diff --git a/platform/repo/metadata_retention.go b/platform/repo/metadata_retention.go index a3f47ed..3541c42 100644 --- a/platform/repo/metadata_retention.go +++ b/platform/repo/metadata_retention.go @@ -17,7 +17,6 @@ func normalizeStoreSnapshot(snapshot StoreSnapshot) StoreSnapshot { snapshot.PluginLifecycles = rewriteSnapshotPluginLifecycles(snapshot.PluginLifecycles, replacements) snapshot.AIConfigDiffs = rewriteSnapshotAIConfigDiffs(snapshot.AIConfigDiffs, replacements) snapshot.PluginDataRecords = retainedSnapshotPluginData(snapshot.PluginDataRecords, replacements) - snapshot.Jobs = retainedSnapshotJobs(snapshot.Jobs, snapshot.ServerInstances, plugins) return snapshot } @@ -119,38 +118,6 @@ func retainedSnapshotPluginData(values []domain.PluginDataRecord, replacements m return out } -func retainedSnapshotJobs(values []domain.Job, instances []domain.ServerInstance, plugins []domain.GamePlugin) []domain.Job { - pluginByID := map[string]domain.GamePlugin{} - for _, plugin := range plugins { - pluginByID[plugin.ID] = plugin - } - serverPlugin := map[string]domain.GamePlugin{} - for _, instance := range instances { - if plugin, ok := pluginByID[instance.PluginID]; ok { - serverPlugin[instance.ID] = plugin - } - } - out := make([]domain.Job, 0, len(values)) - for _, job := range values { - if isLegacySCUMSQLiteQueryJob(job, serverPlugin) { - continue - } - out = append(out, domain.CopyJob(job)) - } - return out -} - -func isLegacySCUMSQLiteQueryJob(job domain.Job, serverPlugin map[string]domain.GamePlugin) bool { - if job.Capability != domain.JobCapabilityRemoteRunDBSQLiteQuery { - return false - } - plugin, ok := serverPlugin[job.ServerInstanceID] - if !ok { - return false - } - return isSCUMGamePlugin(plugin) -} - func isLegacySCUMPluginData(pluginID string, collection string) bool { if !isSCUMPluginID(pluginID) { return false @@ -163,10 +130,6 @@ func isLegacySCUMPluginData(pluginID string, collection string) bool { } } -func isSCUMGamePlugin(plugin domain.GamePlugin) bool { - return strings.EqualFold(strings.TrimSpace(plugin.ServerType), "scum") || isSCUMPluginID(plugin.ID) -} - func isSCUMPluginID(pluginID string) bool { canonical := domain.CanonicalGamePluginID(pluginID) return canonical == "game.scum" || canonical == "server.scum" diff --git a/platform/repo/resources_test.go b/platform/repo/resources_test.go index 9e9d190..ca5645d 100644 --- a/platform/repo/resources_test.go +++ b/platform/repo/resources_test.go @@ -306,8 +306,8 @@ func TestFileStoreNormalizesPluginVersionsAndPrunesLegacySCUMSnapshotData(t *tes t.Fatalf("expected legacy SCUM projection data to be pruned, data=%+v err=%v", legacy, err) } jobs, err := store.Jobs().List(domain.JobFilter{ServerInstanceID: "server-scum"}) - if err != nil || len(jobs) != 1 || jobs[0].ID != "job-scum-start" { - t.Fatalf("expected legacy SCUM sqlite query job to be pruned, jobs=%+v err=%v", jobs, err) + if err != nil || len(jobs) != 2 { + t.Fatalf("expected plugin-declared SCUM sqlite query job to survive the snapshot, jobs=%+v err=%v", jobs, err) } } diff --git a/platform/service/resources.go b/platform/service/resources.go index 4069eed..7adc88c 100644 --- a/platform/service/resources.go +++ b/platform/service/resources.go @@ -336,6 +336,9 @@ func NewCoreServiceWithDurableStores(store repo.Store, logStore LogBodyStore, ar if err := service.recoverJobLogStreams(); err != nil { return nil, err } + if err := service.pruneOrphanJobLogStreams(); err != nil { + return nil, err + } if err := service.recoverLogCursors(); err != nil { return nil, err } @@ -375,6 +378,52 @@ func (svc *CoreService) recoverJobLogStreams() error { return nil } +// pruneOrphanJobLogStreams drops job-scoped log streams whose job no longer +// exists. Deleting a job removes its own streams, so this only clears leftovers +// from producers that are gone. It runs once at startup because the stream +// table is otherwise written one job at a time. +func (svc *CoreService) pruneOrphanJobLogStreams() error { + jobs, err := svc.store.Jobs().List(domain.JobFilter{}) + if err != nil { + return err + } + known := make(map[string]struct{}, len(jobs)) + for _, job := range jobs { + known[job.ID] = struct{}{} + } + streams, err := svc.store.LogStreams().List(domain.LogStreamFilter{}) + if err != nil { + return err + } + for _, stream := range streams { + jobID, ok := jobIDFromJobLogStream(stream) + if !ok || strings.HasPrefix(jobID, "autonomous-") { + continue + } + if _, exists := known[jobID]; exists { + continue + } + if err := svc.store.LogStreams().Delete(stream.ID); err != nil && !errors.Is(err, repo.ErrNotFound) { + return err + } + } + return nil +} + +func jobIDFromJobLogStream(stream domain.LogStream) (string, bool) { + streamKey := strings.TrimSpace(stream.StreamKey) + if streamKey == "" || !strings.HasPrefix(stream.ID, "job.") { + return "", false + } + body := strings.TrimPrefix(stream.ID, "job.") + suffix := "." + streamKey + if !strings.HasSuffix(body, suffix) { + return "", false + } + jobID := strings.TrimSuffix(body, suffix) + return jobID, strings.TrimSpace(jobID) != "" +} + func (svc *CoreService) recoverLogCursors() error { store, ok := svc.logStore.(interface{ LatestSeq(string) (uint64, error) }) if !ok { diff --git a/platform/service/resources_test.go b/platform/service/resources_test.go index 372d280..26743a6 100644 --- a/platform/service/resources_test.go +++ b/platform/service/resources_test.go @@ -282,6 +282,79 @@ func TestCoreServiceStartupRecoversLegacyJobLogStreams(t *testing.T) { } } +func TestCoreServiceStartupPrunesOrphanJobLogStreams(t *testing.T) { + store := repo.NewMemoryStore() + seed := newCoreService(store, func() time.Time { return fixedTime }) + plugin, endpoint := createPluginAndRunEndpoint(t, seed) + plugin.RequiredRunCapabilities = append(plugin.RequiredRunCapabilities, domain.JobCapabilityRemoteRunProgram) + if err := seed.store.GamePlugins().Update(plugin); err != nil { + t.Fatalf("update plugin capabilities: %v", err) + } + endpoint.Capabilities = append(endpoint.Capabilities, domain.JobCapabilityRemoteRunProgram) + if err := seed.store.RunEndpoints().Update(endpoint); err != nil { + t.Fatalf("update endpoint capabilities: %v", err) + } + if _, err := seed.CreateServerInstance(domain.ServerInstance{ID: "orphan-stream-server", PluginID: plugin.ID, RunEndpointID: endpoint.ID, Name: "Orphan Streams"}); err != nil { + t.Fatalf("create server instance: %v", err) + } + if err := store.Jobs().Create(domain.Job{ + ID: "live-terminal-job", + ServerInstanceID: "orphan-stream-server", + RunEndpointID: endpoint.ID, + Capability: domain.JobCapabilityRemoteRunProgram, + TargetKey: "protected-program", + InputRef: "input://protected-program/live-terminal-job", + IdempotencyKey: "live-terminal", + State: domain.JobStateQueued, + CreatedAt: fixedTime, + UpdatedAt: fixedTime, + }); err != nil { + t.Fatalf("seed live job: %v", err) + } + orphanStream := domain.LogStream{ + ID: jobLogStreamID("job-plugin-query-poll-gone", "stdout"), + ServerInstanceID: "orphan-stream-server", + Source: domain.LogStreamSourceProcess, + StreamKey: "stdout", + StorageBackend: domain.LogStorageBackendLocalSegments, + RetentionPolicy: "default", + } + autonomousStream := domain.LogStream{ + ID: jobLogStreamID("autonomous-bootstrap-start", "scum.console.stdout"), + ServerInstanceID: "orphan-stream-server", + Source: domain.LogStreamSourceProcess, + StreamKey: "scum.console.stdout", + StorageBackend: domain.LogStorageBackendLocalSegments, + RetentionPolicy: "default", + } + for _, stream := range []domain.LogStream{orphanStream, autonomousStream} { + if err := store.LogStreams().Create(stream); err != nil { + t.Fatalf("seed log stream %s: %v", stream.ID, err) + } + } + recovered, err := NewCoreServiceWithDurableStores(store, NewMemoryLogBodyStore(), NewMemoryArtifactBodyStore()) + if err != nil { + t.Fatalf("recover durable service: %v", err) + } + streams, err := recovered.ListLogStreams(domain.LogStreamFilter{ServerInstanceID: "orphan-stream-server"}) + if err != nil { + t.Fatalf("list recovered streams: %v", err) + } + byID := map[string]domain.LogStream{} + for _, stream := range streams { + byID[stream.ID] = stream + } + if _, exists := byID[orphanStream.ID]; exists { + t.Fatalf("expected orphan job log stream to be pruned, streams=%+v", byID) + } + if _, exists := byID[autonomousStream.ID]; !exists { + t.Fatalf("expected autonomous job log stream to survive, streams=%+v", byID) + } + if _, exists := byID[jobLogStreamID("live-terminal-job", "stdout")]; !exists { + t.Fatalf("expected live job log stream to be recovered, streams=%+v", byID) + } +} + func TestCoreServiceRejectsInvalidServerDependencies(t *testing.T) { svc := newTestCoreService() plugin, endpoint := createPluginAndRunEndpoint(t, svc)