From 3ea02b2f5dcf5ccb2614076fd8df75b7bf90874f Mon Sep 17 00:00:00 2001 From: Joseph Schorr Date: Mon, 5 Oct 2026 18:06:05 -0400 Subject: [PATCH] Bound operator memory: GC terminal sessions and drop the inmem audit mirror MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The operator OOMKilled on long-lived clusters from two unbounded, retain-forever-by-design sources, neither of which serves its running work. 1. Terminal AgentSessions and the ToolCalls they own were never deleted. The pod reap and PVC reclaim free pods and scratch volumes but keep the CR, so every session and every ToolCall ever created stayed resident in etcd and in the operator's informer cache (observed: 10k+ ToolCalls, oldest months old). Add a session-lifetime GC (--session-gc-after, default 30d): a terminal session past the window is deleted, cascade-collecting its ToolCalls, pods and Secrets. finalize keeps the append-only audit log (it lives in the durable backend; DeleteScope refuses append-only kinds), so the record survives — only the CR and its owned ephemera go. This is the retention the reap already anticipated ("until retention GCs the AgentSession"). 2. Under MEMORY_BACKEND=postgres the ShadowBackend dual-wrote every append-only entry into an in-memory backend as well, while serving reads only from postgres. That inmem mirror was never read and scope-delete can never free it (append-only scopes are undeletable), so it grew with the whole cluster's audit history for the operator's lifetime. Use the postgres backend directly; reads are byte-identical to what the shadow already served from its secondary. --- internal/cmd/operator/main.go | 58 +++++---- internal/cmd/operator/memory_backend_test.go | 20 +++ pkg/controllers/agentsession/controller.go | 33 +++++ pkg/controllers/agentsession/sessiongc.go | 88 ++++++++++++++ .../agentsession/sessiongc_test.go | 114 ++++++++++++++++++ 5 files changed, 291 insertions(+), 22 deletions(-) create mode 100644 pkg/controllers/agentsession/sessiongc.go create mode 100644 pkg/controllers/agentsession/sessiongc_test.go diff --git a/internal/cmd/operator/main.go b/internal/cmd/operator/main.go index c4d61da7..c37be848 100644 --- a/internal/cmd/operator/main.go +++ b/internal/cmd/operator/main.go @@ -158,7 +158,6 @@ import ( searchinmem "github.com/authzed/openagentprimitives/pkg/memory/search/inmem" pgsearch "github.com/authzed/openagentprimitives/pkg/memory/search/postgres" searchsqlite "github.com/authzed/openagentprimitives/pkg/memory/search/sqlite" - memshadow "github.com/authzed/openagentprimitives/pkg/memory/shadow" "github.com/authzed/openagentprimitives/pkg/memory/spicedbauthorizer" memsqlite "github.com/authzed/openagentprimitives/pkg/memory/sqlite" "github.com/authzed/openagentprimitives/pkg/memory/tokens" @@ -739,6 +738,7 @@ type config struct { defaultChannelArchiveAfter time.Duration defaultSessionSleepAfter time.Duration sessionStorageRetention time.Duration + sessionGCAfter time.Duration nodePinnedStorageReclaimGrace time.Duration idleStorageReclaimAfter time.Duration failedSandboxReapGrace time.Duration @@ -857,6 +857,8 @@ func newCommand() *cobra.Command { "how long an idle channel session keeps its pods warm before the operator reaps them (0 = never sleep)") fs.DurationVar(&cfg.sessionStorageRetention, "session-storage-retention", 72*time.Hour, "How long a terminal (Succeeded/Failed) AgentSession's workspace + snapshot-store PVCs are kept past finishedAt before deletion. The session record itself is kept. Per-class override on AgentClass.spec.channels.storageRetention. Set to 0 to disable the sweep.") + fs.DurationVar(&cfg.sessionGCAfter, "session-gc-after", 720*time.Hour, + "How long a terminal (Succeeded/Failed) AgentSession CR is kept past finishedAt before the WHOLE record is deleted — garbage-collecting the ToolCalls/pods/Secrets it owns (which otherwise accumulate unbounded in etcd and the operator's informer cache). The append-only audit log is NOT deleted (it lives in the durable memory backend and finalize keeps it). Must exceed --session-storage-retention. Set to 0 to disable session GC (keep every record forever).") fs.DurationVar(&cfg.nodePinnedStorageReclaimGrace, "node-pinned-storage-reclaim-grace", 15*time.Minute, "When the workspace class is the node-local bundled class (ap-workspace-rwx), CAP a terminal session's workspace + snapshot-store PVC retention at this short grace instead of --session-storage-retention. Bounds the node-local disk-pressure deadlock: local-path PV bytes live on one node's disk with no quota. Kept as a debug window (not instant), and trades away restart-from-here/fork for such sessions. Set to 0 to disable the cap (use the normal retention for node-local too).") fs.DurationVar(&cfg.idleStorageReclaimAfter, "idle-storage-reclaim-after", 30*time.Minute, @@ -979,6 +981,21 @@ func validateMemoryBackend(backend, postgresURI string) (string, error) { } } +// postgresMemoryBackend builds the memory Backend for MEMORY_BACKEND=postgres: +// the postgres backend DIRECTLY, with no in-memory shadow. +// +// The old shadow dual-wrote every append-only entry (transcript, audit, +// authz-decision, tool-session) into an inmem backend AS WELL as postgres, while +// serving reads only from postgres (reads=secondary). That inmem half was a +// write-only mirror nothing ever read, and scope-delete can never free it — +// append-only scopes are structurally undeletable — so it grew with the whole +// cluster's audit history for the operator's entire lifetime and was a primary +// OOM driver. Reading through the postgres backend directly is byte-for-byte the +// same data the shadow already served from its secondary, at none of the RAM. +func postgresMemoryBackend(pg *mempostgres.Client) memorypkg.Backend { + return mempostgres.NewBackend(pg) +} + // validateArtifactStoreURL enforces the fail-closed artifact store contract: // no silent in-memory fallback. `oap install` always injects a value; a bare // `kubectl apply` of the bundle intentionally lands here. @@ -1215,6 +1232,7 @@ func run(cfg *config) { "natsURL", cfg.natsURL, "defaultChannelArchiveAfter", cfg.defaultChannelArchiveAfter.String(), "defaultSessionSleepAfter", cfg.defaultSessionSleepAfter.String(), + "sessionGCAfter", cfg.sessionGCAfter.String(), "failedSandboxReapGrace", cfg.failedSandboxReapGrace.String(), "oauthRefreshThreshold", cfg.oauthRefreshThreshold.String(), "channelsdTokenFile", cfg.channelsdTokenFile, @@ -1734,28 +1752,23 @@ func run(cfg *config) { if err := pgsearch.Migrate(context.Background(), pgClient.Pool()); err != nil { log.Info("PostgreSQL search migration partial (pgvector may not be installed, text search still works)", "err", err.Error()) } - // Reads default to SECONDARY (durable postgres): the shadow still - // dual-writes inmem+postgres as a safety net, but the read path serves - // postgres so a pod roll doesn't wipe what admin UI / agents read. - // MEMORY_READ_SOURCE stays an explicit override — honored when set (e.g. - // 'primary' to temporarily read ephemeral inmem for debugging). - readSource := cfg.memoryReadSource - if readSource == "" { - readSource = memshadow.ReadFromSecondary - } - shadowBackend, err := memshadow.New( - memoryinmem.NewBackend(), - mempostgres.NewBackend(pgClient), - readSource, - log.WithName("shadow-backend"), - ) - if err != nil { - log.Error(err, "shadow backend config", "readSource", readSource) - os.Exit(1) + // Use the postgres backend DIRECTLY — no in-memory shadow. The shadow + // used to dual-write inmem+postgres and read from postgres (secondary), + // which made its inmem half a write-only mirror nothing read and that + // scope-delete could never free (append-only scopes are undeletable), so + // it grew with the whole cluster's audit history for the operator's + // lifetime and drove it OOM. See postgresMemoryBackend. + memBackend = postgresMemoryBackend(pgClient) + // MEMORY_READ_SOURCE selected the shadow's read side; with no shadow it is + // inert. Warn rather than fail so a Deployment still carrying the flag + // starts — but say plainly that 'primary' (read the ephemeral inmem) can + // no longer be honored, since there is no inmem to read. + if cfg.memoryReadSource != "" { + log.Info("MEMORY_READ_SOURCE is ignored under the direct postgres backend (no inmem shadow to read); remove it", + "memoryReadSource", cfg.memoryReadSource) } - memBackend = shadowBackend - log.Info("using postgres memory backend (shadow inmem+postgres, reads=secondary)", - "postgresURI", pgCfg.URI, "readSource", readSource) + log.Info("using postgres memory backend (direct, no inmem shadow)", + "postgresURI", pgCfg.URI) case memoryBackendInmem: memBackend = memoryinmem.NewBackend() log.Info("using in-memory backend (MEMORY_BACKEND=inmem)") @@ -2568,6 +2581,7 @@ func run(cfg *config) { DefaultSessionSleepAfter: cfg.defaultSessionSleepAfter, FailedSandboxReapGrace: cfg.failedSandboxReapGrace, DefaultSessionStorageRetention: cfg.sessionStorageRetention, + SessionGCAfter: cfg.sessionGCAfter, NodePinnedStorageReclaimGrace: cfg.nodePinnedStorageReclaimGrace, IdleStorageReclaimAfter: cfg.idleStorageReclaimAfter, SpiceDBDeleter: spiceDBClient, diff --git a/internal/cmd/operator/memory_backend_test.go b/internal/cmd/operator/memory_backend_test.go index 749d6417..bda39306 100644 --- a/internal/cmd/operator/memory_backend_test.go +++ b/internal/cmd/operator/memory_backend_test.go @@ -7,6 +7,8 @@ import ( "github.com/stretchr/testify/require" memoryinmem "github.com/authzed/openagentprimitives/pkg/memory/inmem" + mempostgres "github.com/authzed/openagentprimitives/pkg/memory/postgres" + memshadow "github.com/authzed/openagentprimitives/pkg/memory/shadow" memsqlite "github.com/authzed/openagentprimitives/pkg/memory/sqlite" ) @@ -85,6 +87,24 @@ func TestValidateMemoryBackend(t *testing.T) { } } +// TestPostgresMemoryBackendIsDirect pins Contributor-2's fix: under +// MEMORY_BACKEND=postgres the operator builds the postgres backend DIRECTLY, +// with no in-memory shadow. The shadow's inmem half was a write-only mirror +// (reads already came from postgres) that scope-delete can never free, so it +// grew with the whole cluster's audit history for the operator's lifetime and +// drove it OOM. A nil client is fine: NewBackend only wraps it, and this asserts +// the construction's TYPE, not any query. +func TestPostgresMemoryBackendIsDirect(t *testing.T) { + be := postgresMemoryBackend(nil) + require.NotNil(t, be, "postgres mode must construct a backend") + + _, isShadow := be.(*memshadow.Backend) + assert.False(t, isShadow, "postgres mode must NOT dual-write to an in-memory shadow") + + _, isPostgres := be.(*mempostgres.Backend) + assert.True(t, isPostgres, "postgres mode must use the postgres backend directly") +} + // TestInmemBackendConstructs asserts the inmem selector's construction path // yields a usable Backend (the postgres connect path is not unit-testable, so // the validation above covers its fail-closed rules). diff --git a/pkg/controllers/agentsession/controller.go b/pkg/controllers/agentsession/controller.go index 0f577a75..c0d5f5e5 100644 --- a/pkg/controllers/agentsession/controller.go +++ b/pkg/controllers/agentsession/controller.go @@ -578,6 +578,17 @@ type Reconciler struct { // 0 disables the sweep. From --session-storage-retention (72h). DefaultSessionStorageRetention time.Duration + // SessionGCAfter is how long a TERMINAL AgentSession is kept past + // status.finishedAt before the whole CR is deleted — the retention the pod + // reap comment anticipates ("until retention GCs the AgentSession"). Deleting + // the session cascade-removes its owned ToolCalls/pods/Secrets (owner-ref GC) + // and runs finalize, which revokes the token and SpiceDB grants but KEEPS the + // append-only audit log (DeleteScope refuses it). Measured from finishedAt + // (creation as a fallback). 0 or negative disables GC entirely. From + // --session-gc-after (720h / 30d). Must exceed DefaultSessionStorageRetention + // so storage is reclaimed well before the CR disappears. + SessionGCAfter time.Duration + // FailedSandboxReapGrace is how long a Failed session's sandbox pods are kept // for debugging before teardown. Measured from status.finishedAt; past it the // bundle SpiceboxSessions and runner pod are reaped so a dead session stops @@ -1223,6 +1234,22 @@ func (r *Reconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Resu // reapable and wakeable — a Failed one is never wakeable — so this narrows to // exactly the resume case. if !shouldWake(&sess) { + // Session GC supersedes the reap: a terminal session past its FULL + // retention window is deleted outright — the retention this block's reap + // comment anticipates ("until retention GCs the AgentSession"). Deleting + // the CR cascade-collects the ToolCalls/pods/Secrets it owns (which + // otherwise accumulate unbounded in etcd and the informer cache) and runs + // finalize, which keeps the append-only audit log. When it fires there is + // nothing left to reap or reclaim, so stop here. See reconcileSessionGC. + gcAfter, deleted, gerr := r.reconcileSessionGC(ctx, &sess) + if gerr != nil { + log.FromContext(ctx).Info("session GC failed; retried on the next reconcile", + "session", sess.Namespace+"/"+sess.Name, "err", gerr.Error()) + } + if deleted { + return ctrl.Result{}, nil + } + // Storage retention runs alongside the pod reap and must thread its // requeue through every branch below: after the pods are reaped, a // terminal session's reconciles keep short-circuiting in this block, @@ -1234,6 +1261,12 @@ func (r *Reconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Resu log.FromContext(ctx).Info("storage reclaim failed; retried on the next reconcile", "session", sess.Namespace+"/"+sess.Name, "err", rerr.Error()) } + // Fold the GC deadline into the same retention wakeup: once storage is + // reclaimed (72h) the GC deadline (far later) is the only event left for a + // terminal session, and nothing else in this block would schedule it. + if gcAfter > 0 && (reclaimAfter == 0 || gcAfter < reclaimAfter) { + reclaimAfter = gcAfter + } switch action, remaining := terminalReapAction(sess.Status.Phase, sess.Status.FinishedAt, sess.CreationTimestamp, r.FailedSandboxReapGrace, r.now()); action { case reapActionRequeue: if reclaimAfter > 0 && reclaimAfter < remaining { diff --git a/pkg/controllers/agentsession/sessiongc.go b/pkg/controllers/agentsession/sessiongc.go new file mode 100644 index 00000000..4209ee06 --- /dev/null +++ b/pkg/controllers/agentsession/sessiongc.go @@ -0,0 +1,88 @@ +// pkg/controllers/agentsession/sessiongc.go +// +// reconcileSessionGC implements the operator's session-lifetime GC: a TERMINAL +// (Succeeded/Failed) AgentSession whose SessionGCAfter window has elapsed past +// status.finishedAt has its whole CR deleted. This is the retention the pod +// reap anticipates — see controller.go's terminal block, "until retention GCs +// the AgentSession" — and the reason downstream consumers already tolerate a +// GC'd session (admind stamps such a session "Gone"). +// +// Deleting the CR is safe for the tamper-evident audit log. The append-only +// kinds (transcript, audit, authz-decision, tool-session) live in the durable +// memory backend, NOT on the CR, and finalize never deletes them +// (memory.Local.DeleteScope structurally refuses append-only kinds). The CR +// carries only the K8s-witnessed trust root (status.auditPublicKey/auditKeyID), +// which witnessAuditKey has already mirrored into the session's own durable +// memory scope precisely so records still verify after a delete-and-recreate. +// What a GC'd CR does give up is the truncation anchor (status.auditChainHeads); +// that is a verification aid, and losing it on long-dead sessions is the +// accepted cost of bounding unbounded CR growth. +// +// The bulk of what this reclaims is the owned ToolCall population: ToolCalls +// carry an AgentSession owner-ref with blockOwnerDeletion, so a session that is +// never deleted strands every ToolCall it ever spawned in etcd and in the +// operator's informer cache. Deleting the session cascade-collects them. +package agentsession + +import ( + "context" + "fmt" + "time" + + apierrors "k8s.io/apimachinery/pkg/api/errors" + "sigs.k8s.io/controller-runtime/pkg/log" + + spiceboxv1alpha1 "github.com/authzed/openagentprimitives/pkg/apis/v1alpha1" +) + +// reconcileSessionGC evaluates the session-lifetime deadline and deletes the +// whole AgentSession CR once past it. Returns: +// - requeueAfter: time until the GC deadline (0 = nothing further to +// schedule, e.g. disabled, non-terminal, or already deleted). +// - deleted: true when the session has been deleted (or was already +// deleting) and the caller should stop reconciling it this pass. +// - error: a delete failure (the caller logs; the next reconcile retries — +// the delete is idempotent). +func (r *Reconciler) reconcileSessionGC(ctx context.Context, sess *spiceboxv1alpha1.AgentSession) (time.Duration, bool, error) { + if r.SessionGCAfter <= 0 { + return 0, false, nil // disabled + } + if !isTerminalPhase(sess.Status.Phase) { + return 0, false, nil + } + // Already deleting (finalize in flight): nothing to schedule, and the caller + // should stop reconciling this pass. + if !sess.DeletionTimestamp.IsZero() { + return 0, true, nil + } + // Deadline from finishedAt, with creation as the fallback — the same + // reference reconcileStorageReclaim uses, so a terminal session that reached + // its phase without a finishedAt stamp is still collected instead of parking + // forever. + ref := sess.CreationTimestamp.Time + if sess.Status.FinishedAt != nil { + ref = sess.Status.FinishedAt.Time + } + deadline := ref.Add(r.SessionGCAfter) + now := r.now() + if now.Before(deadline) { + return deadline.Sub(now), false, nil + } + + // Past deadline: delete the whole CR. The resulting deletionTimestamp routes + // the next reconcile to finalize, which revokes the memory token and SpiceDB + // grants and KEEPS the append-only audit log; owner-ref GC cascade-collects + // the ToolCalls, pods and Secrets the session owns. NotFound is success — a + // prior pass (or another actor) already deleted it. + if err := r.Client.Delete(ctx, sess); err != nil { + if apierrors.IsNotFound(err) { + return 0, true, nil + } + return 0, false, fmt.Errorf("session GC: delete terminal session %s/%s: %w", sess.Namespace, sess.Name, err) + } + log.FromContext(ctx).Info("session GC: deleted terminal session past retention", + "session", sess.Namespace+"/"+sess.Name, + "phase", sess.Status.Phase, + "gcAfter", r.SessionGCAfter.String()) + return 0, true, nil +} diff --git a/pkg/controllers/agentsession/sessiongc_test.go b/pkg/controllers/agentsession/sessiongc_test.go new file mode 100644 index 00000000..82617c06 --- /dev/null +++ b/pkg/controllers/agentsession/sessiongc_test.go @@ -0,0 +1,114 @@ +// pkg/controllers/agentsession/sessiongc_test.go +// +// These tests pin the session-lifetime GC: a terminal AgentSession past its +// SessionGCAfter window has its whole CR deleted (cascade-collecting the +// ToolCalls it owns), while a disabled sweep, a non-terminal session, and a +// still-within-window session are all left untouched. The reconcile-level test +// proves the sweep is actually wired into the terminal block rather than only +// callable in isolation. +package agentsession + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + apierrors "k8s.io/apimachinery/pkg/api/errors" + "sigs.k8s.io/controller-runtime/pkg/client" + + spiceboxv1alpha1 "github.com/authzed/openagentprimitives/pkg/apis/v1alpha1" +) + +// sessionDeletionRequested reports whether the fixture's AgentSession has been +// deleted — either fully gone, or (with the finalizer still attached, as the +// fixture seeds it) marked with a deletionTimestamp that routes the next +// reconcile to finalize. +func sessionDeletionRequested(t *testing.T, f reapFixture) bool { + t.Helper() + var got spiceboxv1alpha1.AgentSession + err := f.c.Get(context.Background(), client.ObjectKey{Namespace: f.sess.Namespace, Name: f.sess.Name}, &got) + if apierrors.IsNotFound(err) { + return true + } + require.NoError(t, err, "get session to inspect deletion state") + return !got.DeletionTimestamp.IsZero() +} + +func TestReconcileSessionGC(t *testing.T) { + const day = 24 * time.Hour + cases := []struct { + name string + phase string + finishedDelta time.Duration // status.finishedAt relative to the fixture clock + gcAfter time.Duration + wantDeleted bool + wantRequeue bool // requeueAfter > 0 (a deadline is scheduled) + }{ + { + name: "disabled (gcAfter=0): a long-dead terminal session is never collected", + phase: spiceboxv1alpha1.AgentSessionPhaseSucceeded, + finishedDelta: -60 * day, gcAfter: 0, + wantDeleted: false, wantRequeue: false, + }, + { + name: "non-terminal (Idle): never collected, even past the window", + phase: spiceboxv1alpha1.AgentSessionPhaseIdle, + finishedDelta: -60 * day, gcAfter: 30 * day, + wantDeleted: false, wantRequeue: false, + }, + { + name: "terminal within window: kept, a deadline is scheduled", + phase: spiceboxv1alpha1.AgentSessionPhaseSucceeded, + finishedDelta: -1 * time.Hour, gcAfter: 30 * day, + wantDeleted: false, wantRequeue: true, + }, + { + name: "Succeeded past window: CR deleted", + phase: spiceboxv1alpha1.AgentSessionPhaseSucceeded, + finishedDelta: -31 * day, gcAfter: 30 * day, + wantDeleted: true, wantRequeue: false, + }, + { + name: "Failed past window: CR deleted", + phase: spiceboxv1alpha1.AgentSessionPhaseFailed, + finishedDelta: -31 * day, gcAfter: 30 * day, + wantDeleted: true, wantRequeue: false, + }, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + f := newReapFixture(t, tc.phase, time.Minute, tc.finishedDelta) + f.r.SessionGCAfter = tc.gcAfter + + requeue, deleted, err := f.r.reconcileSessionGC(context.Background(), f.sess) + require.NoError(t, err, "reconcileSessionGC") + + assert.Equal(t, tc.wantDeleted, deleted, "deleted result") + if tc.wantRequeue { + assert.Positive(t, requeue, "a within-window session must schedule its GC deadline") + } else { + assert.Zero(t, requeue, "no deadline expected") + } + assert.Equal(t, tc.wantDeleted, sessionDeletionRequested(t, f), + "the CR is deleted iff GC fired") + }) + } +} + +// TestReconcile_TerminalSession_GCdPastRetention proves the sweep is wired into +// the terminal reconcile path, not merely callable: a full Reconcile of a +// terminal session past its GC window must request the CR's deletion. +func TestReconcile_TerminalSession_GCdPastRetention(t *testing.T) { + const day = 24 * time.Hour + f := newReapFixture(t, spiceboxv1alpha1.AgentSessionPhaseSucceeded, time.Minute, -31*day) + f.r.SessionGCAfter = 30 * day + + require.False(t, sessionDeletionRequested(t, f), "precondition: the session is live") + + f.reconcile(t) + + assert.True(t, sessionDeletionRequested(t, f), + "a terminal session past its GC window must be deleted through Reconcile") +}