Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
dd7bfae
fix(runnerhub): ignore lifecycle frames that arrive after ERRORED (RI…
rigel-mintaka Oct 6, 2026
2e8c963
fix(runnerhub): key the ERRORED guard on RunnerSeq, not recovery admi…
rigel-mintaka Oct 6, 2026
f2db0e3
fix(runner): share RunnerSeq across Gateways and fence ERRORED by enr…
rigel-mintaka Oct 6, 2026
8a2ac88
fix(runnerhub): fence session frames by stream enrollment; let late s…
rigel-mintaka Oct 6, 2026
02d5f5f
fix(runnerhub): serialize session delivery with enroll; scope gap tra…
rigel-mintaka Oct 6, 2026
5a2733b
fix(runnerhub): drop any lifecycle frame older than the session's lat…
rigel-mintaka Oct 6, 2026
1ae2220
docs(gateway): state when a burned RunnerSeq shows as a gap (RIG-4673)
rigel-mintaka Oct 6, 2026
76b67cb
fix(runnerhub): admit lifecycle frames in seq order; pair router and …
rigel-mintaka Oct 6, 2026
ca21ef3
fix(runnerhub): serialize lifecycle frames per session, not hub-wide …
rigel-mintaka Oct 6, 2026
73c1eb4
fix(runnerhub): ERRORED cleanup releases only the binding it saw (RIG…
rigel-mintaka Oct 6, 2026
6a71170
fix(store): version session bindings with a UUID column, not xmin (RI…
rigel-mintaka Oct 6, 2026
ce396ae
fix(runnerhub): a version-less cleanup deletes no durable row (RIG-4742)
rigel-mintaka Oct 6, 2026
a0d3c6a
fix(runnerhub): a re-bound version-less cleanup is not a loss; backfi…
rigel-mintaka Oct 7, 2026
c1e204a
fix(runnerhub): drop the stale account entry when a peer took the ses…
rigel-mintaka Oct 7, 2026
5822a45
test(store): check usage events on a versioned binding delete (RIG-4742)
rigel-mintaka Oct 7, 2026
260e55e
fix(runnerhub): version reverse read-through bindings; order Stop wit…
rigel-mintaka Oct 7, 2026
d0567c7
Merging 260e55e414d9bca06e6717e4faccaa34fd748d33 into trunk-temp/pr-1…
trunk-io[bot] Oct 8, 2026
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
16 changes: 11 additions & 5 deletions go/internal/runner/gateway/gateway.go
Original file line number Diff line number Diff line change
Expand Up @@ -174,13 +174,12 @@ type Gateway struct {
// across the whole event stream and the hub only flags seq > lastSeq+1, so
// replayed low seqs are accepted and loss in that range stops being detectable.

// Scope is per-Runner-link: relay.go's eventPublisher owns a SECOND counter and
// both feed the hub's one high-water mark, so gap detection is meaningful only
// while exactly one is live. This path replaced the stdout relay for gateway
// traffic, so that holds today; unifying the two is T9.
// In production every Gateway of one Runner link shares deps.Seq, so RunnerSeq
// stays monotonic across container sockets. The hub's stale-ERRORED guard
// relies on that: a resumed session's new socket must not restart at 1.
pubMu sync.Mutex
pub *sessionPublisher
seq seqCounter
seq *SeqCounter

// publishMu fences telemetry admission while a lifecycle terminal state is sent.
publishMu sync.Mutex
Expand Down Expand Up @@ -222,6 +221,8 @@ type Deps struct {
Events EventRelay
// Committer forwards a durable conversation frame to the Server for commit (CommitConversationFrame).
Committer ConversationCommitter
// Seq is the Runner link's RunnerSeq allocator; nil gives this Gateway its own.
Seq *SeqCounter
}

// NewGateway builds the AgentGateway handler for the container's socket:
Expand All @@ -240,6 +241,10 @@ type Deps struct {
// the stream outlives any one agent request. A caller with no distinct socket
// scope (a hermetic test) passes context.Background().
func NewGateway(baseCtx context.Context, containerName string, deps Deps) *Gateway {
seq := deps.Seq
if seq == nil {
seq = &SeqCounter{}
}
return &Gateway{
baseCtx: baseCtx,
containerName: containerName,
Expand All @@ -250,6 +255,7 @@ func NewGateway(baseCtx context.Context, containerName string, deps Deps) *Gatew
board: deps.Board,
events: deps.Events,
committer: deps.Committer,
seq: seq,
control: noopControlRouter{},
// ttl=0: no expiry, a pure size-bounded LRU (committedKeysMax). The cache
// is advisory, so eviction is safe — it never drops the store's boundary.
Expand Down
34 changes: 16 additions & 18 deletions go/internal/runner/gateway/publisher.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,9 @@ package gateway
// emission order, else the hub's gap detector records a false gap. So a publisher
// holds its stream mutex across BOTH the seq allocation AND the Send.

// Scope: the counter is Gateway-scoped (per socket / per Runner link) and survives
// a publisher replacement, since a per-publisher counter would restart the
// sequence on a swap. relay.go's eventPublisher owns a SECOND counter feeding the
// same high-water mark, so gap detection holds only while one is live; unifying is T9.
// Scope: the counter is Runner-link-wide (shared by every container's Gateway) and
// survives a publisher replacement, since a per-publisher counter would restart the
// sequence on a swap.

import (
"context"
Expand All @@ -32,27 +31,26 @@ type EventRelay interface {
PublishEvents(ctx context.Context) *connect.ClientStreamForClient[compassv1internal.PublishEventsRequest, compassv1internal.PublishEventsResponse]
}

// seqCounter is the RunnerSeq allocator shared by every publisher a Gateway
// builds for its socket. It carries ONLY the counter and the lock that guards
// the counter — deliberately not the publishers' stream lock.
// SeqCounter is the RunnerSeq allocator shared by every publisher of every
// Gateway on one Runner link. It carries ONLY the counter and the lock that
// guards the counter — deliberately not the publishers' stream lock.
//
// Sharing the counter is required: a publisher is replaceable within one session,
// and a per-publisher counter restarts the sequence on that swap. Sharing the
// STREAM lock as well is not, and is actively harmful: acquirePublisher installs
// a replacement and then closes the stale publisher outside pubMu, so a
// CloseAndReceive round-trip against an unresponsive-but-connected Server would
// block every forward on the live replacement's separate upstream stream. Two
// distinct streams need no mutual ordering — the hub keeps one global high-water
// mark and cannot observe an interleaving between them — so the coupling would
// buy nothing and cost unbounded liveness.
type seqCounter struct {
// distinct streams may therefore deliver out of seq order; the hub tolerates a
// late lower seq, so the coupling would buy nothing and cost unbounded liveness.
type SeqCounter struct {
mu sync.Mutex
n uint64
}

// next allocates and returns the next sequence. Called with the publisher's own
// stream lock held, so allocation order still equals emission order.
func (c *seqCounter) next() uint64 {
func (c *SeqCounter) next() uint64 {
c.mu.Lock()
defer c.mu.Unlock()
c.n++
Expand All @@ -64,9 +62,9 @@ func (c *seqCounter) next() uint64 {
// permanent hole, and the hub flags a skipped number as in-transit loss
// (runnerhub/hub.go:230). A durable frame erring back to the agent is correct,
// expected behaviour — it must not make the Server report a loss that did not
// happen. Safe because the caller holds its stream lock across allocate-and-send,
// so no other goroutine can have sent past this value.
func (c *seqCounter) rollback(seq uint64) {
// happen. Only the latest seq is reclaimed: if another Gateway allocated since,
// the number stays burned and shows as a gap once the hub has a baseline.
func (c *SeqCounter) rollback(seq uint64) {
c.mu.Lock()
defer c.mu.Unlock()
if c.n == seq {
Expand Down Expand Up @@ -99,7 +97,7 @@ type sessionPublisher struct {
// seq is the Gateway's shared allocator, carried across publishers so the
// sequence survives a replacement. Only the counter is shared; its lock is
// held just long enough to allocate.
seq *seqCounter
seq *SeqCounter
// admit, when set, allocates under the Gateway's publish gate so a sealed
// session refuses frames instead of sequencing them after its ERRORED report.
admit func() (uint64, error)
Expand All @@ -122,7 +120,7 @@ type sessionPublisher struct {
// life: the Publish handler uses the socket-lifetime context, while lifecycle
// reports use a bounded caller ctx. seq is the Gateway counter, carried across
// publishers so the sequence never restarts.
func newSessionPublisher(ctx context.Context, relay EventRelay, sessionID string, seq *seqCounter) *sessionPublisher {
func newSessionPublisher(ctx context.Context, relay EventRelay, sessionID string, seq *SeqCounter) *sessionPublisher {
ctx, cancel := context.WithCancel(ctx)
return &sessionPublisher{
sessionID: sessionID,
Expand Down Expand Up @@ -230,7 +228,7 @@ func (g *Gateway) acquirePublisher(sessionID string) *sessionPublisher {
g.pub = nil
}
if g.pub == nil {
g.pub = newSessionPublisher(g.baseCtx, g.events, sessionID, &g.seq)
g.pub = newSessionPublisher(g.baseCtx, g.events, sessionID, g.seq)
g.pub.admit = func() (uint64, error) { return g.admitFrame(sessionID) }
g.pub.afterAdmit = g.afterAdmit
}
Expand Down
2 changes: 1 addition & 1 deletion go/internal/runner/gateway/session_state.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,7 @@ func (l *SocketListener) PublishSessionState(ctx context.Context, sessionID stri
}
sendCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), stateSendTimeout)
defer cancel()
pub := newSessionPublisher(sendCtx, g.events, sessionID, &g.seq)
pub := newSessionPublisher(sendCtx, g.events, sessionID, g.seq)
forwardErr := pub.forward(frame)
if forwardErr == nil {
sharedErr = nil
Expand Down
24 changes: 23 additions & 1 deletion go/internal/runner/gateway/telemetry_ingest_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1169,7 +1169,7 @@ func (s *seqSink) seqs() []uint64 {
// a variant defect that restarts at a value which happens not to collide would
// pass a duplicates-only check. [1 2] pins the contract.
//
// RED: give newSessionPublisher its own &seqCounter{} instead of the Gateway's
// RED: give newSessionPublisher its own &SeqCounter{} instead of the Gateway's
// -> seqs = [1 1], and this fails every run.
func TestSequenceSurvivesPublisherReplacement(t *testing.T) {
sink := &seqSink{}
Expand Down Expand Up @@ -1204,6 +1204,28 @@ func TestSequenceSurvivesPublisherReplacement(t *testing.T) {
}
}

// A resumed session gets a new container socket, so a per-Gateway counter would
// restart at 1 and fall under the hub's stale-ERRORED boundary.
//
// RED: ignore deps.Seq in NewGateway -> seqs = [1 1].
func TestSequenceSharedAcrossGateways(t *testing.T) {
sink := &seqSink{}
events := newRunnerServiceServer(t, sink)
shared := &SeqCounter{}
for _, name := range []string{"cont-1", "cont-2"} {
g := NewGateway(context.Background(), name, Deps{Sessions: boundSessions(), Events: events, Seq: shared})
if err := g.acquirePublisher("sess-1").forward(traceFrame(name)); err != nil {
t.Fatalf("forward on %s = %v, want success", name, err)
}
if err := releaseCurrentPublisher(g); err != nil {
t.Fatalf("release on %s = %v", name, err)
}
}
if got, want := sink.seqs(), []uint64{1, 2}; !slices.Equal(got, want) {
t.Fatalf("RunnerSeq sequence = %v, want %v across Gateways sharing one counter", got, want)
}
}

// A failed forward must not burn a sequence number. The counter is
// socket-lifetime, so an allocated-but-unsent number is a permanent hole, and
// the hub reads a skipped number as in-transit loss (runnerhub/hub.go:230) — so
Expand Down
7 changes: 5 additions & 2 deletions go/internal/runner/host.go
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,9 @@ type agentHost struct {
mu sync.Mutex
sessions map[string]*liveSession
sockets map[string]*gateway.SocketListener
// runnerSeq is shared by every container Gateway so RunnerSeq never restarts
// within one enrollment; the hub drops lifecycle frames at or below ERRORED's.
runnerSeq gateway.SeqCounter
// afterExitCheck is a test seam between detecting exit and acquiring the
// container lock; its returned func runs when retireOnExit returns. Nil in production.
afterExitCheck func() func()
Expand Down Expand Up @@ -871,7 +874,7 @@ func (h *agentHost) provisionVsockGateway(ctx context.Context, spec runtime.Agen
h.teardownContainer(ctx, name)
return "", fmt.Errorf("resolving vsock gateway endpoint for container %q: backend reports no session", name)
}
deps := gateway.Deps{Sessions: h, Relay: h.link.client, Lifecycle: h.link.client, Events: h.link.client, Committer: h.link.client, Forge: h.link.client, Board: h.link.client}
deps := gateway.Deps{Sessions: h, Relay: h.link.client, Lifecycle: h.link.client, Events: h.link.client, Committer: h.link.client, Forge: h.link.client, Board: h.link.client, Seq: &h.runnerSeq}
h.log.InfoContext(ctx, "serving agent gateway over vsock path",
slog.String("container", name), slog.String("path", endpoint))
listener, err := gateway.Serve(ctx, endpoint, name, deps)
Expand Down Expand Up @@ -1250,7 +1253,7 @@ func (h *agentHost) serveSocket(ctx context.Context, containerName string) (*gat
// serveSocketAt is serveSocket with an explicit socket path: container tiers pass the
// fixed RuntimeDir socket, the host tier a path in the handle's state dir (no mount).
func (h *agentHost) serveSocketAt(ctx context.Context, containerName, path string) (*gateway.SocketListener, error) {
listener, err := gateway.Serve(ctx, path, containerName, gateway.Deps{Sessions: h, Relay: h.link.client, Lifecycle: h.link.client, Events: h.link.client, Committer: h.link.client, Forge: h.link.client, Board: h.link.client})
listener, err := gateway.Serve(ctx, path, containerName, gateway.Deps{Sessions: h, Relay: h.link.client, Lifecycle: h.link.client, Events: h.link.client, Committer: h.link.client, Forge: h.link.client, Board: h.link.client, Seq: &h.runnerSeq})
if err != nil {
return nil, fmt.Errorf("serving agent socket for container %q: %w", containerName, err)
}
Expand Down
Loading
Loading