Files
browser/platform/service/scum_live_data.go
T

398 lines
21 KiB
Go

package service
import (
"errors"
"strings"
"browser.local/platform/domain"
"browser.local/platform/repo"
"browser.local/platform/validator"
)
const (
scumSchemaProbeExecutionKind = "sqlite.schema-probe"
scumSQLiteTemplateExecutionKind = "sqlite.template-query"
scumRCONTemplateExecutionKind = "rcon.template-command"
scumGuardedMutationExecutionKind = "sqlite.guarded-mutation"
scumParsedLogBatchExecutionKind = "log.parsed-events"
)
func (svc *CoreService) NegotiateSCUMCapabilitiesForSession(sessionID, serverInstanceID string) (domain.SCUMCapabilityNegotiation, error) {
instance, err := svc.GetServerInstanceForSession(sessionID, serverInstanceID)
if err != nil {
return domain.SCUMCapabilityNegotiation{}, err
}
plugin, err := svc.store.GamePlugins().Get(instance.PluginID)
if err != nil {
return domain.SCUMCapabilityNegotiation{}, err
}
endpoint, err := svc.store.RunEndpoints().Get(instance.RunEndpointID)
if err != nil {
return domain.SCUMCapabilityNegotiation{}, err
}
adapterVersion := scumSchemaProbeAdapterVersion(plugin.SCUMLiveData)
active := domain.SCUMBindingIdentity{ServerInstanceID: instance.ID, RunEndpointID: instance.RunEndpointID, PluginID: plugin.ID, PluginVersion: plugin.Version, AdapterVersion: adapterVersion, DatabaseIdentity: scumSchemaProbeDatabaseIdentity(plugin.SCUMLiveData.Probe.TargetKey)}
binding, bindingErr := svc.runtimeBindingForServer(instance.ID)
if bindingErr == nil {
binding, bindingErr = normalizeRuntimeBinding(plugin, binding)
}
if bindingErr == nil {
active.RunBindingID = binding.ID
}
probeExecutorAvailable := scumRunCapabilityAvailable(endpoint, domain.SCUMDataCapabilitySchemaProbe)
negotiation := domain.SCUMCapabilityNegotiation{ServerInstanceID: instance.ID, RunEndpointID: instance.RunEndpointID, RunBindingID: active.RunBindingID, PluginID: plugin.ID, PluginVersion: plugin.Version, AdapterVersion: adapterVersion, GameVersion: active.GameVersion, DatabaseIdentity: active.DatabaseIdentity, ProbeExecutorAvailable: probeExecutorAvailable, EvaluatedAt: svc.now()}
evidenceByCapability := svc.latestSCUMCapabilityEvidenceByCapability(instance.ID)
for _, declaration := range plugin.SCUMLiveData.CapabilityGates {
gate := domain.SCUMCapabilityGate{Capability: declaration.Capability, State: domain.SCUMCapabilityGateDisabled, ReasonCode: domain.SCUMSafeErrorProbeMissing, Reason: safeSCUMGateReason(declaration.SafeReason, "current-service evidence is required before this SCUM capability can run")}
if requiredRunCapability := scumRequiredRunCapabilityForDataCapability(declaration.Capability); requiredRunCapability != "" && !containsString(endpoint.Capabilities, requiredRunCapability) {
gate.ReasonCode = domain.SCUMSafeErrorProbeExecutorAbsent
gate.Reason = "bound Run does not expose the generic executor required for this SCUM capability"
negotiation.Gates = append(negotiation.Gates, gate)
continue
}
if bindingErr != nil || binding.Status != domain.RuntimeBindingStatusComplete || binding.PluginVersion != plugin.Version {
gate.ReasonCode = domain.SCUMSafeErrorBindingMismatch
gate.Reason = "active runtime binding is missing, incomplete, or stale for this plugin version"
negotiation.Gates = append(negotiation.Gates, gate)
continue
}
if declaration.Gate != domain.SCUMCapabilityGateEnabled {
gate.ReasonCode = scumReasonCodeForEvidenceStatus(declaration.EvidenceStatus)
negotiation.Gates = append(negotiation.Gates, gate)
continue
}
requirement := domain.SCUMCapabilityRequirement{Capability: declaration.Capability, AdapterVersion: declaration.AdapterVersion, SchemaFingerprint: declaration.RequiredSchemaFingerprint, AssetDigests: domain.CopyStringSlice(declaration.RequiredAssetDigests)}
gate = domain.EvaluateSCUMCapabilityGate(requirement, evidenceByCapability[declaration.Capability], active, probeExecutorAvailable, negotiation.EvaluatedAt)
negotiation.Gates = append(negotiation.Gates, gate)
}
return negotiation, nil
}
func (svc *CoreService) latestSCUMCapabilityEvidenceByCapability(serverInstanceID string) map[domain.SCUMDataCapability]domain.SCUMCapabilityEvidence {
jobs, err := svc.store.Jobs().List(domain.JobFilter{ServerInstanceID: serverInstanceID})
if err != nil {
return map[domain.SCUMDataCapability]domain.SCUMCapabilityEvidence{}
}
evidence := map[domain.SCUMDataCapability]domain.SCUMCapabilityEvidence{}
for _, job := range jobs {
for _, candidate := range scumCapabilityEvidenceFromJob(job) {
current, exists := evidence[candidate.Capability]
if !exists || current.ObservedAt.Before(candidate.ObservedAt) {
evidence[candidate.Capability] = candidate
}
}
}
return evidence
}
func scumCapabilityEvidenceFromJob(job domain.Job) []domain.SCUMCapabilityEvidence {
var evidence []domain.SCUMCapabilityEvidence
if result := job.ExecutionResult.SQLiteSchemaProbe; result != nil {
status := domain.SCUMCapabilityEvidenceFailed
if result.Status == domain.SCUMSchemaProbeStatusSucceeded || result.Status == domain.SCUMCapabilityEvidenceCompatible {
status = domain.SCUMCapabilityEvidenceCompatible
} else if result.Status == domain.SCUMCapabilityEvidenceIncompatible {
status = domain.SCUMCapabilityEvidenceIncompatible
}
evidence = append(evidence, domain.SCUMCapabilityEvidence{Capability: domain.SCUMDataCapabilitySchemaProbe, Status: status, Binding: result.Binding, AdapterVersion: result.Binding.AdapterVersion, SchemaFingerprint: result.SchemaFingerprint, ProbeResultDigest: result.ResultDigest, ObservedAt: result.ObservedAt, SafeError: result.SafeError})
}
if result := job.ExecutionResult.SQLiteTemplate; result != nil {
evidence = append(evidence, domain.SCUMCapabilityEvidence{Capability: result.Capability, Status: scumTerminalEvidenceStatus(result.Status), Binding: result.Binding, AdapterVersion: result.AdapterVersion, SchemaFingerprint: result.SchemaFingerprint, ProbeResultDigest: result.ResultDigest, AssetDigests: []string{result.AssetDigest}, ObservedAt: result.ObservedAt, SafeError: result.SafeError})
}
if result := job.ExecutionResult.RCONTemplate; result != nil {
evidence = append(evidence, domain.SCUMCapabilityEvidence{Capability: result.Capability, Status: scumTerminalEvidenceStatus(result.Status), Binding: result.Binding, AdapterVersion: result.AdapterVersion, SchemaFingerprint: result.SchemaFingerprint, ProbeResultDigest: result.ResultDigest, AssetDigests: []string{result.AssetDigest}, ObservedAt: result.ObservedAt, SafeError: result.SafeError})
}
if result := job.ExecutionResult.GuardedMutation; result != nil {
evidence = append(evidence, domain.SCUMCapabilityEvidence{Capability: result.Capability, Status: scumTerminalEvidenceStatus(result.Status), Binding: result.Binding, AdapterVersion: result.AdapterVersion, SchemaFingerprint: result.SchemaFingerprint, ProbeResultDigest: result.ResultDigest, AssetDigests: []string{result.AssetDigest}, ObservedAt: result.ObservedAt, SafeError: result.SafeError})
}
return evidence
}
func scumTerminalEvidenceStatus(status domain.SCUMTerminalResultStatus) domain.SCUMCapabilityEvidenceStatus {
if status == domain.SCUMTerminalResultSucceeded {
return domain.SCUMCapabilityEvidenceCompatible
}
return domain.SCUMCapabilityEvidenceFailed
}
func scumRunCapabilityAvailable(endpoint domain.RunEndpoint, capability domain.SCUMDataCapability) bool {
return containsString(endpoint.Capabilities, scumRequiredRunCapabilityForDataCapability(capability))
}
func scumRequiredRunCapabilityForDataCapability(capability domain.SCUMDataCapability) string {
switch capability {
case domain.SCUMDataCapabilitySchemaProbe:
return domain.JobCapabilityRemoteRunDBSQLiteProbe
case domain.SCUMDataCapabilityPlayerRead, domain.SCUMDataCapabilityPlayerDetailRead, domain.SCUMDataCapabilitySquadRead, domain.SCUMDataCapabilitySquadMemberRead, domain.SCUMDataCapabilityVehicleRead, domain.SCUMDataCapabilityFlagRead, domain.SCUMDataCapabilityPositionRead:
return domain.JobCapabilityRemoteRunDBSQLiteQuery
case domain.SCUMDataCapabilityEconomyCommand, domain.SCUMDataCapabilityGiftCommand:
return domain.JobCapabilityRemoteRunProtectedRCON
case domain.SCUMDataCapabilityProfileXMLWrite:
return domain.JobCapabilityRemoteRunProtectedSQL
default:
return ""
}
}
func scumReasonCodeForEvidenceStatus(status domain.SCUMCapabilityEvidenceStatus) domain.SCUMSafeErrorCode {
switch status {
case domain.SCUMCapabilityEvidenceFailed:
return domain.SCUMSafeErrorProbeFailed
case domain.SCUMCapabilityEvidenceIncompatible:
return domain.SCUMSafeErrorSchemaIncompatible
default:
return domain.SCUMSafeErrorProbeMissing
}
}
func safeSCUMGateReason(value, fallback string) string {
if strings.TrimSpace(value) == "" {
return fallback
}
return value
}
func (svc *CoreService) RequestSCUMSchemaProbeForSession(sessionID, serverInstanceID, idempotencyKey string) (domain.SCUMSchemaProbeRequest, domain.RemoteAdapterResult, error) {
idempotencyKey = strings.TrimSpace(idempotencyKey)
if idempotencyKey == "" {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, validationError("idempotencyKey is required")
}
instance, err := svc.GetServerInstanceForSession(sessionID, serverInstanceID)
if err != nil {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, err
}
plugin, err := svc.store.GamePlugins().Get(instance.PluginID)
if err != nil {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, err
}
probe := plugin.SCUMLiveData.Probe
if plugin.SCUMLiveData.SchemaVersion == "" || probe.Capability != domain.JobCapabilityRemoteRunDBSQLiteProbe || strings.TrimSpace(probe.TargetKey) == "" {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, validationError("SCUM schema probe is not declared by the plugin")
}
if !scumSchemaProbeHasDataTarget(plugin.RuntimeProfiles, probe.TargetKey) {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, validationError("SCUM schema probe target is not declared as a generated Run data target")
}
endpoint, err := svc.store.RunEndpoints().Get(instance.RunEndpointID)
if err != nil {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, err
}
if endpoint.Status != domain.RunEndpointStatusOnline && endpoint.Status != domain.RunEndpointStatusDegraded {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, validationError("bound Run is not online for SCUM schema probe")
}
if !containsString(endpoint.Capabilities, domain.JobCapabilityRemoteRunDBSQLiteProbe) {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, validationError("bound Run does not expose the generic SQLite schema-probe executor")
}
binding, err := svc.runtimeBindingForServer(instance.ID)
if errors.Is(err, repo.ErrNotFound) {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, validationError("runtime binding is required before SCUM schema probe")
}
if err != nil {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, err
}
adapterVersion := scumSchemaProbeAdapterVersion(plugin.SCUMLiveData)
if adapterVersion == "" {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, validationError("SCUM schema-probe adapter version is not declared")
}
bounds := probe.Bounds
if bounds.MaxObjects == 0 {
bounds = domain.DefaultSCUMSchemaProbeBounds()
}
jobID := jobIDFromParts("job-remote-adapter", instance.ID, idempotencyKey)
request := domain.SCUMSchemaProbeRequest{
RequestID: jobID,
JobID: jobID,
Binding: domain.SCUMBindingIdentity{
ServerInstanceID: instance.ID,
RunBindingID: binding.ID,
RunEndpointID: instance.RunEndpointID,
PluginID: plugin.ID,
PluginVersion: plugin.Version,
AdapterVersion: adapterVersion,
DatabaseIdentity: scumSchemaProbeDatabaseIdentity(probe.TargetKey),
},
Bounds: bounds,
RequestedAt: svc.now(),
}
if err := validator.ValidateSCUMSchemaProbeRequest(request); err != nil {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, err
}
result, err := svc.RequestRemoteAdapterForSession(sessionID, domain.RemoteAdapterRequest{
ServerInstanceID: instance.ID,
DeclarationKey: probe.TargetKey,
TargetKey: probe.TargetKey,
Capability: domain.JobCapabilityRemoteRunDBSQLiteProbe,
TimeoutSeconds: scumSchemaProbeTimeoutSeconds(bounds),
MaxAttempts: 1,
IdempotencyKey: idempotencyKey,
PlatformScheduled: true,
SQLiteSchemaProbe: &request,
})
if err != nil {
return domain.SCUMSchemaProbeRequest{}, domain.RemoteAdapterResult{}, err
}
return request, result, nil
}
func scumSchemaProbeDatabaseIdentity(targetKey string) string {
trimmed := strings.TrimSpace(targetKey)
return strings.TrimPrefix(trimmed, "databases/")
}
func scumSchemaProbeHasDataTarget(profiles domain.GamePluginRuntimeProfiles, targetKey string) bool {
expectedWorkspaceKey := "databases/" + scumSchemaProbeDatabaseIdentity(targetKey)
for _, target := range profiles.DataTargets {
if target.Key == targetKey && target.Kind == "sqlite.snapshot" && target.WorkspaceKey == expectedWorkspaceKey && target.RefreshPolicy == "on-demand-snapshot" {
return true
}
}
return false
}
func scumSchemaProbeAdapterVersion(manifest domain.SCUMLiveDataManifest) string {
for _, gate := range manifest.CapabilityGates {
if gate.Capability == domain.SCUMDataCapabilitySchemaProbe {
return strings.TrimSpace(gate.AdapterVersion)
}
}
for _, gate := range manifest.CapabilityGates {
if strings.TrimSpace(gate.AdapterVersion) != "" {
return strings.TrimSpace(gate.AdapterVersion)
}
}
return ""
}
func scumSchemaProbeTimeoutSeconds(bounds domain.SCUMSchemaProbeBounds) int {
if bounds.TimeoutMS <= 0 {
return 1
}
seconds := (bounds.TimeoutMS + 999) / 1000
if seconds <= 0 {
return 1
}
return seconds
}
func validateSCUMSchemaProbeResultForJob(job domain.Job, result domain.SCUMSchemaProbeResult) error {
if err := validator.ValidateSCUMSchemaProbeResult(result); err != nil {
return err
}
expected := job.ExecutionInput.SQLiteSchemaProbe
if expected == nil {
return validationError("SQLite schema probe request is missing from leased job")
}
if result.JobID != job.ID || result.JobID != expected.JobID || result.RequestID != expected.RequestID {
return validationError("SQLite schema probe result does not match leased job identity")
}
if !sameSCUMSchemaProbeBinding(result.Binding, expected.Binding) {
return validationError("SQLite schema probe result does not match leased binding identity")
}
return nil
}
func validateSCUMSQLiteTemplateResultForJob(job domain.Job, result domain.SCUMSQLiteTemplateResult) error {
if err := validator.ValidateSCUMSQLiteTemplateResult(result); err != nil {
return err
}
expected := job.ExecutionInput.SQLiteTemplate
if expected == nil {
return validationError("SQLite template request is missing from leased job")
}
if result.JobID != job.ID || result.JobID != expected.JobID || result.RequestID != expected.RequestID {
return validationError("SQLite template result does not match leased job identity")
}
if !sameSCUMSchemaProbeBinding(result.Binding, expected.Binding) {
return validationError("SQLite template result does not match leased binding identity")
}
if result.Capability != expected.Capability || result.TargetKey != expected.TargetKey || result.TemplateKey != expected.TemplateKey || result.AdapterVersion != expected.AdapterVersion || result.SchemaFingerprint != expected.RequiredSchemaFingerprint || result.AssetDigest != expected.AssetDigest || result.ParameterDigest != expected.ParameterDigest {
return validationError("SQLite template result does not match leased template, adapter, digest, or parameter identity")
}
return nil
}
func validateSCUMTypedRCONTemplateResultForJob(job domain.Job, result domain.SCUMTypedRCONTemplateResult) error {
if err := validator.ValidateSCUMTypedRCONTemplateResult(result); err != nil {
return err
}
expected := job.ExecutionInput.RCONTemplate
if expected == nil {
return validationError("typed RCON template request is missing from leased job")
}
if result.JobID != job.ID || result.JobID != expected.JobID || result.RequestID != expected.RequestID {
return validationError("typed RCON template result does not match leased job identity")
}
if !sameSCUMSchemaProbeBinding(result.Binding, expected.Binding) {
return validationError("typed RCON template result does not match leased binding identity")
}
if result.Capability != expected.Capability || result.TransportKey != expected.TransportKey || result.TargetKey != expected.TargetKey || result.TemplateKey != expected.TemplateKey || result.AdapterVersion != expected.AdapterVersion || result.AssetDigest != expected.AssetDigest || result.PayloadDigest != expected.PayloadDigest || result.ConfirmationDigest != expected.ConfirmationDigest || result.TargetIdentityDigest != expected.TargetIdentityDigest {
return validationError("typed RCON template result does not match leased template, target, digest, or payload identity")
}
if expected.RequiredSchemaFingerprint != "" && result.SchemaFingerprint != expected.RequiredSchemaFingerprint {
return validationError("typed RCON template result does not match leased schema fingerprint")
}
return nil
}
func validateSCUMGuardedMutationResultForJob(job domain.Job, result domain.SCUMGuardedMutationResult) error {
if err := validator.ValidateSCUMGuardedMutationResult(result); err != nil {
return err
}
expected := job.ExecutionInput.GuardedMutation
if expected == nil {
return validationError("guarded mutation request is missing from leased job")
}
if result.JobID != job.ID || result.JobID != expected.JobID || result.RequestID != expected.RequestID {
return validationError("guarded mutation result does not match leased job identity")
}
if !sameSCUMSchemaProbeBinding(result.Binding, expected.Binding) {
return validationError("guarded mutation result does not match leased binding identity")
}
if result.Capability != expected.Capability || result.TargetKey != expected.TargetKey || result.TemplateKey != expected.TemplateKey || result.AdapterVersion != expected.AdapterVersion || result.SchemaFingerprint != expected.RequiredSchemaFingerprint || result.AssetDigest != expected.AssetDigest || result.TargetIdentityDigest != expected.TargetIdentityDigest || result.ExpectedRowDigest != expected.ExpectedRowDigest || result.ExpectedValueDigest != expected.ExpectedValueDigest || result.ExpectedXMLDigest != expected.ExpectedXMLDigest || result.PatchDigest != expected.PatchDigest || result.BackupEvidenceDigest != expected.BackupEvidenceDigest || result.OfflineEvidenceDigest != expected.OfflineEvidenceDigest || result.DangerConfirmationDigest != expected.DangerConfirmationDigest || result.ReadbackExpectationDigest != expected.ReadbackExpectationDigest {
return validationError("guarded mutation result does not match leased template, target, guard, digest, or readback identity")
}
return nil
}
func validateSCUMParsedLogBatchResultForJob(job domain.Job, result domain.SCUMParsedLogBatchResult) error {
if err := validator.ValidateSCUMParsedLogBatchResult(result); err != nil {
return err
}
expected := job.ExecutionInput.LogSource
if expected == nil {
return validationError("parsed log batch request is missing from leased job")
}
if result.JobID != job.ID {
return validationError("parsed log batch result does not match leased job identity")
}
if result.Binding.ServerInstanceID != job.ServerInstanceID || result.Binding.RunEndpointID != job.RunEndpointID {
return validationError("parsed log batch result does not match leased server or Run endpoint")
}
if result.SourceKey != expected.Key || result.StreamKey != expected.StreamKey {
return validationError("parsed log batch result does not match leased log source identity")
}
for key, value := range map[string]string{
"parserKey": result.ParserKey,
"parserVersion": result.ParserVersion,
"parserDigest": result.ParserDigest,
"adapterVersion": result.AdapterVersion,
} {
if expectedValue := strings.TrimSpace(job.ExecutionInput.Inputs[key]); expectedValue != "" && value != expectedValue {
return validationError("parsed log batch result does not match leased parser identity")
}
}
if job.ExecutionInput.PluginID != "" && result.Binding.PluginID != job.ExecutionInput.PluginID {
return validationError("parsed log batch result does not match leased plugin identity")
}
if job.ExecutionInput.TargetVersion != "" && result.Binding.PluginVersion != job.ExecutionInput.TargetVersion {
return validationError("parsed log batch result does not match leased plugin version")
}
if result.FirstCursor.SourceIdentityDigest != result.LastCursor.SourceIdentityDigest || result.FirstCursor.StreamGeneration != result.LastCursor.StreamGeneration {
return validationError("parsed log batch result crosses source identity or generation boundaries")
}
return nil
}
func sameSCUMSchemaProbeBinding(a, b domain.SCUMBindingIdentity) bool {
return a.ServerInstanceID == b.ServerInstanceID && a.RunBindingID == b.RunBindingID && a.RunEndpointID == b.RunEndpointID && a.PluginID == b.PluginID && a.PluginVersion == b.PluginVersion && a.AdapterVersion == b.AdapterVersion && a.GameVersion == b.GameVersion && a.DatabaseIdentity == b.DatabaseIdentity
}