package service import ( "crypto/sha256" "encoding/hex" "errors" "sort" "strings" "time" "browser.local/platform/domain" "browser.local/platform/repo" "browser.local/platform/validator" ) const ( capacityHeartbeatStaleAfter = 2 * time.Minute capacityRetryAfterSeconds = 30 capacityLogBacklogLimit = 256 capacityArtifactBacklogLimit = 128 aiConfigDiffTTL = 30 * time.Minute ) func (svc *CoreService) GetProductionCapacityForSession(sessionID string) (domain.ProductionCapacitySummary, error) { user, err := svc.GetCurrentUser(sessionID) if err != nil { return domain.ProductionCapacitySummary{}, err } endpoints, err := svc.store.RunEndpoints().List(domain.RunEndpointFilter{}) if err != nil { return domain.ProductionCapacitySummary{}, err } visibleEndpointIDs, err := svc.visibleEndpointIDs(user) if err != nil { return domain.ProductionCapacitySummary{}, err } alerts, err := svc.store.Alerts().List(domain.AlertFilter{}) if err != nil { return domain.ProductionCapacitySummary{}, err } summary := domain.ProductionCapacitySummary{GeneratedAt: svc.now()} for _, endpoint := range endpoints { if !isPlatformAdmin(user) { if _, visible := visibleEndpointIDs[endpoint.ID]; !visible { continue } } running, queued, err := svc.capacityJobCounts(endpoint.ID) if err != nil { return domain.ProductionCapacitySummary{}, err } projection := svc.capacityProjection(endpoint, running, queued) for _, alert := range alerts { if alert.SourceKind == "run-endpoint" && alert.SourceID == endpoint.ID && alert.RuleKey == "capacity.pressure" && alert.State != domain.AlertStateResolved { projection.LastAdmissionDecision = domain.CapacityAdmissionDeferred projection.LastAdmissionReason = alert.Message projection.LastAdmissionCheckedAt = alert.LastSeenAt } } summary.Endpoints = append(summary.Endpoints, projection) summary.TotalMaxJobs += projection.MaxJobs summary.TotalRunningJobs += projection.RunningJobs summary.TotalQueuedJobs += projection.QueuedJobs } for _, alert := range alerts { if alert.State != domain.AlertStateResolved && svc.canAccessAlert(user, alert) { summary.ActiveAlerts++ } } sort.Slice(summary.Endpoints, func(i, j int) bool { return summary.Endpoints[i].RunEndpointID < summary.Endpoints[j].RunEndpointID }) return domain.CopyProductionCapacitySummary(summary), nil } func (svc *CoreService) CheckCapacityAdmissionForSession(sessionID string, request domain.CapacityAdmissionRequest) (domain.CapacityAdmissionDecision, error) { if err := validator.ValidateCapacityAdmissionRequest(request); err != nil { return domain.CapacityAdmissionDecision{}, err } user, err := svc.GetCurrentUser(sessionID) if err != nil { return domain.CapacityAdmissionDecision{}, err } request, err = svc.authorizeCapacityRequest(user, request) if err != nil { return domain.CapacityAdmissionDecision{}, err } svc.productionMu.Lock() defer svc.productionMu.Unlock() return svc.checkCapacityAdmission(user.ID, request) } func (svc *CoreService) checkCapacityAdmission(actorID string, request domain.CapacityAdmissionRequest) (domain.CapacityAdmissionDecision, error) { endpoint, err := svc.store.RunEndpoints().Get(request.RunEndpointID) if err != nil { return domain.CapacityAdmissionDecision{}, err } running, queued, err := svc.capacityJobCounts(endpoint.ID) if err != nil { return domain.CapacityAdmissionDecision{}, err } projection := svc.capacityProjection(endpoint, running, queued) decision := domain.CapacityAdmissionDecision{ Accepted: true, State: domain.CapacityAdmissionAccepted, Reason: "capacity available", ServerInstanceID: request.ServerInstanceID, RunEndpointID: endpoint.ID, Capability: request.Capability, TargetKey: request.TargetKey, MaxJobs: projection.MaxJobs, RunningJobs: projection.RunningJobs, QueuedJobs: projection.QueuedJobs, CheckedAt: svc.now(), } pressure := append([]domain.CapacityPressureCode(nil), projection.PressureCodes...) if len(validator.MissingCapabilities(endpoint.Capabilities, []string{request.Capability})) > 0 { pressure = appendCapacityPressure(pressure, domain.CapacityPressureCapabilityGap) } decision.PressureCodes = pressure hardDenied := containsCapacityPressure(pressure, domain.CapacityPressureEndpointOffline) || containsCapacityPressure(pressure, domain.CapacityPressureCapabilityGap) if hardDenied { decision.Accepted = false decision.State = domain.CapacityAdmissionDenied decision.Reason = "endpoint is unavailable or missing the required capability" } else if len(pressure) > 0 { decision.Accepted = false decision.State = domain.CapacityAdmissionDeferred decision.Reason = "endpoint capacity is temporarily under pressure" decision.RetryAfterSeconds = capacityRetryAfterSeconds } auditResult := domain.AuditResultSuccess auditAction := "capacity.admission.accepted" if !decision.Accepted { auditResult = domain.AuditResultDenied auditAction = "capacity.admission.denied" } auditID, err := svc.recordAuditEventWithID(actorID, auditAction, "run-endpoint", endpoint.ID, auditResult, decision.Reason) if err != nil { return domain.CapacityAdmissionDecision{}, err } decision.AuditEventID = auditID if !decision.Accepted { severity := domain.AlertSeverityWarning if hardDenied { severity = domain.AlertSeverityCritical } alert, err := svc.upsertAlert(domain.AlertRecord{ SourceKind: "run-endpoint", SourceID: endpoint.ID, RuleKey: "capacity.pressure", Severity: severity, Title: "Run endpoint capacity admission blocked", Message: decision.Reason, Retryable: true, RetryAfterSeconds: decision.RetryAfterSeconds, LastAuditEventID: auditID, }) if err != nil { return domain.CapacityAdmissionDecision{}, err } decision.AlertID = alert.ID } else if err := svc.resolveAlertForSource("run-endpoint", endpoint.ID, "capacity.pressure", actorID, "capacity returned to an admissible state", auditID); err != nil { return domain.CapacityAdmissionDecision{}, err } if err := validator.ValidateCapacityAdmissionDecision(decision); err != nil { return domain.CapacityAdmissionDecision{}, err } return domain.CopyCapacityAdmissionDecision(decision), nil } func (svc *CoreService) ListAlertsForSession(sessionID string, filter domain.AlertFilter) ([]domain.AlertRecord, error) { user, err := svc.GetCurrentUser(sessionID) if err != nil { return nil, err } alerts, err := svc.store.Alerts().List(filter) if err != nil { return nil, err } visible := make([]domain.AlertRecord, 0, len(alerts)) for _, alert := range alerts { if svc.canAccessAlert(user, alert) { visible = append(visible, alert) } } sort.Slice(visible, func(i, j int) bool { return visible[i].UpdatedAt.After(visible[j].UpdatedAt) }) return domain.CopyAlertRecords(visible), nil } func (svc *CoreService) AcknowledgeAlertForSession(sessionID string, request domain.AlertAcknowledgeRequest) (domain.AlertRecord, error) { if err := validator.ValidateAlertAcknowledgeRequest(request); err != nil { return domain.AlertRecord{}, err } user, err := svc.GetCurrentUser(sessionID) if err != nil { return domain.AlertRecord{}, err } svc.productionMu.Lock() defer svc.productionMu.Unlock() alert, err := svc.store.Alerts().Get(request.AlertID) if err != nil { return domain.AlertRecord{}, err } if !svc.canAccessAlert(user, alert) { return domain.AlertRecord{}, ErrForbidden } if alert.State == domain.AlertStateResolved { return domain.AlertRecord{}, validationError("resolved alerts cannot be acknowledged") } stamp := svc.now() auditID, err := svc.recordAuditEventWithID(user.ID, "alert.acknowledge", "alert", alert.ID, domain.AuditResultSuccess, defaultAlertNote(request.Note, "alert acknowledged")) if err != nil { return domain.AlertRecord{}, err } alert.State = domain.AlertStateAcknowledged alert.AcknowledgedBy = user.ID alert.AcknowledgedAt = stamp alert.LastAuditEventID = auditID alert.UpdatedAt = stamp if err := validator.ValidateAlertRecord(alert); err != nil { return domain.AlertRecord{}, err } if err := svc.store.Alerts().Update(alert); err != nil { return domain.AlertRecord{}, err } return domain.CopyAlertRecord(alert), nil } func (svc *CoreService) ResolveAlertForSession(sessionID string, request domain.AlertResolveRequest) (domain.AlertRecord, error) { if err := validator.ValidateAlertResolveRequest(request); err != nil { return domain.AlertRecord{}, err } user, err := svc.GetCurrentUser(sessionID) if err != nil { return domain.AlertRecord{}, err } svc.productionMu.Lock() defer svc.productionMu.Unlock() alert, err := svc.store.Alerts().Get(request.AlertID) if err != nil { return domain.AlertRecord{}, err } if !svc.canAccessAlert(user, alert) { return domain.AlertRecord{}, ErrForbidden } if alert.State == domain.AlertStateResolved { return domain.CopyAlertRecord(alert), nil } stamp := svc.now() note := defaultAlertNote(request.Note, "alert resolved after operator review") auditID, err := svc.recordAuditEventWithID(user.ID, "alert.resolve", "alert", alert.ID, domain.AuditResultSuccess, note) if err != nil { return domain.AlertRecord{}, err } alert.State = domain.AlertStateResolved alert.ResolvedBy = user.ID alert.ResolvedAt = stamp alert.ResolutionNote = note alert.LastAuditEventID = auditID alert.UpdatedAt = stamp if err := validator.ValidateAlertRecord(alert); err != nil { return domain.AlertRecord{}, err } if err := svc.store.Alerts().Update(alert); err != nil { return domain.AlertRecord{}, err } return domain.CopyAlertRecord(alert), nil } func (svc *CoreService) RetryAlertForSession(sessionID string, request domain.AlertRetryRequest) (domain.AlertRetryResult, error) { if err := validator.ValidateAlertRetryRequest(request); err != nil { return domain.AlertRetryResult{}, err } user, err := svc.GetCurrentUser(sessionID) if err != nil { return domain.AlertRetryResult{}, err } alert, err := svc.store.Alerts().Get(request.AlertID) if err != nil { return domain.AlertRetryResult{}, err } if !svc.canAccessAlert(user, alert) { return domain.AlertRetryResult{}, ErrForbidden } if !alert.Retryable { return domain.AlertRetryResult{}, validationError("alert source is not retryable") } switch alert.SourceKind { case "run-endpoint": endpoint, err := svc.store.RunEndpoints().Get(alert.SourceID) if err != nil { return domain.AlertRetryResult{}, err } capability := firstCapacityCapability(endpoint.Capabilities) decision, err := svc.CheckCapacityAdmissionForSession(sessionID, domain.CapacityAdmissionRequest{RunEndpointID: endpoint.ID, Capability: capability, IdempotencyKey: request.IdempotencyKey}) if err != nil { return domain.AlertRetryResult{}, err } updated, err := svc.store.Alerts().Get(alert.ID) if err != nil { return domain.AlertRetryResult{}, err } return domain.CopyAlertRetryResult(domain.AlertRetryResult{Alert: updated, Decision: decision, Status: string(decision.State)}), nil case "plugin-lifecycle": installation, err := svc.store.PluginLifecycles().Get(alert.SourceID) if err != nil { return domain.AlertRetryResult{}, err } result, err := svc.RunPluginLifecycleForSession(sessionID, domain.PluginLifecycleRequest{PluginID: installation.PluginID, ServerInstanceID: installation.ServerInstanceID, Operation: installation.LastOperation, TargetVersion: installation.TargetVersion, IdempotencyKey: request.IdempotencyKey, Confirmed: true}) if err != nil { return domain.AlertRetryResult{}, err } return domain.CopyAlertRetryResult(domain.AlertRetryResult{Alert: alert, Decision: result.Decision, Status: result.Status}), nil default: return domain.AlertRetryResult{}, validationError("alert source does not support scoped retry") } } func (svc *CoreService) ListPluginLifecyclesForSession(sessionID string, filter domain.PluginLifecycleFilter) ([]domain.PluginLifecycleInstallation, error) { user, err := svc.GetCurrentUser(sessionID) if err != nil { return nil, err } installations, err := svc.store.PluginLifecycles().List(filter) if err != nil { return nil, err } visible := make([]domain.PluginLifecycleInstallation, 0, len(installations)) for _, installation := range installations { instance, err := svc.store.ServerInstances().Get(installation.ServerInstanceID) if err == nil && canAccessServer(user, instance) { visible = append(visible, installation) } } sort.Slice(visible, func(i, j int) bool { return visible[i].UpdatedAt.After(visible[j].UpdatedAt) }) return domain.CopyPluginLifecycleInstallations(visible), nil } func (svc *CoreService) RunPluginLifecycleForSession(sessionID string, request domain.PluginLifecycleRequest) (domain.PluginLifecycleResult, error) { if err := validator.ValidatePluginLifecycleRequest(request); err != nil { return domain.PluginLifecycleResult{}, err } user, instance, err := svc.requireServerOwner(sessionID, request.ServerInstanceID) if err != nil { return domain.PluginLifecycleResult{}, err } if instance.PluginID != request.PluginID { return domain.PluginLifecycleResult{}, validationError("pluginId must match the server plugin") } plugin, err := svc.store.GamePlugins().Get(request.PluginID) if err != nil { return domain.PluginLifecycleResult{}, err } if !containsString(plugin.ProductionLifecycle.Operations, string(request.Operation)) { return domain.PluginLifecycleResult{}, validationError("plugin lifecycle operation is not declared by the manifest") } endpoint, err := svc.store.RunEndpoints().Get(instance.RunEndpointID) if err != nil { return domain.PluginLifecycleResult{}, err } capability, targetKey, err := pluginLifecycleDispatchMetadata(plugin, request.Operation) if err != nil { return domain.PluginLifecycleResult{}, err } if request.TargetVersion == "" { request.TargetVersion = plugin.Version } if !containsString(plugin.SupportedOS, endpoint.Platform) { return svc.pluginLifecycleDenied(user.ID, instance, plugin, request, "plugin is not compatible with the assigned endpoint platform") } svc.productionMu.Lock() defer svc.productionMu.Unlock() installationID := pluginLifecycleInstallationID(request.PluginID, request.ServerInstanceID) installation, getErr := svc.store.PluginLifecycles().Get(installationID) if errors.Is(getErr, repo.ErrNotFound) { stamp := svc.now() installation = domain.PluginLifecycleInstallation{ID: installationID, PluginID: plugin.ID, ServerInstanceID: instance.ID, TargetVersion: request.TargetVersion, DesiredState: domain.PluginLifecycleStatePending, CurrentState: domain.PluginLifecycleStatePending, Compatibility: "compatible", DependencyState: domain.DependencyStateUnknown, CreatedAt: stamp, UpdatedAt: stamp} } else if getErr != nil { return domain.PluginLifecycleResult{}, getErr } if err := validatePluginLifecycleTransition(installation, request); err != nil { return svc.pluginLifecycleDeniedLocked(user.ID, installation, request, err.Error()) } existingJob, jobErr := svc.store.Jobs().GetByIdempotency(endpoint.ID, request.IdempotencyKey) if jobErr == nil { if existingJob.ServerInstanceID != instance.ID || existingJob.Capability != capability || existingJob.TargetKey != targetKey || existingJob.ExecutionInput.PluginID != plugin.ID || existingJob.ExecutionInput.LifecycleOperation != string(request.Operation) || existingJob.ExecutionInput.TargetVersion != request.TargetVersion { return domain.PluginLifecycleResult{}, validationError("idempotencyKey is already used for different plugin lifecycle inputs") } installation.JobID = existingJob.ID return domain.CopyPluginLifecycleResult(domain.PluginLifecycleResult{Installation: installation, Job: existingJob, Status: "queued"}), nil } if !errors.Is(jobErr, repo.ErrNotFound) { return domain.PluginLifecycleResult{}, jobErr } decision, err := svc.checkCapacityAdmission(user.ID, domain.CapacityAdmissionRequest{ServerInstanceID: instance.ID, RunEndpointID: endpoint.ID, Capability: capability, TargetKey: targetKey, IdempotencyKey: request.IdempotencyKey}) if err != nil { return domain.PluginLifecycleResult{}, err } if !decision.Accepted { return domain.CopyPluginLifecycleResult(domain.PluginLifecycleResult{Installation: installation, Decision: decision, Status: string(decision.State)}), nil } job, err := svc.CreateJob(domain.Job{ ID: jobIDFromParts("job-plugin-lifecycle", installation.ID, request.IdempotencyKey), ServerInstanceID: instance.ID, RunEndpointID: endpoint.ID, Capability: capability, TargetKey: targetKey, IdempotencyKey: request.IdempotencyKey, ExecutionInput: domain.JobExecutionInput{WorkspaceScope: svc.runtimeProfileScope(instance.ID), PluginID: plugin.ID, LifecycleOperation: string(request.Operation), TargetVersion: request.TargetVersion}, Progress: domain.JobProgress{Percent: 0, Message: "plugin lifecycle operation queued"}, }) if err != nil { return domain.PluginLifecycleResult{}, err } stamp := svc.now() installation.PreviousVersion = installation.CurrentVersion installation.TargetVersion = request.TargetVersion installation.DesiredState = desiredPluginLifecycleState(request.Operation, installation) installation.CurrentState = dispatchedPluginLifecycleState(request.Operation, installation.CurrentState) installation.LastOperation = request.Operation installation.JobID = job.ID installation.IdempotencyKey = request.IdempotencyKey installation.FailureReason = "" installation.UpdatedAt = stamp auditID, err := svc.recordAuditEventWithID(user.ID, "plugin.lifecycle."+string(request.Operation), "plugin-lifecycle", installation.ID, domain.AuditResultQueued, "plugin lifecycle operation admitted and queued") if err != nil { return domain.PluginLifecycleResult{}, err } installation.AuditEventID = auditID if err := validator.ValidatePluginLifecycleInstallation(installation); err != nil { return domain.PluginLifecycleResult{}, err } if errors.Is(getErr, repo.ErrNotFound) { err = svc.store.PluginLifecycles().Create(installation) } else { err = svc.store.PluginLifecycles().Update(installation) } if err != nil { return domain.PluginLifecycleResult{}, err } return domain.CopyPluginLifecycleResult(domain.PluginLifecycleResult{Installation: installation, Job: job, Decision: decision, Status: "queued"}), nil } func (svc *CoreService) ListAIConfigDiffsForSession(sessionID string, filter domain.AIConfigDiffFilter) ([]domain.AIConfigDiffPreview, error) { user, err := svc.GetCurrentUser(sessionID) if err != nil { return nil, err } previews, err := svc.store.AIConfigDiffs().List(filter) if err != nil { return nil, err } visible := make([]domain.AIConfigDiffPreview, 0, len(previews)) for _, preview := range previews { instance, err := svc.store.ServerInstances().Get(preview.ServerInstanceID) if err == nil && canAccessServer(user, instance) { visible = append(visible, preview) } } sort.Slice(visible, func(i, j int) bool { return visible[i].CreatedAt.After(visible[j].CreatedAt) }) return domain.CopyAIConfigDiffPreviews(visible), nil } func (svc *CoreService) ApproveAIConfigDiffForSession(sessionID string, request domain.AIConfigDiffApprovalRequest) (domain.AIConfigDiffApprovalResult, error) { if err := validator.ValidateAIConfigDiffApprovalRequest(request); err != nil { return domain.AIConfigDiffApprovalResult{}, err } user, err := svc.GetCurrentUser(sessionID) if err != nil { return domain.AIConfigDiffApprovalResult{}, err } svc.productionMu.Lock() defer svc.productionMu.Unlock() preview, err := svc.store.AIConfigDiffs().Get(request.DiffID) if err != nil { return domain.AIConfigDiffApprovalResult{}, err } instance, err := svc.store.ServerInstances().Get(preview.ServerInstanceID) if err != nil { return domain.AIConfigDiffApprovalResult{}, err } if !canAccessServer(user, instance) || (!isPlatformAdmin(user) && preview.CreatedBy != user.ID) { return domain.AIConfigDiffApprovalResult{}, ErrForbidden } if preview.State == domain.AIConfigDiffStateApproved { if preview.ApprovalIdempotencyKey != request.IdempotencyKey { return domain.AIConfigDiffApprovalResult{}, validationError("AI config diff is already approved with another idempotency key") } job, err := svc.store.Jobs().Get(preview.JobID) if err != nil { return domain.AIConfigDiffApprovalResult{}, err } dispatch := domain.ServerConfigWriteDispatch{Job: job, Status: "queued"} return domain.CopyAIConfigDiffApprovalResult(domain.AIConfigDiffApprovalResult{Preview: preview, Dispatch: dispatch}), nil } if preview.State != domain.AIConfigDiffStatePending { return domain.AIConfigDiffApprovalResult{}, validationError("AI config diff is not pending approval") } if !preview.ExpiresAt.After(svc.now()) { preview.State = domain.AIConfigDiffStateExpired preview.UpdatedAt = svc.now() _ = svc.store.AIConfigDiffs().Update(preview) return domain.AIConfigDiffApprovalResult{}, validationError("AI config diff has expired") } dispatch, err := svc.ApproveServerConfigWriteForSession(sessionID, domain.ServerConfigWriteApproval{ServerInstanceID: preview.ServerInstanceID, ExpectedConfigVersion: preview.ConfigVersion, ExpectedChecksum: preview.CurrentConfigChecksum, Key: preview.Key, ProposedContent: preview.ProposedConfig, IdempotencyKey: request.IdempotencyKey}) if err != nil { return domain.AIConfigDiffApprovalResult{}, err } stamp := svc.now() preview.State = domain.AIConfigDiffStateApproved preview.ApprovedBy = user.ID preview.ApprovedAt = stamp preview.ApprovalIdempotencyKey = request.IdempotencyKey preview.JobID = dispatch.Job.ID preview.UpdatedAt = stamp if err := svc.store.AIConfigDiffs().Update(preview); err != nil { return domain.AIConfigDiffApprovalResult{}, err } if _, err := svc.recordAuditEventWithID(user.ID, "ai.config-diff.approve", "ai-config-diff", preview.ID, domain.AuditResultQueued, "approved reviewed AI config diff and queued one config write job"); err != nil { return domain.AIConfigDiffApprovalResult{}, err } return domain.CopyAIConfigDiffApprovalResult(domain.AIConfigDiffApprovalResult{Preview: preview, Dispatch: dispatch}), nil } func (svc *CoreService) projectProductionOpsJobResult(job domain.Job, stamp time.Time) error { if !strings.HasPrefix(job.ID, "job-plugin-lifecycle-") || job.ExecutionInput.LifecycleOperation == "" || job.ExecutionInput.PluginID == "" { return nil } svc.productionMu.Lock() defer svc.productionMu.Unlock() installation, err := svc.store.PluginLifecycles().Get(pluginLifecycleInstallationID(job.ExecutionInput.PluginID, job.ServerInstanceID)) if err != nil { return err } if installation.JobID != job.ID { return nil } if job.State == domain.JobStateSucceeded { applyPluginLifecycleSuccess(&installation) installation.FailureReason = "" if err := svc.resolveAlertForSource("plugin-lifecycle", installation.ID, "plugin.lifecycle.failed", "run:"+job.RunEndpointID, "plugin lifecycle job completed", installation.AuditEventID); err != nil { return err } } else { installation.CurrentState = domain.PluginLifecycleStateFailed installation.FailureReason = "plugin lifecycle job did not complete successfully" alert, err := svc.upsertAlert(domain.AlertRecord{SourceKind: "plugin-lifecycle", SourceID: installation.ID, RuleKey: "plugin.lifecycle.failed", Severity: domain.AlertSeverityWarning, Title: "Plugin lifecycle operation failed", Message: installation.FailureReason, Retryable: true, LastJobID: job.ID, LastAuditEventID: installation.AuditEventID}) if err != nil { return err } installation.AlertID = alert.ID } installation.UpdatedAt = stamp return svc.store.PluginLifecycles().Update(installation) } func (svc *CoreService) persistAIConfigDiff(actorID string, provider domain.AIProvider, request domain.AIInvocationRequest, result domain.AIProviderInvocationResult) (domain.AIConfigDiffPreview, error) { config, err := svc.getServerConfigForUser(actorID, request.ServerInstanceID) if err != nil { return domain.AIConfigDiffPreview{}, err } stamp := svc.now() preview := domain.AIConfigDiffPreview{ID: aiConfigDiffID(request.RequestID, request.ServerInstanceID), RequestID: request.RequestID, CreatedBy: actorID, ServerInstanceID: request.ServerInstanceID, PluginID: request.PluginID, ProviderID: provider.ID, Model: result.Usage.Model, Key: config.Key, ConfigVersion: config.ConfigVersion, CurrentConfigChecksum: config.Checksum, ProposedConfig: result.SuggestedConfig, DiffSummary: "review required before config write dispatch", State: domain.AIConfigDiffStatePending, ExpiresAt: stamp.Add(aiConfigDiffTTL), CreatedAt: stamp, UpdatedAt: stamp} if err := validator.ValidateAIConfigDiffPreview(preview); err != nil { return domain.AIConfigDiffPreview{}, err } if err := svc.store.AIConfigDiffs().Create(preview); err != nil { if !errors.Is(err, repo.ErrDuplicate) { return domain.AIConfigDiffPreview{}, err } existing, getErr := svc.store.AIConfigDiffs().Get(preview.ID) if getErr != nil { return domain.AIConfigDiffPreview{}, getErr } if existing.CreatedBy != actorID || existing.ServerInstanceID != request.ServerInstanceID || existing.PluginID != request.PluginID || existing.ProposedConfig != result.SuggestedConfig { return domain.AIConfigDiffPreview{}, validationError("requestId is already used for a different AI config recommendation") } return existing, nil } return preview, nil } func (svc *CoreService) getServerConfigForUser(userID, serverInstanceID string) (domain.ServerConfig, error) { instance, err := svc.store.ServerInstances().Get(serverInstanceID) if err != nil { return domain.ServerConfig{}, err } user, err := svc.store.Users().Get(userID) if err != nil { return domain.ServerConfig{}, err } if !canAccessServer(user, instance) { return domain.ServerConfig{}, ErrForbidden } config := domain.ServerConfig{ServerInstanceID: instance.ID, ConfigVersion: instance.ConfigVersion, Format: "properties", Key: instance.ConfigKey, Content: instance.ConfigContent, Checksum: instance.ConfigChecksum, Source: "platform-derived", UpdatedAt: instance.ConfigUpdatedAt} if config.ConfigVersion <= 0 { config.ConfigVersion = 1 } if config.Key == "" { config.Key = "server.properties" } if config.Content == "" { config.Content = buildLogicalServerConfig(instance) } if config.Checksum == "" { config.Checksum = validator.BytesChecksum([]byte(config.Content)) } if config.UpdatedAt.IsZero() { config.UpdatedAt = svc.now() } return config, validator.ValidateServerConfig(config) } func (svc *CoreService) authorizeCapacityRequest(user domain.User, request domain.CapacityAdmissionRequest) (domain.CapacityAdmissionRequest, error) { if request.ServerInstanceID != "" { instance, err := svc.store.ServerInstances().Get(request.ServerInstanceID) if err != nil { return request, err } if !canAccessServer(user, instance) { return request, ErrForbidden } if request.RunEndpointID != "" && request.RunEndpointID != instance.RunEndpointID { return request, validationError("runEndpointId must match server instance") } request.RunEndpointID = instance.RunEndpointID return request, nil } if request.RunEndpointID == "" { return request, validationError("serverInstanceId or runEndpointId is required") } if isPlatformAdmin(user) { return request, nil } visible, err := svc.visibleEndpointIDs(user) if err != nil { return request, err } if _, ok := visible[request.RunEndpointID]; !ok { return request, ErrForbidden } return request, nil } func (svc *CoreService) visibleEndpointIDs(user domain.User) (map[string]struct{}, error) { instances, err := svc.store.ServerInstances().List(domain.ServerInstanceFilter{}) if err != nil { return nil, err } ids := map[string]struct{}{} for _, instance := range instances { if isPlatformAdmin(user) || canAccessServer(user, instance) { ids[instance.RunEndpointID] = struct{}{} } } return ids, nil } func (svc *CoreService) capacityJobCounts(endpointID string) (int, int, error) { jobs, err := svc.store.Jobs().List(domain.JobFilter{RunEndpointID: endpointID}) if err != nil { return 0, 0, err } running, queued := 0, 0 for _, job := range jobs { switch job.State { case domain.JobStateAccepted, domain.JobStateRunning: running++ case domain.JobStateQueued, domain.JobStateRetrying: queued++ } } return running, queued, nil } func (svc *CoreService) capacityProjection(endpoint domain.RunEndpoint, durableRunning, durableQueued int) domain.EndpointCapacityProjection { running := maxInt(endpoint.Capacity.RunningJobs, durableRunning) queued := maxInt(endpoint.Capacity.QueuedJobs, durableQueued) projection := domain.EndpointCapacityProjection{RunEndpointID: endpoint.ID, DisplayName: endpoint.DisplayName, Status: endpoint.Status, Capabilities: endpoint.Capabilities, MaxJobs: endpoint.Capacity.MaxJobs, RunningJobs: running, QueuedJobs: queued, LogBacklogBatches: endpoint.Capacity.LogBacklogBatches, ArtifactBacklogChunks: endpoint.Capacity.ArtifactBacklogChunks, Summary: safeBridgeReason(endpoint.Capacity.Summary), LastHeartbeatAt: endpoint.LastHeartbeatAt} if endpoint.Status != domain.RunEndpointStatusOnline && endpoint.Status != domain.RunEndpointStatusDegraded { projection.PressureCodes = appendCapacityPressure(projection.PressureCodes, domain.CapacityPressureEndpointOffline) } if endpoint.LastHeartbeatAt.IsZero() || svc.now().Sub(endpoint.LastHeartbeatAt) > capacityHeartbeatStaleAfter { projection.PressureCodes = appendCapacityPressure(projection.PressureCodes, domain.CapacityPressureEndpointStale) } if projection.MaxJobs <= 0 || running >= projection.MaxJobs { projection.PressureCodes = appendCapacityPressure(projection.PressureCodes, domain.CapacityPressureJobLimit) } queueLimit := maxInt(4, projection.MaxJobs*2) if queued >= queueLimit { projection.PressureCodes = appendCapacityPressure(projection.PressureCodes, domain.CapacityPressureQueueLimit) } if projection.LogBacklogBatches >= capacityLogBacklogLimit || projection.ArtifactBacklogChunks >= capacityArtifactBacklogLimit || len(endpoint.Capacity.PressureCodes) > 0 { projection.PressureCodes = appendCapacityPressure(projection.PressureCodes, domain.CapacityPressureBacklog) } return projection } func (svc *CoreService) upsertAlert(candidate domain.AlertRecord) (domain.AlertRecord, error) { stamp := svc.now() candidate.ID = alertIDForSource(candidate.SourceKind, candidate.SourceID, candidate.RuleKey) existing, err := svc.store.Alerts().Get(candidate.ID) if err == nil { existing.Severity = candidate.Severity existing.State = domain.AlertStateActive existing.Title = candidate.Title existing.Message = safeBridgeReason(candidate.Message) existing.OccurrenceCount++ existing.Retryable = candidate.Retryable existing.RetryAfterSeconds = candidate.RetryAfterSeconds existing.LastJobID = candidate.LastJobID existing.LastAuditEventID = candidate.LastAuditEventID existing.LastSeenAt = stamp existing.ResolvedBy = "" existing.ResolvedAt = time.Time{} existing.ResolutionNote = "" existing.UpdatedAt = stamp if err := validator.ValidateAlertRecord(existing); err != nil { return domain.AlertRecord{}, err } if err := svc.store.Alerts().Update(existing); err != nil { return domain.AlertRecord{}, err } return existing, nil } if !errors.Is(err, repo.ErrNotFound) { return domain.AlertRecord{}, err } candidate.State = domain.AlertStateActive candidate.Message = safeBridgeReason(candidate.Message) candidate.OccurrenceCount = 1 candidate.LastSeenAt = stamp candidate.CreatedAt = stamp candidate.UpdatedAt = stamp if err := validator.ValidateAlertRecord(candidate); err != nil { return domain.AlertRecord{}, err } if err := svc.store.Alerts().Create(candidate); err != nil { return domain.AlertRecord{}, err } return candidate, nil } func (svc *CoreService) resolveAlertForSource(sourceKind, sourceID, ruleKey, actorID, note, auditID string) error { alert, err := svc.store.Alerts().Get(alertIDForSource(sourceKind, sourceID, ruleKey)) if errors.Is(err, repo.ErrNotFound) { return nil } if err != nil { return err } if alert.State == domain.AlertStateResolved { return nil } stamp := svc.now() alert.State = domain.AlertStateResolved alert.ResolvedBy = actorID alert.ResolvedAt = stamp alert.ResolutionNote = note alert.LastAuditEventID = auditID alert.UpdatedAt = stamp return svc.store.Alerts().Update(alert) } func (svc *CoreService) canAccessAlert(user domain.User, alert domain.AlertRecord) bool { if isPlatformAdmin(user) { return true } switch alert.SourceKind { case "server-instance": instance, err := svc.store.ServerInstances().Get(alert.SourceID) return err == nil && canAccessServer(user, instance) case "run-endpoint": instances, err := svc.store.ServerInstances().List(domain.ServerInstanceFilter{RunEndpointID: alert.SourceID}) if err != nil { return false } for _, instance := range instances { if canAccessServer(user, instance) { return true } } case "plugin-lifecycle": installation, err := svc.store.PluginLifecycles().Get(alert.SourceID) if err != nil { return false } instance, err := svc.store.ServerInstances().Get(installation.ServerInstanceID) return err == nil && canAccessServer(user, instance) case "ai-config-diff": preview, err := svc.store.AIConfigDiffs().Get(alert.SourceID) if err != nil { return false } instance, err := svc.store.ServerInstances().Get(preview.ServerInstanceID) return err == nil && canAccessServer(user, instance) } return false } func (svc *CoreService) pluginLifecycleDenied(actorID string, instance domain.ServerInstance, plugin domain.GamePlugin, request domain.PluginLifecycleRequest, reason string) (domain.PluginLifecycleResult, error) { svc.productionMu.Lock() defer svc.productionMu.Unlock() stamp := svc.now() installation := domain.PluginLifecycleInstallation{ID: pluginLifecycleInstallationID(plugin.ID, instance.ID), PluginID: plugin.ID, ServerInstanceID: instance.ID, TargetVersion: request.TargetVersion, DesiredState: domain.PluginLifecycleStatePending, CurrentState: domain.PluginLifecycleStateFailed, LastOperation: request.Operation, Compatibility: "incompatible", DependencyState: domain.DependencyStateUnknown, FailureReason: safeBridgeReason(reason), IdempotencyKey: request.IdempotencyKey, CreatedAt: stamp, UpdatedAt: stamp} if existing, err := svc.store.PluginLifecycles().Get(installation.ID); err == nil { installation.CreatedAt = existing.CreatedAt } return svc.pluginLifecycleDeniedLocked(actorID, installation, request, reason) } func (svc *CoreService) pluginLifecycleDeniedLocked(actorID string, installation domain.PluginLifecycleInstallation, request domain.PluginLifecycleRequest, reason string) (domain.PluginLifecycleResult, error) { stamp := svc.now() auditID, err := svc.recordAuditEventWithID(actorID, "plugin.lifecycle.denied", "plugin-lifecycle", installation.ID, domain.AuditResultDenied, reason) if err != nil { return domain.PluginLifecycleResult{}, err } installation.CurrentState = domain.PluginLifecycleStateFailed installation.LastOperation = request.Operation installation.TargetVersion = request.TargetVersion installation.FailureReason = safeBridgeReason(reason) installation.AuditEventID = auditID installation.UpdatedAt = stamp alert, err := svc.upsertAlert(domain.AlertRecord{SourceKind: "plugin-lifecycle", SourceID: installation.ID, RuleKey: "plugin.lifecycle.compatibility", Severity: domain.AlertSeverityWarning, Title: "Plugin lifecycle compatibility check failed", Message: installation.FailureReason, Retryable: true, LastAuditEventID: auditID}) if err != nil { return domain.PluginLifecycleResult{}, err } installation.AlertID = alert.ID if err := validator.ValidatePluginLifecycleInstallation(installation); err != nil { return domain.PluginLifecycleResult{}, err } if _, err := svc.store.PluginLifecycles().Get(installation.ID); errors.Is(err, repo.ErrNotFound) { err = svc.store.PluginLifecycles().Create(installation) } else if err == nil { err = svc.store.PluginLifecycles().Update(installation) } if err != nil { return domain.PluginLifecycleResult{}, err } return domain.CopyPluginLifecycleResult(domain.PluginLifecycleResult{Installation: installation, Alert: &alert, Status: "denied"}), nil } func pluginLifecycleDispatchMetadata(plugin domain.GamePlugin, operation domain.PluginLifecycleOperation) (string, string, error) { switch operation { case domain.PluginLifecycleOperationInstall, domain.PluginLifecycleOperationUpgrade, domain.PluginLifecycleOperationRollback: if plugin.LifecycleActions.Install == "" { return "", "", validationError("plugin install action is not declared") } return domain.LifecycleCapabilityInstall, plugin.LifecycleActions.Install, nil case domain.PluginLifecycleOperationEnable: if plugin.LifecycleActions.Start == "" { return "", "", validationError("plugin start action is not declared") } return domain.LifecycleCapabilityStart, plugin.LifecycleActions.Start, nil case domain.PluginLifecycleOperationDisable, domain.PluginLifecycleOperationRetire: if plugin.LifecycleActions.Stop == "" { return "", "", validationError("plugin stop action is not declared") } return domain.LifecycleCapabilityStop, plugin.LifecycleActions.Stop, nil case domain.PluginLifecycleOperationDependencyCheck: return domain.JobCapabilityDependenciesCheck, "dependencies/plugin", nil default: return "", "", validationError("plugin lifecycle operation is invalid") } } func validatePluginLifecycleTransition(installation domain.PluginLifecycleInstallation, request domain.PluginLifecycleRequest) error { switch request.Operation { case domain.PluginLifecycleOperationInstall: if installation.CurrentState == domain.PluginLifecycleStateInstalled || installation.CurrentState == domain.PluginLifecycleStateEnabled || installation.CurrentState == domain.PluginLifecycleStateDisabled { return validationError("plugin is already installed") } case domain.PluginLifecycleOperationEnable: if installation.CurrentState != domain.PluginLifecycleStateInstalled && installation.CurrentState != domain.PluginLifecycleStateDisabled { return validationError("plugin must be installed or disabled before enable") } case domain.PluginLifecycleOperationDisable: if installation.CurrentState != domain.PluginLifecycleStateEnabled { return validationError("plugin must be enabled before disable") } case domain.PluginLifecycleOperationUpgrade: if installation.CurrentVersion == "" || request.TargetVersion == installation.CurrentVersion { return validationError("upgrade requires a different target version") } case domain.PluginLifecycleOperationRollback: if installation.PreviousVersion == "" { return validationError("rollback requires a retained previous version") } case domain.PluginLifecycleOperationRetire: if installation.CurrentState == domain.PluginLifecycleStateRetired { return validationError("plugin is already retired") } } return nil } func desiredPluginLifecycleState(operation domain.PluginLifecycleOperation, installation domain.PluginLifecycleInstallation) domain.PluginLifecycleState { switch operation { case domain.PluginLifecycleOperationInstall: return domain.PluginLifecycleStateInstalled case domain.PluginLifecycleOperationEnable: return domain.PluginLifecycleStateEnabled case domain.PluginLifecycleOperationDisable: return domain.PluginLifecycleStateDisabled case domain.PluginLifecycleOperationRetire: return domain.PluginLifecycleStateRetired default: return installation.DesiredState } } func dispatchedPluginLifecycleState(operation domain.PluginLifecycleOperation, current domain.PluginLifecycleState) domain.PluginLifecycleState { switch operation { case domain.PluginLifecycleOperationUpgrade: return domain.PluginLifecycleStateUpgrading case domain.PluginLifecycleOperationRollback: return domain.PluginLifecycleStateRollingBack case domain.PluginLifecycleOperationInstall: return domain.PluginLifecycleStatePending default: return current } } func applyPluginLifecycleSuccess(installation *domain.PluginLifecycleInstallation) { switch installation.LastOperation { case domain.PluginLifecycleOperationInstall: installation.CurrentVersion = installation.TargetVersion installation.CurrentState = domain.PluginLifecycleStateInstalled case domain.PluginLifecycleOperationEnable: installation.CurrentState = domain.PluginLifecycleStateEnabled case domain.PluginLifecycleOperationDisable: installation.CurrentState = domain.PluginLifecycleStateDisabled case domain.PluginLifecycleOperationUpgrade: installation.CurrentVersion = installation.TargetVersion installation.CurrentState = installation.DesiredState if installation.CurrentState != domain.PluginLifecycleStateEnabled && installation.CurrentState != domain.PluginLifecycleStateDisabled { installation.CurrentState = domain.PluginLifecycleStateInstalled } case domain.PluginLifecycleOperationRollback: current := installation.CurrentVersion installation.CurrentVersion = installation.PreviousVersion installation.PreviousVersion = current installation.TargetVersion = installation.CurrentVersion installation.CurrentState = domain.PluginLifecycleStateInstalled case domain.PluginLifecycleOperationRetire: installation.CurrentState = domain.PluginLifecycleStateRetired case domain.PluginLifecycleOperationDependencyCheck: installation.DependencyState = domain.DependencyStatePresent } } func alertIDForSource(sourceKind, sourceID, ruleKey string) string { sum := sha256.Sum256([]byte(sourceKind + "\x00" + sourceID + "\x00" + ruleKey)) return "alert-" + hex.EncodeToString(sum[:12]) } func pluginLifecycleInstallationID(pluginID, serverInstanceID string) string { sum := sha256.Sum256([]byte(pluginID + "\x00" + serverInstanceID)) return "plugin-lifecycle-" + hex.EncodeToString(sum[:12]) } func aiConfigDiffID(requestID, serverInstanceID string) string { sum := sha256.Sum256([]byte(requestID + "\x00" + serverInstanceID)) return "ai-config-diff-" + hex.EncodeToString(sum[:12]) } func appendCapacityPressure(codes []domain.CapacityPressureCode, code domain.CapacityPressureCode) []domain.CapacityPressureCode { if !containsCapacityPressure(codes, code) { return append(codes, code) } return codes } func containsCapacityPressure(codes []domain.CapacityPressureCode, target domain.CapacityPressureCode) bool { for _, code := range codes { if code == target { return true } } return false } func firstCapacityCapability(capabilities []string) string { for _, capability := range capabilities { if strings.TrimSpace(capability) != "" { return capability } } return "control.heartbeat" } func defaultAlertNote(note, fallback string) string { if strings.TrimSpace(note) == "" { return fallback } return safeBridgeReason(note) } func maxInt(a, b int) int { if a > b { return a } return b }