From 342dba1e3a912d526867a189e0791c2da6730c55 Mon Sep 17 00:00:00 2001 From: mintaka Date: Thu, 1 Oct 2026 20:43:10 -0400 Subject: [PATCH] feat(server): archive the transcript tail when a session ends without Stop (RIG-4135) When the hub sees a session end without Stop (the Runner reports a bound session unknown, or a re-enroll reaps its binding) it now runs the same session_end flush Stop does, detached and under the session's tenant. The PG tail is kept, so a later resume is unaffected; the manifest dedups an identical range. Co-authored-by: Matt Wilkinson --- go/e2e/leg_s3_durability_test.go | 56 +++++++-------------- go/internal/runnerhub/hub.go | 20 +++++++- go/internal/runnerhub/lost_session_test.go | 45 +++++++++++++++++ go/internal/runnerhub/relay_comms.go | 22 ++++++++ go/server/service.go | 16 +----- go/server/service_sessionend_pgtest_test.go | 52 +++++++++++++++++++ go/server/session_end_archive.go | 36 +++++++++++++ go/server/sinks.go | 1 + 8 files changed, 196 insertions(+), 52 deletions(-) create mode 100644 go/server/session_end_archive.go diff --git a/go/e2e/leg_s3_durability_test.go b/go/e2e/leg_s3_durability_test.go index 3118471c7..eb4dd7344 100644 --- a/go/e2e/leg_s3_durability_test.go +++ b/go/e2e/leg_s3_durability_test.go @@ -31,7 +31,7 @@ func TestLegS3DurabilityValves(t *testing.T) { {"valve_server_restart", runS3ServerRestart}, {"valve_runner_restart", runS3RunnerRestart}, {"s3_outage", runS3Outage}, - {"crash_leaves_no_session_end", runS3CrashNoSessionEnd}, + {"crash_archives_session_end", runS3CrashArchivesSessionEnd}, } { t.Run(tc.name, func(t *testing.T) { tc.run(t, ctx) }) } @@ -112,7 +112,6 @@ func runS3AgentCrash(t *testing.T, ctx context.Context) { if err := podmanRemoveForce(ctx, container); err != nil { t.Fatalf("crash agent container %q: %v", container, err) } - assertNoSessionEnd(t, ctx, f.DSN(), sessionID) resumeS3AndVerify(t, ctx, f, st, "s3-valve-agent-crash", sessionID, channel) } @@ -247,21 +246,26 @@ func runS3TurnSettled(ctx context.Context, f *Fixture, sessionID, channel, marke } func waitSafetyValveSegment(t *testing.T, ctx context.Context, dsn, sessionID string) { + t.Helper() + waitSegmentKind(t, ctx, dsn, sessionID, "safety_valve") +} + +func waitSegmentKind(t *testing.T, ctx context.Context, dsn, sessionID, kind string) { t.Helper() deadline := time.Now().Add(settleTimeout) ticker := time.NewTicker(transcriptPollInterval) defer ticker.Stop() for { - if segmentCountByKind(t, ctx, dsn, sessionID, "safety_valve") > 0 { + if segmentCountByKind(t, ctx, dsn, sessionID, kind) > 0 { return } if !time.Now().Before(deadline) { assertArchiveAndCheckpointDiagnostics(t, ctx, dsn, sessionID) - t.Fatalf("safety_valve segment did not appear within %s", settleTimeout) + t.Fatalf("%s segment did not appear within %s", kind, settleTimeout) } select { case <-ctx.Done(): - t.Fatalf("waiting for safety_valve segment: %v", ctx.Err()) + t.Fatalf("waiting for %s segment: %v", kind, ctx.Err()) case <-ticker.C: } } @@ -382,10 +386,11 @@ func assertSafetyValveEviction(t *testing.T, ctx context.Context, f *Fixture, se } } -// Crash does not archive the tail; archive-on-crash remains a design decision. -func runS3CrashNoSessionEnd(t *testing.T, ctx context.Context) { +// A lost container is detected on the next deliver; the hub then archives the tail +// as session_end without pruning it, so the woken session still resumes from PG. +func runS3CrashArchivesSessionEnd(t *testing.T, ctx context.Context) { t.Helper() - f, st, _, container, joined := runS3DurabilityBase(t, ctx, "s3-crash-no-end", nil) + f, _, _, container, joined := runS3DurabilityBase(t, ctx, "s3-crash-end", nil) sessionID, channel := splitDurabilityIDs(t, joined) if err := runS3TurnSettled(ctx, f, sessionID, channel, "durability-first"); err != nil { t.Fatalf("first turn: %v", err) @@ -393,21 +398,14 @@ func runS3CrashNoSessionEnd(t *testing.T, ctx context.Context) { if err := podmanRemoveForce(ctx, container); err != nil { t.Fatalf("crash agent container: %v", err) } - assertNoSessionEnd(t, ctx, f.DSN(), sessionID) - var minHotSeq int64 - pool, err := pgxpool.New(ctx, f.DSN()) - if err != nil { - t.Fatalf("open crash tail pool: %v", err) + if _, err := f.PostMessage(ctx, channel, "general", "marker durability-after-crash"); err != nil { + t.Fatalf("post after crash: %v", err) } - if err := pool.QueryRow(ctx, `SELECT COALESCE(MIN(entry_seq), 0) FROM agent_session_transcript_entries WHERE session_id=$1`, sessionID).Scan(&minHotSeq); err != nil { - pool.Close() - t.Fatalf("query crash hot tail min: %v", err) - } - pool.Close() - if minHotSeq == 0 { - t.Fatal("crash left no transcript entries in PostgreSQL hot tail") + waitSegmentKind(t, ctx, f.DSN(), sessionID, "session_end") + assertArchiveHasReply(t, ctx, f, sessionID, "durability first reply") + if segmentCountByKind(t, ctx, f.DSN(), sessionID, "superseded") != 0 { + t.Fatal("detected-end archive pruned the tail as superseded") } - resumeS3AndVerify(t, ctx, f, st, "s3-crash-no-end", sessionID, channel) } func resumeS3AndVerify(t *testing.T, ctx context.Context, f *Fixture, st *store.Store, accountHandle, sessionID, channel string) { @@ -438,22 +436,6 @@ func resumeS3AndVerify(t *testing.T, ctx context.Context, f *Fixture, st *store. assertArchiveHasReply(t, ctx, f, sessionID, "durability first reply") } -func assertNoSessionEnd(t *testing.T, ctx context.Context, dsn, sessionID string) { - t.Helper() - pool, err := pgxpool.New(ctx, dsn) - if err != nil { - t.Fatalf("open assertion pool: %v", err) - } - defer pool.Close() - var count int - if err := pool.QueryRow(ctx, `SELECT count(*) FROM agent_session_archive_segments WHERE session_id=$1 AND kind='session_end'`, sessionID).Scan(&count); err != nil { - t.Fatalf("query session_end segments: %v", err) - } - if count != 0 { - t.Fatalf("session_end segments after crash = %d, want 0", count) - } -} - // assertArchiveHasReply requires want in an archived S3 object, not just in the PG tail. func assertArchiveHasReply(t *testing.T, ctx context.Context, f *Fixture, sessionID, want string) { t.Helper() diff --git a/go/internal/runnerhub/hub.go b/go/internal/runnerhub/hub.go index 19a50858c..6b602d6d8 100644 --- a/go/internal/runnerhub/hub.go +++ b/go/internal/runnerhub/hub.go @@ -130,6 +130,13 @@ type SessionLostSink interface { OnSessionLost(sessionID string, account store.AccountID) } +// SessionEndSink archives the transcript tail of a session the hub saw end without +// Stop (a lost container or a re-enroll reap). ctx is scoped to the session tenant. +// It may block on object-store I/O, so the hub calls it off its locks and loops. +type SessionEndSink interface { + OnSessionEnded(ctx context.Context, sessionID string) +} + // PresenceSink is notified of the two hub-side edges the RIG-1569 T8 presence // projection (design record D4) is fed by: a session lifecycle transition at the // deliverSession arm (the SAME arm SettleSink rides, right after the @@ -351,7 +358,8 @@ type Hub struct { // a pre-T9 Runner link loss left behind. Nil until SetSessionReapSink; read under mu. reap SessionReapSink // lost is notified after dropLostSession releases a dead session. Read under mu. - lost SessionLostSink + lost SessionLostSink + ended SessionEndSink // presence is the RIG-1569 T8 presence projection's sink, notified at // deliverSession (lifecycle transition) and promoteSession (reconciliation). Nil // until SetPresenceSink; read under mu. Nil-safe (today's behavior). @@ -518,6 +526,13 @@ func (h *Hub) SetSessionReapSink(reap SessionReapSink) { h.reap = reap } +// SetSessionEndSink wires the transcript archive for sessions that end without Stop. +func (h *Hub) SetSessionEndSink(end SessionEndSink) { + h.mu.Lock() + defer h.mu.Unlock() + h.ended = end +} + // SetSessionLostSink wires the delivery consumer to wake an agent whose session died. func (h *Hub) SetSessionLostSink(lost SessionLostSink) { h.mu.Lock() @@ -1029,6 +1044,9 @@ func (h *Hub) enroll(ctx context.Context, id string, subject store.Subject, tier if reap != nil { reap.OnSessionsReaped(reapedSessions) } + for _, sessionID := range reapedSessions { + h.archiveEnded(ctx, sessionID) + } return reattached } diff --git a/go/internal/runnerhub/lost_session_test.go b/go/internal/runnerhub/lost_session_test.go index 544415f60..806e1c540 100644 --- a/go/internal/runnerhub/lost_session_test.go +++ b/go/internal/runnerhub/lost_session_test.go @@ -5,6 +5,7 @@ package runnerhub import ( "context" "testing" + "time" compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" "github.com/RigelBuild/compass/go/internal/store" @@ -45,3 +46,47 @@ func TestDropLostSessionOnlyForOwningRunner(t *testing.T) { t.Fatalf("lost = %v, want [%s]", sink.lost, testAgentAccount) } } + +type chanEndSink chan string + +func (c chanEndSink) OnSessionEnded(_ context.Context, sessionID string) { c <- sessionID } + +func recvEnded(t *testing.T, c chanEndSink) string { + t.Helper() + select { + case s := <-c: + return s + case <-time.After(10 * time.Second): + t.Fatal("no session-end archive within 10s") + return "" + } +} + +// A session that ends without Stop is archived: on a lost-session drop by its owning +// Runner, and on each binding the re-enroll reap removes. A foreign refusal is not. +func TestSessionEndedWithoutStopIsArchived(t *testing.T) { + ctx := t.Context() + hub := newHubOnly() + bindings := newFakeBindingStore() + hub.SetSessionBindingStore(bindings) + ended := make(chanEndSink, 4) + hub.SetSessionEndSink(ended) + subj := runnerSubject() + tier, egress := compassv1.RuntimeTier_RUNTIME_TIER_UNSPECIFIED, compassv1.EgressPosture_EGRESS_POSTURE_UNSPECIFIED + hub.enroll(ctx, "runner-1", subj, tier, egress) + + hub.bindContainer("cont-1", testAgentAccount, "runner-1") + hub.promoteSession(ctx, "cont-1", "sess-lost") + hub.dropLostSession(ctx, "runner-other", "sess-lost") + hub.dropLostSession(ctx, "runner-1", "sess-lost") + if got := recvEnded(t, ended); got != "sess-lost" { + t.Fatalf("archived %q after the lost drop, want sess-lost", got) + } + + hub.bindContainer("cont-2", testAgentAccount, "runner-1") + hub.promoteSession(ctx, "cont-2", "sess-reaped") + hub.enroll(ctx, "runner-1", subj, tier, egress) + if got := recvEnded(t, ended); got != "sess-reaped" { + t.Fatalf("archived %q after the re-enroll reap, want sess-reaped (the foreign refusal must not archive)", got) + } +} diff --git a/go/internal/runnerhub/relay_comms.go b/go/internal/runnerhub/relay_comms.go index 6dcb27a8f..cfea960e0 100644 --- a/go/internal/runnerhub/relay_comms.go +++ b/go/internal/runnerhub/relay_comms.go @@ -13,6 +13,7 @@ import ( "context" "errors" "maps" + "time" "connectrpc.com/connect" "go.opentelemetry.io/otel/attribute" @@ -878,6 +879,7 @@ func (h *Hub) dropLostSession(ctx context.Context, runnerID, sessionID string) { h.mu.Lock() lost := h.lost h.mu.Unlock() + h.archiveEnded(ctx, sessionID) h.log.Warn("runner reports bound session unknown; released binding to wake the agent", "session_id", sessionID, "agent_account_id", account) if lost != nil { @@ -915,3 +917,23 @@ func (h *Hub) runnerSessionCtx(ctx context.Context, runnerID, sessionID string) } return ctx, true } + +// archiveEnded hands a session that ended without Stop to the transcript archive. +// It runs detached and bounded: the caller is a Runner stream loop, and an +// object-store PUT must not stall it nor die with the stream. +func (h *Hub) archiveEnded(ctx context.Context, sessionID string) { + h.mu.Lock() + ended := h.ended + h.mu.Unlock() + if ended == nil || sessionID == "" { + return + } + actx, cancel := context.WithTimeout(context.WithoutCancel(ctx), sessionEndArchiveTimeout) + go func() { + defer cancel() + ended.OnSessionEnded(actx, sessionID) + }() +} + +// sessionEndArchiveTimeout bounds one detached session-end archive. +const sessionEndArchiveTimeout = 2 * time.Minute diff --git a/go/server/service.go b/go/server/service.go index e4abe571c..31f061174 100644 --- a/go/server/service.go +++ b/go/server/service.go @@ -211,21 +211,9 @@ func (s *service) StopAgentSession( if err != nil { return nil, err } - // RIG-1667 T4 session-end flush: archive the remaining hot-tail as one session_end - // segment so history is COMPLETE for analytics (not pruned, never read on resume). // BEST-EFFORT: the Stop relay already killed the agent, so a flush failure must - // NEVER convert a successful Stop into a failure. - sessionID := req.Msg.GetSessionId() - if maxSeq, seqErr := s.store.SessionMaxEntrySeq(ctx, sessionID); seqErr != nil { - slog.ErrorContext(ctx, "session-end transcript flush skipped: could not read max entry seq", - "session_id", sessionID, "error", seqErr) - } else if maxSeq > 0 { - // maxSeq == 0 means no transcript rows: nothing to archive, skip the flush. - if flushErr := s.store.FlushSuperseded(ctx, sessionID, maxSeq, store.SegmentKindSessionEnd); flushErr != nil { - slog.ErrorContext(ctx, "session-end transcript flush failed; Stop still succeeded", - "session_id", sessionID, "upto_entry_seq", maxSeq, "error", flushErr) - } - } + // never convert a successful Stop into a failure. + archiveSessionEnd(ctx, s.store, req.Msg.GetSessionId()) return connect.NewResponse(resp), nil } diff --git a/go/server/service_sessionend_pgtest_test.go b/go/server/service_sessionend_pgtest_test.go index 145c7f87a..a7561f248 100644 --- a/go/server/service_sessionend_pgtest_test.go +++ b/go/server/service_sessionend_pgtest_test.go @@ -9,6 +9,7 @@ package server import ( "context" + "runtime" "testing" "connectrpc.com/connect" @@ -177,3 +178,54 @@ func (f *serverFakeObjectStore) GetSegment(_ context.Context, key string) ([]byt } return body, nil } + +// TestReenrollReapArchivesSessionEnd pins archive-on-detected-end: a session the +// re-enroll reap removes (it ended without Stop) gets one session_end segment over +// its whole tail, and the PG tail is kept for a later resume. +func TestReenrollReapArchivesSessionEnd(t *testing.T) { + f := newPlacementFixture(t) + ctx := context.Background() // the test root context + f.store.SetObjectStore(newServerFakeObjectStore()) + // Provision binds the container to the agent, so Start binds the session and the + // reap has a row to remove. + if _, _, err := f.hub.Provision(ctx, "", f.agentID, &compassv1.ProvisionAgentWorkspaceRequest{}); err != nil { + t.Fatalf("Provision: %v", err) + } + startBoundSession(t, f, ctx) + if got, ok := f.hub.SessionForAccount(ctx, f.agentID); !ok || got != fakeSessionID { + t.Fatalf("precondition: SessionForAccount = (%q, %v), want (%q, true)", got, ok, fakeSessionID) + } + for seq := uint64(1); seq <= 2; seq++ { + if err := f.store.AppendTranscriptEntry(ctx, fakeSessionID, seq, seq == 1, `{"e":1}`, "reap-k"+string(rune('0'+seq))); err != nil { + t.Fatalf("append %d: %v", seq, err) + } + } + + // The same Runner id dials again: a re-enroll, which reaps every binding it held. + attachFakeRunner(t, f.store, f.hub, false) + + conn := connectPG(t, ctx, f.dsn) + deadline := timeAfter() + for { + var n int + if err := conn.QueryRow(ctx, `SELECT COUNT(*) FROM agent_session_archive_segments WHERE session_id = $1 AND kind = 'session_end'`, fakeSessionID).Scan(&n); err != nil { + t.Fatalf("count session_end segments: %v", err) + } + if n > 0 { + segs := sessionEndSegments(t, ctx, f.dsn, fakeSessionID) + if len(segs) != 1 || segs[0].minSeq != 1 || segs[0].maxSeq != 2 { + t.Fatalf("session_end segments = %+v, want one span [1..2]", segs) + } + break + } + select { + case <-deadline: + t.Fatalf("session_end segments = %d after the re-enroll reap, want 1 within %s", n, testTimeout) + default: + } + runtime.Gosched() + } + if got := transcriptRowCount(t, ctx, f.dsn, fakeSessionID); got != 2 { + t.Fatalf("transcript rows = %d, want 2 (session_end never prunes)", got) + } +} diff --git a/go/server/session_end_archive.go b/go/server/session_end_archive.go new file mode 100644 index 000000000..cd9ddb11e --- /dev/null +++ b/go/server/session_end_archive.go @@ -0,0 +1,36 @@ +package server + +import ( + "context" + "log/slog" + + "github.com/RigelBuild/compass/go/internal/store" +) + +// archiveSessionEnd archives a session's remaining hot tail as one session_end +// segment, so history is complete for analytics. The PG tail is not pruned and +// stays authoritative for resume. Best-effort: failures are logged. +func archiveSessionEnd(ctx context.Context, st *store.Store, sessionID string) { + maxSeq, err := st.SessionMaxEntrySeq(ctx, sessionID) + if err != nil { + slog.ErrorContext(ctx, "session-end transcript flush skipped: could not read max entry seq", + "session_id", sessionID, "error", err) + return + } + if maxSeq == 0 { + return // no transcript rows: nothing to archive + } + if err := st.FlushSuperseded(ctx, sessionID, maxSeq, store.SegmentKindSessionEnd); err != nil { + slog.ErrorContext(ctx, "session-end transcript flush failed", + "session_id", sessionID, "upto_entry_seq", maxSeq, "error", err) + } +} + +// sessionEndArchiver archives sessions the hub saw end without Stop: a lost +// container or a re-enroll reap. A later resume and Stop archive again; the +// manifest dedups an identical range. +type sessionEndArchiver struct{ st *store.Store } + +func (a sessionEndArchiver) OnSessionEnded(ctx context.Context, sessionID string) { + archiveSessionEnd(ctx, a.st, sessionID) +} diff --git a/go/server/sinks.go b/go/server/sinks.go index 033d1893f..a98d33b1d 100644 --- a/go/server/sinks.go +++ b/go/server/sinks.go @@ -63,6 +63,7 @@ func newRunnerHub(st *store.Store, brd *board.Projection, tail runnerhub.Session // seam (SessionResumeSnapshot + ReadArchiveSegment), wired here beside the // write seam so the one store instance serves both legs. hub.SetTranscriptReader(st) + hub.SetSessionEndSink(sessionEndArchiver{st: st}) return hub }