Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
58 changes: 36 additions & 22 deletions internal/cmd/operator/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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)")
Expand Down Expand Up @@ -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,
Expand Down
20 changes: 20 additions & 0 deletions internal/cmd/operator/memory_backend_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand Down Expand Up @@ -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).
Expand Down
33 changes: 33 additions & 0 deletions pkg/controllers/agentsession/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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 {
Expand Down
88 changes: 88 additions & 0 deletions pkg/controllers/agentsession/sessiongc.go
Original file line number Diff line number Diff line change
@@ -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
}
Loading
Loading