From 137b6fd3df18efa1d03958bf392caebe731b8418 Mon Sep 17 00:00:00 2001 From: shrutiyam-glitch Date: Mon, 24 Aug 2026 13:57:11 -0700 Subject: [PATCH] atepg: deliver worker watch events via a transactional change feed (#934) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Fixes #920 partially. Replaces per-write pg_notify with a worker_changes outbox table written in the same transaction, and LISTEN with a 100ms polling watcher. #### Motivation * **Scalability Bottleneck**: `pg_notify` serializes the commits of all notifying transactions through a global lock held across the commit (including `fsync`). This artificially caps worker writes at ~600/s on Cloud SQL regardless of instance size, whereas our target is O(10K) worker updates/s. * **Payload Limits**: Bypasses the 8KB `NOTIFY` payload limit that previously caused writes to fail. * **Reliability**: A cursor-based polling watcher survives reconnects and failovers without missing events, which was a known flaw with the ephemeral `LISTEN` approach. *(Known Postgres pathology prior art: [Recall.ai](https://www.recall.ai/blog/postgres-listen-notify-does-not-scale), [DBOS](https://www.dbos.dev/blog/postgres-listen-notify-scalability)).* ### Performance Improvement WorkerUpdate @ 1,000 QPS , 1M workers (preloaded) — before vs after the change feed: | | p50 | p90 | p95 | p99 | |---|---|---|---|---| | Before (per-update pg_notify) | 40.3s |55.6s | 61.2s | 63.8s | | After (change-feed table) | 7.08 ms | 8.04 ms | 8.52 ms | 27.5 ms | #### Changes Made * **Schema**: Added transactional outbox table `worker_outbox`. * **Write Path**: Worker writes now append to the `worker_outbox` feed inside the same transaction instead of calling `pg_notify()`. * **Watch Path**: Replaced `LISTEN` in `WatchWorkers` with a polling watcher that queries the feed every 50ms. * **Cleanup**: Implemented a janitor process during polling to periodically delete old feed rows. * **Tests**: Updated atomicity tests to verify feed inserts instead of `pg_notify` payloads. For full architecture: https://docs.google.com/document/d/10K0wB6aTeFkJCL4HN3NbLJCdFGoLYdhIcqkFnqFHkKc/edit?tab=t.txqcjwvmhp3v - [x] Tests pass - [x] Appropriate changes to documentation are included in the PR --- .../internal/actoridentity/actoridentity.go | 50 +- .../actoridentity/actoridentity_test.go | 95 ++ cmd/ateapi/internal/store/atepg/atepg.go | 210 ++-- cmd/ateapi/internal/store/atepg/atepg_test.go | 102 +- cmd/ateapi/internal/store/atepg/outbox.go | 571 +++++++++ .../internal/store/atepg/outbox_test.go | 1019 +++++++++++++++++ cmd/ateapi/internal/store/atepg/schema.go | 41 + cmd/ateapi/main.go | 5 + 8 files changed, 1853 insertions(+), 240 deletions(-) create mode 100644 cmd/ateapi/internal/store/atepg/outbox.go create mode 100644 cmd/ateapi/internal/store/atepg/outbox_test.go diff --git a/cmd/ateapi/internal/actoridentity/actoridentity.go b/cmd/ateapi/internal/actoridentity/actoridentity.go index 473e79ccc..a186ca276 100644 --- a/cmd/ateapi/internal/actoridentity/actoridentity.go +++ b/cmd/ateapi/internal/actoridentity/actoridentity.go @@ -317,23 +317,55 @@ func validateWorkerRef(worker *ateapipb.ObjectRef) error { // authorizeActor resolves the actor from the authenticated worker and verifies // that the worker and actor still point at one another. Actor identity supplied // by the requester never participates in this authorization decision. +// The worker is resolved from cache first (hot path), but denials fall back +// to the authoritative store to handle watch-delivery lag right after ResumeActor. func (s *Server) authorizeActor(ctx context.Context, caller *ateletCaller, req *ateapipb.MintCertRequest) (*ateapipb.Actor, resources.ActorRef, error) { - // Denials are deliberately indistinguishable from each other: a caller that - // is not entitled to a worker should not learn its assignment. - deny := func(reason string, args ...any) error { - slog.WarnContext(ctx, "ActorIdentity denied: "+reason, - append([]any{slog.String("worker", req.GetWorker().GetName()), slog.String("callerPod", caller.podName), slog.String("callerNode", caller.nodeName)}, args...)...) - return status.Errorf(codes.PermissionDenied, "caller is not permitted to mint credentials for this actor") - } - worker, err := s.workers.Worker(req.GetWorker().GetName()) if err != nil { if errors.Is(err, store.ErrNotFound) { - return nil, resources.ActorRef{}, deny("worker not found") + // Cache misses cannot read through without the pool name; worker rows exist for the pod's lifetime. + return nil, resources.ActorRef{}, s.denyMint(ctx, caller, req, "worker not found") } slog.ErrorContext(ctx, "ActorIdentity: failed to read worker", slog.Any("err", err)) return nil, resources.ActorRef{}, status.Error(codes.Internal, "failed to look up worker") } + + actor, actorRef, err := s.authorizeWithWorker(ctx, worker, caller, req) + if status.Code(err) != codes.PermissionDenied { + return actor, actorRef, err // authorized, or a non-denial error + } + + // Read-through: re-check the authoritative worker from the store on denial. + fresh, ferr := s.store.GetWorker(ctx, req.GetWorker().GetName()) + if ferr != nil { + if !errors.Is(ferr, store.ErrNotFound) { + slog.ErrorContext(ctx, "ActorIdentity: read-through worker lookup failed", slog.Any("err", ferr)) + } + return nil, resources.ActorRef{}, err // the original denial stands + } + actor, actorRef, retryErr := s.authorizeWithWorker(ctx, fresh, caller, req) + if retryErr == nil { + slog.InfoContext(ctx, "ActorIdentity: authorized via store read-through; worker cache was stale", + slog.String("worker", req.GetWorker().GetName())) + } + return actor, actorRef, retryErr +} + +// denyMint logs the internal reason and returns a uniform PermissionDenied. +// Denials are deliberately indistinguishable from each other: a caller that +// is not entitled to a worker should not learn its assignment. +func (s *Server) denyMint(ctx context.Context, caller *ateletCaller, req *ateapipb.MintCertRequest, reason string, args ...any) error { + slog.WarnContext(ctx, "ActorIdentity denied: "+reason, + append([]any{slog.String("worker", req.GetWorker().GetName()), slog.String("callerPod", caller.podName), slog.String("callerNode", caller.nodeName)}, args...)...) + return status.Errorf(codes.PermissionDenied, "caller is not permitted to mint credentials for this actor: %s", reason) +} + +// authorizeWithWorker runs the worker↔actor mutual-pointer checks against one +// specific worker record (cached or freshly read — the caller decides). +func (s *Server) authorizeWithWorker(ctx context.Context, worker *ateapipb.Worker, caller *ateletCaller, req *ateapipb.MintCertRequest) (*ateapipb.Actor, resources.ActorRef, error) { + deny := func(reason string, args ...any) error { + return s.denyMint(ctx, caller, req, reason, args...) + } if worker.GetNodeName() != caller.nodeName { return nil, resources.ActorRef{}, deny("worker is hosted on a different node", slog.String("workerNode", worker.GetNodeName())) } diff --git a/cmd/ateapi/internal/actoridentity/actoridentity_test.go b/cmd/ateapi/internal/actoridentity/actoridentity_test.go index 171f969da..b54f38a9e 100644 --- a/cmd/ateapi/internal/actoridentity/actoridentity_test.go +++ b/cmd/ateapi/internal/actoridentity/actoridentity_test.go @@ -166,6 +166,101 @@ func newTestServer(t *testing.T, st store.Interface) *Server { return New("issuer", "", poolFile, st, workers) } +// staleWatchStore wraps a store with a WatchWorkers that never delivers, +// freezing any workercache built over it at its seed-time state — the unit +// analog of the watch's delivery latency. +type staleWatchStore struct{ store.Interface } + +func (s staleWatchStore) WatchWorkers(context.Context) (*store.WorkerWatch, error) { + return store.NewWorkerWatch(make(chan store.WorkerEvent), func() {}), nil +} + +// TestMintCertReadsThroughStaleWorkerCache pins the authorization +// read-through: an atelet minting immediately after ResumeActor committed the +// worker's assignment must be authorized from the store even though this +// replica's cache has not yet seen the assignment. The control case keeps the +// store unassigned too and must still deny — only fresh data may authorize, +// and only fresh data may deny. +func TestMintCertReadsThroughStaleWorkerCache(t *testing.T) { + for name, assignInStore := range map[string]bool{ + "assignment committed but not yet in cache: authorized via read-through": true, + "unassigned in cache and store: denial stands": false, + } { + t.Run(name, func(t *testing.T) { + ctx := context.Background() + st, cleanup := storetest.SetupTestStore(t) + defer cleanup() + + // Phase 1: worker exists, unassigned; the cache seeds this view and + // (via the inert watch) never learns anything newer. + seedActor(t, ctx, st, actorFixture{state: ateapipb.ActorState_ACTOR_STATE_RUNNING, workerNode: testNode, unassigned: true}) + workers := workercache.New(staleWatchStore{st}, time.Hour) + cacheCtx, cancel := context.WithCancel(ctx) + t.Cleanup(cancel) + if err := workers.Start(cacheCtx); err != nil { + t.Fatalf("start worker cache: %v", err) + } + + actor, err := st.GetActor(ctx, resources.ActorRef{Atespace: testAtespace, Name: testActorName}) + if err != nil { + t.Fatalf("read seeded actor: %v", err) + } + if assignInStore { + // Phase 2: commit the assignment to the store only, as + // AssignWorker does (possibly on another replica). + worker, err := st.GetWorker(ctx, testWorkerName) + if err != nil { + t.Fatalf("read seeded worker: %v", err) + } + if worker.Status == nil { + worker.Status = &ateapipb.WorkerStatus{} + } + worker.Status.Assignment = &ateapipb.ActorAssignment{ + Actor: (resources.ActorRef{Atespace: testAtespace, Name: testActorName}).ToObjectRef(), + ActorUid: actor.GetMetadata().GetUid(), + } + if err := st.UpdateWorker(ctx, worker, worker.GetMetadata().GetVersion()); err != nil { + t.Fatalf("assign worker in store: %v", err) + } + } + + srv := newTestServerWithCache(t, st, workers) + resp, err := srv.MintCert(ctxWithCert(ateletCertOn(t, testNode)), mintCertRequest(t, actor.GetMetadata().GetUid())) + + wantCode := codes.PermissionDenied + if assignInStore { + wantCode = codes.OK + } + if got := status.Code(err); got != wantCode { + t.Fatalf("MintCert() code = %v (err = %v), want %v", got, err, wantCode) + } + if assignInStore && len(resp.GetActorCertificates()) == 0 { + t.Fatal("MintCert() returned no certificates") + } + }) + } +} + +// newTestServerWithCache is newTestServer with a caller-controlled worker +// cache (e.g. one frozen at a stale state). +func newTestServerWithCache(t *testing.T, st store.Interface, workers *workercache.Cache) *Server { + t.Helper() + + ca, err := localca.GenerateED25519CA("test-actor-ca") + if err != nil { + t.Fatalf("generate CA: %v", err) + } + poolBytes, err := localca.Marshal(&localca.Pool{CAs: []*localca.CA{ca}}) + if err != nil { + t.Fatalf("marshal CA pool: %v", err) + } + poolFile := filepath.Join(t.TempDir(), "actor-ca-pool.json") + if err := os.WriteFile(poolFile, poolBytes, 0o600); err != nil { + t.Fatalf("write CA pool: %v", err) + } + return New("issuer", "", poolFile, st, workers) +} + func TestMintJWTRequiresConfiguredJWTProvider(t *testing.T) { srv := &Server{actorIdentityJWTIssuer: "https://kubernetes.example"} for _, tt := range []struct { diff --git a/cmd/ateapi/internal/store/atepg/atepg.go b/cmd/ateapi/internal/store/atepg/atepg.go index dcccfb54f..80b094c52 100644 --- a/cmd/ateapi/internal/store/atepg/atepg.go +++ b/cmd/ateapi/internal/store/atepg/atepg.go @@ -23,7 +23,6 @@ package atepg import ( "context" - "encoding/json" "errors" "fmt" "log/slog" @@ -36,21 +35,37 @@ import ( "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgconn" "github.com/jackc/pgx/v5/pgxpool" - "google.golang.org/protobuf/encoding/protojson" "google.golang.org/protobuf/proto" "google.golang.org/protobuf/types/known/timestamppb" ) // Persistence is a service that stores ate state in PostgreSQL. +// watchPoolMaxConns sizes the dedicated outbox watch pool: one connection +// for the WatchWorkers poller, one for the maintenance loop, and one of headroom +// so a transiently slow poll can never gate a maintenance pass. +const ( + watchPoolMaxConns = 3 + watchPoolMinConns = 1 +) + type Persistence struct { - pool *pgxpool.Pool - lockTTL time.Duration + pool *pgxpool.Pool + // watchPool serves the outbox side only: the WatchWorkers pollers + // and the partition-maintenance loop. + watchPool *pgxpool.Pool + ownsWatchPool bool + lockTTL time.Duration + pollFailureCloseAfter time.Duration + stopMaintenance context.CancelFunc + maintenanceDone chan struct{} } var _ store.Interface = (*Persistence)(nil) // Connect opens a pgxpool against dsn, verifies connectivity, and applies the -// embedded schema. Startup fails if the database cannot be reached. +// embedded schema. Startup fails if the database cannot be reached. A second, +// two-connection watch pool (owned by the Persistence, closed by Close) isolates +// outbox polling and maintenance from write traffic. func Connect(ctx context.Context, dsn string) (*Persistence, error) { pool, err := pgxpool.New(ctx, dsn) if err != nil { @@ -60,22 +75,71 @@ func Connect(ctx context.Context, dsn string) (*Persistence, error) { pool.Close() return nil, fmt.Errorf("pinging PostgreSQL: %w", err) } - p, err := NewPersistence(ctx, pool) + + watchCfg, err := pgxpool.ParseConfig(dsn) if err != nil { + pool.Close() + return nil, fmt.Errorf("parsing watch pool config: %w", err) + } + watchCfg.MaxConns = watchPoolMaxConns + watchCfg.MinConns = watchPoolMinConns + watchPool, err := pgxpool.NewWithConfig(ctx, watchCfg) + if err != nil { + pool.Close() + return nil, fmt.Errorf("opening PostgreSQL watch pool: %w", err) + } + + p, err := newPersistence(ctx, pool, watchPool) + if err != nil { + watchPool.Close() pool.Close() return nil, err } + p.ownsWatchPool = true return p, nil } // NewPersistence wraps an already-open pool, applying the idempotent schema. -// Callers that already hold a pool (e.g. tests using -// testcontainers) use this directly instead of Connect. +// Callers that already hold a pool (e.g. tests using testcontainers) use +// this directly instead of Connect; outbox watch traffic shares the given pool. func NewPersistence(ctx context.Context, pool *pgxpool.Pool) (*Persistence, error) { + return newPersistence(ctx, pool, pool) +} + +func newPersistence(ctx context.Context, pool, watchPool *pgxpool.Pool) (*Persistence, error) { if err := applySchema(ctx, pool); err != nil { return nil, err } - return &Persistence{pool: pool, lockTTL: defaultLockTTL}, nil + maintenanceCtx, stopMaintenance := context.WithCancel(context.Background()) + p := &Persistence{pool: pool, watchPool: watchPool, lockTTL: defaultLockTTL, pollFailureCloseAfter: outboxPollFailureCloseAfter, stopMaintenance: stopMaintenance, maintenanceDone: make(chan struct{})} + // Cover the partition lead before accepting writes; from then on the + // maintenance loop keeps partitions ahead of the clock (and the + // DEFAULT partition catches writes if it ever falls behind). + bootNow, err := p.outboxNow(ctx) + if err != nil { + stopMaintenance() + return nil, err + } + if err := p.createWorkerOutboxPartitions(ctx, outboxPartitionLeadTimes(bootNow)...); err != nil { + stopMaintenance() + return nil, err + } + go func() { + defer close(p.maintenanceDone) + p.outboxMaintenance(maintenanceCtx) + }() + return p, nil +} + +// Close stops the outbox maintenance loop and waits for it to exit, +// then closes the watch pool if Connect created one. It does not close the +// main pool, which the caller owns. +func (p *Persistence) Close() { + p.stopMaintenance() + <-p.maintenanceDone + if p.ownsWatchPool { + p.watchPool.Close() + } } // querier is satisfied by both *pgxpool.Pool and pgx.Tx, letting read helpers @@ -1074,77 +1138,6 @@ func (p *Persistence) DeleteActorSnapshotTag(ctx context.Context, atespace, name // --- Workers --- -const ( - // workerChangeChannel is the fixed LISTEN/NOTIFY channel for worker changes. - workerChangeChannel = "worker_changes" - // maxNotifyPayloadBytes reflects PostgreSQL's NOTIFY payload size limit. - // Writes fail rather than silently omit a notification if exceeded. - maxNotifyPayloadBytes = 8000 -) - -type workerEventEnvelope struct { - Type int `json:"t"` - Worker string `json:"w"` // protojson-encoded Worker -} - -func marshalWorkerEvent(eventType store.WorkerEventType, worker *ateapipb.Worker) ([]byte, error) { - workerJSON, err := protojson.Marshal(worker) - if err != nil { - return nil, fmt.Errorf("in protojson.Marshal: %w", err) - } - msg, err := json.Marshal(workerEventEnvelope{Type: int(eventType), Worker: string(workerJSON)}) - if err != nil { - return nil, fmt.Errorf("in json.Marshal: %w", err) - } - return msg, nil -} - -func unmarshalWorkerEvent(payload string) (store.WorkerEvent, error) { - var env workerEventEnvelope - if err := json.Unmarshal([]byte(payload), &env); err != nil { - return store.WorkerEvent{}, fmt.Errorf("in json.Unmarshal: %w", err) - } - worker := &ateapipb.Worker{} - if err := protojson.Unmarshal([]byte(env.Worker), worker); err != nil { - return store.WorkerEvent{}, fmt.Errorf("in protojson.Unmarshal: %w", err) - } - return store.WorkerEvent{Type: store.WorkerEventType(env.Type), Worker: worker}, nil -} - -// writeAndNotify runs fn inside a transaction, then--only if fn reports a -// change worth notifying--calls pg_notify in the same transaction so -// delivery happens if and only if the transaction commits. -func (p *Persistence) writeAndNotify(ctx context.Context, eventType store.WorkerEventType, worker *ateapipb.Worker, fn func(ctx context.Context, tx pgx.Tx) (notify bool, err error)) error { - tx, err := p.pool.Begin(ctx) - if err != nil { - return fmt.Errorf("beginning transaction: %w", err) - } - defer tx.Rollback(ctx) //nolint:errcheck // no-op once committed - - notify, err := fn(ctx, tx) - if err != nil { - return err - } - - if notify { - payload, err := marshalWorkerEvent(eventType, worker) - if err != nil { - return fmt.Errorf("marshaling worker event: %w", err) - } - if len(payload) > maxNotifyPayloadBytes { - return fmt.Errorf("worker event payload of %d bytes exceeds PostgreSQL NOTIFY limit of %d bytes", len(payload), maxNotifyPayloadBytes) - } - if _, err := tx.Exec(ctx, `SELECT pg_notify($1, $2)`, workerChangeChannel, string(payload)); err != nil { - return fmt.Errorf("notifying worker change: %w", err) - } - } - - if err := tx.Commit(ctx); err != nil { - return fmt.Errorf("committing transaction: %w", err) - } - return nil -} - func (p *Persistence) CreateWorker(ctx context.Context, worker *ateapipb.Worker) error { dbWorker := proto.Clone(worker).(*ateapipb.Worker) // Workers are global-scoped, so the atespace is always empty. @@ -1155,7 +1148,7 @@ func (p *Persistence) CreateWorker(ctx context.Context, worker *ateapipb.Worker) return fmt.Errorf("marshaling worker: %w", err) } - err = p.writeAndNotify(ctx, store.WorkerEventCreated, dbWorker, func(ctx context.Context, tx pgx.Tx) (bool, error) { + err = p.writeAndAppendEvent(ctx, store.WorkerEventCreated, dbWorker, func(ctx context.Context, tx pgx.Tx) (bool, error) { _, err := tx.Exec(ctx, ` INSERT INTO workers (name, uid, version, proto) VALUES ($1, $2, $3, $4)`, @@ -1206,7 +1199,7 @@ func (p *Persistence) UpdateWorker(ctx context.Context, worker *ateapipb.Worker, return fmt.Errorf("marshaling worker: %w", err) } - return p.writeAndNotify(ctx, store.WorkerEventUpdated, dbWorker, func(ctx context.Context, tx pgx.Tx) (bool, error) { + return p.writeAndAppendEvent(ctx, store.WorkerEventUpdated, dbWorker, func(ctx context.Context, tx pgx.Tx) (bool, error) { var returned []byte err := tx.QueryRow(ctx, ` UPDATE workers @@ -1235,7 +1228,7 @@ func (p *Persistence) UpdateWorker(ctx context.Context, worker *ateapipb.Worker, func (p *Persistence) DeleteWorker(ctx context.Context, name string) error { deletedEvent := &ateapipb.Worker{Metadata: &ateapipb.ResourceMetadata{Name: name}} - return p.writeAndNotify(ctx, store.WorkerEventDeleted, deletedEvent, func(ctx context.Context, tx pgx.Tx) (bool, error) { + return p.writeAndAppendEvent(ctx, store.WorkerEventDeleted, deletedEvent, func(ctx context.Context, tx pgx.Tx) (bool, error) { var protoBytes []byte err := tx.QueryRow(ctx, ` DELETE FROM workers @@ -1243,7 +1236,7 @@ func (p *Persistence) DeleteWorker(ctx context.Context, name string) error { RETURNING proto`, name).Scan(&protoBytes) if err != nil { if errors.Is(err, pgx.ErrNoRows) { - // Idempotent: nothing existed, so nothing to notify either. + // Idempotent: nothing existed, so no event to publish either. return false, nil } return false, fmt.Errorf("deleting worker %s: %w", name, err) @@ -1304,53 +1297,6 @@ func (p *Persistence) ListWorkers(ctx context.Context, opts store.ListOptions) ( return store.ListResponse[*ateapipb.Worker]{Items: result, NextPageToken: nextToken}, nil } -// WatchWorkers acquires a dedicated connection (hijacked out of the pool, so -// it's never handed back for unrelated queries), LISTENs on the fixed -// worker-change channel, and forwards decoded notifications until the -// context is cancelled or the caller closes the watch. -func (p *Persistence) WatchWorkers(ctx context.Context) (*store.WorkerWatch, error) { - watchCtx, cancel := context.WithCancel(ctx) - - poolConn, err := p.pool.Acquire(watchCtx) - if err != nil { - cancel() - return nil, fmt.Errorf("acquiring watch connection: %w", err) - } - conn := poolConn.Hijack() - - if _, err := conn.Exec(watchCtx, "LISTEN "+workerChangeChannel); err != nil { - conn.Close(watchCtx) //nolint:errcheck - cancel() - return nil, fmt.Errorf("listening for worker changes: %w", err) - } - - ch := make(chan store.WorkerEvent, 128) - go func() { - defer close(ch) - defer conn.Close(context.Background()) //nolint:errcheck - for { - notification, err := conn.WaitForNotification(watchCtx) - if err != nil { - // Context cancelled (caller closed the watch) or the - // connection was lost. Either way, the caller must - // re-subscribe; matches ateredis's WatchWorkers contract. - return - } - event, err := unmarshalWorkerEvent(notification.Payload) - if err != nil { - slog.ErrorContext(ctx, "worker event unmarshal failed; closing watch", slog.Any("err", err)) - return - } - select { - case ch <- event: - case <-watchCtx.Done(): - return - } - } - }() - return store.NewWorkerWatch(ch, cancel), nil -} - // --- Workflow locks --- // defaultLockTTL is how long a lock may go unrenewed before another client @@ -1507,7 +1453,7 @@ func (p *Persistence) releaseLease(ctx context.Context, key, token string) error // --- Debug --- func (p *Persistence) DebugClearAll(ctx context.Context) error { - if _, err := p.pool.Exec(ctx, `TRUNCATE atespaces, actors, actor_templates, actor_snapshots, actor_snapshot_tags, workers, leases`); err != nil { + if _, err := p.pool.Exec(ctx, `TRUNCATE atespaces, actors, actor_templates, actor_snapshots, actor_snapshot_tags, workers, leases, worker_outbox, worker_outbox_trim`); err != nil { return fmt.Errorf("truncating tables: %w", err) } return nil diff --git a/cmd/ateapi/internal/store/atepg/atepg_test.go b/cmd/ateapi/internal/store/atepg/atepg_test.go index 88dffeba1..607db78d0 100644 --- a/cmd/ateapi/internal/store/atepg/atepg_test.go +++ b/cmd/ateapi/internal/store/atepg/atepg_test.go @@ -26,7 +26,6 @@ import ( "github.com/google/go-cmp/cmp" "github.com/jackc/pgx/v5/pgxpool" "github.com/testcontainers/testcontainers-go/modules/postgres" - "google.golang.org/protobuf/proto" "google.golang.org/protobuf/testing/protocmp" "github.com/agent-substrate/substrate/cmd/ateapi/internal/store" @@ -40,6 +39,7 @@ import ( var ( containerOnce sync.Once containerPool *pgxpool.Pool + containerDSN string containerPG *postgres.PostgresContainer containerErr error ) @@ -79,6 +79,7 @@ func requirePool(t *testing.T) *pgxpool.Pool { containerErr = err return } + containerDSN = dsn pool, err := pgxpool.New(ctx, dsn) if err != nil { containerErr = err @@ -115,6 +116,7 @@ func setupPostgresPersistence(t *testing.T) *Persistence { if err != nil { t.Fatalf("NewPersistence failed: %v", err) } + t.Cleanup(p.Close) if err := p.DebugClearAll(ctx); err != nil { t.Fatalf("DebugClearAll failed: %v", err) } @@ -354,104 +356,6 @@ func TestCreateActor_MissingAtespace_FailedPrecondition(t *testing.T) { } } -// TestWorkerNotification_OnlyAfterCommit proves the doc's atomicity claim: a -// worker write's pg_notify shares the write's transaction, so a rolled-back -// write never notifies, while a committed write always does. -func TestWorkerNotification_OnlyAfterCommit(t *testing.T) { - s := setupPostgresStore(t).(*Persistence) - ctx := context.Background() - - watch, err := s.WatchWorkers(ctx) - if err != nil { - t.Fatalf("WatchWorkers failed: %v", err) - } - defer watch.Close() - - const workerName = "6e4d2f81-b3a9-4c05-8e72-1f9d4a0c7b63" - worker := &ateapipb.Worker{ - Metadata: &ateapipb.ResourceMetadata{Name: workerName}, - WorkerNamespace: "ns", - WorkerPool: "pool", - WorkerPod: "pod", - WorkerPodUid: workerName, - } - protoBytes, err := proto.Marshal(worker) - if err != nil { - t.Fatalf("marshaling worker: %v", err) - } - - // Write the row and roll back instead of committing: no notification - // should ever arrive, proving pg_notify's effect is undone with the rest - // of the transaction. - tx, err := s.pool.Begin(ctx) - if err != nil { - t.Fatalf("Begin failed: %v", err) - } - if _, err := tx.Exec(ctx, ` - INSERT INTO workers (name, uid, version, proto) - VALUES ($1, $2, $3, $4)`, - workerName, "rolled-back-uid", int64(1), protoBytes); err != nil { - t.Fatalf("insert failed: %v", err) - } - if _, err := tx.Exec(ctx, `SELECT pg_notify($1, $2)`, workerChangeChannel, "rolled-back-payload"); err != nil { - t.Fatalf("pg_notify failed: %v", err) - } - if err := tx.Rollback(ctx); err != nil { - t.Fatalf("Rollback failed: %v", err) - } - - select { - case event := <-watch.Events: - t.Fatalf("received event %+v from a rolled-back transaction; NOTIFY should not survive rollback", event) - case <-time.After(500 * time.Millisecond): - // Expected: nothing arrives. - } - - // The equivalent committed write must notify. - if err := s.CreateWorker(ctx, worker); err != nil { - t.Fatalf("CreateWorker failed: %v", err) - } - select { - case event := <-watch.Events: - if event.Type != store.WorkerEventCreated { - t.Errorf("expected WorkerEventCreated, got %v", event.Type) - } - // CreateWorker assigns the uid, version and timestamps server-side. - want := proto.Clone(worker).(*ateapipb.Worker) - want.Metadata.Version = 1 - if diff := cmp.Diff(want, event.Worker, protocmp.Transform(), - protocmp.IgnoreFields(&ateapipb.ResourceMetadata{}, "uid", "create_time", "update_time")); diff != "" { - t.Errorf("event worker mismatch (-want +got):\n%s", diff) - } - case <-time.After(2 * time.Second): - t.Fatal("timed out waiting for event from a committed write") - } -} - -func TestWatchWorkers_MalformedNotificationClosesWatch(t *testing.T) { - s := setupPostgresStore(t).(*Persistence) - ctx := context.Background() - - watch, err := s.WatchWorkers(ctx) - if err != nil { - t.Fatalf("WatchWorkers failed: %v", err) - } - defer watch.Close() - - if _, err := s.pool.Exec(ctx, `SELECT pg_notify($1, $2)`, workerChangeChannel, "not-json"); err != nil { - t.Fatalf("pg_notify failed: %v", err) - } - - select { - case event, ok := <-watch.Events: - if ok { - t.Fatalf("received event %+v from malformed notification; want closed watch", event) - } - case <-time.After(2 * time.Second): - t.Fatal("timed out waiting for malformed notification to close watch") - } -} - func TestListActors_InvalidPageToken(t *testing.T) { s := setupPostgresStore(t).(*Persistence) ctx := context.Background() diff --git a/cmd/ateapi/internal/store/atepg/outbox.go b/cmd/ateapi/internal/store/atepg/outbox.go new file mode 100644 index 000000000..89d2d27f2 --- /dev/null +++ b/cmd/ateapi/internal/store/atepg/outbox.go @@ -0,0 +1,571 @@ +// Copyright 2026 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +// The worker outbox: worker writes append one event row to the +// range-partitioned, UNLOGGED worker_outbox table in the same transaction +// (writeAndAppendEvent), per-replica watchers poll it with an xmin-fenced +// xid cursor (WatchWorkers), and a background loop pre-creates partitions +// and retires old ones by dropping them, recording a trim high-water mark +// that lets lagging watchers detect loss and resync (outboxMaintenance). + +package atepg + +import ( + "context" + "fmt" + "log/slog" + "strings" + "time" + + "github.com/agent-substrate/substrate/cmd/ateapi/internal/store" + "github.com/agent-substrate/substrate/pkg/proto/ateapipb" + "github.com/jackc/pgx/v5" + "google.golang.org/protobuf/proto" +) + +// Outbox payload format: one event-type byte followed by the binary Worker proto. +// The tag byte is read by other replicas during rolling deploys, so +// store.WorkerEventType values must stay append-only stable and fit a byte. +func marshalWorkerEvent(eventType store.WorkerEventType, worker *ateapipb.Worker) ([]byte, error) { + b, err := proto.Marshal(worker) + if err != nil { + return nil, fmt.Errorf("in proto.Marshal: %w", err) + } + return append([]byte{byte(eventType)}, b...), nil +} + +func unmarshalWorkerEvent(payload []byte) (store.WorkerEvent, error) { + if len(payload) == 0 { + return store.WorkerEvent{}, fmt.Errorf("empty worker event payload") + } + // Assert invariants at the boundary. Corrupted payloads or unknown types + // must fail here to trigger a loud resync, rather than falling through + // downstream as silent no-ops. + eventType := store.WorkerEventType(payload[0]) + switch eventType { + case store.WorkerEventCreated, store.WorkerEventUpdated, store.WorkerEventDeleted: + default: + return store.WorkerEvent{}, fmt.Errorf("unknown worker event type byte %d", payload[0]) + } + worker := &ateapipb.Worker{} + if err := proto.Unmarshal(payload[1:], worker); err != nil { + return store.WorkerEvent{}, fmt.Errorf("in proto.Unmarshal: %w", err) + } + if worker.GetMetadata().GetName() == "" { + return store.WorkerEvent{}, fmt.Errorf("worker event payload has no worker name") + } + return store.WorkerEvent{Type: eventType, Worker: worker}, nil +} + +// writeAndAppendEvent runs fn inside a transaction, then--only if fn +// reports a change worth publishing--appends the event to the worker_outbox +// table in the same transaction, so watchers see it if and only if the +// transaction commits. +func (p *Persistence) writeAndAppendEvent(ctx context.Context, eventType store.WorkerEventType, worker *ateapipb.Worker, fn func(ctx context.Context, tx pgx.Tx) (changed bool, err error)) error { + tx, err := p.pool.Begin(ctx) + if err != nil { + return fmt.Errorf("beginning transaction: %w", err) + } + defer tx.Rollback(ctx) //nolint:errcheck // no-op once committed + + changed, err := fn(ctx, tx) + if err != nil { + return err + } + + if changed { + payload, err := marshalWorkerEvent(eventType, worker) + if err != nil { + return fmt.Errorf("marshaling worker event: %w", err) + } + if _, err := tx.Exec(ctx, `INSERT INTO worker_outbox (payload) VALUES ($1)`, payload); err != nil { + return fmt.Errorf("appending worker outbox: %w", err) + } + } + + if err := tx.Commit(ctx); err != nil { + return fmt.Errorf("committing transaction: %w", err) + } + return nil +} + +const ( + // Bound worker-event delivery latency in the absence of an xmin stall. + outboxPollInterval = 50 * time.Millisecond + + // Cap rows fetched per poll; a burst beyond it carries over to the next poll + // (events are delayed, never dropped). + outboxBatch = 1024 + + // Minimum time retention keeps outbox rows. + outboxRetentionAge = 15 * time.Minute + + // Paces partition maintenance. + outboxMaintenanceInterval = time.Minute + + // The outbox partition range width. + outboxPartitionInterval = 15 * time.Minute + + // Bounds a maintenance pass to prevent indefinite hangs (e.g., from lock waits) + // which would permanently starve partition creation. Stalls abort and retry. + outboxMaintenancePassTimeout = 5 * time.Minute + + // How many intervals ahead partitions are pre-created: creation must stall past + // lead-1 intervals before any write detours into the DEFAULT partition backstop. + outboxPartitionLead = 2 +) + +// Bounds stale-serving during polling outages: after this duration of +// uninterrupted failures, the watch closes and forces a full cache relist. +// Balances riding out transient blips vs. failing fast on real outages. +// Configured per-Persistence-instance to prevent data races in tests. +const outboxPollFailureCloseAfter = 30 * time.Second + +// outboxNow returns the database's clock_timestamp — the clock rows route +// by, and therefore the one partition bounds and expiry must use. +func (p *Persistence) outboxNow(ctx context.Context) (time.Time, error) { + var now time.Time + if err := p.watchPool.QueryRow(ctx, `SELECT clock_timestamp()`).Scan(&now); err != nil { + return time.Time{}, fmt.Errorf("reading database clock: %w", err) + } + return now.UTC(), nil +} + +// Maintains worker_outbox partitions on a fixed timer. +func (p *Persistence) outboxMaintenance(ctx context.Context) { + ticker := time.NewTicker(outboxMaintenanceInterval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + } + passCtx, cancel := context.WithTimeout(ctx, outboxMaintenancePassTimeout) + if err := p.maintainWorkerOutboxPartitions(passCtx); err != nil && ctx.Err() == nil { + slog.WarnContext(ctx, "worker outbox maintenance failed", slog.Any("err", err)) + } + cancel() + } +} + +// Database-scoped advisory lock used to elect a single replica to run the retention transaction (drops + trim). +const outboxMaintenanceLockKey = "atepg-outbox-maintenance" + +// pollWorkerOutboxSQL is the watch's batch query. The xid::text cast MUST +// carry an alias so that ORDER BY xid sorts numerically by the table's xid8 +// column instead of alphabetically by the string output. +const pollWorkerOutboxSQL = ` + SELECT xid::text AS xid_text, payload FROM worker_outbox + WHERE xid > $1::xid8 + AND xid < pg_snapshot_xmin(pg_current_snapshot()) + ORDER BY xid LIMIT $2` + +// pollSafetySQL returns cheap safety scalars fetched on every poll: +// 1. A fell-behind check (trim mark is past both cursor and baseline). +// 2. The postmaster start time (to detect database restarts). +const pollSafetySQL = ` + SELECT EXISTS( + SELECT 1 FROM worker_outbox_trim + WHERE xid > $1::xid8 AND xid > $2::xid8), + pg_postmaster_start_time()::text` + +// maintainWorkerOutboxPartitions runs one maintenance pass. Partition creation +// is unelected. DEFAULT truncate and partition drops run in SEPARATE elected +// transactions to prevent AB/BA deadlocks against writers. +func (p *Persistence) maintainWorkerOutboxPartitions(ctx context.Context) error { + // Must use the database's clock. App-sourced time would let a fast-clocked + // replica accidentally drift partition bounds and shorten retention. + now, err := p.outboxNow(ctx) + if err != nil { + return err + } + if err := p.createWorkerOutboxPartitions(ctx, outboxPartitionLeadTimes(now)...); err != nil { + return err + } + + if err := p.retireStrayedOutboxDefault(ctx); err != nil { + return err + } + return p.dropExpiredOutboxRetention(ctx, now) +} + +// retireStrayedOutboxDefault truncates a non-empty DEFAULT partition in its +// own elected transaction. Locks touched: DEFAULT child only — never the +// parent (see maintainWorkerOutboxPartitions on deadlock ordering). +func (p *Persistence) retireStrayedOutboxDefault(ctx context.Context) error { + tx, err := p.watchPool.Begin(ctx) + if err != nil { + return fmt.Errorf("beginning outbox stray-cleanup transaction: %w", err) + } + defer tx.Rollback(ctx) //nolint:errcheck // no-op once committed + + var elected bool + if err := tx.QueryRow(ctx, `SELECT pg_try_advisory_xact_lock(hashtext(current_database() || ':' || $1))`, outboxMaintenanceLockKey).Scan(&elected); err != nil { + return fmt.Errorf("electing outbox maintenance: %w", err) + } + if !elected { + return nil // another replica is maintaining; next tick retries + } + // A non-empty DEFAULT partition means partition creation stalled and writes + // detoured here. Watchers that lose events will detect the trim mark and resync. + var strays bool + if err := tx.QueryRow(ctx, `SELECT EXISTS(SELECT 1 FROM worker_outbox_default)`).Scan(&strays); err != nil { + return fmt.Errorf("checking outbox default partition: %w", err) + } + if !strays { + // Nothing to clean; end the election transaction explicitly rather + // than leaning on the deferred rollback for a read-only exit. + if err := tx.Commit(ctx); err != nil { + return fmt.Errorf("committing outbox stray-cleanup transaction: %w", err) + } + return nil + } + slog.WarnContext(ctx, "outbox DEFAULT partition is non-empty; partition creation has stalled and writes are detouring") + if err := p.truncateWorkerOutboxDefault(ctx, tx); err != nil { + return err + } + if err := tx.Commit(ctx); err != nil { + return fmt.Errorf("committing outbox stray-cleanup transaction: %w", err) + } + return nil +} + +// dropExpiredOutboxRetention drops aged-out partitions in its own elected +// transaction. Locks touched: parent (ACCESS EXCLUSIVE, blocking every +// worker write's outbox append) plus the dropped children — never the +// DEFAULT while waiting on the parent. +func (p *Persistence) dropExpiredOutboxRetention(ctx context.Context, now time.Time) error { + tx, err := p.watchPool.Begin(ctx) + if err != nil { + return fmt.Errorf("beginning outbox retention transaction: %w", err) + } + defer tx.Rollback(ctx) //nolint:errcheck // no-op once committed + + var elected bool + if err := tx.QueryRow(ctx, `SELECT pg_try_advisory_xact_lock(hashtext(current_database() || ':' || $1))`, outboxMaintenanceLockKey).Scan(&elected); err != nil { + return fmt.Errorf("electing outbox maintenance: %w", err) + } + if !elected { + return nil // another replica is maintaining; next tick retries + } + // Drops run last in the pass with commit immediately after, keeping the + // writer-blocking window minimal. + if err := p.dropExpiredWorkerOutboxPartitions(ctx, tx, now); err != nil { + return err + } + if err := tx.Commit(ctx); err != nil { + return fmt.Errorf("committing outbox retention transaction: %w", err) + } + return nil +} + +// workerOutboxPartitionName names the partition covering the given instant. +func workerOutboxPartitionName(at time.Time) string { + // it truncates to the partition boundary itself so callers can pass any moment within the range. + return "worker_outbox_p" + at.UTC().Truncate(outboxPartitionInterval).Format("200601021504") +} + +// outboxPartitionLeadTimes lists instants covering now through the creation lead, one per partition interval. +func outboxPartitionLeadTimes(now time.Time) []time.Time { + times := make([]time.Time, outboxPartitionLead+1) + for i := range times { + times[i] = now.UTC().Add(time.Duration(i) * outboxPartitionInterval) + } + return times +} + +// createWorkerOutboxPartitions idempotently creates partitions. +// If writes spilled into DEFAULT, CREATE PARTITION fails (23514). We catch +// this, TRUNCATE the strays (triggering watcher resyncs), and run CREATE +// PARTITION inside the SAME transaction to safely un-wedge the system. +func (p *Persistence) createWorkerOutboxPartitions(ctx context.Context, instants ...time.Time) error { + err := p.tryCreateWorkerOutboxPartitions(ctx, false, instants...) + if err == nil || !isCheckViolation(err) { + return err + } + slog.WarnContext(ctx, "outbox DEFAULT partition holds rows in a range being created; truncating strays to un-wedge partition creation", + slog.Any("err", err)) + return p.tryCreateWorkerOutboxPartitions(ctx, true, instants...) +} + +// isCheckViolation matches SQLSTATE 23514, which CREATE ... PARTITION OF +// raises when the DEFAULT partition holds rows inside the new range. +func isCheckViolation(err error) bool { return pgErrCode(err) == "23514" } + +func (p *Persistence) tryCreateWorkerOutboxPartitions(ctx context.Context, truncateStrays bool, instants ...time.Time) error { + tx, err := p.watchPool.Begin(ctx) + if err != nil { + return fmt.Errorf("beginning outbox partition transaction: %w", err) + } + defer tx.Rollback(ctx) //nolint:errcheck // no-op once committed + + if _, err := tx.Exec(ctx, `SELECT pg_advisory_xact_lock(hashtext('agent-substrate-atepg-outbox-partitions'))`); err != nil { + return fmt.Errorf("locking outbox partition DDL: %w", err) + } + if truncateStrays { + // Parent BEFORE child: Writers lock parent-then-child, so locking DEFAULT + // first (and later needing parent for CREATE PARTITION) causes deadlocks. + // Locking the parent first matches writer order, and holding it through + // the CREATEs prevents concurrent writes from re-seeding DEFAULT mid-rescue. + if _, err := tx.Exec(ctx, `LOCK TABLE worker_outbox IN ACCESS EXCLUSIVE MODE`); err != nil { + return fmt.Errorf("locking outbox parent for stray rescue: %w", err) + } + if err := p.truncateWorkerOutboxDefault(ctx, tx); err != nil { + return err + } + } + for _, at := range instants { + start := at.UTC().Truncate(outboxPartitionInterval) + // UNLOGGED: see schema comment for the durability trade-off. + // autovacuum off: partitions are insert-only and discarded whole, + // so autovacuum is unnecessary and its scans would cause latency spikes. + stmt := fmt.Sprintf(`CREATE UNLOGGED TABLE IF NOT EXISTS %s PARTITION OF worker_outbox FOR VALUES FROM ('%s') TO ('%s') WITH (autovacuum_enabled = off)`, + workerOutboxPartitionName(start), start.Format(time.RFC3339), start.Add(outboxPartitionInterval).Format(time.RFC3339)) + if _, err := tx.Exec(ctx, stmt); err != nil { + return fmt.Errorf("creating outbox partition for %s: %w", start, err) + } + } + if err := tx.Commit(ctx); err != nil { + return fmt.Errorf("committing outbox partition transaction: %w", err) + } + return nil +} + +// dropExpiredWorkerOutboxPartitions drops every outbox partition whose +// entire range is older than retention. +func (p *Persistence) dropExpiredWorkerOutboxPartitions(ctx context.Context, q querier, now time.Time) error { + rows, err := q.Query(ctx, ` + SELECT c.relname FROM pg_inherits i + JOIN pg_class c ON c.oid = i.inhrelid + JOIN pg_class parent ON parent.oid = i.inhparent + WHERE parent.relname = 'worker_outbox'`) + if err != nil { + return fmt.Errorf("listing outbox partitions: %w", err) + } + var names []string + for rows.Next() { + var name string + if err := rows.Scan(&name); err != nil { + rows.Close() + return fmt.Errorf("scanning outbox partition name: %w", err) + } + names = append(names, name) + } + rows.Close() + if err := rows.Err(); err != nil { + return fmt.Errorf("listing outbox partitions: %w", err) + } + + for _, name := range names { + // The DEFAULT partition (worker_outbox_default) doesn't match + // the range prefix and is skipped here naturally. + suffix, ok := strings.CutPrefix(name, "worker_outbox_p") + if !ok { + continue + } + start, err := time.Parse("200601021504", suffix) + if err != nil { + continue // not a partition this maintenance loop manages + } + if now.Sub(start.Add(outboxPartitionInterval)) < outboxRetentionAge { + continue + } + if err := p.dropWorkerOutboxPartition(ctx, q, name); err != nil { + return err + } + } + return nil +} + +// dropWorkerOutboxPartition records the trim mark and drops the +// partition on the caller's (elected, single) retention transaction. +func (p *Persistence) dropWorkerOutboxPartition(ctx context.Context, q querier, name string) error { + ident := pgx.Identifier{name}.Sanitize() + // The mark is the partition's greatest xid. + if _, err := q.Exec(ctx, fmt.Sprintf(` + INSERT INTO worker_outbox_trim (xid) + SELECT xid FROM %s ORDER BY xid DESC LIMIT 1 + ON CONFLICT (id) DO UPDATE SET xid = EXCLUDED.xid + WHERE EXCLUDED.xid > worker_outbox_trim.xid`, ident)); err != nil { + return fmt.Errorf("recording trim mark for outbox partition %s: %w", name, err) + } + if _, err := q.Exec(ctx, `DROP TABLE `+ident); err != nil { + return fmt.Errorf("dropping outbox partition %s: %w", name, err) + } + return nil +} + +// truncateWorkerOutboxDefault discards the DEFAULT partition wholesale to +// un-stall partition creation. +func (p *Persistence) truncateWorkerOutboxDefault(ctx context.Context, q querier) error { + // Lock BEFORE reading the trim mark to block concurrent writers. This ensures + // our snapshot sees exactly what TRUNCATE will destroy, preventing silent data loss. + if _, err := q.Exec(ctx, `LOCK TABLE worker_outbox_default IN ACCESS EXCLUSIVE MODE`); err != nil { + return fmt.Errorf("locking outbox default partition: %w", err) + } + // Highest xid is recorded as a trim mark in the same transaction so lagging watchers detect the loss and resync. + if _, err := q.Exec(ctx, ` + INSERT INTO worker_outbox_trim (xid) + SELECT xid FROM worker_outbox_default ORDER BY xid DESC LIMIT 1 + ON CONFLICT (id) DO UPDATE SET xid = EXCLUDED.xid + WHERE EXCLUDED.xid > worker_outbox_trim.xid`); err != nil { + return fmt.Errorf("recording trim mark for outbox default partition: %w", err) + } + if _, err := q.Exec(ctx, `TRUNCATE worker_outbox_default`); err != nil { + return fmt.Errorf("truncating outbox default partition: %w", err) + } + return nil +} + +// WatchWorkers subscribes by polling the worker_outbox table using an xid cursor. +// It fences reads behind pg_snapshot_xmin (the oldest in-flight transaction) +// to guarantee gap-free delivery. Note that a long-running transaction anywhere +// in the database will stall delivery. +// +// Events are delivered in xid order, so consumers must reconcile worker versions. +// If the watcher detects missed events—either by lagging behind retention drops +// or if a database restart truncates the UNLOGGED partitions—it closes the channel +// to force the consumer to resync from the primary tables. +func (p *Persistence) WatchWorkers(ctx context.Context) (*store.WorkerWatch, error) { + watchCtx, cancel := context.WithCancel(ctx) + + // cursor starts at xmin-1. baseline records pre-subscribe history (xid < xmin) + // so past drops aren't mistaken for losses, while preventing artificially + // inflated baselines from masking the future drop of slow, in-flight txs. + var cursorXid, baselineXid, baselineStart string + if err := p.watchPool.QueryRow(watchCtx, ` + SELECT (pg_snapshot_xmin(pg_current_snapshot())::text::numeric - 1)::text, + GREATEST( + COALESCE((SELECT max(xid) FROM worker_outbox + WHERE xid < pg_snapshot_xmin(pg_current_snapshot())), '0'::xid8), + COALESCE((SELECT xid FROM worker_outbox_trim), '0'::xid8))::text, + pg_postmaster_start_time()::text`).Scan(&cursorXid, &baselineXid, &baselineStart); err != nil { + cancel() + return nil, fmt.Errorf("reading worker outbox cursor: %w", err) + } + + ch := make(chan store.WorkerEvent, 128) + go func() { + defer close(ch) + ticker := time.NewTicker(outboxPollInterval) + defer ticker.Stop() + // failingSince limits how long consumers serve stale state during an outage. + // Past outboxPollFailureCloseAfter, the channel closes. Unlike the postmaster + // restart check (which fires when connectivity returns), this surfaces the + // actual outage in real-time. + var failingSince time.Time + for { + select { + case <-watchCtx.Done(): + return + case <-ticker.C: + } + // Drain until a batch is partial. Sleeping between full batches would + // cap throughput and cause unrecoverable lag during bursts. + for { + // Safety checks share the batch round trip but must remain separate + // queries so we can detect gaps and restarts even when no rows match. + b := &pgx.Batch{} + b.Queue(pollWorkerOutboxSQL, cursorXid, outboxBatch) + b.Queue(pollSafetySQL, cursorXid, baselineXid) + br := p.watchPool.SendBatch(watchCtx, b) + + type feedRow struct { + xid string + payload []byte + } + var batch []feedRow + rows, err := br.Query() + if err == nil { + for rows.Next() { + var r feedRow + if err = rows.Scan(&r.xid, &r.payload); err != nil { + batch = nil + break + } + batch = append(batch, r) + } + rows.Close() + } + var fellBehind bool + var pmStart string + if err == nil { + err = br.QueryRow().Scan(&fellBehind, &pmStart) + } + if closeErr := br.Close(); err == nil { + err = closeErr + } + if err != nil { + if watchCtx.Err() != nil { + return + } + // Transient failure: retry next tick. Persistent failure (past the + // deadline): close the watch to flip workercache not-ready, forcing + // callers to fail fast instead of serving a frozen fleet view. + if failingSince.IsZero() { + failingSince = time.Now() + } else if time.Since(failingSince) > p.pollFailureCloseAfter { + slog.WarnContext(watchCtx, "worker outbox polling has failed persistently; closing watch", + slog.Duration("failing_for", time.Since(failingSince)), slog.Any("err", err)) + return + } + slog.WarnContext(watchCtx, "worker outbox poll failed", slog.Any("err", err)) + break + } + failingSince = time.Time{} + // A restarted postmaster truncated the UNLOGGED outbox: + // committed-but-undelivered events may be gone, so close + // before the cursor can skip past them; consumers resync + // with a full relist. + if pmStart != baselineStart { + slog.WarnContext(watchCtx, "database restarted under the outbox; closing watch for resync", + slog.String("was", baselineStart), slog.String("now", pmStart)) + return + } + // Retention safety: if retention's recorded trim high-water + // mark is ahead of everything this watcher has seen, a row + // it never consumed was discarded. Close before delivering + // anything past the gap. + if fellBehind { + slog.WarnContext(watchCtx, "worker watch fell behind outbox retention; closing for resync", + slog.String("cursor_xid", cursorXid)) + return + } + + for _, r := range batch { + event, err := unmarshalWorkerEvent(r.payload) + if err != nil { + // Close to force a relist. Skipping it would cause silent + // data loss. The fresh watch starts at the current xmin, + // naturally bypassing the corrupt row to prevent a boot loop. + slog.ErrorContext(watchCtx, "worker event unmarshal failed; closing watch for resync", + slog.String("xid", r.xid), slog.Any("err", err)) + return + } + select { + case ch <- event: + cursorXid = r.xid + case <-watchCtx.Done(): + return + } + } + if len(batch) < outboxBatch { + break // caught up; wait for the next tick + } + } + } + }() + return store.NewWorkerWatch(ch, cancel), nil +} diff --git a/cmd/ateapi/internal/store/atepg/outbox_test.go b/cmd/ateapi/internal/store/atepg/outbox_test.go new file mode 100644 index 000000000..8c778a162 --- /dev/null +++ b/cmd/ateapi/internal/store/atepg/outbox_test.go @@ -0,0 +1,1019 @@ +// Copyright 2026 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +// Tests for the worker outbox (outbox.go): transactional append, watch +// delivery and its safety fences, partition maintenance/retention, and the +// dedicated watch pool. Shared container fixtures live in atepg_test.go. + +package atepg + +import ( + "context" + "fmt" + "strings" + "sync" + "testing" + "time" + + "github.com/google/go-cmp/cmp" + "google.golang.org/protobuf/proto" + "google.golang.org/protobuf/testing/protocmp" + + "github.com/agent-substrate/substrate/cmd/ateapi/internal/store" + "github.com/agent-substrate/substrate/pkg/proto/ateapipb" +) + +// TestConnect_DedicatedWatchPool covers the dual-pool path only Connect +// takes (the rest of the suite uses NewPersistence, where feed traffic +// shares the caller's pool): the watch pool must be distinct and owned, and +// the watch must deliver through it end to end. +func TestConnect_DedicatedWatchPool(t *testing.T) { + requirePool(t) // ensures the container is up and containerDSN is set + ctx := context.Background() + + p, err := Connect(ctx, containerDSN) + if err != nil { + t.Fatalf("Connect failed: %v", err) + } + defer p.pool.Close() + defer p.Close() + + if p.watchPool == p.pool { + t.Fatal("Connect did not create a dedicated watch pool") + } + if !p.ownsWatchPool { + t.Fatal("Connect must own the watch pool so Close releases it") + } + if got := p.watchPool.Config().MaxConns; got != watchPoolMaxConns { + t.Fatalf("watch pool MaxConns = %d, want %d", got, watchPoolMaxConns) + } + + if err := p.DebugClearAll(ctx); err != nil { + t.Fatalf("DebugClearAll failed: %v", err) + } + watch, err := p.WatchWorkers(ctx) + if err != nil { + t.Fatalf("WatchWorkers failed: %v", err) + } + defer watch.Close() + worker := &ateapipb.Worker{ + Metadata: &ateapipb.ResourceMetadata{Name: "watchpool-worker"}, + WorkerNamespace: "ns", + WorkerPool: "pool", + WorkerPod: "watchpool-pod", + } + if err := p.CreateWorker(ctx, worker); err != nil { + t.Fatalf("CreateWorker failed: %v", err) + } + select { + case event := <-watch.Events: + if event.Type != store.WorkerEventCreated || event.Worker.GetWorkerPod() != "watchpool-pod" { + t.Fatalf("unexpected event %v %s", event.Type, event.Worker.GetWorkerPod()) + } + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for event through the watch pool") + } +} + +// TestWorkerEvent_OnlyAfterCommit proves the doc's atomicity claim: a +// worker write's outbox insert shares the write's transaction, so a +// rolled-back write never produces an event, while a committed write always +// does. +func TestWorkerEvent_OnlyAfterCommit(t *testing.T) { + s := setupPostgresStore(t).(*Persistence) + ctx := context.Background() + + watch, err := s.WatchWorkers(ctx) + if err != nil { + t.Fatalf("WatchWorkers failed: %v", err) + } + defer watch.Close() + + const workerName = "6e4d2f81-b3a9-4c05-8e72-1f9d4a0c7b63" + worker := &ateapipb.Worker{ + Metadata: &ateapipb.ResourceMetadata{Name: workerName}, + WorkerNamespace: "ns", + WorkerPool: "pool", + WorkerPod: "pod", + WorkerPodUid: workerName, + } + protoBytes, err := proto.Marshal(worker) + if err != nil { + t.Fatalf("marshaling worker: %v", err) + } + + // Write the row and roll back instead of committing: no event should + // ever arrive, proving the outbox insert is undone with the rest of the + // transaction. + tx, err := s.pool.Begin(ctx) + if err != nil { + t.Fatalf("Begin failed: %v", err) + } + if _, err := tx.Exec(ctx, ` + INSERT INTO workers (name, uid, version, proto) + VALUES ($1, $2, $3, $4)`, + workerName, "rolled-back-uid", int64(1), protoBytes); err != nil { + t.Fatalf("insert failed: %v", err) + } + if _, err := tx.Exec(ctx, `INSERT INTO worker_outbox (payload) VALUES ($1)`, []byte("rolled-back-payload")); err != nil { + t.Fatalf("outbox insert failed: %v", err) + } + if err := tx.Rollback(ctx); err != nil { + t.Fatalf("Rollback failed: %v", err) + } + + select { + case event := <-watch.Events: + t.Fatalf("received event %+v from a rolled-back transaction; the outbox insert must be undone with the rest of the transaction", event) + case <-time.After(500 * time.Millisecond): + // Expected: nothing arrives. + } + + // The equivalent committed write must produce an event. + if err := s.CreateWorker(ctx, worker); err != nil { + t.Fatalf("CreateWorker failed: %v", err) + } + select { + case event := <-watch.Events: + if event.Type != store.WorkerEventCreated { + t.Errorf("expected WorkerEventCreated, got %v", event.Type) + } + // CreateWorker assigns the uid, version and timestamps server-side. + want := proto.Clone(worker).(*ateapipb.Worker) + want.Metadata.Version = 1 + if diff := cmp.Diff(want, event.Worker, protocmp.Transform(), + protocmp.IgnoreFields(&ateapipb.ResourceMetadata{}, "uid", "create_time", "update_time")); diff != "" { + t.Errorf("event worker mismatch (-want +got):\n%s", diff) + } + case <-time.After(2 * time.Second): + t.Fatal("timed out waiting for event from a committed write") + } +} + +// TestWatchWorkers_OutOfOrderCommitNotSkipped reproduces the commit-order +// gap: xids are assigned at a transaction's first write but rows appear at +// COMMIT, so a transaction holding a lower xid can commit after a +// higher-xid sibling. A watcher that advanced past every visible row would +// skip the in-flight one and lose its event permanently. The xmin fence +// must instead hold the committed sibling back until the older +// transaction resolves, then deliver both in order. +func TestWatchWorkers_OutOfOrderCommitNotSkipped(t *testing.T) { + s := setupPostgresStore(t).(*Persistence) + ctx := context.Background() + + watch, err := s.WatchWorkers(ctx) + if err != nil { + t.Fatalf("WatchWorkers failed: %v", err) + } + defer watch.Close() + + mkPayload := func(pod string) []byte { + payload, err := marshalWorkerEvent(store.WorkerEventCreated, + &ateapipb.Worker{ + Metadata: &ateapipb.ResourceMetadata{Name: pod}, + WorkerNamespace: "ns", WorkerPool: "pool", WorkerPod: pod, + }) + if err != nil { + t.Fatalf("marshaling event for %q: %v", pod, err) + } + return payload + } + + // tx1 appends first (lower xid) and stays open. + tx1, err := s.pool.Begin(ctx) + if err != nil { + t.Fatalf("Begin tx1 failed: %v", err) + } + defer tx1.Rollback(ctx) //nolint:errcheck // no-op once committed + if _, err := tx1.Exec(ctx, `INSERT INTO worker_outbox (payload) VALUES ($1)`, mkPayload("first-xid-late-commit")); err != nil { + t.Fatalf("tx1 outbox insert failed: %v", err) + } + + // tx2 appends second (higher xid) and commits immediately. + tx2, err := s.pool.Begin(ctx) + if err != nil { + t.Fatalf("Begin tx2 failed: %v", err) + } + if _, err := tx2.Exec(ctx, `INSERT INTO worker_outbox (payload) VALUES ($1)`, mkPayload("second-xid-early-commit")); err != nil { + t.Fatalf("tx2 outbox insert failed: %v", err) + } + if err := tx2.Commit(ctx); err != nil { + t.Fatalf("tx2 Commit failed: %v", err) + } + + // While tx1 is in flight, tx2's committed event must be held back by + // the xmin fence — otherwise the cursor has already skipped tx1's row. + select { + case event := <-watch.Events: + t.Fatalf("event %q delivered while an older feed transaction was still in flight; its sibling event is now unreachable", event.Worker.GetWorkerPod()) + case <-time.After(500 * time.Millisecond): + // Expected: fence holds both events back. + } + + if err := tx1.Commit(ctx); err != nil { + t.Fatalf("tx1 Commit failed: %v", err) + } + + var got []string + for len(got) < 2 { + select { + case event := <-watch.Events: + got = append(got, event.Worker.GetWorkerPod()) + case <-time.After(5 * time.Second): + t.Fatalf("timed out waiting for both events; delivered so far: %v", got) + } + } + want := []string{"first-xid-late-commit", "second-xid-early-commit"} + if diff := cmp.Diff(want, got); diff != "" { + t.Errorf("event delivery order mismatch (-want +got):\n%s", diff) + } +} + +// TestWorkerOutboxPartitionRetention verifies partition-based +// retention: an hourly partition wholly past outboxRetentionAge is +// dropped (with its greatest xid recorded in worker_outbox_trim), fresh +// rows survive, and aged strays in the DEFAULT partition are trimmed by the +// row-wise fallback. +func TestWorkerOutboxPartitionRetention(t *testing.T) { + s := setupPostgresPersistence(t) + ctx := context.Background() + + // A partition two hours back, holding one aged event. + stale := time.Now().UTC().Add(-2 * time.Hour) + if err := s.createWorkerOutboxPartitions(ctx, stale); err != nil { + t.Fatalf("creating stale partition failed: %v", err) + } + var staleXid string + if err := s.pool.QueryRow(ctx, `INSERT INTO worker_outbox (payload, created_at) VALUES ($1, $2) RETURNING xid::text`, + []byte("old"), stale).Scan(&staleXid); err != nil { + t.Fatalf("inserting aged row failed: %v", err) + } + // An aged stray in the DEFAULT partition (no hourly partition covers a + // day ago), and a fresh row in the current partition. + if _, err := s.pool.Exec(ctx, `INSERT INTO worker_outbox (payload, created_at) VALUES ($1, now() - interval '1 day')`, []byte("stray")); err != nil { + t.Fatalf("inserting default-partition stray failed: %v", err) + } + if _, err := s.pool.Exec(ctx, `INSERT INTO worker_outbox (payload) VALUES ($1)`, []byte("fresh")); err != nil { + t.Fatalf("inserting fresh row failed: %v", err) + } + + if err := s.maintainWorkerOutboxPartitions(ctx); err != nil { + t.Fatalf("maintainWorkerOutboxPartitions failed: %v", err) + } + + var staleExists bool + if err := s.pool.QueryRow(ctx, `SELECT to_regclass($1) IS NOT NULL`, + workerOutboxPartitionName(stale)).Scan(&staleExists); err != nil { + t.Fatalf("checking stale partition failed: %v", err) + } + if staleExists { + t.Errorf("stale partition %s still exists, want dropped", workerOutboxPartitionName(stale)) + } + var remaining int + if err := s.pool.QueryRow(ctx, `SELECT count(*) FROM worker_outbox`).Scan(&remaining); err != nil { + t.Fatalf("counting remaining rows failed: %v", err) + } + if remaining != 1 { + t.Errorf("%d rows remain, want 1 (the fresh row; aged partition row and default stray gone)", remaining) + } + var trim bool + if err := s.pool.QueryRow(ctx, `SELECT (SELECT xid FROM worker_outbox_trim) >= $1::xid8`, staleXid).Scan(&trim); err != nil { + t.Fatalf("reading trim mark failed: %v", err) + } + if !trim { + t.Errorf("trim mark does not cover dropped partition's xid %s", staleXid) + } +} + +// TestOutboxMaintenance_SingleMaintainer verifies the retention +// election: while another replica holds the advisory lock, a pass skips +// retention cleanly (no error, nothing dropped); once released, the next +// pass does the work. (Partition creation is deliberately unelected.) +func TestOutboxMaintenance_SingleMaintainer(t *testing.T) { + s := setupPostgresPersistence(t) + ctx := context.Background() + + stale := time.Now().UTC().Add(-2 * time.Hour) + if err := s.createWorkerOutboxPartitions(ctx, stale); err != nil { + t.Fatalf("creating stale partition failed: %v", err) + } + + // Another "replica" mid-pass: hold the advisory lock in an open + // transaction of our own. + holder, err := s.pool.Begin(ctx) + if err != nil { + t.Fatalf("Begin holder failed: %v", err) + } + defer holder.Rollback(ctx) //nolint:errcheck // released below + if _, err := holder.Exec(ctx, `SELECT pg_advisory_xact_lock(hashtext(current_database() || ':' || $1))`, outboxMaintenanceLockKey); err != nil { + t.Fatalf("taking maintenance lock failed: %v", err) + } + + if err := s.maintainWorkerOutboxPartitions(ctx); err != nil { + t.Fatalf("pass with lock held must skip cleanly, got: %v", err) + } + var staleExists bool + if err := s.pool.QueryRow(ctx, `SELECT to_regclass($1) IS NOT NULL`, workerOutboxPartitionName(stale)).Scan(&staleExists); err != nil { + t.Fatalf("checking stale partition failed: %v", err) + } + if !staleExists { + t.Fatal("stale partition was dropped by a pass that lost the election") + } + + if err := holder.Rollback(ctx); err != nil { + t.Fatalf("releasing maintenance lock failed: %v", err) + } + if err := s.maintainWorkerOutboxPartitions(ctx); err != nil { + t.Fatalf("pass after lock release failed: %v", err) + } + if err := s.pool.QueryRow(ctx, `SELECT to_regclass($1) IS NOT NULL`, workerOutboxPartitionName(stale)).Scan(&staleExists); err != nil { + t.Fatalf("re-checking stale partition failed: %v", err) + } + if staleExists { + t.Error("stale partition survived a pass that held the election") + } +} + +// TestWorkerEvents_OneRowPerTransaction pins the invariant the xid-only +// watch cursor rests on: writeAndAppendEvent appends exactly one outbox row +// per transaction, so xids are distinct across the outbox and a poll batch +// can never split a same-xid group. +func TestWorkerEvents_OneRowPerTransaction(t *testing.T) { + s := setupPostgresPersistence(t) + ctx := context.Background() + + worker := &ateapipb.Worker{ + Metadata: &ateapipb.ResourceMetadata{Name: "one-row-worker"}, + WorkerNamespace: "ns", + WorkerPool: "pool", + WorkerPod: "pod", + } + if err := s.CreateWorker(ctx, worker); err != nil { + t.Fatalf("CreateWorker failed: %v", err) + } + for i := 0; i < 10; i++ { + stored, err := s.GetWorker(ctx, "one-row-worker") + if err != nil { + t.Fatalf("GetWorker failed: %v", err) + } + if err := s.UpdateWorker(ctx, stored, stored.GetMetadata().GetVersion()); err != nil { + t.Fatalf("UpdateWorker %d failed: %v", i, err) + } + } + if err := s.DeleteWorker(ctx, "one-row-worker"); err != nil { + t.Fatalf("DeleteWorker failed: %v", err) + } + + var total, distinct int + if err := s.pool.QueryRow(ctx, `SELECT count(*), count(DISTINCT xid) FROM worker_outbox`).Scan(&total, &distinct); err != nil { + t.Fatalf("counting outbox rows failed: %v", err) + } + if total == 0 || total != distinct { + t.Errorf("feed has %d rows but %d distinct xids; the one-row-per-transaction invariant is broken", total, distinct) + } +} + +// TestWatchWorkers_DeliveryFencedByOldestTransaction documents the xmin +// fence's real bound: one old transaction anywhere holds back delivery of +// everything committed after it, for as long as it lives. +func TestWatchWorkers_DeliveryFencedByOldestTransaction(t *testing.T) { + s := setupPostgresPersistence(t) + ctx := context.Background() + + watch, err := s.WatchWorkers(ctx) + if err != nil { + t.Fatalf("WatchWorkers failed: %v", err) + } + defer watch.Close() + + // An unrelated transaction that merely holds an xid. + blocker, err := s.pool.Begin(ctx) + if err != nil { + t.Fatalf("Begin blocker failed: %v", err) + } + defer blocker.Rollback(ctx) //nolint:errcheck // released below + if _, err := blocker.Exec(ctx, `SELECT pg_current_xact_id()`); err != nil { + t.Fatalf("assigning blocker xid failed: %v", err) + } + + worker := &ateapipb.Worker{ + Metadata: &ateapipb.ResourceMetadata{Name: "fenced-worker"}, + WorkerNamespace: "ns", + WorkerPool: "pool", + WorkerPod: "fenced", + } + if err := s.CreateWorker(ctx, worker); err != nil { + t.Fatalf("CreateWorker failed: %v", err) + } + + select { + case event := <-watch.Events: + t.Fatalf("event %+v delivered through the fence while an older transaction was in flight", event) + case <-time.After(600 * time.Millisecond): + // Expected: committed but fenced behind the blocker's xid. + } + + if err := blocker.Rollback(ctx); err != nil { + t.Fatalf("ending blocker failed: %v", err) + } + select { + case event := <-watch.Events: + if got := event.Worker.GetWorkerPod(); got != "fenced" { + t.Errorf("delivered %q, want %q", got, "fenced") + } + case <-time.After(5 * time.Second): + t.Fatal("event not delivered after the fencing transaction ended") + } +} + +// TestClose_StopsMaintenance pins that Close ends the background +// maintenance goroutine (main.go defers it for exactly this): Close blocks +// on the loop's done channel, so its return IS the assertion. +func TestClose_StopsMaintenance(t *testing.T) { + ctx := context.Background() + p, err := NewPersistence(ctx, requirePool(t)) + if err != nil { + t.Fatalf("NewPersistence failed: %v", err) + } + closed := make(chan struct{}) + go func() { + p.Close() + close(closed) + }() + select { + case <-closed: + case <-time.After(5 * time.Second): + t.Fatal("Close did not stop the maintenance loop") + } +} + +// TestPollQueryPlanStaysOnIndex pins the poll's plan shape against the +// output-column shadowing bug: an unaliased xid::text captures the bare +// ORDER BY name, sorting xids as text — which both diverges from the +// cursor predicate's xid8 order (silently skipping events across digit +// boundaries) and forces full scans with a top-N sort. Behavioral tests +// cannot see this; the plan can. +func TestPollQueryPlanStaysOnIndex(t *testing.T) { + s := setupPostgresPersistence(t) + ctx := context.Background() + + // Seed enough rows (and stats) for the planner to have a real choice: + // on empty partitions it costs bitmap scans plus an explicit Sort as + // cheapest regardless of the index, which would make the Merge Append + // assertion below vacuously unreachable. + if _, err := s.pool.Exec(ctx, `INSERT INTO worker_outbox (payload) SELECT 'x'::bytea FROM generate_series(1, 3000)`); err != nil { + t.Fatalf("seeding outbox rows failed: %v", err) + } + if _, err := s.pool.Exec(ctx, `ANALYZE worker_outbox`); err != nil { + t.Fatalf("ANALYZE failed: %v", err) + } + + rows, err := s.pool.Query(ctx, "EXPLAIN "+pollWorkerOutboxSQL, "100", outboxBatch) + if err != nil { + t.Fatalf("EXPLAIN failed: %v", err) + } + defer rows.Close() + var plan strings.Builder + for rows.Next() { + var line string + if err := rows.Scan(&line); err != nil { + t.Fatalf("scanning plan line: %v", err) + } + plan.WriteString(line) + plan.WriteString("\n") + } + got := plan.String() + if strings.Contains(got, "::text") { + t.Errorf("poll plan sorts by a text expression (output-column shadowing is back):\n%s", got) + } + if !strings.Contains(got, "Merge Append") { + t.Errorf("poll plan is not an index-ordered Merge Append:\n%s", got) + } +} + +// TestOutboxMaintenance_ConcurrentPassesAreHarmless backs the doc's +// claim: two replicas racing a maintenance pass produce no errors and the +// correct end state (one wins the election, the loser skips). +func TestOutboxMaintenance_ConcurrentPassesAreHarmless(t *testing.T) { + s := setupPostgresPersistence(t) + ctx := context.Background() + + replica, err := NewPersistence(ctx, s.pool) + if err != nil { + t.Fatalf("second Persistence failed: %v", err) + } + t.Cleanup(replica.Close) + + stale := time.Now().UTC().Add(-2 * time.Hour) + if err := s.createWorkerOutboxPartitions(ctx, stale); err != nil { + t.Fatalf("creating stale partition failed: %v", err) + } + + var wg sync.WaitGroup + errs := make([]error, 2) + for i, p := range []*Persistence{s, replica} { + wg.Add(1) + go func(i int, p *Persistence) { + defer wg.Done() + errs[i] = p.maintainWorkerOutboxPartitions(ctx) + }(i, p) + } + wg.Wait() + for i, err := range errs { + if err != nil { + t.Errorf("concurrent pass %d returned error: %v", i, err) + } + } + var staleExists bool + if err := s.pool.QueryRow(ctx, `SELECT to_regclass($1) IS NOT NULL`, workerOutboxPartitionName(stale)).Scan(&staleExists); err != nil { + t.Fatalf("checking stale partition failed: %v", err) + } + if staleExists { + t.Error("stale partition survived both concurrent passes") + } +} + +// TestWorkerOutboxPartitionsAreUnlogged pins the maintenance profile the +// schema documents: every outbox partition must be UNLOGGED (relpersistence +// 'u') with autovacuum disabled (all are insert-only and discarded whole, +// by drop or truncate — an in-window insert-autovacuum is a measured p99 +// spike); and worker_outbox_trim — the loss-detection high-water mark — +// must remain logged so it survives a crash. +func TestWorkerOutboxPartitionsAreUnlogged(t *testing.T) { + s := setupPostgresPersistence(t) + ctx := context.Background() + + rows, err := s.pool.Query(ctx, ` + SELECT c.relname, c.relpersistence, COALESCE(array_to_string(c.reloptions, ','), '') FROM pg_inherits i + JOIN pg_class c ON c.oid = i.inhrelid + JOIN pg_class parent ON parent.oid = i.inhparent + WHERE parent.relname = 'worker_outbox'`) + if err != nil { + t.Fatalf("listing outbox partitions: %v", err) + } + defer rows.Close() + checked := 0 + for rows.Next() { + var name, persistence, options string + if err := rows.Scan(&name, &persistence, &options); err != nil { + t.Fatalf("scanning partition row: %v", err) + } + if persistence != "u" { + t.Errorf("partition %s has relpersistence %q, want 'u' (unlogged)", name, persistence) + } + if !strings.Contains(options, "autovacuum_enabled=off") { + t.Errorf("partition %s does not disable autovacuum (reloptions %q)", name, options) + } + checked++ + } + if checked == 0 { + t.Fatal("no outbox partitions found to check") + } + var trimPersistence string + if err := s.pool.QueryRow(ctx, `SELECT relpersistence FROM pg_class WHERE relname = 'worker_outbox_trim'`).Scan(&trimPersistence); err != nil { + t.Fatalf("checking worker_outbox_trim persistence: %v", err) + } + if trimPersistence != "p" { + t.Errorf("worker_outbox_trim has relpersistence %q, want 'p' (logged) — the trim mark must survive a crash", trimPersistence) + } +} + +// The restart escape hatch (a changed pg_postmaster_start_time() closes +// the watch, because a restart truncates the UNLOGGED feed) has no e2e +// test here: restarting the testcontainer remaps its host port, severing +// the pool permanently — unlike production, where the database endpoint is +// stable across restarts. The comparison itself is four lines in +// WatchWorkers' poll loop; the trimmed-past-cursor test below covers the +// shared close-for-resync path. + +// TestWatchWorkers_ClosesWhenTrimmedPastCursor verifies the retention +// escape hatch: when rows a watcher has not consumed are deleted out from +// under it (a retention trim on a badly lagging watcher), the watcher must +// close its channel — the signal consumers treat as resync-and-relist — +// rather than silently skip the gap. +func TestWatchWorkers_ClosesWhenTrimmedPastCursor(t *testing.T) { + s := setupPostgresPersistence(t) + ctx := context.Background() + + watch, err := s.WatchWorkers(ctx) + if err != nil { + t.Fatalf("WatchWorkers failed: %v", err) + } + defer watch.Close() + + // Deliver one event normally so the cursor is established. + worker := &ateapipb.Worker{ + Metadata: &ateapipb.ResourceMetadata{Name: "trim-worker"}, + WorkerNamespace: "ns", + WorkerPool: "pool", + WorkerPod: "pod", + } + if err := s.CreateWorker(ctx, worker); err != nil { + t.Fatalf("CreateWorker failed: %v", err) + } + select { + case <-watch.Events: + case <-time.After(2 * time.Second): + t.Fatal("timed out waiting for the first event") + } + + // Atomically append three events and trim them away unconsumed — + // the watcher never gets a chance to see them, exactly as if + // retention took rows a lagging watcher had not reached. + payload, err := marshalWorkerEvent(store.WorkerEventUpdated, worker) + if err != nil { + t.Fatalf("marshaling event: %v", err) + } + tx, err := s.pool.Begin(ctx) + if err != nil { + t.Fatalf("Begin failed: %v", err) + } + defer tx.Rollback(ctx) //nolint:errcheck // no-op once committed + if _, err := tx.Exec(ctx, `INSERT INTO worker_outbox (payload) VALUES ($1), ($1), ($1)`, payload); err != nil { + t.Fatalf("outbox inserts failed: %v", err) + } + // Mirrors trimWorkerChangesDefault's shape: the mark is the deleted + // set's greatest xid. (The three rows above share one transaction — + // fine here: this test only needs the recorded mark to land past the + // watcher's cursor, and deletes everything the watcher has not seen.) + if _, err := tx.Exec(ctx, ` + WITH doomed AS ( + DELETE FROM worker_outbox WHERE xid > (SELECT COALESCE((SELECT xid FROM worker_outbox_trim), '0'::xid8)) + RETURNING xid + ) + INSERT INTO worker_outbox_trim (xid) + SELECT xid FROM doomed ORDER BY xid DESC LIMIT 1 + ON CONFLICT (id) DO UPDATE SET xid = EXCLUDED.xid + WHERE EXCLUDED.xid > worker_outbox_trim.xid`); err != nil { + t.Fatalf("trim failed: %v", err) + } + if err := tx.Commit(ctx); err != nil { + t.Fatalf("Commit failed: %v", err) + } + + // The watcher must close the channel, not deliver past the gap. + select { + case event, ok := <-watch.Events: + if ok { + t.Fatalf("received event %+v past a trimmed gap; expected the channel to close for resync", event) + } + // Expected: channel closed. + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for the watch channel to close after a trim past the cursor") + } +} + +// TestPartitionCreation_UnwedgesFromStrayedDefault reproduces the +// self-wedging order caught in review: creation stalls past the lead while +// writes continue, strays land in DEFAULT with current-range timestamps, and +// CREATE ... PARTITION OF for that range then fails outright — permanently, +// since maintenance returns before its strays cleanup and boot fails the +// same way. The fix retries creation with a same-transaction truncate of the +// strays; this test drills the exact scenario. +func TestPartitionCreation_UnwedgesFromStrayedDefault(t *testing.T) { + s := setupPostgresPersistence(t) + ctx := context.Background() + + // Stay clear of a partition boundary so the stray and the re-created + // partition land in the same range. + now := time.Now().UTC() + if rem := now.Truncate(outboxPartitionInterval).Add(outboxPartitionInterval).Sub(now); rem < 5*time.Second { + time.Sleep(rem + time.Second) + now = time.Now().UTC() + } + + // Simulate the stall: the current range's partition does not exist. + if _, err := s.pool.Exec(ctx, `DROP TABLE `+workerOutboxPartitionName(now)); err != nil { + t.Fatalf("dropping current partition: %v", err) + } + // A write during the stall detours into DEFAULT. + worker := &ateapipb.Worker{ + Metadata: &ateapipb.ResourceMetadata{Name: "wedge-worker"}, + WorkerNamespace: "ns", + WorkerPool: "pool", + WorkerPod: "wedge-pod", + } + if err := s.CreateWorker(ctx, worker); err != nil { + t.Fatalf("CreateWorker during stall: %v", err) + } + var strays bool + if err := s.pool.QueryRow(ctx, `SELECT EXISTS(SELECT 1 FROM worker_outbox_default)`).Scan(&strays); err != nil { + t.Fatalf("checking DEFAULT: %v", err) + } + if !strays { + t.Fatal("expected the stall-window write to land in the DEFAULT partition") + } + + // The stall clears: creation must dig itself out (pre-fix this returned + // SQLSTATE 23514 forever, and NewPersistence failed the same way). + if err := s.createWorkerOutboxPartitions(ctx, outboxPartitionLeadTimes(now)...); err != nil { + t.Fatalf("createWorkerOutboxPartitions did not un-wedge: %v", err) + } + + if err := s.pool.QueryRow(ctx, `SELECT EXISTS(SELECT 1 FROM worker_outbox_default)`).Scan(&strays); err != nil { + t.Fatalf("re-checking DEFAULT: %v", err) + } + if strays { + t.Fatal("DEFAULT partition still holds strays after the rescue") + } + var partitionExists bool + if err := s.pool.QueryRow(ctx, `SELECT to_regclass($1) IS NOT NULL`, workerOutboxPartitionName(now)).Scan(&partitionExists); err != nil { + t.Fatalf("checking recreated partition: %v", err) + } + if !partitionExists { + t.Fatal("current-range partition was not recreated") + } + // The truncated stray was never trim-covered before the fix; the rescue + // must have recorded it so lagging watchers resync rather than skip. + var trimSet bool + if err := s.pool.QueryRow(ctx, `SELECT EXISTS(SELECT 1 FROM worker_outbox_trim WHERE xid > '0'::xid8)`).Scan(&trimSet); err != nil { + t.Fatalf("checking trim mark: %v", err) + } + if !trimSet { + t.Fatal("rescue did not record a trim mark for the truncated strays") + } +} + +// TestOutboxMaintenance_NoDeadlockWithConcurrentWriters drills the measured +// child-vs-parent lock-order deadlock: writers route into DEFAULT exactly +// when the stray cleanup runs (that's the truncate path's precondition, not +// a coincidence), and writers lock parent-then-child while a single +// truncate+drops transaction locked child-then-parent. With retention split +// into two transactions and the wedge rescue locking the parent first, no +// interleaving can cycle. Pre-fix this test deadlocked (SQLSTATE 40P01) +// within an iteration or two. +func TestOutboxMaintenance_NoDeadlockWithConcurrentWriters(t *testing.T) { + s := setupPostgresPersistence(t) + ctx := context.Background() + + now := time.Now().UTC() + if rem := now.Truncate(outboxPartitionInterval).Add(outboxPartitionInterval).Sub(now); rem < 10*time.Second { + time.Sleep(rem + time.Second) + now = time.Now().UTC() + } + + // An aged-out partition with a row, so the drops leg has real work. + old := now.Add(-2*outboxRetentionAge - 2*outboxPartitionInterval) + if err := s.createWorkerOutboxPartitions(ctx, old); err != nil { + t.Fatalf("creating expired partition: %v", err) + } + if _, err := s.pool.Exec(ctx, `INSERT INTO worker_outbox (created_at, payload) VALUES ($1, $2)`, + old, []byte{byte(store.WorkerEventUpdated)}); err != nil { + t.Fatalf("seeding expired partition: %v", err) + } + + // Writers hammering the outbox for the whole test. + writerCtx, stopWriters := context.WithCancel(ctx) + defer stopWriters() + var wg sync.WaitGroup + writerErr := make(chan error, 4) + for g := 0; g < 4; g++ { + wg.Add(1) + go func(g int) { + defer wg.Done() + for i := 0; writerCtx.Err() == nil; i++ { + w := &ateapipb.Worker{ + Metadata: &ateapipb.ResourceMetadata{Name: fmt.Sprintf("ddl-worker-%d-%d", g, i)}, + WorkerNamespace: "ns", WorkerPool: "pool", WorkerPod: fmt.Sprintf("ddl-pod-%d-%d", g, i), + } + if err := s.CreateWorker(writerCtx, w); err != nil && writerCtx.Err() == nil { + select { + case writerErr <- fmt.Errorf("writer %d iteration %d: %w", g, i, err): + default: + } + return + } + } + }(g) + } + + // Repeatedly re-enter the degraded state under live writers: drop the + // current partition (parent-first, or this helper itself would ABBA), + // let writes detour into DEFAULT, then run a full maintenance pass — + // rescue-create, stray cleanup, and drops, all with writers in flight. + for i := 0; i < 4; i++ { + if _, err := s.pool.Exec(ctx, + `DO $$ BEGIN LOCK TABLE worker_outbox IN ACCESS EXCLUSIVE MODE; EXECUTE 'DROP TABLE IF EXISTS `+workerOutboxPartitionName(time.Now().UTC())+`'; END $$`); err != nil { + t.Fatalf("dropping current partition (iteration %d): %v", i, err) + } + time.Sleep(150 * time.Millisecond) // let writers detour into DEFAULT + if err := s.maintainWorkerOutboxPartitions(ctx); err != nil { + t.Fatalf("maintenance pass %d failed under concurrent writers: %v", i, err) + } + } + + stopWriters() + wg.Wait() + select { + case err := <-writerErr: + t.Fatalf("concurrent writer failed: %v", err) + default: + } + + var strays bool + if err := s.pool.QueryRow(ctx, `SELECT EXISTS(SELECT 1 FROM worker_outbox_default)`).Scan(&strays); err != nil { + t.Fatalf("checking DEFAULT: %v", err) + } + if strays { + t.Fatal("DEFAULT partition still holds strays after maintenance") + } + var oldExists bool + if err := s.pool.QueryRow(ctx, `SELECT to_regclass($1) IS NOT NULL`, workerOutboxPartitionName(old)).Scan(&oldExists); err != nil { + t.Fatalf("checking expired partition: %v", err) + } + if oldExists { + t.Fatal("expired partition survived the drops leg") + } +} + +// TestWatchWorkers_ClosesAfterPersistentPollFailure pins the watch +// contract's loss signal for polling outages: transient failures retry with +// the cursor kept (edge case 12), but a persistent failure must close the +// channel within outboxPollFailureCloseAfter so workercache flips not-ready +// and callers fail fast — instead of serving a frozen fleet view for as +// long as the outage lasts. (The postmaster-restart check cannot cover +// this: it fires only after connectivity returns.) +func TestWatchWorkers_ClosesAfterPersistentPollFailure(t *testing.T) { + requirePool(t) + ctx := context.Background() + + // Connect so the watcher has its own pool: killing it simulates a + // persistent outage without touching the shared container pool. + p, err := Connect(ctx, containerDSN) + if err != nil { + t.Fatalf("Connect failed: %v", err) + } + p.pollFailureCloseAfter = 300 * time.Millisecond // before WatchWorkers starts the poller + defer p.pool.Close() + defer p.Close() + + watch, err := p.WatchWorkers(ctx) + if err != nil { + t.Fatalf("WatchWorkers failed: %v", err) + } + defer watch.Close() + + p.watchPool.Close() // every subsequent poll now fails + + deadline := time.After(5 * time.Second) + for { + select { + case _, ok := <-watch.Events: + if !ok { + return // closed: the loss signal fired + } + case <-deadline: + t.Fatal("watch did not close after persistent poll failure") + } + } +} + +// TestWatchWorkers_BaselineDoesNotMaskOwedTrims pins the subscribe-baseline +// corner from review: a transaction that took its xid BEFORE this watch +// subscribed but commits AFTER it is a legitimate owed event — yet with a +// baseline of "highest existing row xid", an already-committed HIGHER xid +// (out-of-order commit) raised the baseline above the owed xid, so when the +// owed row was trimmed while the xmin fence stalled, trim ≤ baseline and the +// fell-behind close never fired: silent loss. The baseline must only cover +// settled history below the cursor's own xmin snapshot. +func TestWatchWorkers_BaselineDoesNotMaskOwedTrims(t *testing.T) { + s := setupPostgresStore(t).(*Persistence) + ctx := context.Background() + + // L: an old in-flight transaction pinning xmin — the fence stall. + lTx, err := s.pool.Begin(ctx) + if err != nil { + t.Fatalf("Begin L: %v", err) + } + defer lTx.Rollback(ctx) //nolint:errcheck + var xidL string + if err := lTx.QueryRow(ctx, `SELECT pg_current_xact_id()::text`).Scan(&xidL); err != nil { + t.Fatalf("assigning L's xid: %v", err) + } + + // W: takes the next xid now, but commits only after the subscribe. + wTx, err := s.pool.Begin(ctx) + if err != nil { + t.Fatalf("Begin W: %v", err) + } + defer wTx.Rollback(ctx) //nolint:errcheck + var xidW string + if err := wTx.QueryRow(ctx, `SELECT pg_current_xact_id()::text`).Scan(&xidW); err != nil { + t.Fatalf("assigning W's xid: %v", err) + } + + // C: a LATER xid that commits BEFORE the subscribe — the baseline poison: + // a visible outbox row whose xid exceeds W's. + if err := s.CreateWorker(ctx, &ateapipb.Worker{ + Metadata: &ateapipb.ResourceMetadata{Name: "baseline-poison"}, + WorkerNamespace: "ns", WorkerPool: "pool", WorkerPod: "baseline-pod", + }); err != nil { + t.Fatalf("CreateWorker (C): %v", err) + } + + watch, err := s.WatchWorkers(ctx) // cursor = xidL-1; baseline must exclude C's xid (>= xmin) + if err != nil { + t.Fatalf("WatchWorkers: %v", err) + } + defer watch.Close() + + // W commits its owed event post-subscribe (fenced behind L, undeliverable). + if _, err := wTx.Exec(ctx, `INSERT INTO worker_outbox (payload) VALUES ($1)`, + []byte{byte(store.WorkerEventUpdated)}); err != nil { + t.Fatalf("W's outbox insert: %v", err) + } + if err := wTx.Commit(ctx); err != nil { + t.Fatalf("committing W: %v", err) + } + + // Retention drops W's partition while the fence still stalls: record the + // trim exactly as dropWorkerOutboxPartition would. + if _, err := s.pool.Exec(ctx, ` + INSERT INTO worker_outbox_trim (xid) VALUES ($1::xid8) + ON CONFLICT (id) DO UPDATE SET xid = EXCLUDED.xid + WHERE EXCLUDED.xid > worker_outbox_trim.xid`, xidW); err != nil { + t.Fatalf("recording trim of W's partition: %v", err) + } + + // The fell-behind close must fire: trim (xidW) > cursor (xidL-1), and the + // baseline may not be raised by C's higher-but-already-visible xid. + // Nothing can be delivered first — everything above xidL is fenced. + select { + case event, ok := <-watch.Events: + if ok { + t.Fatalf("delivered event %+v through the fence", event) + } + // Closed: loss surfaced. + case <-time.After(5 * time.Second): + t.Fatal("watch stayed open: the trimmed owed event was silently lost (baseline masked the trim)") + } +} + +// TestUnmarshalWorkerEvent_BoundaryAssertions pins the write-side invariants +// asserted at the read boundary: known event-type byte and a keyable worker. +func TestUnmarshalWorkerEvent_BoundaryAssertions(t *testing.T) { + valid, err := marshalWorkerEvent(store.WorkerEventUpdated, &ateapipb.Worker{ + Metadata: &ateapipb.ResourceMetadata{Name: "w1"}, + }) + if err != nil { + t.Fatalf("marshal: %v", err) + } + for name, tc := range map[string]struct { + payload []byte + wantErr bool + }{ + "valid": {valid, false}, + "empty": {nil, true}, + "unknown type byte": {[]byte{0xff, 0x00}, true}, + "type byte only": {[]byte{byte(store.WorkerEventCreated)}, true}, // empty proto = nameless + "garbage proto": {append([]byte{byte(store.WorkerEventCreated)}, 0xde, 0xad, 0xbe), true}, + "nameless worker": {func() []byte { b, _ := marshalWorkerEvent(store.WorkerEventDeleted, &ateapipb.Worker{}); return b }(), true}, + } { + t.Run(name, func(t *testing.T) { + _, err := unmarshalWorkerEvent(tc.payload) + if (err != nil) != tc.wantErr { + t.Fatalf("unmarshalWorkerEvent() err = %v, wantErr = %v", err, tc.wantErr) + } + }) + } +} + +// TestWatchWorkers_ClosesOnCorruptPayload pins close-over-skip: a payload +// that fails the boundary checks must close the watch (loss surfaced, +// relist repairs) rather than advance the cursor past it silently. +func TestWatchWorkers_ClosesOnCorruptPayload(t *testing.T) { + s := setupPostgresStore(t).(*Persistence) + ctx := context.Background() + + watch, err := s.WatchWorkers(ctx) + if err != nil { + t.Fatalf("WatchWorkers: %v", err) + } + defer watch.Close() + + if _, err := s.pool.Exec(ctx, `INSERT INTO worker_outbox (payload) VALUES ($1)`, + []byte{0xff, 0xde, 0xad}); err != nil { + t.Fatalf("inserting corrupt payload: %v", err) + } + + select { + case event, ok := <-watch.Events: + if ok { + t.Fatalf("delivered event %+v from a corrupt payload", event) + } + // Closed: the loss signal fired. + case <-time.After(5 * time.Second): + t.Fatal("watch stayed open past a corrupt payload (silent skip)") + } +} diff --git a/cmd/ateapi/internal/store/atepg/schema.go b/cmd/ateapi/internal/store/atepg/schema.go index 48ca4a0f5..fae80f1ac 100644 --- a/cmd/ateapi/internal/store/atepg/schema.go +++ b/cmd/ateapi/internal/store/atepg/schema.go @@ -91,6 +91,36 @@ CREATE TABLE IF NOT EXISTS workers ( proto bytea NOT NULL ); +-- Transactional outbox backing WatchWorkers. +-- +-- 1. Ordering (xid): writeAndAppendEvent guarantees exactly one row per tx, +-- ensuring distinct xids so polling batches never split a transaction. +-- 2. Retention (created_at partitions): outboxMaintenance drops expired +-- partitions to avoid VACUUM I/O debt. A DEFAULT partition catches overflow. +-- 3. Durability (UNLOGGED): Skips WAL overhead. Crash recoveries trigger +-- watchers to rebuild from the primary workers table. worker_outbox_trim +-- remains LOGGED to preserve the high-water mark across restarts. +CREATE TABLE IF NOT EXISTS worker_outbox ( + xid xid8 NOT NULL DEFAULT pg_current_xact_id(), + -- MUST use clock_timestamp() instead of now(). now() freezes at tx start, + -- causing slow transactions to route into expired partitions. + created_at timestamptz NOT NULL DEFAULT clock_timestamp(), + payload bytea NOT NULL +) PARTITION BY RANGE (created_at); + +CREATE INDEX IF NOT EXISTS worker_outbox_xid ON worker_outbox (xid); + +CREATE UNLOGGED TABLE IF NOT EXISTS worker_outbox_default PARTITION OF worker_outbox DEFAULT WITH (autovacuum_enabled = off); + +-- Single-row high-water mark of retention: the greatest xid ever discarded +-- from worker_outbox (dropped with an expired partition, or truncated +-- with the DEFAULT partition). Watchers compare it against their cursor to +-- detect exactly that unconsumed rows were discarded out from under them. +CREATE TABLE IF NOT EXISTS worker_outbox_trim ( + id boolean PRIMARY KEY DEFAULT true CHECK (id), + xid xid8 NOT NULL +); + CREATE TABLE IF NOT EXISTS leases ( key text PRIMARY KEY, token text NOT NULL, @@ -108,6 +138,17 @@ func applySchema(ctx context.Context, pool *pgxpool.Pool) error { } defer tx.Rollback(ctx) //nolint:errcheck // no-op once committed + // The schema needs PostgreSQL 13+ (xid8, pg_current_xact_id, + // pg_current_snapshot); fail with a clear message rather than an + // opaque DDL or function error. + var version int + if err := tx.QueryRow(ctx, `SELECT current_setting('server_version_num')::int`).Scan(&version); err != nil { + return fmt.Errorf("reading PostgreSQL version: %w", err) + } + if version < 130000 { + return fmt.Errorf("atepg requires PostgreSQL 13 or newer (xid8 and pg_current_snapshot); server_version_num is %d", version) + } + // Multiple ateapi replicas can start against an empty database together. // PostgreSQL's IF NOT EXISTS does not eliminate every concurrent-DDL race, // so serialize schema application with a transaction-scoped advisory lock. diff --git a/cmd/ateapi/main.go b/cmd/ateapi/main.go index 345ea631f..6c9b7b9f8 100644 --- a/cmd/ateapi/main.go +++ b/cmd/ateapi/main.go @@ -138,6 +138,11 @@ func main() { if err != nil { serverboot.Fatal(ctx, "Failed to set up persistence backend", err) } + // Backends may run background maintenance rooted in their own context + // (atepg's outbox maintenance loop); stop it on shutdown. + if closer, ok := persistence.(interface{ Close() }); ok { + defer closer.Close() + } clientset, ateClient, err := newKubeClients() if err != nil {