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") +}