Files
browser/platform/service/resources.go
T

2461 lines
100 KiB
Go

package service
import (
"crypto/pbkdf2"
"crypto/rand"
"crypto/sha256"
"crypto/subtle"
"encoding/base64"
"errors"
"fmt"
"strconv"
"strings"
"sync"
"time"
"browser.local/platform/domain"
"browser.local/platform/repo"
"browser.local/platform/validator"
)
var (
ErrUnauthorized = errors.New("unauthorized")
ErrForbidden = errors.New("forbidden")
)
type ForbiddenError struct {
Reason string
}
func (err ForbiddenError) Error() string {
if strings.TrimSpace(err.Reason) == "" {
return ErrForbidden.Error()
}
return ErrForbidden.Error() + ": " + err.Reason
}
func (err ForbiddenError) Is(target error) bool {
return target == ErrForbidden
}
func forbiddenError(reason string) error {
if strings.TrimSpace(reason) == "" {
return ErrForbidden
}
return ForbiddenError{Reason: reason}
}
type Core interface {
CreateUser(domain.User) (domain.User, error)
UpdateUser(string, domain.User) (domain.User, error)
GetUser(string) (domain.User, error)
ListUsers(domain.UserFilter) ([]domain.User, error)
RegisterUser(domain.UserRegistration) (domain.AuthSession, error)
LoginUser(domain.UserLogin) (domain.AuthSession, error)
LogoutUser(string) error
RotateUserSession(string) (domain.AuthSession, error)
GetCurrentUser(string) (domain.User, error)
UpdateCurrentUserProfile(string, domain.UserProfile) (domain.User, error)
UpdateCurrentUserTheme(string, domain.UserThemePreference) (domain.UserThemePreference, error)
CreateAIProvider(domain.AIProvider) (domain.AIProvider, error)
UpdateAIProvider(string, domain.AIProvider) (domain.AIProvider, error)
SetAIProviderStatus(string, domain.AIProviderStatus) (domain.AIProvider, error)
TestAIProvider(string) (domain.AIProviderTestResult, error)
ListAIProviderModels(string) (domain.AIProviderModels, error)
GetAIProvider(string) (domain.AIProvider, error)
ListAIProviders(domain.AIProviderFilter) ([]domain.AIProvider, error)
InvokeAIForSession(string, domain.AIInvocationRequest) (domain.AIInvocationResponse, error)
CreateGamePlugin(domain.GamePlugin) (domain.GamePlugin, error)
RegisterGamePluginManifest(domain.GamePluginManifestRegistration) (domain.GamePlugin, error)
GetGamePlugin(string) (domain.GamePlugin, error)
ListGamePlugins(domain.GamePluginFilter) ([]domain.GamePlugin, error)
ListMarketplacePlugins(domain.PluginMarketplaceFilter) ([]domain.PluginMarketplacePlugin, error)
GetMarketplacePlugin(string) (domain.PluginMarketplacePlugin, error)
SetMarketplacePluginState(string, domain.PluginMarketplaceStateAction) (domain.PluginMarketplacePlugin, error)
AuthorizePluginBridgeAction(domain.PluginBridgeAuthorizeRequest) (domain.PluginBridgeAuthorization, error)
AuthorizePluginBridgeActionForSession(string, domain.PluginBridgeAuthorizeRequest) (domain.PluginBridgeAuthorization, error)
ExecutePluginBridgeAction(string, domain.PluginBridgeExecuteRequest) (domain.PluginBridgeExecuteResponse, error)
CreateRunEndpoint(domain.RunEndpoint) (domain.RunEndpoint, error)
GetRunEndpoint(string) (domain.RunEndpoint, error)
ListRunEndpoints(domain.RunEndpointFilter) ([]domain.RunEndpoint, error)
RegisterRunHello(domain.RunControlHello) (domain.RunControlHelloResult, error)
AcceptRunHeartbeat(domain.RunControlHeartbeat) (domain.RunControlHeartbeatResult, error)
AuthorizeRunRequestSignature(domain.RunRequestSignature) error
CreateServerInstance(domain.ServerInstance) (domain.ServerInstance, error)
CreateServerInstanceForSession(string, domain.ServerInstance) (domain.ServerInstance, error)
CreateServerInstanceWorkflow(domain.ServerLifecycleCreate) (domain.ServerLifecycleResult, error)
CreateServerInstanceWorkflowForSession(string, domain.ServerLifecycleCreate) (domain.ServerLifecycleResult, error)
GetServerDeploymentForSession(string, string) (domain.ServerDeploymentView, error)
RevealServerDeploymentForSession(string, string) (domain.ServerDeploymentReveal, error)
UpdateServerDeploymentForSession(string, string, domain.ServerDeploymentUpdate) (domain.ServerDeploymentView, error)
DeployServerInstanceForSession(string, domain.ServerLifecycleCommand) (domain.ServerLifecycleResult, error)
StartServerInstance(domain.ServerLifecycleCommand) (domain.ServerLifecycleResult, error)
StartServerInstanceForSession(string, domain.ServerLifecycleCommand) (domain.ServerLifecycleResult, error)
StopServerInstance(domain.ServerLifecycleCommand) (domain.ServerLifecycleResult, error)
StopServerInstanceForSession(string, domain.ServerLifecycleCommand) (domain.ServerLifecycleResult, error)
QueryServerInstanceProcessForSession(string, domain.ServerLifecycleCommand) (domain.ServerLifecycleResult, error)
GetServerInstance(string) (domain.ServerInstance, error)
GetServerInstanceForSession(string, string) (domain.ServerInstance, error)
UpdateServerInstanceForSession(string, string, domain.ServerInstanceUpdate) (domain.ServerInstance, error)
ListServerInstances(domain.ServerInstanceFilter) ([]domain.ServerInstance, error)
ListServerInstancesForSession(string, domain.ServerInstanceFilter) ([]domain.ServerInstance, error)
ListServerAdministratorCandidates(string, string) ([]domain.User, error)
AddServerAdministrator(string, string, string) (domain.ServerInstance, error)
RemoveServerAdministrator(string, string, string) (domain.ServerInstance, error)
DeleteServerInstanceForSession(string, string, string) (domain.ServerInstance, error)
GetPlatformResourceUsage() (domain.PlatformResourceUsage, error)
ListServerMetricsForSession(string) ([]domain.ServerMetrics, error)
GetProductionCapacityForSession(string) (domain.ProductionCapacitySummary, error)
CheckCapacityAdmissionForSession(string, domain.CapacityAdmissionRequest) (domain.CapacityAdmissionDecision, error)
ListAlertsForSession(string, domain.AlertFilter) ([]domain.AlertRecord, error)
AcknowledgeAlertForSession(string, domain.AlertAcknowledgeRequest) (domain.AlertRecord, error)
ResolveAlertForSession(string, domain.AlertResolveRequest) (domain.AlertRecord, error)
RetryAlertForSession(string, domain.AlertRetryRequest) (domain.AlertRetryResult, error)
ListPluginLifecyclesForSession(string, domain.PluginLifecycleFilter) ([]domain.PluginLifecycleInstallation, error)
RunPluginLifecycleForSession(string, domain.PluginLifecycleRequest) (domain.PluginLifecycleResult, error)
ListAIConfigDiffsForSession(string, domain.AIConfigDiffFilter) ([]domain.AIConfigDiffPreview, error)
ApproveAIConfigDiffForSession(string, domain.AIConfigDiffApprovalRequest) (domain.AIConfigDiffApprovalResult, error)
IngestMetricBatch(domain.MetricBatchIngest) (domain.MetricBatchIngestResult, error)
ListMetricSamplesForSession(string, domain.MetricSampleFilter) ([]domain.MetricSample, error)
CreateBackupForSession(string, domain.BackupRecord) (domain.BackupRecord, error)
GetBackupForSession(string, string) (domain.BackupRecord, error)
ListBackupsForSession(string, domain.BackupFilter) ([]domain.BackupRecord, error)
ListRemoteAdapterDeclarationsForSession(string, string) ([]domain.RemoteAdapterDeclaration, error)
RequestRemoteAdapterForSession(string, domain.RemoteAdapterRequest) (domain.RemoteAdapterResult, error)
GetServerConfigForSession(string, string) (domain.ServerConfig, error)
PreviewServerConfigWriteForSession(string, domain.ServerConfigDiffRequest) (domain.ServerConfigDiffPreview, error)
ApproveServerConfigWriteForSession(string, domain.ServerConfigWriteApproval) (domain.ServerConfigWriteDispatch, error)
DispatchFileOperationForSession(string, domain.FileOperationDispatchRequest) (domain.FileOperationDispatchResult, error)
CreateJob(domain.Job) (domain.Job, error)
GetJob(string) (domain.Job, error)
ListJobs(domain.JobFilter) ([]domain.Job, error)
GetJobForSession(string, string) (domain.Job, error)
ListJobsForSession(string, domain.JobFilter) ([]domain.Job, error)
RequestRunJobCancelForSession(string, domain.RunJobCancelRequest) (domain.RunJobCancelRequestResult, error)
ClaimRunJob(domain.RunJobClaim) (domain.RunJobClaimResult, error)
AckRunJob(domain.RunJobAck) (domain.RunJobAckResult, error)
UpdateRunJobProgress(domain.RunJobProgress) (domain.RunJobProgressResult, error)
CompleteRunJob(domain.RunJobResult) (domain.RunJobResultResult, error)
GetDistributionBuildInput(domain.DistributionBuildInputRequest) (domain.DistributionBuildInput, error)
GetDependencyExecutionInput(domain.DependencyExecutionInputRequest) (domain.DependencyExecutionInput, error)
DispatchSourceRCONCommandForSession(string, domain.SourceRCONCommandRequest) (domain.SourceRCONCommandDispatch, error)
GetSourceRCONExecutionInput(domain.SourceRCONExecutionInputRequest) (domain.SourceRCONExecutionInput, error)
GetRunUpdateInput(domain.RunUpdateInputRequest) (domain.RunUpdateInput, error)
ReadRunUpdateChunk(domain.RunUpdateChunkRequest) (domain.RunUpdateChunk, error)
ReportRunUpdateHealth(domain.RunUpdateHealthReport) (domain.RunUpdateHealthResult, error)
RequestRunJobCancel(domain.RunJobCancelRequest) (domain.RunJobCancelRequestResult, error)
PollRunJobCancel(domain.RunJobCancelPoll) (domain.RunJobCancelPollResult, error)
ReconcileRunJobs(domain.RunJobReconcile) (domain.RunJobReconcileResult, error)
CreateArtifact(domain.Artifact) (domain.Artifact, error)
GetArtifact(string) (domain.Artifact, error)
ListArtifacts(domain.ArtifactFilter) ([]domain.Artifact, error)
ListArtifactsForSession(string, domain.ArtifactFilter) ([]domain.Artifact, error)
GetArtifactForSession(string, string) (domain.Artifact, error)
OpenArtifactDownloadForSession(string, domain.ArtifactDownloadReferenceRequest) (domain.ArtifactDownloadReference, error)
ReadArtifactContentForSession(string, domain.ArtifactContentRequest) (domain.ArtifactContent, error)
GetServerRuntimeActionsForSession(string, string) (domain.ServerRuntimeActions, error)
GetServerRuntimeBindingForSession(string, string) (domain.RuntimeBindingView, error)
UpdateServerRuntimeBindingForSession(string, string, domain.RuntimeBindingUpdate) (domain.RuntimeBindingView, error)
GenerateRunDistributionForSession(string, domain.RunDistributionGenerateRequest) (domain.RunDistribution, error)
GenerateClientManagerDistributionForSession(string, domain.ClientManagerBuildRequest) (domain.ClientManagerDistribution, error)
OpenLatestRunDistributionDownloadForSession(string, string) (domain.ArtifactDownloadReference, error)
OpenLatestClientManagerDistributionDownloadForSession(string, string, string) (domain.ArtifactDownloadReference, error)
ResetComponentKeyForSession(string, domain.ComponentKeyResetRequest) (domain.EncryptedComponentKey, error)
AuthenticateComponent(domain.ComponentAuthenticationRequest) (domain.ComponentAuthenticationResult, error)
DeployClientManagerForSession(string, domain.ClientManagerDeployRequest) (domain.ClientManagerLifecycleView, error)
ControlClientManagerForSession(string, domain.ClientManagerControlRequest) (domain.ClientManagerLifecycleView, error)
UpdateClientManagerForSession(string, domain.ClientManagerUpdateRequest) (domain.ClientManagerLifecycleView, error)
UninstallClientManagerForSession(string, domain.ClientManagerUninstallRequest) (domain.ClientManagerLifecycleView, error)
RetryClientManagerLifecycleForSession(string, domain.ClientManagerRetryRequest) (domain.ClientManagerLifecycleView, error)
RevokeClientManagerSessionForSession(string, domain.ClientManagerRevokeSessionRequest) (domain.ClientManagerLifecycleView, error)
GetClientManagerLifecycleForSession(string, string, string) (domain.ClientManagerLifecycleView, error)
ListClientManagerLifecyclesForSession(string, string) ([]domain.ClientManagerLifecycleView, error)
GetClientManagerLifecycleInput(domain.ClientManagerLifecycleInputRequest) (domain.ClientManagerLifecycleInput, error)
ReadClientManagerLifecycleChunk(domain.RunUpdateChunkRequest) (domain.RunUpdateChunk, error)
RegisterClientManager(domain.ClientManagerRegisterRequest) (domain.ClientManagerRegisterResult, error)
AcceptClientManagerHeartbeat(domain.ClientManagerHeartbeat) (domain.ClientManagerHeartbeatResult, error)
ReconcileClientManagerLifecycle() error
QueueGameClientBridgeCommandForSession(string, domain.GameClientBridgeQueueRequest) (domain.GameClientBridgeCommand, error)
ClaimGameClientBridgeCommands(domain.GameClientBridgeClaimRequest) ([]domain.GameClientBridgeCommand, error)
AckGameClientBridgeCommand(domain.GameClientBridgeAckRequest) (domain.GameClientBridgeCommand, error)
CompleteGameClientBridgeCommand(domain.GameClientBridgeResultRequest) (domain.GameClientBridgeCommand, error)
CancelGameClientBridgeCommandForSession(string, domain.GameClientBridgeCancelRequest) (domain.GameClientBridgeCommand, error)
UploadGameClientBridgeSnapshot(domain.GameClientBridgeSnapshotIngestRequest) (domain.GameClientBridgeSnapshot, error)
ReconcileGameClientBridgeCommands() error
GetGameClientBridgeStatusForSession(string, string) (domain.GameClientBridgeStatus, error)
ListGameClientBridgeCommandsForSession(string, domain.GameClientBridgeCommandFilter) ([]domain.GameClientBridgeCommand, error)
GetGameClientBridgeCommandForSession(string, string) (domain.GameClientBridgeCommand, error)
QueryGameClientBridgeSnapshotsForSession(string, domain.GameClientBridgeSnapshotQuery) ([]domain.GameClientBridgeSnapshot, error)
PushRunUpdateForSession(string, domain.RunUpdateRequest) (domain.RunUpdateJob, error)
ListRunUpdateJobsForSession(string, string) ([]domain.RunUpdateJob, error)
GetDependencyCatalogForSession(string, string) (domain.DependencyCatalog, error)
QueueDependencyJobForSession(string, domain.DependencyJobRequest) (domain.Job, error)
QueueLogBackfillForSession(string, domain.LogBackfillRequest) (domain.Job, error)
OpenArtifactTransfer(domain.ArtifactTransferOpen) (domain.ArtifactTransferOpenResult, error)
UploadArtifactChunk(domain.ArtifactChunkUpload) (domain.ArtifactChunkUploadResult, error)
QueryArtifactTransferStatus(domain.ArtifactTransferStatusQuery) (domain.ArtifactTransferStatusResult, error)
CompleteArtifactTransfer(domain.ArtifactTransferComplete) (domain.ArtifactTransferCompleteResult, error)
CreateLogStream(domain.LogStream) (domain.LogStream, error)
GetLogStream(string) (domain.LogStream, error)
ListLogStreams(domain.LogStreamFilter) ([]domain.LogStream, error)
GetLogStreamForSession(string, string) (domain.LogStream, error)
ListLogStreamsForSession(string, domain.LogStreamFilter) ([]domain.LogStream, error)
QueryLogStreamForSession(string, domain.LogStreamCursorQuery) (domain.LogStreamCursorResult, error)
IngestLogBatch(domain.LogBatchIngest) (domain.LogBatchIngestResult, error)
QueryLogStream(domain.LogStreamCursorQuery) (domain.LogStreamCursorResult, error)
ListGamePlayersForSession(string, domain.GamePlayerFilter) ([]domain.GamePlayer, error)
GetGamePlayerProfileForSession(string, string) (domain.GamePlayerProfile, error)
GetGamePlayerStateForSession(string, string) (domain.GamePlayerStateSnapshot, error)
RequestGamePlayerStatePatchForSession(string, string, domain.GamePlayerStatePatchRequest) (domain.GamePlayerStatePatch, error)
ApproveGamePlayerStatePatchForSession(string, string) (domain.GamePlayerStatePatch, error)
ListGamePlayerStatePatchesForSession(string, string) ([]domain.GamePlayerStatePatch, error)
CreateAuditEvent(domain.AuditEvent) (domain.AuditEvent, error)
GetAuditEvent(string) (domain.AuditEvent, error)
ListAuditEvents(domain.AuditEventFilter) ([]domain.AuditEvent, error)
SeedPlatformAdmin(string, string) error
}
type CoreService struct {
store repo.Store
now func() time.Time
authMu sync.Mutex
authSessions map[string]string
controlMu sync.Mutex
runSessions map[string]domain.RunControlSession
runSessionSeq uint64
jobMu sync.Mutex
bridgeMu sync.Mutex
bridgeSeq uint64
logStore LogBodyStore
artifactStore ArtifactBodyStore
artifactMu sync.Mutex
artifactTransfers map[string]domain.ArtifactTransferSession
artifactPayloads map[string][]byte
artifactTransferSeq uint64
auditMu sync.Mutex
auditSeq uint64
productionMu sync.Mutex
sourceRCONCommands *sourceRCONCommandBroker
aiProviderClient AIProviderClient
secretEnvelope SecretEnvelope
networkFingerprintKey []byte
}
var _ Core = (*CoreService)(nil)
func NewCoreService(store repo.Store) *CoreService {
return newCoreService(store, func() time.Time { return time.Now().UTC() })
}
func NewCoreServiceWithLogStore(store repo.Store, logStore LogBodyStore) *CoreService {
return newCoreServiceWithLogStore(store, logStore, func() time.Time { return time.Now().UTC() })
}
func newCoreService(store repo.Store, now func() time.Time) *CoreService {
return newCoreServiceWithLogStore(store, NewMemoryLogBodyStore(), now)
}
func newCoreServiceWithLogStore(store repo.Store, logStore LogBodyStore, now func() time.Time) *CoreService {
if logStore == nil {
logStore = NewMemoryLogBodyStore()
}
artifactStore := NewMemoryArtifactBodyStore()
service := &CoreService{
store: store,
now: now,
authSessions: map[string]string{},
runSessions: map[string]domain.RunControlSession{},
logStore: logStore,
artifactStore: artifactStore,
artifactTransfers: map[string]domain.ArtifactTransferSession{},
artifactPayloads: map[string][]byte{},
sourceRCONCommands: newSourceRCONCommandBroker(now),
aiProviderClient: MockAIProviderClient{},
secretEnvelope: newSecretEnvelope(developmentSecretEnvelopeKey),
networkFingerprintKey: []byte(developmentSecretEnvelopeKey),
}
return service
}
func NewCoreServiceWithDurableStores(store repo.Store, logStore LogBodyStore, artifactStore ArtifactBodyStore) (*CoreService, error) {
if artifactStore == nil {
artifactStore = NewMemoryArtifactBodyStore()
}
service := newCoreServiceWithLogStore(store, logStore, func() time.Time { return time.Now().UTC() })
service.artifactStore = artifactStore
sessions, err := artifactStore.LoadTransfers()
if err != nil {
return nil, err
}
for _, session := range sessions {
service.artifactTransfers[session.TransferID] = domain.CopyArtifactTransferSession(session)
}
if err := service.recoverLogCursors(); err != nil {
return nil, err
}
if err := service.RecoverIncompleteBackups(); err != nil {
return nil, err
}
if err := service.ReconcileClientManagerLifecycle(); err != nil {
return nil, err
}
if err := service.ReconcileGameClientBridgeCommands(); err != nil {
return nil, err
}
return service, nil
}
func (svc *CoreService) recoverLogCursors() error {
store, ok := svc.logStore.(interface{ LatestSeq(string) (uint64, error) })
if !ok {
return nil
}
streams, err := svc.store.LogStreams().List(domain.LogStreamFilter{})
if err != nil {
return err
}
for _, stream := range streams {
latest, latestErr := store.LatestSeq(stream.ID)
if latestErr != nil || latest <= stream.LatestSeq {
if latestErr != nil {
return latestErr
}
continue
}
stream.LatestSeq = latest
stream.UpdatedAt = svc.now()
if err := svc.store.LogStreams().Update(stream); err != nil {
return err
}
}
return nil
}
func (svc *CoreService) CreateUser(user domain.User) (domain.User, error) {
if strings.TrimSpace(user.ID) == "" {
id, err := svc.nextUserID(user)
if err != nil {
return domain.User{}, err
}
user.ID = id
}
if user.Status == "" {
user.Status = domain.UserStatusActive
}
if len(user.Roles) == 0 {
user.Roles = []string{"server-admin"}
}
if strings.TrimSpace(user.PasswordHash) != "" && !strings.Contains(user.PasswordHash, "$") {
hash, err := hashPassword(user.PasswordHash)
if err != nil {
return domain.User{}, err
}
user.PasswordHash = hash
}
stamp := svc.now()
if user.CreatedAt.IsZero() {
user.CreatedAt = stamp
}
if user.UpdatedAt.IsZero() {
user.UpdatedAt = stamp
}
if err := validator.ValidateUser(user); err != nil {
return domain.User{}, err
}
if err := svc.store.Users().Create(user); err != nil {
return domain.User{}, err
}
return domain.CopyUser(user), nil
}
func (svc *CoreService) UpdateUser(id string, user domain.User) (domain.User, error) {
existing, err := svc.store.Users().Get(id)
if err != nil {
return domain.User{}, err
}
user.ID = id
if user.PasswordHash == "" {
user.PasswordHash = existing.PasswordHash
} else if !strings.Contains(user.PasswordHash, "$") {
hash, err := hashPassword(user.PasswordHash)
if err != nil {
return domain.User{}, err
}
user.PasswordHash = hash
}
if user.CreatedAt.IsZero() {
user.CreatedAt = existing.CreatedAt
}
user.UpdatedAt = svc.now()
if user.Status == "" {
user.Status = existing.Status
}
if user.Roles == nil {
user.Roles = domain.CopyStringSlice(existing.Roles)
}
if err := validator.ValidateUser(user); err != nil {
return domain.User{}, err
}
if err := svc.store.Users().Update(user); err != nil {
return domain.User{}, err
}
if user.Status != domain.UserStatusActive {
if err := svc.revokeUserSessions(user.ID); err != nil {
return domain.User{}, err
}
}
return domain.CopyUser(user), nil
}
func (svc *CoreService) GetUser(id string) (domain.User, error) {
return svc.store.Users().Get(id)
}
func (svc *CoreService) ListUsers(filter domain.UserFilter) ([]domain.User, error) {
return svc.store.Users().List(filter)
}
func (svc *CoreService) RegisterUser(registration domain.UserRegistration) (domain.AuthSession, error) {
if len([]rune(registration.Password)) < 6 {
return domain.AuthSession{}, validationError("password must be at least 6 characters")
}
hash, err := hashPassword(registration.Password)
if err != nil {
return domain.AuthSession{}, err
}
users, err := svc.store.Users().List(domain.UserFilter{})
if err != nil {
return domain.AuthSession{}, err
}
firstUser := len(users) == 0
user := domain.User{
ID: userIDFromEmail(registration.Email),
DisplayName: registration.DisplayName,
Email: registration.Email,
Status: domain.UserStatusPending,
Roles: []string{"server-admin"},
PasswordHash: hash,
Profile: registration.Profile,
}
if firstUser {
user.Status = domain.UserStatusActive
user.Roles = []string{"platform-admin"}
}
created, err := svc.CreateUser(user)
if err != nil {
return domain.AuthSession{}, err
}
if firstUser {
return svc.issueAuthSession(created, "首个账号已创建为平台管理员。")
}
return domain.AuthSession{
User: created,
Status: "pending",
Message: "注册申请已提交,等待平台管理员审核。",
}, nil
}
func (svc *CoreService) LoginUser(login domain.UserLogin) (domain.AuthSession, error) {
var matched domain.User
users, err := svc.store.Users().List(domain.UserFilter{})
if err != nil {
return domain.AuthSession{}, err
}
account := strings.ToLower(strings.TrimSpace(login.Account))
for _, user := range users {
if strings.ToLower(user.ID) == account || strings.ToLower(strings.TrimSpace(user.Email)) == account {
matched = user
break
}
}
if matched.ID == "" || !verifyPassword(matched.PasswordHash, login.Password) {
return domain.AuthSession{}, ErrUnauthorized
}
if matched.Status == domain.UserStatusPending {
return domain.AuthSession{}, ErrForbidden
}
if matched.Status == domain.UserStatusDisabled {
return domain.AuthSession{}, ErrForbidden
}
return svc.issueAuthSession(matched, "登录成功")
}
func (svc *CoreService) LogoutUser(sessionID string) error {
return svc.revokeAuthSession(sessionID)
}
func (svc *CoreService) GetCurrentUser(sessionID string) (domain.User, error) {
userID, err := svc.userIDForSession(sessionID)
if err != nil {
return domain.User{}, err
}
return svc.store.Users().Get(userID)
}
func (svc *CoreService) UpdateCurrentUserProfile(sessionID string, profile domain.UserProfile) (domain.User, error) {
user, err := svc.GetCurrentUser(sessionID)
if err != nil {
return domain.User{}, err
}
user.Profile = profile
return svc.UpdateUser(user.ID, user)
}
func (svc *CoreService) UpdateCurrentUserTheme(sessionID string, preference domain.UserThemePreference) (domain.UserThemePreference, error) {
user, err := svc.GetCurrentUser(sessionID)
if err != nil {
return domain.UserThemePreference{}, err
}
preference.UserID = user.ID
preference.Persistence = "api"
preference.UpdatedAt = svc.now()
user.Theme = preference
if _, err := svc.UpdateUser(user.ID, user); err != nil {
return domain.UserThemePreference{}, err
}
return preference, nil
}
func (svc *CoreService) SeedLocalPlatformAdmin() error {
return svc.SeedPlatformAdmin("operator.local@example.test", "operator-local")
}
func (svc *CoreService) SeedPlatformAdmin(email string, password string) error {
email = strings.TrimSpace(email)
if email == "" {
email = "operator.local@example.test"
}
if len([]rune(password)) < 12 {
return validationError("bootstrap admin password must be at least 12 characters")
}
_, err := svc.store.Users().Get("user-admin")
if err == nil {
return nil
}
if !errors.Is(err, repo.ErrNotFound) {
return err
}
return svc.store.Users().Create(domain.User{
ID: "user-admin",
DisplayName: "Operator",
Email: email,
Status: domain.UserStatusActive,
Roles: []string{"platform-admin"},
PasswordHash: mustHashPassword(password),
Profile: domain.UserProfile{ContactNote: "bootstrap platform admin"},
CreatedAt: svc.now(),
UpdatedAt: svc.now(),
})
}
func (svc *CoreService) CreateAIProvider(provider domain.AIProvider) (domain.AIProvider, error) {
if provider.Status == "" {
provider.Status = domain.AIProviderStatusActive
}
if err := validator.ValidateAIProvider(provider); err != nil {
return domain.AIProvider{}, err
}
if err := svc.store.AIProviders().Create(provider); err != nil {
return domain.AIProvider{}, err
}
return domain.CopyAIProvider(provider), nil
}
func (svc *CoreService) UpdateAIProvider(id string, provider domain.AIProvider) (domain.AIProvider, error) {
existing, err := svc.store.AIProviders().Get(id)
if err != nil {
return domain.AIProvider{}, err
}
provider.ID = id
provider.Status = existing.Status
if err := validator.ValidateAIProvider(provider); err != nil {
return domain.AIProvider{}, err
}
if err := svc.store.AIProviders().Update(provider); err != nil {
return domain.AIProvider{}, err
}
return domain.CopyAIProvider(provider), nil
}
func (svc *CoreService) SetAIProviderStatus(id string, status domain.AIProviderStatus) (domain.AIProvider, error) {
if status != domain.AIProviderStatusActive && status != domain.AIProviderStatusDisabled {
return domain.AIProvider{}, validationError("status must be active or disabled")
}
provider, err := svc.store.AIProviders().Get(id)
if err != nil {
return domain.AIProvider{}, err
}
provider.Status = status
if err := validator.ValidateAIProvider(provider); err != nil {
return domain.AIProvider{}, err
}
if err := svc.store.AIProviders().Update(provider); err != nil {
return domain.AIProvider{}, err
}
return domain.CopyAIProvider(provider), nil
}
func (svc *CoreService) TestAIProvider(id string) (domain.AIProviderTestResult, error) {
provider, err := svc.store.AIProviders().Get(id)
if err != nil {
return domain.AIProviderTestResult{}, err
}
result := domain.AIProviderTestResult{
ProviderID: provider.ID,
Mode: "provider",
Success: true,
Message: "provider invocation passed",
}
if err := validator.ValidateAIProvider(provider); err != nil {
result.Success = false
result.Message = "provider validation failed"
var validationErr validator.ValidationError
if errors.As(err, &validationErr) {
result.Violations = append(result.Violations, validationErr.Violations...)
} else {
result.Violations = append(result.Violations, err.Error())
}
}
if provider.Status != domain.AIProviderStatusActive {
result.Success = false
result.Message = "provider validation failed"
result.Violations = append(result.Violations, "provider must be active")
}
if !result.Success {
return domain.CopyAIProviderTestResult(result), nil
}
_, invokeErr := svc.aiProviderClient.Invoke(provider, domain.AIInvocationRequest{RequestID: "provider-test-" + provider.ID, Purpose: "provider.health", Prompt: "Return a short health acknowledgement.", Model: provider.DefaultModel})
if invokeErr != nil {
result.Success = false
result.Message = "provider invocation failed safely"
result.Violations = []string{"provider invocation failed safely"}
auditID, auditErr := svc.recordAuditEventWithID("platform", "ai.provider.test.failed", "ai-provider", provider.ID, domain.AuditResultFailed, result.Message)
if auditErr != nil {
return domain.AIProviderTestResult{}, auditErr
}
svc.productionMu.Lock()
_, alertErr := svc.upsertAlert(domain.AlertRecord{SourceKind: "ai-provider", SourceID: provider.ID, RuleKey: "ai.provider.failed", Severity: domain.AlertSeverityWarning, Title: "AI provider health check failed", Message: result.Message, Retryable: false, LastAuditEventID: auditID})
svc.productionMu.Unlock()
if alertErr != nil {
return domain.AIProviderTestResult{}, alertErr
}
} else {
svc.productionMu.Lock()
resolveErr := svc.resolveAlertForSource("ai-provider", provider.ID, "ai.provider.failed", "platform", "AI provider health check passed", "")
svc.productionMu.Unlock()
if resolveErr != nil {
return domain.AIProviderTestResult{}, resolveErr
}
}
return domain.CopyAIProviderTestResult(result), nil
}
func (svc *CoreService) ListAIProviderModels(id string) (domain.AIProviderModels, error) {
provider, err := svc.store.AIProviders().Get(id)
if err != nil {
return domain.AIProviderModels{}, err
}
return domain.CopyAIProviderModels(domain.AIProviderModels{
ProviderID: provider.ID,
DefaultModel: provider.DefaultModel,
Models: provider.Models,
}), nil
}
func (svc *CoreService) GetAIProvider(id string) (domain.AIProvider, error) {
return svc.store.AIProviders().Get(id)
}
func (svc *CoreService) ListAIProviders(filter domain.AIProviderFilter) ([]domain.AIProvider, error) {
return svc.store.AIProviders().List(filter)
}
func (svc *CoreService) CreateGamePlugin(plugin domain.GamePlugin) (domain.GamePlugin, error) {
if plugin.Status == "" {
plugin.Status = domain.GamePluginStatusInstalled
}
plugin.ProductionLifecycle = normalizedProductionLifecycle(plugin.ProductionLifecycle)
if err := validator.ValidateGamePlugin(plugin); err != nil {
return domain.GamePlugin{}, err
}
if err := svc.store.GamePlugins().Create(plugin); err != nil {
return domain.GamePlugin{}, err
}
return domain.CopyGamePlugin(plugin), nil
}
func (svc *CoreService) RegisterGamePluginManifest(registration domain.GamePluginManifestRegistration) (domain.GamePlugin, error) {
if err := validator.ValidateGamePluginManifestRegistration(registration); err != nil {
return domain.GamePlugin{}, err
}
return svc.CreateGamePlugin(gamePluginFromManifestRegistration(registration))
}
func gamePluginFromManifestRegistration(registration domain.GamePluginManifestRegistration) domain.GamePlugin {
registration = domain.CopyGamePluginManifestRegistration(registration)
manifest := registration.Manifest
return domain.GamePlugin{
ID: manifest.ID,
Name: manifest.Name,
Description: manifest.Description,
Version: manifest.Version,
ServerType: manifest.Server.Type,
ServerDisplayName: manifest.Server.DisplayName,
SupportedOS: manifest.Server.SupportedOS,
ManifestRef: registration.ManifestRef,
CreateFormSchemaRef: manifest.Server.CreateFormSchema,
CreateFields: manifest.Server.CreateFields,
RequiredRunCapabilities: manifest.Capabilities,
DeclaredPermissions: manifest.Permissions,
Permissions: pluginPermissionsFromManifest(manifest.Permissions),
LifecycleActions: manifest.Actions,
BridgeActions: manifest.Bridge.Actions,
Pages: manifest.Pages,
FileWorkspace: manifest.FileWorkspace,
Tags: manifest.Tags,
AIPurposes: manifest.AI.Purposes,
ProductionLifecycle: manifest.ProductionLifecycle,
RemoteAccess: manifest.RemoteAccess,
RuntimeProfiles: manifest.RuntimeProfiles,
GameClientBridge: manifest.GameClientBridge,
Status: domain.GamePluginStatusInstalled,
}
}
func (svc *CoreService) AuthorizePluginBridgeAction(request domain.PluginBridgeAuthorizeRequest) (domain.PluginBridgeAuthorization, error) {
if err := validator.ValidatePluginBridgeAuthorizeRequest(request); err != nil {
return domain.PluginBridgeAuthorization{}, err
}
plugin, err := svc.store.GamePlugins().Get(request.PluginID)
if err != nil {
return domain.PluginBridgeAuthorization{}, err
}
result, err := validator.AuthorizePluginBridgeAction(plugin, request)
if err != nil {
return domain.PluginBridgeAuthorization{}, err
}
return domain.CopyPluginBridgeAuthorization(result), nil
}
func (svc *CoreService) ExecutePluginBridgeAction(sessionID string, request domain.PluginBridgeExecuteRequest) (domain.PluginBridgeExecuteResponse, error) {
request = domain.CopyPluginBridgeExecuteRequest(request)
if err := validator.ValidatePluginBridgeExecuteRequest(request); err != nil {
return domain.PluginBridgeExecuteResponse{}, err
}
base := domain.PluginBridgeExecuteResponse{
RequestID: request.RequestID,
PluginID: request.PluginID,
RouteKey: request.RouteKey,
ServerInstanceID: request.ServerInstanceID,
Action: request.Action,
}
plugin, err := svc.store.GamePlugins().Get(request.PluginID)
if err != nil {
return domain.PluginBridgeExecuteResponse{}, err
}
authorization, err := validator.AuthorizePluginBridgeAction(plugin, domain.PluginBridgeAuthorizeRequest{
PluginID: request.PluginID,
RouteKey: request.RouteKey,
ServerInstanceID: request.ServerInstanceID,
Action: request.Action,
AIPurpose: request.AIPurpose,
})
if err != nil {
return domain.PluginBridgeExecuteResponse{}, err
}
if !authorization.Allowed {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "permission_denied", Message: safeBridgeReason(authorization.Reason)}
return domain.CopyPluginBridgeExecuteResponse(base), nil
}
var instance domain.ServerInstance
if request.ServerInstanceID != "" {
instance, err = svc.GetServerInstanceForSession(sessionID, request.ServerInstanceID)
if err != nil {
return domain.PluginBridgeExecuteResponse{}, err
}
if instance.PluginID != request.PluginID {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "server_scope_denied", Message: "server instance is outside plugin scope"}
return domain.CopyPluginBridgeExecuteResponse(base), nil
}
}
switch request.Action {
case domain.PluginBridgeActionServerInstancesRead:
base.Status = "ok"
base.Result = map[string]string{
"serverInstanceId": instance.ID,
"pluginId": instance.PluginID,
"pluginVersion": instance.PluginVersion,
"runEndpointId": instance.RunEndpointID,
"state": string(instance.State),
"configVersion": strconv.Itoa(instance.ConfigVersion),
}
case domain.PluginBridgeActionJobsDispatch:
base = svc.executeBridgeJobDispatch(sessionID, base, plugin, instance, request.Payload)
case domain.PluginBridgeActionLogsQuery:
base = svc.executeBridgeLogsQuery(base, instance, request.Payload)
case domain.PluginBridgeActionFilesRequest:
base = svc.executeBridgeFileRequest(sessionID, base, request)
case domain.PluginBridgeActionRemoteAccessRequest:
base = svc.executeBridgeRemoteAccessRequest(sessionID, base, plugin, instance, request.Payload)
case domain.PluginBridgeActionRunDistribution:
base = svc.executeBridgeRunDistribution(sessionID, base, request)
case domain.PluginBridgeActionDependenciesRequest:
base = svc.executeBridgeDependenciesRequest(sessionID, base, plugin, instance, request.Payload)
case domain.PluginBridgeActionLogsBackfillRequest:
base = svc.executeBridgeLogsBackfillRequest(base, plugin, instance, request.Payload)
case domain.PluginBridgeActionClientManager:
base = svc.executeBridgeClientManager(sessionID, base, request)
case domain.PluginBridgeActionPluginLifecycle:
base = svc.executeBridgePluginLifecycle(sessionID, base, request)
case domain.PluginBridgeActionArtifactsOpen:
base = svc.executeBridgeArtifactOpen(sessionID, base, request)
case domain.PluginBridgeActionAIInvoke:
base = svc.executeBridgeAIInvoke(sessionID, base, request)
default:
base.Status = "unsupported"
base.Error = &domain.PluginBridgeSafeError{Code: "unsupported_action", Message: "bridge action is not supported"}
}
return domain.CopyPluginBridgeExecuteResponse(base), nil
}
func (svc *CoreService) executeBridgeArtifactOpen(sessionID string, base domain.PluginBridgeExecuteResponse, request domain.PluginBridgeExecuteRequest) domain.PluginBridgeExecuteResponse {
artifactID := strings.TrimSpace(request.Payload["artifactId"])
if artifactID == "" {
base.Status = "error"
base.Error = &domain.PluginBridgeSafeError{Code: "validation", Message: "artifactId is required"}
return base
}
reference, err := svc.OpenArtifactDownloadForSession(sessionID, domain.ArtifactDownloadReferenceRequest{ArtifactID: artifactID})
if err != nil {
return bridgeExecutionError(base, err)
}
if reference.OwnerKind == domain.ArtifactOwnerKindJob {
job, err := svc.GetJob(reference.OwnerID)
if err != nil {
return bridgeExecutionError(base, err)
}
if job.ServerInstanceID != request.ServerInstanceID {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "server_scope_denied", Message: "artifact is outside server scope"}
return base
}
}
if reference.OwnerKind == domain.ArtifactOwnerKindServerInstance && reference.OwnerID != request.ServerInstanceID {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "server_scope_denied", Message: "artifact is outside server scope"}
return base
}
base.Status = "ok"
base.Result = map[string]string{
"artifactId": reference.ArtifactID,
"filename": reference.Filename,
"contentType": reference.ContentType,
"sizeBytes": strconv.FormatInt(reference.SizeBytes, 10),
"checksum": reference.Checksum,
"downloadUrl": reference.DownloadURL,
"expiresAt": reference.ExpiresAt.Format(time.RFC3339),
"rangeSupported": strconv.FormatBool(reference.RangeSupported),
"chunkSizeBytes": strconv.Itoa(reference.ChunkSizeBytes),
"storageBehavior": reference.StorageBehavior,
}
return base
}
func (svc *CoreService) executeBridgeAIInvoke(sessionID string, base domain.PluginBridgeExecuteResponse, request domain.PluginBridgeExecuteRequest) domain.PluginBridgeExecuteResponse {
response, err := svc.InvokeAIForSession(sessionID, domain.AIInvocationRequest{
RequestID: request.RequestID,
PluginID: request.PluginID,
RouteKey: request.RouteKey,
ServerInstanceID: request.ServerInstanceID,
Purpose: request.AIPurpose,
Model: request.Payload["model"],
Prompt: defaultBridgeValue(request.Payload["prompt"], "Review the current server context and provide a safe recommendation."),
CurrentConfig: request.Payload["currentConfig"],
ContextRefs: map[string]string{"server": "server://" + request.ServerInstanceID},
})
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = response.Status
base.Result = map[string]string{
"purpose": response.Purpose,
"recommendation": response.Recommendation,
"providerId": response.ProviderID,
"model": response.Model,
"mocked": strconv.FormatBool(response.Usage.Mocked),
}
if response.ConfigRecommendation != nil {
base.Result["suggestedConfig"] = response.ConfigRecommendation.SuggestedConfig
base.Result["diffSummary"] = response.ConfigRecommendation.DiffSummary
base.Result["diffId"] = response.ConfigRecommendation.DiffID
base.Result["key"] = response.ConfigRecommendation.Key
base.Result["expiresAt"] = response.ConfigRecommendation.ExpiresAt
}
if response.Error != nil {
base.Error = &domain.PluginBridgeSafeError{Code: response.Error.Code, Message: response.Error.Message, Details: response.Error.Details}
}
return base
}
func (svc *CoreService) executeBridgePluginLifecycle(sessionID string, base domain.PluginBridgeExecuteResponse, request domain.PluginBridgeExecuteRequest) domain.PluginBridgeExecuteResponse {
confirmed, err := strconv.ParseBool(defaultBridgeValue(request.Payload["confirmed"], "false"))
if err != nil {
base.Status = "error"
base.Error = &domain.PluginBridgeSafeError{Code: "validation", Message: "confirmed must be true or false"}
return base
}
result, err := svc.RunPluginLifecycleForSession(sessionID, domain.PluginLifecycleRequest{PluginID: request.PluginID, ServerInstanceID: request.ServerInstanceID, Operation: domain.PluginLifecycleOperation(request.Payload["operation"]), TargetVersion: request.Payload["targetVersion"], IdempotencyKey: defaultBridgeValue(request.Payload["idempotencyKey"], request.RequestID), Confirmed: confirmed})
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = result.Status
base.Result = map[string]string{"installationId": result.Installation.ID, "currentState": string(result.Installation.CurrentState), "desiredState": string(result.Installation.DesiredState), "jobId": result.Job.ID, "alertId": result.Installation.AlertID, "auditEventId": result.Installation.AuditEventID, "admissionState": string(result.Decision.State), "admissionReason": result.Decision.Reason}
return base
}
func (svc *CoreService) executeBridgeJobDispatch(sessionID string, base domain.PluginBridgeExecuteResponse, plugin domain.GamePlugin, instance domain.ServerInstance, payload map[string]string) domain.PluginBridgeExecuteResponse {
capability := strings.TrimSpace(payload["capability"])
if capability == "" {
capability = "process.start"
}
if !containsString(plugin.RequiredRunCapabilities, capability) {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "capability_denied", Message: "requested job capability is not declared by plugin"}
return base
}
lifecycleAction := domain.ServerLifecycleAction(strings.TrimSpace(payload["lifecycleAction"]))
if lifecycleAction == domain.ServerLifecycleActionStart || lifecycleAction == domain.ServerLifecycleActionStop {
expectedVersion, _ := strconv.Atoi(payload["expectedConfigVersion"])
command := domain.ServerLifecycleCommand{
ServerInstanceID: instance.ID,
ExpectedConfigVersion: expectedVersion,
IdempotencyKey: defaultBridgeValue(payload["idempotencyKey"], base.RequestID),
}
var result domain.ServerLifecycleResult
var err error
switch lifecycleAction {
case domain.ServerLifecycleActionStart:
if capability != domain.LifecycleCapabilityStart {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "capability_denied", Message: "lifecycle action must match requested capability"}
return base
}
result, err = svc.StartServerInstanceForSession(sessionID, command)
case domain.ServerLifecycleActionStop:
if capability != domain.LifecycleCapabilityStop {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "capability_denied", Message: "lifecycle action must match requested capability"}
return base
}
result, err = svc.StopServerInstanceForSession(sessionID, command)
}
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = "queued"
base.Result = map[string]string{
"jobId": result.Job.ID,
"state": string(result.Job.State),
"capability": result.Job.Capability,
"lifecycleAction": string(result.Action),
"serverInstanceId": result.Instance.ID,
}
return base
}
job, err := svc.CreateJob(domain.Job{
ID: jobIDFromParts("job-bridge", base.RequestID, capability),
ServerInstanceID: instance.ID,
RunEndpointID: instance.RunEndpointID,
Capability: capability,
IdempotencyKey: defaultBridgeValue(payload["idempotencyKey"], base.RequestID),
Progress: domain.JobProgress{Percent: 0, Message: "plugin bridge job queued"},
})
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = "queued"
base.Result = map[string]string{"jobId": job.ID, "state": string(job.State), "capability": job.Capability}
return base
}
func (svc *CoreService) executeBridgeLogsQuery(base domain.PluginBridgeExecuteResponse, instance domain.ServerInstance, payload map[string]string) domain.PluginBridgeExecuteResponse {
streamID := strings.TrimSpace(payload["logStreamId"])
if streamID == "" {
base.Status = "error"
base.Error = &domain.PluginBridgeSafeError{Code: "validation", Message: "logStreamId is required"}
return base
}
stream, err := svc.GetLogStream(streamID)
if err != nil {
return bridgeExecutionError(base, err)
}
if stream.ServerInstanceID != instance.ID {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "server_scope_denied", Message: "log stream is outside server scope"}
return base
}
afterSeq, _ := strconv.ParseUint(payload["afterSeq"], 10, 64)
limit, _ := strconv.Atoi(payload["limit"])
result, err := svc.QueryLogStream(domain.LogStreamCursorQuery{LogStreamID: streamID, AfterSeq: afterSeq, Limit: limit})
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = "ok"
base.Result = map[string]string{
"logStreamId": stream.ID,
"entryCount": strconv.Itoa(len(result.Entries)),
"nextSeq": strconv.FormatUint(result.NextSeq, 10),
"latestSeq": strconv.FormatUint(result.LatestSeq, 10),
}
return base
}
func (svc *CoreService) executeBridgeFileRequest(sessionID string, base domain.PluginBridgeExecuteResponse, request domain.PluginBridgeExecuteRequest) domain.PluginBridgeExecuteResponse {
payload := request.Payload
operation := domain.FileOperationKind(defaultBridgeValue(payload["operation"], string(domain.FileOperationRead)))
expectedVersion, _ := strconv.Atoi(payload["expectedConfigVersion"])
result, err := svc.DispatchFileOperationForSession(sessionID, domain.FileOperationDispatchRequest{
ServerInstanceID: request.ServerInstanceID,
PluginID: request.PluginID,
Operation: operation,
Key: payload["key"],
InputRef: payload["inputRef"],
ExpectedConfigVersion: expectedVersion,
IdempotencyKey: defaultBridgeValue(payload["idempotencyKey"], request.RequestID),
})
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = "queued"
base.Result = map[string]string{
"jobId": result.Job.ID,
"state": string(result.Job.State),
"capability": result.Job.Capability,
"targetKey": result.Key,
}
return base
}
func (svc *CoreService) executeBridgeRemoteAccessRequest(sessionID string, base domain.PluginBridgeExecuteResponse, plugin domain.GamePlugin, instance domain.ServerInstance, payload map[string]string) domain.PluginBridgeExecuteResponse {
capability := strings.TrimSpace(payload["capability"])
if capability == "" {
base.Status = "error"
base.Error = &domain.PluginBridgeSafeError{Code: "validation", Message: "capability is required"}
return base
}
if !containsString(plugin.RequiredRunCapabilities, capability) || !containsString(plugin.RemoteAccess.RunCapabilities, capability) {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "capability_denied", Message: "requested remote capability is not declared by plugin"}
return base
}
declarationKey := strings.TrimSpace(payload["declarationKey"])
if declarationKey == "" {
for _, profile := range plugin.RuntimeProfiles.TransportProfiles {
if profile.TargetKey == payload["targetKey"] && containsString(profile.Capabilities, capability) {
declarationKey = profile.Key
break
}
}
}
if declarationKey == "" {
declarationKey = "legacy-" + string(remoteAdapterKindForCapability(capability))
}
timeoutSeconds, _ := strconv.Atoi(payload["timeoutSeconds"])
maxAttempts, _ := strconv.Atoi(payload["maxAttempts"])
inputs := map[string]string{}
for key, value := range payload {
if strings.HasPrefix(key, "input.") {
inputs[strings.TrimPrefix(key, "input.")] = value
}
}
if capability == domain.JobCapabilityRemoteRunDBSQLiteQuery {
templateKey := strings.TrimSpace(inputs["templateKey"])
if templateKey == "" {
base.Status = "error"
base.Error = &domain.PluginBridgeSafeError{Code: "validation", Message: "input.templateKey is required for sqlite query requests"}
return base
}
template, reason := findBridgeQueryTemplate(plugin, base.RouteKey, templateKey)
if reason != "" {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "query_template_denied", Message: reason}
return base
}
if template.Engine != "sqlite" || template.TransportKey != declarationKey || template.TargetKey != payload["targetKey"] {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "query_template_denied", Message: "query template transport or target is not approved"}
return base
}
if timeoutSeconds == 0 {
timeoutSeconds = template.TimeoutSeconds
} else if timeoutSeconds > template.TimeoutSeconds {
base.Status = "error"
base.Error = &domain.PluginBridgeSafeError{Code: "validation", Message: "query template timeout limit exceeded"}
return base
}
maxRows := template.MaxRows
if requestedRows, ok := inputs["maxRows"]; ok && strings.TrimSpace(requestedRows) != "" {
parsedRows, parseErr := strconv.Atoi(requestedRows)
if parseErr != nil || parsedRows <= 0 {
base.Status = "error"
base.Error = &domain.PluginBridgeSafeError{Code: "validation", Message: "input.maxRows must be a positive integer"}
return base
}
if parsedRows < maxRows {
maxRows = parsedRows
}
}
inputs["templateKey"] = template.Key
inputs["maxRows"] = strconv.Itoa(maxRows)
}
result, err := svc.RequestRemoteAdapterForSession(sessionID, domain.RemoteAdapterRequest{ServerInstanceID: instance.ID, DeclarationKey: declarationKey, TargetKey: payload["targetKey"], Capability: capability, TimeoutSeconds: timeoutSeconds, MaxAttempts: maxAttempts, IdempotencyKey: defaultBridgeValue(payload["idempotencyKey"], base.RequestID), InputRef: payload["inputRef"], Inputs: inputs})
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = "queued"
base.Result = map[string]string{
"jobId": result.RequestID,
"state": result.Status,
"capability": capability,
"targetKey": result.TargetKey,
"serverInstanceId": result.ServerInstanceID,
"adapterKind": string(result.Kind),
}
return base
}
func findBridgeQueryTemplate(plugin domain.GamePlugin, routeKey string, templateKey string) (domain.GameClientBridgeQueryTemplateDeclaration, string) {
pageFound := false
pageAllowsTemplate := false
for _, page := range plugin.GameClientBridge.Pages {
if page.PageKey != routeKey {
continue
}
pageFound = true
if containsString(page.QueryTemplateKeys, templateKey) {
pageAllowsTemplate = true
}
}
if !pageFound || !pageAllowsTemplate {
return domain.GameClientBridgeQueryTemplateDeclaration{}, "query template is not declared by the bridge page"
}
var selected domain.GameClientBridgeQueryTemplateDeclaration
for _, template := range plugin.GameClientBridge.QueryTemplates {
if template.Key == templateKey {
selected = template
break
}
}
if selected.Key == "" {
return domain.GameClientBridgeQueryTemplateDeclaration{}, "query template is not declared by the plugin"
}
for _, page := range plugin.Pages {
if page.Key != routeKey {
continue
}
if !containsString(page.Permissions, selected.Permission) {
return domain.GameClientBridgeQueryTemplateDeclaration{}, "query template permission is not declared by the plugin page"
}
if !containsString(page.BridgeActions, string(domain.PluginBridgeActionRemoteAccessRequest)) {
return domain.GameClientBridgeQueryTemplateDeclaration{}, "query template page does not declare remote access"
}
return selected, ""
}
return domain.GameClientBridgeQueryTemplateDeclaration{}, "query template plugin page is not declared"
}
func (svc *CoreService) executeBridgeRunDistribution(sessionID string, base domain.PluginBridgeExecuteResponse, request domain.PluginBridgeExecuteRequest) domain.PluginBridgeExecuteResponse {
distribution, err := svc.GenerateRunDistributionForSession(sessionID, domain.RunDistributionGenerateRequest{
ServerInstanceID: request.ServerInstanceID,
TargetOS: defaultBridgeValue(request.Payload["targetOs"], "linux"),
TargetArch: defaultBridgeValue(request.Payload["targetArch"], "amd64"),
IdempotencyKey: defaultBridgeValue(request.Payload["idempotencyKey"], request.RequestID),
})
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = "ok"
base.Result = map[string]string{
"distributionId": distribution.ID,
"artifactId": distribution.ArtifactID,
"checksum": distribution.Checksum,
"keyGeneration": strconv.Itoa(distribution.KeyGeneration),
"status": string(distribution.Status),
}
return base
}
func (svc *CoreService) executeBridgeClientManager(sessionID string, base domain.PluginBridgeExecuteResponse, request domain.PluginBridgeExecuteRequest) domain.PluginBridgeExecuteResponse {
profileKey := request.Payload["profileKey"]
operation := defaultBridgeValue(request.Payload["operation"], "status")
idempotencyKey := defaultBridgeValue(request.Payload["idempotencyKey"], request.RequestID)
generation, _ := strconv.Atoi(request.Payload["expectedDeploymentGeneration"])
instance, err := svc.GetServerInstanceForSession(sessionID, request.ServerInstanceID)
if err != nil {
return bridgeExecutionError(base, err)
}
plugin, err := svc.store.GamePlugins().Get(instance.PluginID)
if err != nil {
return bridgeExecutionError(base, err)
}
profile, err := findRuntimeClientManagerProfile(plugin, profileKey)
if err != nil {
return bridgeExecutionError(base, ErrForbidden)
}
var view domain.ClientManagerLifecycleView
switch operation {
case "generate":
distribution, err := svc.GenerateClientManagerDistributionForSession(sessionID, domain.ClientManagerBuildRequest{ServerInstanceID: instance.ID, ProfileKey: profileKey, TargetOS: defaultBridgeValue(request.Payload["targetOs"], "windows"), TargetArch: defaultBridgeValue(request.Payload["targetArch"], "amd64"), RepositoryURL: profile.RepositoryURL, SourceRevision: clientManagerProfileRevision(profile), IdempotencyKey: idempotencyKey})
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = "queued"
base.Result = map[string]string{"distributionId": distribution.ID, "buildJobId": distribution.BuildJobID, "artifactId": distribution.ArtifactID, "checksum": distribution.Checksum, "keyGeneration": strconv.Itoa(distribution.KeyGeneration), "version": distribution.Version, "status": string(distribution.Status)}
return base
case "download":
reference, err := svc.OpenLatestClientManagerDistributionDownloadForSession(sessionID, instance.ID, profileKey)
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = "ok"
base.Result = map[string]string{"artifactId": reference.ArtifactID, "downloadUrl": reference.DownloadURL, "checksum": reference.Checksum, "sizeBytes": strconv.FormatInt(reference.SizeBytes, 10), "expiresAt": reference.ExpiresAt.Format(time.RFC3339), "rangeSupported": strconv.FormatBool(reference.RangeSupported), "chunkSizeBytes": strconv.Itoa(reference.ChunkSizeBytes)}
return base
case "reset-key":
key, err := svc.ResetComponentKeyForSession(sessionID, domain.ComponentKeyResetRequest{ServerInstanceID: instance.ID, ComponentKind: domain.DistributionComponentClientManager, ComponentKey: profileKey})
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = "ok"
base.Result = map[string]string{"profileKey": profileKey, "keyGeneration": strconv.Itoa(key.Generation), "status": string(key.Status), "requiresRedeploy": "true"}
return base
case "status":
view, err = svc.GetClientManagerLifecycleForSession(sessionID, instance.ID, profileKey)
case "deploy":
distributionID, resolveErr := svc.resolveClientManagerDistributionID(instance.ID, profileKey, request.Payload["artifactId"])
if resolveErr != nil {
return bridgeExecutionError(base, resolveErr)
}
view, err = svc.DeployClientManagerForSession(sessionID, domain.ClientManagerDeployRequest{ServerInstanceID: instance.ID, ProfileKey: profileKey, DistributionID: distributionID, ExpectedDeploymentGeneration: generation, IdempotencyKey: idempotencyKey})
case "start", "stop", "restart", "rollback":
view, err = svc.ControlClientManagerForSession(sessionID, domain.ClientManagerControlRequest{ServerInstanceID: instance.ID, ProfileKey: profileKey, Operation: domain.ClientManagerLifecycleOperation(operation), ExpectedDeploymentGeneration: generation, IdempotencyKey: idempotencyKey})
case "update":
distributionID, resolveErr := svc.resolveClientManagerDistributionID(instance.ID, profileKey, request.Payload["artifactId"])
if resolveErr != nil {
return bridgeExecutionError(base, resolveErr)
}
view, err = svc.UpdateClientManagerForSession(sessionID, domain.ClientManagerUpdateRequest{ServerInstanceID: instance.ID, ProfileKey: profileKey, DistributionID: distributionID, ExpectedDeploymentGeneration: generation, Approved: true, IdempotencyKey: idempotencyKey})
case "retry":
view, err = svc.RetryClientManagerLifecycleForSession(sessionID, domain.ClientManagerRetryRequest{ServerInstanceID: instance.ID, ProfileKey: profileKey, ExpectedDeploymentGeneration: generation, IdempotencyKey: idempotencyKey})
case "revoke-session":
view, err = svc.RevokeClientManagerSessionForSession(sessionID, domain.ClientManagerRevokeSessionRequest{ServerInstanceID: instance.ID, ProfileKey: profileKey, Reason: "plugin bridge operator request"})
case "uninstall":
view, err = svc.UninstallClientManagerForSession(sessionID, domain.ClientManagerUninstallRequest{ServerInstanceID: instance.ID, ProfileKey: profileKey, ExpectedDeploymentGeneration: generation, Confirmed: true, IdempotencyKey: idempotencyKey})
default:
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "invalid_client_manager_operation", Message: "client-manager operation is not supported"}
return base
}
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = "ok"
if view.Job.ID != "" && !isTerminalJobState(view.Job.State) {
base.Status = "queued"
}
base.Result = safeClientManagerBridgeResult(view)
return base
}
func (svc *CoreService) executeBridgeDependenciesRequest(sessionID string, base domain.PluginBridgeExecuteResponse, plugin domain.GamePlugin, instance domain.ServerInstance, payload map[string]string) domain.PluginBridgeExecuteResponse {
action := defaultBridgeValue(payload["operation"], "check")
capability := domain.JobCapabilityDependenciesCheck
if action == "install" {
capability = domain.JobCapabilityDependenciesInstall
} else if action != "check" {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "invalid_dependency_operation", Message: "dependency operation must be check or install"}
return base
}
if !containsString(plugin.RequiredRunCapabilities, capability) {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "capability_denied", Message: "dependency capability is not declared by plugin"}
return base
}
job, err := svc.QueueDependencyJobForSession(sessionID, domain.DependencyJobRequest{
ServerInstanceID: instance.ID,
ProbeKey: payload["probeKey"],
InstallPlanKey: payload["planKey"],
PlanDigest: payload["planDigest"],
TargetOS: payload["targetOS"],
TargetArch: payload["targetArch"],
IdempotencyKey: defaultBridgeValue(payload["idempotencyKey"], base.RequestID),
Install: action == "install",
})
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = "queued"
base.Result = map[string]string{"jobId": job.ID, "state": string(job.State), "capability": job.Capability, "targetKey": job.TargetKey}
return base
}
func (svc *CoreService) executeBridgeLogsBackfillRequest(base domain.PluginBridgeExecuteResponse, plugin domain.GamePlugin, instance domain.ServerInstance, payload map[string]string) domain.PluginBridgeExecuteResponse {
if !containsString(plugin.RequiredRunCapabilities, domain.JobCapabilityLogsBackfill) {
base.Status = "denied"
base.Error = &domain.PluginBridgeSafeError{Code: "capability_denied", Message: "log backfill capability is not declared by plugin"}
return base
}
job, err := svc.CreateJob(domain.Job{
ID: jobIDFromParts("job-logs-backfill", base.RequestID, payload["sourceKey"]),
ServerInstanceID: instance.ID,
RunEndpointID: instance.RunEndpointID,
Capability: domain.JobCapabilityLogsBackfill,
TargetKey: defaultBridgeValue(payload["sourceKey"], "logs/default"),
InputRef: payload["checkpointRef"],
IdempotencyKey: defaultBridgeValue(payload["idempotencyKey"], base.RequestID),
Progress: domain.JobProgress{Percent: 0, Message: "historical log backfill queued"},
})
if err != nil {
return bridgeExecutionError(base, err)
}
base.Status = "queued"
base.Result = map[string]string{"jobId": job.ID, "state": string(job.State), "capability": job.Capability, "sourceKey": job.TargetKey}
return base
}
func bridgeExecutionError(base domain.PluginBridgeExecuteResponse, err error) domain.PluginBridgeExecuteResponse {
base.Status = "error"
base.Error = &domain.PluginBridgeSafeError{Code: "execution_failed", Message: safeBridgeReason(err.Error())}
return base
}
func defaultBridgeValue(value string, fallback string) string {
if strings.TrimSpace(value) == "" {
return fallback
}
return value
}
func safeBridgeReason(reason string) string {
reason = strings.TrimSpace(reason)
if reason == "" {
return "bridge action is not allowed"
}
for _, forbidden := range []string{"/Users/", "/private/", "unix://", "tcp://", "Bearer ", "sk-", "password=", "api_key=", "apikey="} {
if strings.Contains(strings.ToLower(reason), strings.ToLower(forbidden)) {
return "bridge action failed safely"
}
}
return reason
}
func pluginPermissionsFromManifest(permissions []string) domain.PluginPermissions {
var aggregate domain.PluginPermissions
for _, permission := range permissions {
switch permission {
case "ai.invoke":
aggregate.AI = true
case "server.logs.read":
aggregate.Logs = true
case "server.files.read", "server.files.write":
aggregate.Files = true
case "server.lifecycle", "server.create":
aggregate.Jobs = true
case "server.artifacts.read", "server.artifacts.write":
aggregate.Artifacts = true
case "server.remote.access":
aggregate.RemoteAccess = true
case "server.run.distribution", "server.dependencies.manage", "server.client-manager.manage":
aggregate.Jobs = true
aggregate.Artifacts = true
}
}
return aggregate
}
func (svc *CoreService) GetGamePlugin(id string) (domain.GamePlugin, error) {
return svc.store.GamePlugins().Get(id)
}
func (svc *CoreService) ListGamePlugins(filter domain.GamePluginFilter) ([]domain.GamePlugin, error) {
return svc.store.GamePlugins().List(filter)
}
func (svc *CoreService) ListMarketplacePlugins(filter domain.PluginMarketplaceFilter) ([]domain.PluginMarketplacePlugin, error) {
filter.Keyword = strings.TrimSpace(filter.Keyword)
if err := validator.ValidatePluginMarketplaceFilter(filter); err != nil {
return nil, err
}
plugins, err := svc.store.GamePlugins().List(domain.GamePluginFilter{
ServerType: filter.ServerType,
Status: filter.Status,
})
if err != nil {
return nil, err
}
items := make([]domain.PluginMarketplacePlugin, 0, len(plugins))
for _, plugin := range plugins {
projected := marketplacePluginFromGamePlugin(plugin)
if filter.Capability != "" && !containsString(projected.Capabilities, filter.Capability) && !containsString(projected.BridgeActions, filter.Capability) {
continue
}
if filter.Keyword != "" && !marketplacePluginMatchesKeyword(projected, filter.Keyword) {
continue
}
items = append(items, projected)
}
if err := validator.ValidatePluginMarketplacePlugins(items); err != nil {
return nil, err
}
return domain.CopyPluginMarketplacePluginSlice(items), nil
}
func (svc *CoreService) GetMarketplacePlugin(id string) (domain.PluginMarketplacePlugin, error) {
if strings.TrimSpace(id) == "" {
return domain.PluginMarketplacePlugin{}, validationError("pluginId is required")
}
plugin, err := svc.store.GamePlugins().Get(id)
if err != nil {
return domain.PluginMarketplacePlugin{}, err
}
projected := marketplacePluginFromGamePlugin(plugin)
if err := validator.ValidatePluginMarketplacePlugin(projected); err != nil {
return domain.PluginMarketplacePlugin{}, err
}
return domain.CopyPluginMarketplacePlugin(projected), nil
}
func (svc *CoreService) SetMarketplacePluginState(id string, action domain.PluginMarketplaceStateAction) (domain.PluginMarketplacePlugin, error) {
if strings.TrimSpace(id) == "" {
return domain.PluginMarketplacePlugin{}, validationError("pluginId is required")
}
if err := validator.ValidatePluginMarketplaceStateAction(action); err != nil {
return domain.PluginMarketplacePlugin{}, err
}
plugin, err := svc.store.GamePlugins().Get(id)
if err != nil {
return domain.PluginMarketplacePlugin{}, err
}
switch action {
case domain.PluginMarketplaceStateActionInstall, domain.PluginMarketplaceStateActionEnable:
plugin.Status = domain.GamePluginStatusInstalled
case domain.PluginMarketplaceStateActionDisable:
plugin.Status = domain.GamePluginStatusDisabled
}
if err := validator.ValidateGamePlugin(plugin); err != nil {
return domain.PluginMarketplacePlugin{}, err
}
if err := svc.store.GamePlugins().Update(plugin); err != nil {
return domain.PluginMarketplacePlugin{}, err
}
projected := marketplacePluginFromGamePlugin(plugin)
if err := validator.ValidatePluginMarketplacePlugin(projected); err != nil {
return domain.PluginMarketplacePlugin{}, err
}
return domain.CopyPluginMarketplacePlugin(projected), nil
}
func marketplacePluginFromGamePlugin(plugin domain.GamePlugin) domain.PluginMarketplacePlugin {
plugin = domain.CopyGamePlugin(plugin)
return domain.PluginMarketplacePlugin{
ID: plugin.ID,
Name: plugin.Name,
Description: plugin.Description,
Version: plugin.Version,
ServerType: plugin.ServerType,
ServerDisplayName: plugin.ServerDisplayName,
SupportedOS: plugin.SupportedOS,
ManifestRef: plugin.ManifestRef,
CreateFormSchemaRef: plugin.CreateFormSchemaRef,
CreateFields: plugin.CreateFields,
Capabilities: plugin.RequiredRunCapabilities,
DeclaredPermissions: plugin.DeclaredPermissions,
Permissions: plugin.Permissions,
LifecycleActions: plugin.LifecycleActions,
BridgeActions: plugin.BridgeActions,
Pages: plugin.Pages,
Tags: plugin.Tags,
AIPurposes: plugin.AIPurposes,
ProductionLifecycle: plugin.ProductionLifecycle,
RemoteAccess: plugin.RemoteAccess,
RuntimeProfiles: plugin.RuntimeProfiles,
GameClientBridge: plugin.GameClientBridge,
ValidationViolations: plugin.ValidationViolations,
Status: plugin.Status,
Source: "platform-registry",
}
}
func normalizedProductionLifecycle(lifecycle domain.GamePluginProductionLifecycle) domain.GamePluginProductionLifecycle {
if len(lifecycle.Operations) == 0 {
lifecycle.Operations = []string{"install", "enable", "disable", "upgrade", "rollback", "retire", "dependency-check"}
}
if lifecycle.DependencyPolicy == "" {
lifecycle.DependencyPolicy = "optional"
}
if len(lifecycle.ApprovalRequired) == 0 {
lifecycle.ApprovalRequired = []string{"disable", "rollback", "retire"}
}
return lifecycle
}
func marketplacePluginMatchesKeyword(plugin domain.PluginMarketplacePlugin, keyword string) bool {
keyword = strings.ToLower(strings.TrimSpace(keyword))
if keyword == "" {
return true
}
fields := []string{plugin.ID, plugin.Name, plugin.Description, plugin.ServerType, plugin.ServerDisplayName, plugin.Version}
fields = append(fields, plugin.Tags...)
fields = append(fields, plugin.Capabilities...)
for _, field := range fields {
if strings.Contains(strings.ToLower(field), keyword) {
return true
}
}
return false
}
func (svc *CoreService) CreateRunEndpoint(endpoint domain.RunEndpoint) (domain.RunEndpoint, error) {
if endpoint.Status == "" {
endpoint.Status = domain.RunEndpointStatusOnline
}
if endpoint.LastHeartbeatAt.IsZero() && (endpoint.Status == domain.RunEndpointStatusOnline || endpoint.Status == domain.RunEndpointStatusDegraded) {
endpoint.LastHeartbeatAt = svc.now()
}
if err := validator.ValidateRunEndpoint(endpoint); err != nil {
return domain.RunEndpoint{}, err
}
if err := svc.store.RunEndpoints().Create(endpoint); err != nil {
return domain.RunEndpoint{}, err
}
return domain.CopyRunEndpoint(endpoint), nil
}
func (svc *CoreService) GetRunEndpoint(id string) (domain.RunEndpoint, error) {
return svc.store.RunEndpoints().Get(id)
}
func (svc *CoreService) ListRunEndpoints(filter domain.RunEndpointFilter) ([]domain.RunEndpoint, error) {
return svc.store.RunEndpoints().List(filter)
}
func (svc *CoreService) CreateServerInstance(instance domain.ServerInstance) (domain.ServerInstance, error) {
plugin, err := svc.store.GamePlugins().Get(instance.PluginID)
if err != nil {
return domain.ServerInstance{}, fmt.Errorf("get plugin dependency: %w", err)
}
var endpoint domain.RunEndpoint
if strings.TrimSpace(instance.RunEndpointID) != "" {
endpoint, err = svc.store.RunEndpoints().Get(instance.RunEndpointID)
if err != nil {
return domain.ServerInstance{}, fmt.Errorf("get run endpoint dependency: %w", err)
}
}
if instance.PluginVersion == "" {
instance.PluginVersion = plugin.Version
}
if instance.State == "" {
instance.State = domain.ServerInstanceStateDraft
}
if instance.ConfigVersion == 0 {
instance.ConfigVersion = 1
}
stamp := svc.now()
if instance.CreatedAt.IsZero() {
instance.CreatedAt = stamp
}
if instance.UpdatedAt.IsZero() {
instance.UpdatedAt = stamp
}
if err := validator.ValidateServerInstance(instance); err != nil {
return domain.ServerInstance{}, err
}
if instance.State != domain.ServerInstanceStateDraft && strings.TrimSpace(instance.RunEndpointID) == "" {
return domain.ServerInstance{}, validationError("runEndpointId is required when server is not a draft")
}
if strings.TrimSpace(instance.RunEndpointID) != "" {
if err := validator.ValidateServerInstanceDependencies(instance, plugin, endpoint); err != nil {
return domain.ServerInstance{}, err
}
} else if plugin.Status != domain.GamePluginStatusInstalled || plugin.Version != instance.PluginVersion {
return domain.ServerInstance{}, validationError("plugin must be installed and match the server plugin version")
}
if err := svc.store.ServerInstances().Create(instance); err != nil {
return domain.ServerInstance{}, err
}
return domain.CopyServerInstance(instance), nil
}
func (svc *CoreService) CreateServerInstanceForSession(sessionID string, instance domain.ServerInstance) (domain.ServerInstance, error) {
user, err := svc.GetCurrentUser(sessionID)
if err != nil {
return domain.ServerInstance{}, err
}
if strings.TrimSpace(instance.OwnerUserID) == "" {
instance.OwnerUserID = user.ID
}
if !isPlatformAdmin(user) && instance.OwnerUserID != user.ID {
return domain.ServerInstance{}, ErrForbidden
}
return svc.CreateServerInstance(instance)
}
func (svc *CoreService) GetServerInstance(id string) (domain.ServerInstance, error) {
return svc.store.ServerInstances().Get(id)
}
func (svc *CoreService) GetServerInstanceForSession(sessionID string, id string) (domain.ServerInstance, error) {
user, err := svc.GetCurrentUser(sessionID)
if err != nil {
return domain.ServerInstance{}, err
}
instance, err := svc.store.ServerInstances().Get(id)
if err != nil {
return domain.ServerInstance{}, err
}
if !canAccessServer(user, instance) {
return domain.ServerInstance{}, ErrForbidden
}
return domain.CopyServerInstance(instance), nil
}
func (svc *CoreService) UpdateServerInstanceForSession(sessionID string, id string, update domain.ServerInstanceUpdate) (domain.ServerInstance, error) {
if err := validator.ValidateServerInstanceUpdate(update); err != nil {
return domain.ServerInstance{}, err
}
user, err := svc.GetCurrentUser(sessionID)
if err != nil {
return domain.ServerInstance{}, err
}
instance, err := svc.store.ServerInstances().Get(id)
if err != nil {
return domain.ServerInstance{}, err
}
if !isPlatformAdmin(user) && instance.OwnerUserID != user.ID {
return domain.ServerInstance{}, ErrForbidden
}
if instance.State == domain.ServerInstanceStateDeleted {
return domain.ServerInstance{}, validationError("deleted server instances cannot be edited")
}
if update.Name != nil {
instance.Name = *update.Name
}
instance.UpdatedAt = svc.now()
if err := validator.ValidateStoredServerInstance(instance); err != nil {
return domain.ServerInstance{}, err
}
if err := svc.store.ServerInstances().Update(instance); err != nil {
return domain.ServerInstance{}, err
}
return domain.CopyServerInstance(instance), nil
}
func (svc *CoreService) ListServerInstances(filter domain.ServerInstanceFilter) ([]domain.ServerInstance, error) {
return svc.store.ServerInstances().List(filter)
}
func (svc *CoreService) ListServerInstancesForSession(sessionID string, filter domain.ServerInstanceFilter) ([]domain.ServerInstance, error) {
user, err := svc.GetCurrentUser(sessionID)
if err != nil {
return nil, err
}
if !isPlatformAdmin(user) {
filter.VisibleToUserID = user.ID
}
return svc.store.ServerInstances().List(filter)
}
func (svc *CoreService) ListServerAdministratorCandidates(sessionID string, serverInstanceID string) ([]domain.User, error) {
user, instance, err := svc.requireServerOwner(sessionID, serverInstanceID)
if err != nil {
return nil, err
}
_ = user
users, err := svc.store.Users().List(domain.UserFilter{Status: domain.UserStatusActive})
if err != nil {
return nil, err
}
candidates := make([]domain.User, 0, len(users))
for _, candidate := range users {
if candidate.ID == instance.OwnerUserID || containsString(instance.AdminUserIDs, candidate.ID) || isPlatformAdmin(candidate) {
continue
}
candidates = append(candidates, domain.CopyUser(candidate))
}
return candidates, nil
}
func (svc *CoreService) GetPlatformResourceUsage() (domain.PlatformResourceUsage, error) {
instances, err := svc.store.ServerInstances().List(domain.ServerInstanceFilter{})
if err != nil {
return domain.PlatformResourceUsage{}, err
}
endpoints, err := svc.store.RunEndpoints().List(domain.RunEndpointFilter{})
if err != nil {
return domain.PlatformResourceUsage{}, err
}
jobs, err := svc.store.Jobs().List(domain.JobFilter{})
if err != nil {
return domain.PlatformResourceUsage{}, err
}
runningServers := 0
for _, instance := range instances {
if instance.State == domain.ServerInstanceStateRunning {
runningServers++
}
}
onlineEndpoints := 0
for _, endpoint := range endpoints {
if endpoint.Status == domain.RunEndpointStatusOnline {
onlineEndpoints++
}
}
activeJobs := 0
for _, job := range jobs {
if job.State == domain.JobStateQueued || job.State == domain.JobStateAccepted || job.State == domain.JobStateRunning {
activeJobs++
}
}
usage := domain.PlatformResourceUsage{
CPUPercent: clampPercent(float64(runningServers*18 + activeJobs*6 + onlineEndpoints*4)),
MemoryPercent: clampPercent(float64(runningServers*22 + onlineEndpoints*8 + len(instances)*3)),
DiskPercent: clampPercent(float64(len(instances)*9 + len(jobs)*2)),
Source: "platform-derived",
CollectedAt: svc.now(),
}
if err := validator.ValidatePlatformResourceUsage(usage); err != nil {
return domain.PlatformResourceUsage{}, err
}
return domain.CopyPlatformResourceUsage(usage), nil
}
func (svc *CoreService) ListServerMetricsForSession(sessionID string) ([]domain.ServerMetrics, error) {
instances, err := svc.ListServerInstancesForSession(sessionID, domain.ServerInstanceFilter{})
if err != nil {
return nil, err
}
items := make([]domain.ServerMetrics, 0, len(instances))
for _, instance := range instances {
items = append(items, svc.metricsForServer(instance))
}
if err := validator.ValidateServerMetricsList(items); err != nil {
return nil, err
}
return domain.CopyServerMetricsSlice(items), nil
}
func (svc *CoreService) GetServerConfigForSession(sessionID string, serverInstanceID string) (domain.ServerConfig, error) {
instance, err := svc.GetServerInstanceForSession(sessionID, serverInstanceID)
if err != nil {
return domain.ServerConfig{}, err
}
config := domain.ServerConfig{
ServerInstanceID: instance.ID,
ConfigVersion: instance.ConfigVersion,
Format: "properties",
Key: "server.properties",
Source: "platform-derived",
UpdatedAt: instance.UpdatedAt,
Checksum: instance.ConfigChecksum,
}
if config.UpdatedAt.IsZero() {
config.UpdatedAt = svc.now()
}
config.Key = instance.ConfigKey
if config.Key == "" {
config.Key = "server.properties"
}
config.Content = instance.ConfigContent
if config.Content == "" {
config.Content = buildLogicalServerConfig(instance)
}
if config.Checksum == "" {
config.Checksum = validator.BytesChecksum([]byte(config.Content))
}
if err := validator.ValidateServerConfig(config); err != nil {
return domain.ServerConfig{}, err
}
return domain.CopyServerConfig(config), nil
}
func (svc *CoreService) PreviewServerConfigWriteForSession(sessionID string, request domain.ServerConfigDiffRequest) (domain.ServerConfigDiffPreview, error) {
if request.Key == "" {
request.Key = "server.properties"
}
if err := validator.ValidateServerConfigDiffRequest(request); err != nil {
return domain.ServerConfigDiffPreview{}, err
}
config, err := svc.GetServerConfigForSession(sessionID, request.ServerInstanceID)
if err != nil {
return domain.ServerConfigDiffPreview{}, err
}
if config.ConfigVersion != request.ExpectedConfigVersion {
return domain.ServerConfigDiffPreview{}, validationError("expectedConfigVersion must match server instance")
}
if request.ExpectedChecksum != "" && config.Checksum != request.ExpectedChecksum {
return domain.ServerConfigDiffPreview{}, validationError("expectedChecksum must match server config")
}
if config.Key != request.Key {
return domain.ServerConfigDiffPreview{}, validationError("key must match server config")
}
preview := domain.ServerConfigDiffPreview{
ServerInstanceID: request.ServerInstanceID,
ConfigVersion: config.ConfigVersion,
Checksum: config.Checksum,
Key: request.Key,
CurrentContent: config.Content,
ProposedContent: request.ProposedContent,
ProposedContentInputRef: request.ProposedContentInputRef,
Diff: buildConfigDiffLines(config.Content, request.ProposedContent),
Source: "platform-review",
ReviewedAt: svc.now(),
}
preview.HasChanges = config.Content != request.ProposedContent
return domain.CopyServerConfigDiffPreview(preview), nil
}
func (svc *CoreService) ApproveServerConfigWriteForSession(sessionID string, approval domain.ServerConfigWriteApproval) (domain.ServerConfigWriteDispatch, error) {
if approval.Key == "" {
approval.Key = "server.properties"
}
if approval.ProposedContentInputRef == "" {
approval.ProposedContentInputRef = configWriteInputRef(approval.ServerInstanceID, approval.Key, approval.ExpectedConfigVersion)
}
if err := validator.ValidateServerConfigWriteApproval(approval); err != nil {
return domain.ServerConfigWriteDispatch{}, err
}
preview, err := svc.PreviewServerConfigWriteForSession(sessionID, domain.ServerConfigDiffRequest{
ServerInstanceID: approval.ServerInstanceID,
ExpectedConfigVersion: approval.ExpectedConfigVersion,
ExpectedChecksum: approval.ExpectedChecksum,
Key: approval.Key,
ProposedContent: approval.ProposedContent,
ProposedContentInputRef: approval.ProposedContentInputRef,
})
if err != nil {
return domain.ServerConfigWriteDispatch{}, err
}
if !preview.HasChanges {
return domain.ServerConfigWriteDispatch{}, validationError("config diff has no changes")
}
instance, err := svc.GetServerInstanceForSession(sessionID, approval.ServerInstanceID)
if err != nil {
return domain.ServerConfigWriteDispatch{}, err
}
user, err := svc.GetCurrentUser(sessionID)
if err != nil {
return domain.ServerConfigWriteDispatch{}, err
}
if err := svc.requireCompleteRuntimeBindings(user.ID, instance.ID, "config.write.denied"); err != nil {
return domain.ServerConfigWriteDispatch{}, err
}
job, err := svc.CreateJob(domain.Job{
ID: jobIDFromParts("job-config-write", approval.ServerInstanceID, approval.IdempotencyKey),
ServerInstanceID: instance.ID,
RunEndpointID: instance.RunEndpointID,
Capability: domain.JobCapabilityConfigWrite,
TargetKey: approval.Key,
InputRef: approval.ProposedContentInputRef,
ExecutionInput: domain.JobExecutionInput{WorkspaceScope: svc.runtimeProfileScope(instance.ID), Content: approval.ProposedContent, ExpectedVersion: approval.ExpectedConfigVersion, ExpectedChecksum: preview.Checksum, MaxReadBytes: 64 * 1024},
IdempotencyKey: approval.IdempotencyKey,
Progress: domain.JobProgress{Percent: 0, Message: "config write queued"},
})
if err != nil {
return domain.ServerConfigWriteDispatch{}, err
}
return domain.CopyServerConfigWriteDispatch(domain.ServerConfigWriteDispatch{Preview: preview, Job: job, Status: "queued"}), nil
}
func (svc *CoreService) DispatchFileOperationForSession(sessionID string, request domain.FileOperationDispatchRequest) (domain.FileOperationDispatchResult, error) {
if err := validator.ValidateFileOperationDispatchRequest(request); err != nil {
return domain.FileOperationDispatchResult{}, err
}
instance, err := svc.GetServerInstanceForSession(sessionID, request.ServerInstanceID)
if err != nil {
return domain.FileOperationDispatchResult{}, err
}
if request.ExpectedConfigVersion > 0 && request.ExpectedConfigVersion != instance.ConfigVersion {
return domain.FileOperationDispatchResult{}, validationError("expectedConfigVersion must match server instance")
}
user, err := svc.GetCurrentUser(sessionID)
if err != nil {
return domain.FileOperationDispatchResult{}, err
}
if err := svc.requireCompleteRuntimeBindings(user.ID, instance.ID, "file.operation.denied"); err != nil {
return domain.FileOperationDispatchResult{}, err
}
content := request.Content
if request.Operation == domain.FileOperationWrite && content == "" && strings.HasPrefix(request.InputRef, "artifact://") {
artifactID := strings.TrimPrefix(request.InputRef, "artifact://")
artifact, artifactErr := svc.store.Artifacts().Get(artifactID)
if artifactErr != nil {
return domain.FileOperationDispatchResult{}, artifactErr
}
if artifact.OwnerKind != domain.ArtifactOwnerKindServerInstance || artifact.OwnerID != instance.ID || artifact.State != domain.ArtifactStateAvailable {
return domain.FileOperationDispatchResult{}, ErrForbidden
}
payload, payloadErr := svc.artifactPayload(artifactID)
if payloadErr != nil {
return domain.FileOperationDispatchResult{}, payloadErr
}
content = string(payload)
}
if request.PluginID != "" {
plugin, err := svc.store.GamePlugins().Get(request.PluginID)
if err != nil {
return domain.FileOperationDispatchResult{}, err
}
if plugin.ID != instance.PluginID {
return domain.FileOperationDispatchResult{}, validationError("pluginId must match server instance")
}
if plugin.Status != domain.GamePluginStatusInstalled {
return domain.FileOperationDispatchResult{}, validationError("plugin must be installed")
}
if request.Operation == domain.FileOperationRead && !plugin.Permissions.Files && !containsString(plugin.DeclaredPermissions, "server.files.read") {
return domain.FileOperationDispatchResult{}, ErrForbidden
}
if request.Operation == domain.FileOperationWrite && !containsString(plugin.DeclaredPermissions, "server.files.write") {
return domain.FileOperationDispatchResult{}, ErrForbidden
}
}
capability := domain.JobCapabilityFilesRead
message := "file read queued"
if request.Operation == domain.FileOperationWrite {
capability = domain.JobCapabilityFilesWrite
message = "file write queued"
}
job, err := svc.CreateJob(domain.Job{
ID: jobIDFromParts("job-file", request.ServerInstanceID, request.IdempotencyKey),
ServerInstanceID: instance.ID,
RunEndpointID: instance.RunEndpointID,
Capability: capability,
TargetKey: request.Key,
InputRef: request.InputRef,
ExecutionInput: domain.JobExecutionInput{WorkspaceScope: svc.runtimeProfileScope(instance.ID), Content: content, ExpectedVersion: request.ExpectedConfigVersion, ExpectedChecksum: request.ExpectedChecksum, MaxReadBytes: 64 * 1024},
IdempotencyKey: request.IdempotencyKey,
Progress: domain.JobProgress{Percent: 0, Message: message},
})
if err != nil {
return domain.FileOperationDispatchResult{}, err
}
return domain.CopyFileOperationDispatchResult(domain.FileOperationDispatchResult{
ServerInstanceID: request.ServerInstanceID,
PluginID: request.PluginID,
Operation: request.Operation,
Key: request.Key,
InputRef: request.InputRef,
Job: job,
Status: "queued",
}), nil
}
func (svc *CoreService) runtimeProfileScope(serverInstanceID string) string {
binding, err := svc.runtimeBindingForServer(serverInstanceID)
if err != nil {
return "default"
}
return binding.ProfileKey
}
func (svc *CoreService) metricsForServer(instance domain.ServerInstance) domain.ServerMetrics {
metrics := domain.ServerMetrics{
ServerInstanceID: instance.ID,
Online: instance.State == domain.ServerInstanceStateRunning,
Source: "platform-derived",
CollectedAt: svc.now(),
}
if !metrics.Online {
return metrics
}
seed := len(instance.ID) + len(instance.Name) + instance.ConfigVersion
playerCount := seed % 20
maxPlayers := 20
tps := 18.5 + float64(seed%15)/10
latency := 35.0 + float64(seed%40)
cpu := clampPercent(25 + float64(seed%45))
memory := clampPercent(30 + float64(seed%50))
disk := clampPercent(20 + float64(seed%60))
metrics.PlayerCount = &playerCount
metrics.MaxPlayers = &maxPlayers
metrics.TPS = &tps
metrics.LatencyMS = &latency
metrics.CPUPercent = &cpu
metrics.MemoryPercent = &memory
metrics.DiskPercent = &disk
return metrics
}
func buildLogicalServerConfig(instance domain.ServerInstance) string {
lines := []string{
"# platform logical server config",
"server.id=" + instance.ID,
"server.name=" + instance.Name,
"plugin.id=" + instance.PluginID,
"plugin.version=" + instance.PluginVersion,
fmt.Sprintf("config.version=%d", instance.ConfigVersion),
"state=" + string(instance.State),
}
return strings.Join(lines, "\n") + "\n"
}
func buildConfigDiffLines(current string, proposed string) []domain.ConfigDiffLine {
currentLines := strings.Split(current, "\n")
proposedLines := strings.Split(proposed, "\n")
maxLen := len(currentLines)
if len(proposedLines) > maxLen {
maxLen = len(proposedLines)
}
lines := make([]domain.ConfigDiffLine, 0, maxLen*2)
for i := 0; i < maxLen; i++ {
oldExists := i < len(currentLines)
newExists := i < len(proposedLines)
oldLine := ""
newLine := ""
if oldExists {
oldLine = currentLines[i]
}
if newExists {
newLine = proposedLines[i]
}
if oldExists && newExists && oldLine == newLine {
lines = append(lines, domain.ConfigDiffLine{Kind: "context", OldNumber: i + 1, NewNumber: i + 1, Content: oldLine})
continue
}
if oldExists {
lines = append(lines, domain.ConfigDiffLine{Kind: "removed", OldNumber: i + 1, Content: oldLine})
}
if newExists {
lines = append(lines, domain.ConfigDiffLine{Kind: "added", NewNumber: i + 1, Content: newLine})
}
}
return lines
}
func configWriteInputRef(serverInstanceID string, key string, version int) string {
return fmt.Sprintf("input://server-config/%s/%s/v%d", serverInstanceID, strings.ReplaceAll(key, "/", "-"), version)
}
func jobIDFromParts(prefix string, resourceID string, idempotencyKey string) string {
return fmt.Sprintf("%s-%s-%d", prefix, resourceID, stableStringNumber(idempotencyKey))
}
func stableStringNumber(value string) int {
sum := 0
for _, char := range value {
sum = sum*31 + int(char)
if sum < 0 {
sum = -sum
}
}
return sum
}
func clampPercent(value float64) float64 {
if value < 0 {
return 0
}
if value > 100 {
return 100
}
return value
}
func (svc *CoreService) AddServerAdministrator(sessionID string, serverInstanceID string, userID string) (domain.ServerInstance, error) {
_, instance, err := svc.requireServerOwner(sessionID, serverInstanceID)
if err != nil {
return domain.ServerInstance{}, err
}
member, err := svc.store.Users().Get(userID)
if err != nil {
return domain.ServerInstance{}, err
}
if member.Status != domain.UserStatusActive || isPlatformAdmin(member) || member.ID == instance.OwnerUserID {
return domain.ServerInstance{}, ErrForbidden
}
if !containsString(instance.AdminUserIDs, member.ID) {
instance.AdminUserIDs = append(instance.AdminUserIDs, member.ID)
instance.UpdatedAt = svc.now()
if err := validator.ValidateServerInstance(instance); err != nil {
return domain.ServerInstance{}, err
}
if err := svc.store.ServerInstances().Update(instance); err != nil {
return domain.ServerInstance{}, err
}
}
return domain.CopyServerInstance(instance), nil
}
func (svc *CoreService) RemoveServerAdministrator(sessionID string, serverInstanceID string, userID string) (domain.ServerInstance, error) {
_, instance, err := svc.requireServerOwner(sessionID, serverInstanceID)
if err != nil {
return domain.ServerInstance{}, err
}
member, err := svc.store.Users().Get(userID)
if err != nil {
return domain.ServerInstance{}, err
}
if isPlatformAdmin(member) {
return domain.ServerInstance{}, ErrForbidden
}
nextAdmins := make([]string, 0, len(instance.AdminUserIDs))
for _, adminID := range instance.AdminUserIDs {
if adminID != userID {
nextAdmins = append(nextAdmins, adminID)
}
}
instance.AdminUserIDs = nextAdmins
instance.UpdatedAt = svc.now()
if err := validator.ValidateServerInstance(instance); err != nil {
return domain.ServerInstance{}, err
}
if err := svc.store.ServerInstances().Update(instance); err != nil {
return domain.ServerInstance{}, err
}
return domain.CopyServerInstance(instance), nil
}
func (svc *CoreService) DeleteServerInstanceForSession(sessionID string, serverInstanceID string, password string) (domain.ServerInstance, error) {
user, err := svc.GetCurrentUser(sessionID)
if err != nil {
return domain.ServerInstance{}, err
}
instance, err := svc.GetServerInstanceForSession(sessionID, serverInstanceID)
if err != nil {
return domain.ServerInstance{}, err
}
if !isPlatformAdmin(user) && instance.OwnerUserID != user.ID {
return domain.ServerInstance{}, ErrForbidden
}
if strings.TrimSpace(password) == "" {
return domain.ServerInstance{}, validationError("password is required")
}
if !verifyPassword(user.PasswordHash, password) {
return domain.ServerInstance{}, forbiddenError("password confirmation failed")
}
if instance.State == domain.ServerInstanceStateRunning || instance.State == domain.ServerInstanceStateInstalling {
return domain.ServerInstance{}, validationError("running or installing server instances must be stopped before delete")
}
if instance.State == domain.ServerInstanceStateDeleted {
return domain.CopyServerInstance(instance), nil
}
instance.State = domain.ServerInstanceStateDeleted
instance.UpdatedAt = svc.now()
if err := validator.ValidateStoredServerInstance(instance); err != nil {
return domain.ServerInstance{}, err
}
if err := svc.store.ServerInstances().Update(instance); err != nil {
return domain.ServerInstance{}, err
}
return domain.CopyServerInstance(instance), nil
}
func (svc *CoreService) CreateJob(job domain.Job) (domain.Job, error) {
if job.State == "" {
job.State = domain.JobStateQueued
}
stamp := svc.now()
if job.CreatedAt.IsZero() {
job.CreatedAt = stamp
}
if job.UpdatedAt.IsZero() {
job.UpdatedAt = stamp
}
job = normalizeJobScheduling(job, stamp)
if err := validator.ValidateJob(job); err != nil {
return domain.Job{}, err
}
existing, err := svc.store.Jobs().GetByIdempotency(job.RunEndpointID, job.IdempotencyKey)
if err == nil {
return existing, nil
}
if !errors.Is(err, repo.ErrNotFound) {
return domain.Job{}, err
}
endpoint, err := svc.store.RunEndpoints().Get(job.RunEndpointID)
if err != nil {
return domain.Job{}, fmt.Errorf("get run endpoint dependency: %w", err)
}
if err := svc.validateRunnableEndpoint(endpoint, job.Capability); err != nil {
return domain.Job{}, err
}
if job.ServerInstanceID != "" {
instance, err := svc.store.ServerInstances().Get(job.ServerInstanceID)
if err != nil {
return domain.Job{}, fmt.Errorf("get server instance dependency: %w", err)
}
plugin, err := svc.store.GamePlugins().Get(instance.PluginID)
if err != nil {
return domain.Job{}, fmt.Errorf("get server plugin dependency: %w", err)
}
if err := validateJobServerTarget(job, instance, plugin); err != nil {
return domain.Job{}, err
}
}
if err := svc.store.Jobs().Create(job); err != nil {
return domain.Job{}, err
}
return domain.CopyJob(job), nil
}
func (svc *CoreService) GetJob(id string) (domain.Job, error) {
job, err := svc.store.Jobs().Get(id)
if err != nil {
return domain.Job{}, err
}
return normalizeJobScheduling(job, svc.now()), nil
}
func (svc *CoreService) ListJobs(filter domain.JobFilter) ([]domain.Job, error) {
jobs, err := svc.store.Jobs().List(filter)
if err != nil {
return nil, err
}
for i := range jobs {
jobs[i] = normalizeJobScheduling(jobs[i], svc.now())
}
return jobs, nil
}
func (svc *CoreService) CreateArtifact(artifact domain.Artifact) (domain.Artifact, error) {
if artifact.State == "" {
artifact.State = domain.ArtifactStateUploading
}
stamp := svc.now()
if artifact.CreatedAt.IsZero() {
artifact.CreatedAt = stamp
}
if artifact.UpdatedAt.IsZero() {
artifact.UpdatedAt = stamp
}
if err := validator.ValidateArtifact(artifact); err != nil {
return domain.Artifact{}, err
}
if err := svc.store.Artifacts().Create(artifact); err != nil {
return domain.Artifact{}, err
}
return domain.CopyArtifact(artifact), nil
}
func (svc *CoreService) GetArtifact(id string) (domain.Artifact, error) {
return svc.store.Artifacts().Get(id)
}
func (svc *CoreService) ListArtifacts(filter domain.ArtifactFilter) ([]domain.Artifact, error) {
return svc.store.Artifacts().List(filter)
}
func (svc *CoreService) CreateLogStream(stream domain.LogStream) (domain.LogStream, error) {
instance, err := svc.store.ServerInstances().Get(stream.ServerInstanceID)
if err != nil {
return domain.LogStream{}, fmt.Errorf("get server instance dependency: %w", err)
}
if instance.State == domain.ServerInstanceStateDeleted {
return domain.LogStream{}, validationError("server instance must not be deleted")
}
stamp := svc.now()
if stream.CreatedAt.IsZero() {
stream.CreatedAt = stamp
}
if stream.UpdatedAt.IsZero() {
stream.UpdatedAt = stamp
}
if err := validator.ValidateLogStream(stream); err != nil {
return domain.LogStream{}, err
}
if err := svc.store.LogStreams().Create(stream); err != nil {
return domain.LogStream{}, err
}
return domain.CopyLogStream(stream), nil
}
func (svc *CoreService) GetLogStream(id string) (domain.LogStream, error) {
return svc.store.LogStreams().Get(id)
}
func (svc *CoreService) ListLogStreams(filter domain.LogStreamFilter) ([]domain.LogStream, error) {
return svc.store.LogStreams().List(filter)
}
func (svc *CoreService) CreateAuditEvent(event domain.AuditEvent) (domain.AuditEvent, error) {
if event.CreatedAt.IsZero() {
event.CreatedAt = svc.now()
}
if err := validator.ValidateAuditEvent(event); err != nil {
return domain.AuditEvent{}, err
}
if err := svc.store.AuditEvents().Create(event); err != nil {
return domain.AuditEvent{}, err
}
return domain.CopyAuditEvent(event), nil
}
func (svc *CoreService) GetAuditEvent(id string) (domain.AuditEvent, error) {
return svc.store.AuditEvents().Get(id)
}
func (svc *CoreService) ListAuditEvents(filter domain.AuditEventFilter) ([]domain.AuditEvent, error) {
return svc.store.AuditEvents().List(filter)
}
func (svc *CoreService) validateRunnableEndpoint(endpoint domain.RunEndpoint, capability string) error {
if endpoint.Status != domain.RunEndpointStatusOnline && endpoint.Status != domain.RunEndpointStatusDegraded {
return validationError("run endpoint must be online or degraded")
}
if !svc.runEndpointHeartbeatCurrent(endpoint) {
return validationError("run endpoint heartbeat is stale")
}
if len(validator.MissingCapabilities(endpoint.Capabilities, []string{capability})) > 0 {
return validationError("run endpoint missing required capability: " + capability)
}
return nil
}
func (svc *CoreService) runEndpointHeartbeatCurrent(endpoint domain.RunEndpoint) bool {
if endpoint.LastHeartbeatAt.IsZero() {
return false
}
return !svc.now().After(endpoint.LastHeartbeatAt.Add(capacityHeartbeatStaleAfter))
}
func validateJobServerTarget(job domain.Job, instance domain.ServerInstance, plugin domain.GamePlugin) error {
if instance.State == domain.ServerInstanceStateDeleted {
return validationError("server instance must not be deleted")
}
usesDeploymentTarget := job.Capability == domain.JobCapabilityDistributionBuild && instance.DeploymentTargetID != "" && instance.DeploymentTargetID == job.RunEndpointID
if instance.RunEndpointID != job.RunEndpointID && !usesDeploymentTarget {
return validationError("job runEndpointId must match server instance")
}
if plugin.ID != instance.PluginID {
return validationError("job plugin must match server instance")
}
if job.Capability != domain.JobCapabilityDistributionBuild && !containsString(plugin.RequiredRunCapabilities, job.Capability) {
return validationError("plugin missing required capability: " + job.Capability)
}
return nil
}
func validationError(violation string) error {
return validator.ValidationError{Violations: []string{violation}}
}
func (svc *CoreService) userIDForSession(sessionID string) (string, error) {
session, err := svc.authenticatedSession(sessionID)
if err != nil {
return "", err
}
return session.UserID, nil
}
func userIDFromEmail(email string) string {
email = strings.ToLower(strings.TrimSpace(email))
var b strings.Builder
b.WriteString("user-")
for _, r := range email {
switch {
case r >= 'a' && r <= 'z':
b.WriteRune(r)
case r >= '0' && r <= '9':
b.WriteRune(r)
default:
b.WriteByte('-')
}
}
return strings.Trim(b.String(), "-")
}
func (svc *CoreService) nextUserID(user domain.User) (string, error) {
base := userIDFromEmail(user.Email)
if base == "user" || base == "" {
base = userIDFromEmail(user.DisplayName)
}
if base == "user" || base == "" {
base = "user-account"
}
if _, err := svc.store.Users().Get(base); errors.Is(err, repo.ErrNotFound) {
return base, nil
} else if err != nil {
return "", err
}
token, err := randomToken()
if err != nil {
return "", err
}
suffix := strings.ToLower(strings.TrimRight(token[:8], "-_"))
if suffix == "" {
suffix = "generated"
}
return base + "-" + suffix, nil
}
func hashPassword(password string) (string, error) {
salt := make([]byte, 16)
if _, err := rand.Read(salt); err != nil {
return "", err
}
key, err := pbkdf2.Key(sha256.New, password, salt, 120000, 32)
if err != nil {
return "", err
}
return "pbkdf2-sha256$120000$" + base64.RawStdEncoding.EncodeToString(salt) + "$" + base64.RawStdEncoding.EncodeToString(key), nil
}
func mustHashPassword(password string) string {
hash, err := hashPassword(password)
if err != nil {
panic(err)
}
return hash
}
func randomToken() (string, error) {
token := make([]byte, 32)
if _, err := rand.Read(token); err != nil {
return "", err
}
return base64.RawURLEncoding.EncodeToString(token), nil
}
func randomRunComponentKey() (string, error) {
token := make([]byte, 64)
if _, err := rand.Read(token); err != nil {
return "", err
}
return base64.RawURLEncoding.EncodeToString(token), nil
}
func verifyPassword(hash string, password string) bool {
parts := strings.Split(hash, "$")
if len(parts) != 4 || parts[0] != "pbkdf2-sha256" {
return false
}
var iterations int
if _, err := fmt.Sscanf(parts[1], "%d", &iterations); err != nil || iterations <= 0 {
return false
}
salt, err := base64.RawStdEncoding.DecodeString(parts[2])
if err != nil {
return false
}
want, err := base64.RawStdEncoding.DecodeString(parts[3])
if err != nil {
return false
}
got, err := pbkdf2.Key(sha256.New, password, salt, iterations, len(want))
if err != nil {
return false
}
return subtle.ConstantTimeCompare(got, want) == 1
}