diff --git a/go/internal/delivery/consumer.go b/go/internal/delivery/consumer.go index 77d85f315..4243994b7 100644 --- a/go/internal/delivery/consumer.go +++ b/go/internal/delivery/consumer.go @@ -142,6 +142,9 @@ type DeliveryReads interface { //nolint:interfacebloat // one method per store r type settleEvent struct { sessionID string state compassv1.AgentSessionState + // upTo bounds the commit times this edge fires. A real settle fires all; a + // late hold's replay fires only its settled turn, not a still-streaming one. + upTo int64 } // startEvent is one queued session-start edge handed from the hub's Start (or @@ -169,6 +172,7 @@ type heldEntry struct { messageID string traceparent string tenant store.TenantID + atUnixMs int64 // commit time, matched against settleEvent.upTo } // Consumer consumes message_posted refs and fans posted messages out to @@ -193,26 +197,17 @@ type Consumer struct { mu sync.Mutex // held is the pending-deliver registry (design.md:157-168), keyed by the // AUTHOR's live session id: an agent-authored message posted while its author - // still streams is HELD here until that author's session settles - // (WORKING->READY) or reaches a terminal frame. The value is the ordered set - // of message ids held for that author, in post order, so a settle fires them - // ascending. A no-frame author death (no settle edge ever enqueues) leaves - // its entry here until it is reaped: the reap happens in-process on the next - // Runner (re-)enroll via the hub's SessionReapSink (OnSessionsReaped), which - // drops the entry for every session id whose hub binding enroll just cleared. - // So the common no-frame death is reaped at that next enroll rather than - // persisting until process restart. The reap is best-effort, NOT a hard - // bound: the reaped set is exactly the session ids bound at enroll time, and a - // no-frame-dead session is never re-promoted (its id, once cleared, never - // re-enters the hub's session map), so a narrow race can still strand one - // entry until process restart — a Deliver that resolved the author LIVE an - // instant before enroll cleared the maps can hold(sess) just AFTER that - // enroll's reap, re-adding the dead session's entry; because that id never - // re-enrolls, no later enroll reaps it. Delivery correctness (no-loss) is - // unaffected either way — only the reap (a leak bound, not the delivery - // guarantee) is best-effort: the recipient still receives the message via the - // reconnect cursor sweep, independent of this registry. + // still streams is HELD here until that session settles (WORKING->READY) or + // reaches a terminal frame. Values are in post order, so a settle fires them + // ascending. A no-frame author death never settles; its entry waits for the + // next-enroll reap (OnSessionsReaped), and a reap race can strand one entry. + // The sweeps skip only messages held for a LIVE author, so the cursor sweep + // still delivers a stranded entry. held map[string][]heldEntry + // lastSettle maps an author session id to the unix ms of its latest settle + // edge, so a message whose hold lost the race with that settle fires at once. + // The recovery pass drops entries of dead sessions; the reap drops them too. + lastSettle map[string]int64 // settleQueue buffers author-settle edges the hook enqueues, drained by the // loop under its ctx. A slice (never lost) plus a buffered notify channel // (coalescing wakeups): the hook appends and signals without blocking Deliver. @@ -245,6 +240,9 @@ type Consumer struct { // it drives. newFloorTicker func() (<-chan time.Time, func()) + // now stamps and prunes lastSettle; a test swaps in a fixed clock. + now func() time.Time + // dispatched counts control dispatches (deliver + steer), labelled only by // op kind (compass.op.kind = steer|deliver). Created ONCE at NewConsumer from // the global meter; nil when meter construction failed, in which case the @@ -281,6 +279,7 @@ func NewConsumer(st DeliveryReads, dispatch ControlDispatcher, resolver SessionR resolver: resolver, log: log, held: make(map[string][]heldEntry), + lastSettle: make(map[string]int64), notify: make(chan struct{}, 1), gates: make(map[string]*sync.Mutex), dispatched: dispatched, @@ -288,6 +287,7 @@ func NewConsumer(st DeliveryReads, dispatch ControlDispatcher, resolver SessionR t := time.NewTicker(recoveryFloorInterval) return t.C, t.Stop }, + now: time.Now, } } @@ -357,9 +357,8 @@ func (c *Consumer) Run(ctx context.Context) error { } // onEventRef handles one ref on the fabric goroutine, concurrently with Run's -// drains. Any read failure is logged and acked (see Run for recovery). A settle -// drained before an earlier post is held leaves that message waiting for the -// author's next settle or session edge. +// drains. Any read failure is logged and acked (see Run for recovery). A post +// whose author settled before the hold landed is delivered at once (hold). func (c *Consumer) onEventRef(ctx context.Context, ref fabric.EventRef) { ctx = store.WithTenant(ctx, store.TenantID(ref.Tenant)) m, err := c.st.MessageByID(ctx, ref.RowID) @@ -392,6 +391,16 @@ func (c *Consumer) drainRecovery(ctx context.Context) { if !pending { return } + // A settle time guards a hold at any age, for example when a backlog replays + // after an outage, so only a dead session's entry goes. + live := c.liveSessionIDs() + c.mu.Lock() + for sid := range c.lastSettle { + if _, ok := live[sid]; !ok { + delete(c.lastSettle, sid) + } + } + c.mu.Unlock() c.sweepAllLive(ctx) c.scanMissedMentions(ctx) } diff --git a/go/internal/delivery/consumer_test.go b/go/internal/delivery/consumer_test.go index f44cec6e5..9fd1c049c 100644 --- a/go/internal/delivery/consumer_test.go +++ b/go/internal/delivery/consumer_test.go @@ -8,6 +8,7 @@ package delivery import ( "context" + "slices" "sync" "testing" "time" @@ -629,6 +630,13 @@ func (d *blockCapturingDispatcher) countFor(messageID string) int { return n } +// records returns every recorded deliver, in dispatch order. +func (d *blockCapturingDispatcher) records() []blockRecord { + d.mu.Lock() + defer d.mu.Unlock() + return slices.Clone(d.calls) +} + // The no-live-author path delivers the STORED block set. The ref carries no // blocks, so the deliver must carry what the store holds. func TestAgentAuthoredNoLiveAuthorDeliversStoredBlocks(t *testing.T) { diff --git a/go/internal/delivery/dispatch.go b/go/internal/delivery/dispatch.go index 2b72d6a13..481eb05fa 100644 --- a/go/internal/delivery/dispatch.go +++ b/go/internal/delivery/dispatch.go @@ -63,35 +63,65 @@ func (c *Consumer) onMessagePosted(ctx context.Context, msg *compassv1.Message) } // Agent-authored. If the author has a live session, HOLD until it settles; - // otherwise deliver now, re-reading the settled blocks from the store (no - // live turn to wait on) — mirroring fireHeld, never the posted (possibly - // partial) wire message (design.md:177-178, :306). + // otherwise deliver now, re-reading the settled blocks from the store — + // mirroring fireHeld, never the posted (possibly partial) wire message + // (design.md:177-178, :306). authorSession, live := c.resolver.SessionForAccount(ctx, author) if !live { - wire, channel, author, err := c.storeMessageToWire(ctx, messageID) - if err != nil { - // The message vanished between post and deliver (unexpected): skip it; - // the cursor never advanced, so the sweep still redelivers. - c.log.ErrorContext(ctx, "delivery: re-read message for dead-author deliver", "error", err, "message_id", messageID) - return - } - c.fanOut(ctx, channel, author, wire) + c.fanOutStored(ctx, messageID) + return + } + c.hold(ctx, authorSession, messageID, msg.GetAtUnixMs()) +} + +// fanOutStored delivers a message now from its stored blocks: the author has no +// live turn. +func (c *Consumer) fanOutStored(ctx context.Context, messageID string) { + wire, channel, author, err := c.storeMessageToWire(ctx, messageID) + if err != nil { + // The message vanished between post and deliver (unexpected): skip it; + // the cursor never advanced, so the sweep still redelivers. + c.log.ErrorContext(ctx, "delivery: re-read message for stored-block deliver", "error", err, "message_id", messageID) return } - c.hold(ctx, authorSession, messageID) + c.fanOut(ctx, channel, author, wire) } // hold registers messageID under its author's session for later firing at the // author's settle edge (design.md:157-160), in post order. It captures the origin -// trace and tenant from ctx for fireHeld. -func (c *Consumer) hold(ctx context.Context, authorSession, messageID string) { - entry := heldEntry{messageID: messageID, traceparent: otelx.Traceparent(ctx)} +// trace and tenant from ctx for fireHeld. If the author already settled at or +// after atUnixMs, it also queues a settle edge, so the loop fires it at once and +// still behind any earlier message of that author. +// +// The two clocks come from different instances. A settling clock behind the +// committing one holds a message until the next settle, which is benign. A +// settling clock ahead by more than the gap between turns fires a still-streaming +// message early, with partial blocks. One process stamping both has no skew. +func (c *Consumer) hold(ctx context.Context, authorSession, messageID string, atUnixMs int64) { + entry := heldEntry{messageID: messageID, traceparent: otelx.Traceparent(ctx), atUnixMs: atUnixMs} if tenant, ok := store.TenantFromContext(ctx); ok { entry.tenant = tenant } c.mu.Lock() - defer c.mu.Unlock() c.held[authorSession] = append(c.held[authorSession], entry) + settled, ok := c.lastSettle[authorSession] + fireNow := ok && settled >= atUnixMs + if fireNow { + // lastSettle stays as is: this replays that settle, bounded to its turn so + // a later, still-streaming message of the author stays held. + c.settleQueue = append(c.settleQueue, settleEvent{ + sessionID: authorSession, + state: compassv1.AgentSessionState_AGENT_SESSION_STATE_READY, + upTo: settled, + }) + } + c.mu.Unlock() + if fireNow { + select { + case c.notify <- struct{}{}: + default: + } + } } // fanOut dispatches one settled message. It first routes any `@`-mentions to a diff --git a/go/internal/delivery/held_gaps_test.go b/go/internal/delivery/held_gaps_test.go new file mode 100644 index 000000000..b1a013910 --- /dev/null +++ b/go/internal/delivery/held_gaps_test.go @@ -0,0 +1,345 @@ +//go:build unix + +package delivery + +// The fabric callback holds concurrently with Run's settle drain, and the start +// sweep runs while an author streams. These cases pin both held-message gaps. + +import ( + "sync" + "sync/atomic" + "testing" + "time" + + compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" + "github.com/RigelBuild/compass/go/internal/store" +) + +// newHeldGapConsumer wires a live agent author and a live subscribed recipient +// over a block-capturing dispatcher, with the consumer clock pinned to settleAt. +// It does not start the consumer, so a test can install read hooks first. +func newHeldGapConsumer(t *testing.T, settleAt time.Time) (*Consumer, *blockCapturingDispatcher, *fakeReads) { + t.Helper() + disp := newBlockCapturingDispatcher() + res := newFakeResolver() + reads := newFakeReads() + c := NewConsumer(reads, disp, res, newFakeFabric(), discardLogger()) + c.now = func() time.Time { return settleAt } + reads.subscribers["chan-1"] = []store.AccountID{"agent-recip"} + reads.agents["agent-author"] = true + res.bind("agent-author", "sess-author") + res.bind("agent-recip", "sess-recip") + return c, disp, reads +} + +// heldGapFixture is newHeldGapConsumer, started and subscribed. +func heldGapFixture(t *testing.T, settleAt time.Time) (*Consumer, *blockCapturingDispatcher, *fakeReads) { + t.Helper() + c, disp, reads := newHeldGapConsumer(t, settleAt) + startConsumer(t, c) + fakeFabricOf(c).waitSubscribed(t) + return c, disp, reads +} + +// messageAt builds an author message stamped at the given commit time. +func messageAt(id, body string, at time.Time) store.Message { + m := textMessage(id, "agent-author", body) + m.At = at + return m +} + +// A settle drained before the callback holds a message from that same turn must +// not strand it until the next settle: the hold sees the settle and fires. A +// commit stamped at the settle ms is still that turn, and a terminal settle +// counts like READY. +func TestHoldAfterSettleFiresAtOnce(t *testing.T) { + settleAt := time.UnixMilli(2_000_000) + cases := []struct { + name string + at time.Time + state compassv1.AgentSessionState + }{ + {"before settle", settleAt.Add(-time.Second), compassv1.AgentSessionState_AGENT_SESSION_STATE_READY}, + {"at settle", settleAt, compassv1.AgentSessionState_AGENT_SESSION_STATE_READY}, + {"stopped", settleAt.Add(-time.Second), compassv1.AgentSessionState_AGENT_SESSION_STATE_STOPPED}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + c, disp, reads := heldGapFixture(t, settleAt) + + c.OnSessionSettled("sess-author", tc.state) + c.waitSettleDrained(t) + + postMessage(t, c, reads, messageAt("m1", "settled body", tc.at)) + rec := disp.waitFor(t, "m1") + if rec.sessionID != "sess-recip" || rec.firstText != "settled body" { + t.Fatalf("deliver = %+v, want {sess-recip, m1, settled body}", rec) + } + fakeFabricOf(c).waitAcked(t, "m1") + if n := disp.countFor("m1"); n != 1 { + t.Fatalf("m1 dispatched %d times, want exactly 1", n) + } + if c.isHeld("sess-author", "m1") { + t.Fatal("m1 still held after its fire") + } + }) + } +} + +// A post whose settle already landed must not overtake an earlier message of +// the same author that fireHeld is still sending. The loop is parked inside +// fireHeld's re-read of m1 while m2 arrives, so the order is event-gated. +func TestFireNowKeepsPostOrderBehindHeld(t *testing.T) { + for _, state := range []compassv1.AgentSessionState{ + compassv1.AgentSessionState_AGENT_SESSION_STATE_READY, + compassv1.AgentSessionState_AGENT_SESSION_STATE_STOPPED, + } { + t.Run(state.String(), func(t *testing.T) { + c, disp, reads := newHeldGapConsumer(t, time.UnixMilli(200)) + entered := make(chan struct{}) + release := make(chan struct{}) + var armed atomic.Bool + reads.beforeMessageByID = func(id string) { + if id == "m1" && armed.CompareAndSwap(true, false) { + close(entered) + <-release + } + } + startConsumer(t, c) + // Also runs before startConsumer's cleanup, so a failed assert never + // leaves the loop parked in the hook. + unpark := sync.OnceFunc(func() { close(release) }) + t.Cleanup(unpark) + fakeFabricOf(c).waitSubscribed(t) + + postMessage(t, c, reads, messageAt("m1", "m1 partial", time.UnixMilli(100))) + c.waitHeld(t, "sess-author", 1) + reads.seedMessage(messageAt("m1", "m1 settled", time.UnixMilli(100))) + armed.Store(true) + + c.OnSessionSettled("sess-author", state) + select { + case <-entered: + case <-time.After(testTimeout): + t.Fatal("fireHeld never re-read m1") + } + postMessage(t, c, reads, messageAt("m2", "m2 body", time.UnixMilli(150))) + fakeFabricOf(c).waitAcked(t, "m2") + unpark() + + disp.waitFor(t, "m2") + disp.waitFor(t, "m1") + got := disp.records() + if len(got) != 2 || + got[0] != (blockRecord{sessionID: "sess-recip", messageID: "m1", firstText: "m1 settled"}) || + got[1] != (blockRecord{sessionID: "sess-recip", messageID: "m2", firstText: "m2 body"}) { + t.Fatalf("delivers = %+v, want [m1 settled, m2 body] to sess-recip, each once", got) + } + }) + } +} + +// The edge a late hold queues must fire only that turn's messages. A message +// of the next turn, still streaming, stays held until its own settle. +func TestFireNowSparesNextTurnMessage(t *testing.T) { + c, disp, reads := newHeldGapConsumer(t, time.UnixMilli(200)) + entered := make(chan struct{}) + release := make(chan struct{}) + var armed atomic.Bool + reads.beforeMessageByID = func(id string) { + if id == "m1" && armed.CompareAndSwap(true, false) { + close(entered) + <-release + } + } + startConsumer(t, c) + // Also runs before startConsumer's cleanup, so a failed assert never leaves + // the loop parked in the hook. + unpark := sync.OnceFunc(func() { close(release) }) + t.Cleanup(unpark) + fab := fakeFabricOf(c) + fab.waitSubscribed(t) + + postMessage(t, c, reads, messageAt("m1", "m1 body", time.UnixMilli(100))) + c.waitHeld(t, "sess-author", 1) + armed.Store(true) + c.OnSessionSettled("sess-author", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) + select { + case <-entered: + case <-time.After(testTimeout): + t.Fatal("fireHeld never re-read m1") + } + postMessage(t, c, reads, messageAt("m-late", "late body", time.UnixMilli(150))) + postMessage(t, c, reads, messageAt("m-next", "next partial", time.UnixMilli(300))) + fab.waitAcked(t, "m-late") + fab.waitAcked(t, "m-next") + unpark() + + disp.waitFor(t, "m-late") + // The loop drains edges in order, so once this one is popped every earlier + // fire returned. + c.OnSessionSettled("sess-other", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) + c.waitSettleDrained(t) + got := disp.records() + if len(got) != 2 || got[0].messageID != "m1" || got[1].messageID != "m-late" { + t.Fatalf("delivers = %+v, want exactly [m1, m-late]", got) + } + if !c.isHeld("sess-author", "m-next") { + t.Fatal("m-next is no longer held before its own turn settled") + } + + reads.seedMessage(messageAt("m-next", "next settled", time.UnixMilli(300))) + c.OnSessionSettled("sess-author", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) + rec := disp.waitFor(t, "m-next") + if rec.sessionID != "sess-recip" || rec.firstText != "next settled" { + t.Fatalf("m-next deliver = %+v, want {sess-recip, next settled}", rec) + } + if n := disp.countFor("m-next"); n != 1 { + t.Fatalf("m-next dispatched %d times, want exactly 1", n) + } +} + +// Control: a message committed after the recorded settle belongs to a later +// turn, so it is held until the next settle. +func TestHoldAfterSettleKeepsLaterMessageHeld(t *testing.T) { + settleAt := time.UnixMilli(2_000_000) + c, disp, reads := heldGapFixture(t, settleAt) + + c.OnSessionSettled("sess-author", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) + c.waitSettleDrained(t) + + postMessage(t, c, reads, messageAt("m2", "next turn", settleAt.Add(time.Second))) + c.waitHeld(t, "sess-author", 1) + if n := disp.countFor("m2"); n != 0 { + t.Fatalf("m2 dispatched %d times before the next settle, want 0", n) + } + + c.OnSessionSettled("sess-author", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) + disp.waitFor(t, "m2") + if n := disp.countFor("m2"); n != 1 { + t.Fatalf("m2 dispatched %d times, want exactly 1", n) + } +} + +// The recovery pass keeps a live session's settle time at any age, since a +// backlog replayed after an outage still needs it, and drops a dead session's. +func TestRecoveryPrunesSettleTimesByLiveness(t *testing.T) { + c, _, res, _ := newTestConsumer(t) + now := time.UnixMilli(2_000_000) + c.now = func() time.Time { return now } + res.bind("agent-live", "sess-live") + + c.OnSessionSettled("sess-live", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) + c.OnSessionSettled("sess-dead", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) + now = now.Add(recoveryFloorInterval + time.Hour) + if !c.hasLastSettle("sess-live") || !c.hasLastSettle("sess-dead") { + t.Fatal("precondition: both settles should be recorded") + } + + c.requestRecovery() + c.drainRecovery(t.Context()) + if !c.hasLastSettle("sess-live") { + t.Fatal("an old settle time of a live session was pruned") + } + if c.hasLastSettle("sess-dead") { + t.Fatal("the settle time of a non-live session survived recovery") + } +} + +// A recipient that starts while its author streams must not get the held +// message from partial blocks: dedup would drop the settled deliver. The settle +// then sends it once, carrying the settled blocks. +func TestStartSweepSkipsHeldMessage(t *testing.T) { + disp := newBlockCapturingDispatcher() + res := newFakeResolver() + reads := newFakeReads() + c := NewConsumer(reads, disp, res, newFakeFabric(), discardLogger()) + const ch store.ChannelID = "chan-1" + const recipient store.AccountID = "agent-recip" + + reads.subscribers[ch] = []store.AccountID{recipient} + reads.agents["agent-author"] = true + res.bind("agent-author", "sess-author") + startConsumer(t, c) + fakeFabricOf(c).waitSubscribed(t) + + postMessage(t, c, reads, textMessage("m1", "agent-author", "partial body")) + c.waitHeld(t, "sess-author", 1) + // The plain message follows m1 in the owed order, so its dispatch proves the + // sweep already passed m1. + reads.mu.Lock() + reads.owed[recipient] = map[store.ChannelID][]store.Message{ch: { + textMessage("m1", "agent-author", "partial body"), + textMessage("plain", "human-1", "plain body"), + }} + reads.mu.Unlock() + + res.bind(recipient, "sess-recip") + c.OnSessionStarted("sess-recip", recipient) + disp.waitFor(t, "plain") + if n := disp.countFor("m1"); n != 0 { + t.Fatalf("start sweep sent held m1 %d times, want 0 (fireHeld owns it)", n) + } + + reads.seedMessage(textMessage("m1", "agent-author", "settled body")) + c.OnSessionSettled("sess-author", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) + rec := disp.waitFor(t, "m1") + if rec.sessionID != "sess-recip" || rec.firstText != "settled body" { + t.Fatalf("settled deliver = %+v, want {sess-recip, m1, settled body}", rec) + } + if n := disp.countFor("m1"); n != 1 { + t.Fatalf("m1 dispatched %d times, want exactly 1", n) + } +} + +// deadAuthorFixture holds m1 under an author session that is not live, as a +// reap race or a no-frame death leaves it, and owes it to agent-recip. +func deadAuthorFixture(t *testing.T) (*Consumer, *blockCapturingDispatcher, *fakeResolver, *fakeReads) { + t.Helper() + disp := newBlockCapturingDispatcher() + res := newFakeResolver() + reads := newFakeReads() + c := NewConsumer(reads, disp, res, newFakeFabric(), discardLogger()) + c.hold(t.Context(), "sess-dead", "m1", 0) + reads.owed["agent-recip"] = map[store.ChannelID][]store.Message{"chan-1": { + textMessage("m1", "agent-author", "stored body"), + textMessage("plain", "human-1", "plain body"), + }} + return c, disp, res, reads +} + +// No settle ever fires a dead author's held entry, so the start sweep must +// deliver it rather than skip it. +func TestStartSweepDeliversHeldForDeadAuthor(t *testing.T) { + c, disp, res, _ := deadAuthorFixture(t) + startConsumer(t, c) + fakeFabricOf(c).waitSubscribed(t) + + res.bind("agent-recip", "sess-recip") + c.OnSessionStarted("sess-recip", "agent-recip") + disp.waitFor(t, "plain") // the sweep passed m1 + if n := disp.countFor("m1"); n != 1 { + t.Fatalf("start sweep sent m1 %d times, want 1 (its author is dead)", n) + } +} + +// The recovery sweep must also deliver a held entry stranded under a dead author. +func TestRecoverySweepDeliversHeldForDeadAuthor(t *testing.T) { + c, disp, res, reads := deadAuthorFixture(t) + tick := make(chan time.Time) + c.newFloorTicker = func() (<-chan time.Time, func()) { return tick, func() {} } + res.bind("agent-recip", "sess-recip") + startConsumer(t, c) + fakeFabricOf(c).waitSubscribed(t) + + passes := reads.unroutedCallCount() + select { + case tick <- time.Now(): + case <-time.After(testTimeout): + t.Fatal("consumer loop never read the floor tick") + } + reads.waitUnroutedCalls(t, passes+1) + if n := disp.countFor("m1"); n != 1 { + t.Fatalf("recovery sweep sent m1 %d times, want 1 (its author is dead)", n) + } +} diff --git a/go/internal/delivery/introspect_test.go b/go/internal/delivery/introspect_test.go index 544ef0de8..5939cef24 100644 --- a/go/internal/delivery/introspect_test.go +++ b/go/internal/delivery/introspect_test.go @@ -46,6 +46,14 @@ func (c *Consumer) isHeld(authorSession, messageID string) bool { }) } +// hasLastSettle reports whether a settle time is recorded for authorSession. +func (c *Consumer) hasLastSettle(authorSession string) bool { + c.mu.Lock() + defer c.mu.Unlock() + _, ok := c.lastSettle[authorSession] + return ok +} + // waitSettleDrained blocks until the settle queue is empty, or fails at the // deadline. Paired with an OnSessionSettled for a throwaway session, it is a // deterministic barrier that a prior settle edge was fully processed. diff --git a/go/internal/delivery/reap_test.go b/go/internal/delivery/reap_test.go index 4f3252b21..f58fab897 100644 --- a/go/internal/delivery/reap_test.go +++ b/go/internal/delivery/reap_test.go @@ -10,6 +10,8 @@ package delivery import ( "context" "testing" + + compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" ) // OnSessionsReaped drops exactly the held entries for the reaped session ids and @@ -19,9 +21,9 @@ func TestOnSessionsReapedDropsHeldEntries(t *testing.T) { c, _, _, _ := newTestConsumer(t) //nolint:dogsled // this test needs only the consumer; the fakes (dispatcher/resolver/reads) are unused here — the reap is a pure in-memory delete with no dispatch/resolve/read path. // Two authors hold pending delivers; a no-frame death would strand both. - c.hold(context.Background(), "sess-dead", "m1") - c.hold(context.Background(), "sess-dead", "m2") - c.hold(context.Background(), "sess-live", "m3") + c.hold(context.Background(), "sess-dead", "m1", 0) + c.hold(context.Background(), "sess-dead", "m2", 0) + c.hold(context.Background(), "sess-live", "m3", 0) if !c.isHeld("sess-dead", "m1") || !c.isHeld("sess-dead", "m2") { t.Fatal("precondition: sess-dead should hold m1 and m2") @@ -46,7 +48,7 @@ func TestOnSessionsReapedDropsHeldEntries(t *testing.T) { // no-op that touches no held entry. func TestOnSessionsReapedEmptyIsNoop(t *testing.T) { c, _, _, _ := newTestConsumer(t) //nolint:dogsled // this test needs only the consumer; the fakes are unused — an empty-slice reap touches no dispatch/resolve/read path. - c.hold(context.Background(), "sess-a", "m1") + c.hold(context.Background(), "sess-a", "m1", 0) c.OnSessionsReaped(nil) @@ -60,7 +62,7 @@ func TestOnSessionsReapedEmptyIsNoop(t *testing.T) { // leaves every unrelated held entry intact — delete of an absent map key. func TestOnSessionsReapedAbsentIDIsNoop(t *testing.T) { c, _, _, _ := newTestConsumer(t) //nolint:dogsled // this test needs only the consumer; the fakes are unused — reaping an absent id touches no dispatch/resolve/read path. - c.hold(context.Background(), "sess-live", "m1") + c.hold(context.Background(), "sess-live", "m1", 0) c.OnSessionsReaped([]string{"sess-never-held"}) @@ -68,3 +70,20 @@ func TestOnSessionsReapedAbsentIDIsNoop(t *testing.T) { t.Fatal("reap of an absent id dropped an unrelated held entry, want it intact") } } + +// A reaped session's settle time goes with it, so the settle map is bounded by +// enroll too; an unreaped session keeps its entry. +func TestOnSessionsReapedDropsSettleTimes(t *testing.T) { + c, _, _, _ := newTestConsumer(t) //nolint:dogsled // only the consumer's settle map is exercised. + c.OnSessionSettled("sess-dead", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) + c.OnSessionSettled("sess-live", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) + + c.OnSessionsReaped([]string{"sess-dead"}) + + if c.hasLastSettle("sess-dead") { + t.Fatal("sess-dead settle time survived the reap, want dropped") + } + if !c.hasLastSettle("sess-live") { + t.Fatal("sess-live settle time was dropped by the reap, want it kept") + } +} diff --git a/go/internal/delivery/scan_test.go b/go/internal/delivery/scan_test.go index 41bb5ecae..e003f516a 100644 --- a/go/internal/delivery/scan_test.go +++ b/go/internal/delivery/scan_test.go @@ -59,7 +59,7 @@ func TestScanSkipsHeldMessage(t *testing.T) { reads.members[ch] = []store.AccountID{agentA} reads.handles["aa"] = agentAccount(agentA, "aa") reads.seedUnrouted(textMessage("m1", author, "@aa ping"), ch, 1) - c.hold(context.Background(), "author-sess", "m1") // registered in c.held under its author session + c.hold(context.Background(), "author-sess", "m1", 0) // registered in c.held under its author session c.scanMissedMentions(context.Background()) @@ -253,7 +253,7 @@ func TestScanSkipsMarkWhenHeldAtMarkTime(t *testing.T) { reads.beforeMessageByID = func(id string) { if id == "m1" && !injected { injected = true - c.hold(context.Background(), "sess-author", "m1") + c.hold(context.Background(), "sess-author", "m1", 0) } } c.scanMissedMentions(context.Background()) diff --git a/go/internal/delivery/settle.go b/go/internal/delivery/settle.go index 571376077..6276f1d7e 100644 --- a/go/internal/delivery/settle.go +++ b/go/internal/delivery/settle.go @@ -4,6 +4,7 @@ package delivery import ( "context" + "math" comms "github.com/RigelBuild/compass/go/internal/comms" @@ -27,7 +28,10 @@ func (c *Consumer) OnSessionSettled(sessionID string, state compassv1.AgentSessi return } c.mu.Lock() - c.settleQueue = append(c.settleQueue, settleEvent{sessionID: sessionID, state: state}) + c.settleQueue = append(c.settleQueue, settleEvent{sessionID: sessionID, state: state, upTo: math.MaxInt64}) + // Recorded with the enqueue, so a hold that loses the race to the drain + // still sees this settle. + c.lastSettle[sessionID] = c.now().UnixMilli() c.mu.Unlock() // Coalescing wakeup: a full buffer already signals a pending drain, so a // dropped send loses nothing (the loop drains the whole queue). @@ -84,21 +88,13 @@ func firesHeldDelivers(state compassv1.AgentSessionState) bool { } // drainSettles fires every queued author-settle edge under the loop's ctx. Each -// edge fires the messages held for that author session, in post order, from each -// message's CURRENT (settled) stored blocks (design.md:158-168), then clears the -// registry entry — a no-frame author death never enqueues an edge, so its held -// entry is left in place until it is reaped. The reap happens in-process on the -// next Runner (re-)enroll via the hub's SessionReapSink (OnSessionsReaped), -// which drops the entry for every session id enroll just cleared — so the common -// no-frame death is reaped at that next enroll rather than persisting until -// process restart. The reap is best-effort, not a hard bound: a no-frame-dead -// session is never re-promoted, so a narrow race (a Deliver that resolved the -// author LIVE an instant before enroll cleared the maps re-holds the dead -// session just AFTER that enroll's reap) can still strand one entry until process -// restart, since that id never re-enrolls to be reaped again. No-loss is -// unaffected regardless — the reconnect sweep still delivers the message -// (design.md:168-176); only the leak bound, not the delivery guarantee, is -// best-effort. +// edge fires the messages held for that author session up to its upTo, in post +// order, from each message's CURRENT (settled) stored blocks (design.md:158-168). +// An edge for a session with nothing in range is a no-op. +// +// A no-frame author death never enqueues an edge. No-loss still holds: the +// sweeps skip only messages held for a LIVE author, so the cursor sweep +// delivers an entry stranded under a dead session (design.md:168-176). func (c *Consumer) drainSettles(ctx context.Context) { for { c.mu.Lock() @@ -109,7 +105,7 @@ func (c *Consumer) drainSettles(ctx context.Context) { ev := c.settleQueue[0] c.settleQueue = c.settleQueue[1:] c.mu.Unlock() - c.fireHeld(ctx, ev.sessionID) + c.fireHeld(ctx, ev.sessionID, ev.upTo) } } @@ -131,7 +127,9 @@ func (c *Consumer) drainStarts(ctx context.Context) { ev := c.startQueue[0] c.startQueue = c.startQueue[1:] c.mu.Unlock() - c.sweepSession(ctx, ev.account, ev.sessionID, false) + // Skip messages held for a live author: fireHeld re-resolves recipients + // at settle, so this session still gets them with their settled blocks. + c.sweepSession(ctx, ev.account, ev.sessionID, true) if err := c.sweepPins(ctx, ev.account, ev.sessionID); err != nil { c.log.ErrorContext(ctx, "delivery: sweep pins on session start", "error", err, "account", string(ev.account), "session_id", ev.sessionID) @@ -261,17 +259,28 @@ func (c *Consumer) sweepOwedMentions(ctx context.Context, agent store.AccountID, return nil } -// fireHeld dispatches every message held for authorSession, ascending, and -// clears the registry entry. Each message is re-read under its hold-time tenant, -// so the deliver carries the SETTLED blocks (design.md:158-161), and recipients -// are re-resolved against the then-current subscription + liveness. -func (c *Consumer) fireHeld(ctx context.Context, authorSession string) { +// fireHeld dispatches, ascending, the messages held for authorSession whose +// commit time is at most upTo, and keeps the rest held in order. Each is +// re-read under its hold-time tenant, so the deliver carries the SETTLED blocks +// (design.md:158-161), and recipients are re-resolved at fire time. +func (c *Consumer) fireHeld(ctx context.Context, authorSession string, upTo int64) { c.mu.Lock() - held := c.held[authorSession] - delete(c.held, authorSession) + var fire, keep []heldEntry + for _, e := range c.held[authorSession] { + if e.atUnixMs <= upTo { + fire = append(fire, e) + } else { + keep = append(keep, e) + } + } + if len(keep) == 0 { + delete(c.held, authorSession) + } else { + c.held[authorSession] = keep + } c.mu.Unlock() - for _, entry := range held { + for _, entry := range fire { // The drain ctx carries no tenant; re-read under the one captured at hold. tctx := store.WithTenant(ctx, entry.tenant) wire, channel, author, err := c.storeMessageToWire(tctx, entry.messageID) @@ -290,46 +299,65 @@ func (c *Consumer) fireHeld(ctx context.Context, authorSession string) { } } -// OnSessionsReaped drops the held-deliver registry entries for sessions whose -// hub bindings were cleared at a Runner (re-)enroll (SessionReapSink). A -// no-frame author death emits no terminal frame, so no settle edge ever fires -// fireHeld to clear its entry; this enroll-bounded reap realizes the design's -// promised cleanup (design.md:172-175) so the registry does not leak an entry -// per no-frame death until process restart. No-loss is unaffected: any message -// still owed is redelivered by the recipient's reconnect cursor sweep. Pure -// in-memory work under c.mu — it does not enqueue onto the consumer loop or -// touch the store, so it is safe to run directly on the hub's enroll goroutine. +// OnSessionsReaped drops the held-deliver and settle-time entries for sessions +// whose hub bindings were cleared at a Runner (re-)enroll, so a no-frame +// author death does not leak an entry (design.md:172-175). The cursor sweep +// still delivers what was held, since it skips only live authors' messages. func (c *Consumer) OnSessionsReaped(sessionIDs []string) { c.mu.Lock() defer c.mu.Unlock() for _, sid := range sessionIDs { delete(c.held, sid) + delete(c.lastSettle, sid) } } // sweepAllLive is the recovery pass a fabric reconnect or the floor tick runs: -// it redelivers every owed message to every live agent session, skipping held -// ones so partial blocks never go out ahead of fireHeld (message_id dedup would -// drop the settled deliver). +// it redelivers every owed message to every live agent session, skipping those +// held for a live author so partial blocks never go out ahead of fireHeld +// (message_id dedup would drop the settled deliver). func (c *Consumer) sweepAllLive(ctx context.Context) { for account, sessionID := range c.resolver.LiveAgentSessions() { c.sweepSession(ctx, account, sessionID, true) } } -// heldIDs snapshots every held message id. +// heldIDs snapshots the ids held for LIVE author sessions. An entry stranded +// under a dead session has no settle coming, so the sweeps must deliver it. +// Liveness is read after the held snapshot, so a session that went live in +// between is still skipped; a stale-live one also stays skipped, the safe side. func (c *Consumer) heldIDs() map[string]struct{} { c.mu.Lock() - defer c.mu.Unlock() - ids := make(map[string]struct{}) - for _, entries := range c.held { + bySession := make(map[string][]string, len(c.held)) + for sid, entries := range c.held { for _, e := range entries { - ids[e.messageID] = struct{}{} + bySession[sid] = append(bySession[sid], e.messageID) + } + } + c.mu.Unlock() + live := c.liveSessionIDs() // resolver read stays outside c.mu + ids := make(map[string]struct{}) + for sid, msgIDs := range bySession { + if _, ok := live[sid]; !ok { + continue + } + for _, id := range msgIDs { + ids[id] = struct{}{} } } return ids } +// liveSessionIDs snapshots the set of live agent session ids. +func (c *Consumer) liveSessionIDs() map[string]struct{} { + bound := c.resolver.LiveAgentSessions() + live := make(map[string]struct{}, len(bound)) + for _, sid := range bound { + live[sid] = struct{}{} + } + return live +} + // sweepSession redelivers one agent's owed messages in seq order under the // session's gate held for the whole pass, so live delivers queue behind it // (design.md:220-225). skipHeld leaves messages held at owed-read time to fireHeld.