diff --git a/platform/repo/resources.go b/platform/repo/resources.go index ef8daaa..d47c6d4 100644 --- a/platform/repo/resources.go +++ b/platform/repo/resources.go @@ -71,6 +71,7 @@ type JobRepository interface { GetByIdempotency(runEndpointID string, idempotencyKey string) (domain.Job, error) List(domain.JobFilter) ([]domain.Job, error) Update(domain.Job) error + Delete(id string) error } type ArtifactRepository interface { @@ -120,6 +121,7 @@ type LogStreamRepository interface { Get(id string) (domain.LogStream, error) List(domain.LogStreamFilter) ([]domain.LogStream, error) Update(domain.LogStream) error + Delete(id string) error } type MetricSampleRepository interface { diff --git a/platform/service/scum_query_ingest.go b/platform/service/scum_query_ingest.go index 10bea0d..e043941 100644 --- a/platform/service/scum_query_ingest.go +++ b/platform/service/scum_query_ingest.go @@ -6,6 +6,7 @@ import ( "fmt" "log" "math" + "sort" "strconv" "strings" "time" @@ -24,6 +25,15 @@ const ( const scumQueryTemplateIdempotencyPrefix = "scum-query:" +// Steady-state bounds for the recurring database poll. Every dispatched poll +// creates one durable job, so a new poll retires older terminal jobs of the +// same template together with their job log streams instead of leaving an +// unbounded history behind. +const ( + scumQueryJobRetention = 2 + scumQueryJobScanLimit = 64 +) + // ReconcileSCUMQueryTemplatesForRunEndpoint dispatches due SCUM database read // templates for every live server instance registered to one run endpoint. func (svc *CoreService) ReconcileSCUMQueryTemplatesForRunEndpoint(runEndpointID string) error { @@ -128,9 +138,57 @@ func (svc *CoreService) ReconcileSCUMQueryTemplates(serverInstanceID string) err }, }, } - if _, err := svc.CreateJob(job); err != nil { + created, err := svc.CreateJob(job) + if err != nil { return err } + svc.pruneSCUMQueryJobs(instance.ID, job.InputRef, created.ID) + } + return nil +} + +// pruneSCUMQueryJobs keeps the newest terminal jobs for one dispatched template +// and removes the rest. It runs once per poll interval, so the caller pays one +// bounded job read per dispatch and no work at all while a poll is in flight. +func (svc *CoreService) pruneSCUMQueryJobs(serverInstanceID string, inputRef string, keepJobID string) { + jobs, err := svc.store.Jobs().List(domain.JobFilter{ServerInstanceID: serverInstanceID, Limit: scumQueryJobScanLimit}) + if err != nil { + log.Printf("SCUM query job retention skipped server=%s error=%s", serverInstanceID, err.Error()) + return + } + expired := make([]domain.Job, 0, len(jobs)) + for _, candidate := range jobs { + if candidate.ID == keepJobID || candidate.InputRef != inputRef || !isTerminalJobState(candidate.State) { + continue + } + expired = append(expired, candidate) + } + retained := scumQueryJobRetention - 1 + if retained < 0 { + retained = 0 + } + if len(expired) <= retained { + return + } + sort.SliceStable(expired, func(i, j int) bool { return expired[i].CreatedAt.After(expired[j].CreatedAt) }) + for _, job := range expired[retained:] { + if err := svc.deleteJobWithLogStreams(job); err != nil { + log.Printf("SCUM query job retention failed server=%s job=%s error=%s", serverInstanceID, job.ID, err.Error()) + return + } + } +} + +func (svc *CoreService) deleteJobWithLogStreams(job domain.Job) error { + for _, streamKey := range []string{"stdout", "stderr"} { + err := svc.store.LogStreams().Delete(jobLogStreamID(job.ID, streamKey)) + if err != nil && !errors.Is(err, repo.ErrNotFound) { + return err + } + } + err := svc.store.Jobs().Delete(job.ID) + if err != nil && !errors.Is(err, repo.ErrNotFound) { + return err } return nil } diff --git a/platform/service/scum_query_ingest_test.go b/platform/service/scum_query_ingest_test.go index c76f7db..176efa6 100644 --- a/platform/service/scum_query_ingest_test.go +++ b/platform/service/scum_query_ingest_test.go @@ -1,6 +1,7 @@ package service import ( + "errors" "testing" "time" @@ -225,3 +226,61 @@ func TestSCUMQueryProjectionIgnoresUndeclaredTemplateResults(t *testing.T) { t.Fatalf("expected no projected users for an undeclared template, users=%+v err=%v", users, err) } } + +func TestSCUMQueryJobRetentionKeepsBoundedHistory(t *testing.T) { + svc, _, instance, clock := newSCUMQueryIngestFixture(t) + aligned := time.Date(2026, 9, 16, 8, 0, 0, 0, time.UTC) + *clock = aligned + refreshSCUMQueryRunHeartbeat(t, svc, aligned) + + playersJobs := func() []domain.Job { + jobs := countSCUMQueryJobs(t, svc, instance.ID) + players := make([]domain.Job, 0, len(jobs)) + for _, job := range jobs { + if job.ExecutionInput.Inputs["templateKey"] == "scum.database.players" { + players = append(players, job) + } + } + return players + } + + var firstJob domain.Job + for cycle := 0; cycle < 4; cycle++ { + if cycle > 0 { + *clock = aligned.Add(time.Duration(cycle) * 61 * time.Second) + refreshSCUMQueryRunHeartbeat(t, svc, *clock) + } + if err := svc.ReconcileSCUMQueryTemplates(instance.ID); err != nil { + t.Fatalf("reconcile cycle %d: %v", cycle, err) + } + for _, job := range playersJobs() { + if isTerminalJobState(job.State) { + continue + } + if cycle == 0 { + firstJob = job + } + job.State = domain.JobStateSucceeded + job.TerminalAt = *clock + if err := svc.store.Jobs().Update(job); err != nil { + t.Fatalf("complete query job: %v", err) + } + } + } + + retained := playersJobs() + if len(retained) != scumQueryJobRetention { + t.Fatalf("expected %d retained player query jobs, got %d", scumQueryJobRetention, len(retained)) + } + for _, job := range retained { + if !isTerminalJobState(job.State) { + t.Fatalf("retention must not remove in-flight jobs: %+v", job) + } + } + if _, err := svc.store.Jobs().Get(firstJob.ID); !errors.Is(err, repo.ErrNotFound) { + t.Fatalf("expected the oldest query job to be retired, err=%v", err) + } + if _, err := svc.store.LogStreams().Get(jobLogStreamID(firstJob.ID, "stdout")); !errors.Is(err, repo.ErrNotFound) { + t.Fatalf("expected retired query job log streams to be removed, err=%v", err) + } +}