993 lines
42 KiB
Go
993 lines
42 KiB
Go
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 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
|
|
}
|