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
56 changes: 19 additions & 37 deletions go/e2e/leg_s3_durability_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) })
}
Expand Down Expand Up @@ -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)
}

Expand Down Expand Up @@ -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:
}
}
Expand Down Expand Up @@ -382,32 +386,26 @@ 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)
}
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) {
Expand Down Expand Up @@ -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()
Expand Down
20 changes: 19 additions & 1 deletion go/internal/runnerhub/hub.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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).
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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
}

Expand Down
45 changes: 45 additions & 0 deletions go/internal/runnerhub/lost_session_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
}
}
22 changes: 22 additions & 0 deletions go/internal/runnerhub/relay_comms.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (
"context"
"errors"
"maps"
"time"

"connectrpc.com/connect"
"go.opentelemetry.io/otel/attribute"
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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
16 changes: 2 additions & 14 deletions go/server/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand Down
52 changes: 52 additions & 0 deletions go/server/service_sessionend_pgtest_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ package server

import (
"context"
"runtime"
"testing"

"connectrpc.com/connect"
Expand Down Expand Up @@ -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)
}
}
36 changes: 36 additions & 0 deletions go/server/session_end_archive.go
Original file line number Diff line number Diff line change
@@ -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)
}
1 change: 1 addition & 0 deletions go/server/sinks.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand Down
Loading