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 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 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) 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) 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) 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) 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) ArchiveServerInstanceForSession(string, string) (domain.ServerInstance, error) GetPlatformResourceUsage() (domain.PlatformResourceUsage, error) ListServerMetricsForSession(string) ([]domain.ServerMetrics, 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) ClaimRunJob(domain.RunJobClaim) (domain.RunJobClaimResult, error) AckRunJob(domain.RunJobAck) (domain.RunJobAckResult, error) UpdateRunJobProgress(domain.RunJobProgress) (domain.RunJobProgressResult, error) CompleteRunJob(domain.RunJobResult) (domain.RunJobResultResult, 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) 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) 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) PushRunUpdateForSession(string, domain.RunUpdateRequest) (domain.RunUpdateJob, 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) IngestLogBatch(domain.LogBatchIngest) (domain.LogBatchIngestResult, error) QueryLogStream(domain.LogStreamCursorQuery) (domain.LogStreamCursorResult, error) CreateAuditEvent(domain.AuditEvent) (domain.AuditEvent, error) GetAuditEvent(string) (domain.AuditEvent, error) ListAuditEvents(domain.AuditEventFilter) ([]domain.AuditEvent, 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 jobLeases map[string]domain.RunJobLease jobLeaseSeq uint64 logStore LogBodyStore artifactMu sync.Mutex artifactTransfers map[string]domain.ArtifactTransferSession artifactPayloads map[string][]byte artifactTransferSeq uint64 auditMu sync.Mutex auditSeq uint64 aiProviderClient AIProviderClient } 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() } return &CoreService{ store: store, now: now, authSessions: map[string]string{}, runSessions: map[string]domain.RunControlSession{}, jobLeases: map[string]domain.RunJobLease{}, logStore: logStore, artifactTransfers: map[string]domain.ArtifactTransferSession{}, artifactPayloads: map[string][]byte{}, aiProviderClient: MockAIProviderClient{}, } } 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 } 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 { sessionID, err := randomToken() if err != nil { return domain.AuthSession{}, err } svc.authMu.Lock() svc.authSessions[sessionID] = created.ID svc.authMu.Unlock() return domain.AuthSession{ SessionID: sessionID, User: created, Status: "authenticated", Message: "首个账号已创建为平台管理员。", }, nil } 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 } sessionID, err := randomToken() if err != nil { return domain.AuthSession{}, err } svc.authMu.Lock() svc.authSessions[sessionID] = matched.ID svc.authMu.Unlock() return domain.AuthSession{SessionID: sessionID, User: matched, Status: "authenticated", Message: "登录成功"}, nil } func (svc *CoreService) LogoutUser(sessionID string) error { if strings.TrimSpace(sessionID) == "" { return ErrUnauthorized } svc.authMu.Lock() defer svc.authMu.Unlock() if _, exists := svc.authSessions[sessionID]; !exists { return ErrUnauthorized } delete(svc.authSessions, sessionID) return nil } 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 { const adminEmail = "operator.local@example.test" _, 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: adminEmail, Status: domain.UserStatusActive, Roles: []string{"platform-admin"}, PasswordHash: mustHashPassword("operator-local"), Profile: domain.UserProfile{ContactNote: "local development 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: "metadata", Success: true, Message: "metadata validation passed", } if err := validator.ValidateAIProvider(provider); err != nil { result.Success = false result.Message = "metadata 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 = "metadata validation failed" result.Violations = append(result.Violations, "provider must be active") } 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 } 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, RequiredRunCapabilities: manifest.Capabilities, DeclaredPermissions: manifest.Permissions, Permissions: pluginPermissionsFromManifest(manifest.Permissions), LifecycleActions: manifest.Actions, BridgeActions: manifest.Bridge.Actions, Pages: manifest.Pages, Tags: manifest.Tags, AIPurposes: manifest.AI.Purposes, RemoteAccess: manifest.RemoteAccess, 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(base, plugin, instance, request.Payload) case domain.PluginBridgeActionRunDistribution: base = svc.executeBridgeRunDistribution(sessionID, base, request) case domain.PluginBridgeActionDependenciesRequest: base = svc.executeBridgeDependenciesRequest(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.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, ProviderID: request.Payload["providerId"], 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 } 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) 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(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) { base.Status = "denied" base.Error = &domain.PluginBridgeSafeError{Code: "capability_denied", Message: "requested remote capability is not declared by plugin"} return base } job := domain.Job{ ID: jobIDFromParts("job-remote", base.RequestID, capability), ServerInstanceID: instance.ID, RunEndpointID: instance.RunEndpointID, Capability: capability, TargetKey: payload["targetKey"], InputRef: payload["inputRef"], IdempotencyKey: defaultBridgeValue(payload["idempotencyKey"], base.RequestID), Progress: domain.JobProgress{Percent: 0, Message: "remote access job queued"}, } created, err := svc.CreateJob(job) if err != nil { return bridgeExecutionError(base, err) } base.Status = "queued" base.Result = map[string]string{ "jobId": created.ID, "state": string(created.State), "capability": created.Capability, "targetKey": created.TargetKey, "serverInstanceId": created.ServerInstanceID, } return base } 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), "secretRef": distribution.SecretRef, "status": string(distribution.Status), } return base } func (svc *CoreService) executeBridgeClientManager(sessionID string, base domain.PluginBridgeExecuteResponse, request domain.PluginBridgeExecuteRequest) domain.PluginBridgeExecuteResponse { distribution, err := svc.GenerateClientManagerDistributionForSession(sessionID, domain.ClientManagerBuildRequest{ ServerInstanceID: request.ServerInstanceID, ProfileKey: request.Payload["profileKey"], TargetOS: defaultBridgeValue(request.Payload["targetOs"], "windows"), TargetArch: defaultBridgeValue(request.Payload["targetArch"], "amd64"), RepositoryURL: request.Payload["repositoryUrl"], SourceRevision: request.Payload["sourceRevision"], 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, "buildJobId": distribution.BuildJobID, "artifactId": distribution.ArtifactID, "checksum": distribution.Checksum, "keyGeneration": strconv.Itoa(distribution.KeyGeneration), "secretRef": distribution.SecretRef, "status": string(distribution.Status), } return base } func (svc *CoreService) executeBridgeDependenciesRequest(base domain.PluginBridgeExecuteResponse, plugin domain.GamePlugin, instance domain.ServerInstance, payload map[string]string) domain.PluginBridgeExecuteResponse { action := defaultBridgeValue(payload["action"], "check") capability := domain.JobCapabilityDependenciesCheck message := "dependency check queued" if action == "install" { capability = domain.JobCapabilityDependenciesInstall message = "dependency install queued" } 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.CreateJob(domain.Job{ ID: jobIDFromParts("job-dependencies", base.RequestID, capability), ServerInstanceID: instance.ID, RunEndpointID: instance.RunEndpointID, Capability: capability, TargetKey: defaultBridgeValue(payload["probeKey"], "dependencies/default"), InputRef: payload["inputRef"], IdempotencyKey: defaultBridgeValue(payload["idempotencyKey"], base.RequestID), Progress: domain.JobProgress{Percent: 0, Message: message}, }) 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, Capabilities: plugin.RequiredRunCapabilities, DeclaredPermissions: plugin.DeclaredPermissions, Permissions: plugin.Permissions, LifecycleActions: plugin.LifecycleActions, BridgeActions: plugin.BridgeActions, Pages: plugin.Pages, Tags: plugin.Tags, AIPurposes: plugin.AIPurposes, RemoteAccess: plugin.RemoteAccess, ValidationViolations: plugin.ValidationViolations, Status: plugin.Status, Source: "platform-registry", } } 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 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) } 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 err := validator.ValidateServerInstanceDependencies(instance, plugin, endpoint); err != nil { return domain.ServerInstance{}, err } 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, } if config.UpdatedAt.IsZero() { config.UpdatedAt = svc.now() } config.Content = buildLogicalServerConfig(instance) 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 config.Key != request.Key { return domain.ServerConfigDiffPreview{}, validationError("key must match server config") } preview := domain.ServerConfigDiffPreview{ ServerInstanceID: request.ServerInstanceID, ConfigVersion: config.ConfigVersion, 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, 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 } 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, 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") } 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, 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) 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) ArchiveServerInstanceForSession(sessionID string, serverInstanceID 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 instance.State == domain.ServerInstanceStateRunning || instance.State == domain.ServerInstanceStateInstalling { return domain.ServerInstance{}, validationError("running or installing server instances must be stopped before archive") } 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 } 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 := 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) { return svc.store.Jobs().Get(id) } func (svc *CoreService) ListJobs(filter domain.JobFilter) ([]domain.Job, error) { return svc.store.Jobs().List(filter) } 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 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 len(validator.MissingCapabilities(endpoint.Capabilities, []string{capability})) > 0 { return validationError("run endpoint missing required capability: " + capability) } return nil } 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") } if instance.RunEndpointID != job.RunEndpointID { return validationError("job runEndpointId must match server instance") } if plugin.ID != instance.PluginID { return validationError("job plugin must match server instance") } if !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) { sessionID = strings.TrimSpace(sessionID) if sessionID == "" { return "", ErrUnauthorized } svc.authMu.Lock() defer svc.authMu.Unlock() userID, exists := svc.authSessions[sessionID] if !exists { return "", ErrUnauthorized } return 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 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 }