Keep the recurring SCUM database poll bounded

Each dispatched poll is a durable job with two job log streams, so an
unbounded poll history would grow the platform store by roughly two
thousand jobs a day per server. A new poll now retires older terminal jobs
of the same template together with their job log streams, leaving at most
two rows per template while nothing is in flight.
This commit is contained in:
npc0-hue
2026-09-16 18:08:18 +08:00
parent 300390dc4d
commit 6985f77963
3 changed files with 120 additions and 1 deletions
+2
View File
@@ -71,6 +71,7 @@ type JobRepository interface {
GetByIdempotency(runEndpointID string, idempotencyKey string) (domain.Job, error) GetByIdempotency(runEndpointID string, idempotencyKey string) (domain.Job, error)
List(domain.JobFilter) ([]domain.Job, error) List(domain.JobFilter) ([]domain.Job, error)
Update(domain.Job) error Update(domain.Job) error
Delete(id string) error
} }
type ArtifactRepository interface { type ArtifactRepository interface {
@@ -120,6 +121,7 @@ type LogStreamRepository interface {
Get(id string) (domain.LogStream, error) Get(id string) (domain.LogStream, error)
List(domain.LogStreamFilter) ([]domain.LogStream, error) List(domain.LogStreamFilter) ([]domain.LogStream, error)
Update(domain.LogStream) error Update(domain.LogStream) error
Delete(id string) error
} }
type MetricSampleRepository interface { type MetricSampleRepository interface {
+59 -1
View File
@@ -6,6 +6,7 @@ import (
"fmt" "fmt"
"log" "log"
"math" "math"
"sort"
"strconv" "strconv"
"strings" "strings"
"time" "time"
@@ -24,6 +25,15 @@ const (
const scumQueryTemplateIdempotencyPrefix = "scum-query:" 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 // ReconcileSCUMQueryTemplatesForRunEndpoint dispatches due SCUM database read
// templates for every live server instance registered to one run endpoint. // templates for every live server instance registered to one run endpoint.
func (svc *CoreService) ReconcileSCUMQueryTemplatesForRunEndpoint(runEndpointID string) error { 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 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 return nil
} }
@@ -1,6 +1,7 @@
package service package service
import ( import (
"errors"
"testing" "testing"
"time" "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) 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)
}
}