atepg: deliver worker watch events via a transactional change feed (#934)

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
This commit is contained in:
shrutiyam-glitch
2026-08-24 16:57:11 -04:00
committed by GitHub
parent 62dbb86718
commit 137b6fd3df
8 changed files with 1853 additions and 240 deletions
@@ -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()))
}
@@ -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 {
+78 -132
View File
@@ -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
+3 -99
View File
@@ -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()
+571
View File
@@ -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
}
File diff suppressed because it is too large Load Diff
+41
View File
@@ -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.
+5
View File
@@ -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 {