From a55fdc5cab2c438f6f0c73196a4e1024453196dd Mon Sep 17 00:00:00 2001 From: mintaka Date: Sat, 26 Sep 2026 04:59:40 -0400 Subject: [PATCH 1/3] fix(delivery): fire a hold that lost the race to its settle; skip held in start sweeps (RIG-4031) The fabric callback runs concurrently with the settle drain. A message held after its author settled waited for the next settle, unbounded for an author that stays READY. OnSessionSettled now records a per-session settle time, and hold delivers at once from stored blocks when that settle is at or after the message commit time. Entries are pruned by the recovery pass and dropped on reap. The session-start sweep also skips held messages, so a recipient that starts mid-turn gets the settled blocks from fireHeld instead of a partial deliver that dedup would keep. Spec-impact: none. Refs RIG-4031, RIG-3107 Co-authored-by: Matt Wilkinson --- go/internal/delivery/consumer.go | 23 +++- go/internal/delivery/dispatch.go | 40 +++--- go/internal/delivery/held_gaps_test.go | 154 ++++++++++++++++++++++++ go/internal/delivery/introspect_test.go | 8 ++ go/internal/delivery/reap_test.go | 29 ++++- go/internal/delivery/scan_test.go | 4 +- go/internal/delivery/settle.go | 8 +- 7 files changed, 240 insertions(+), 26 deletions(-) create mode 100644 go/internal/delivery/held_gaps_test.go diff --git a/go/internal/delivery/consumer.go b/go/internal/delivery/consumer.go index 77d85f315..2c1c57908 100644 --- a/go/internal/delivery/consumer.go +++ b/go/internal/delivery/consumer.go @@ -213,6 +213,10 @@ type Consumer struct { // guarantee) is best-effort: the recipient still receives the message via the // reconnect cursor sweep, independent of this registry. 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. + // Pruned by the recovery pass and dropped on reap. + 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 +249,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 +288,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 +296,7 @@ func NewConsumer(st DeliveryReads, dispatch ControlDispatcher, resolver SessionR t := time.NewTicker(recoveryFloorInterval) return t.C, t.Stop }, + now: time.Now, } } @@ -357,9 +366,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) @@ -388,6 +396,15 @@ func (c *Consumer) drainRecovery(ctx context.Context) { c.mu.Lock() pending := c.recoveryPending c.recoveryPending = false + if pending { + // A hold this late is not a same-turn race, so the entry can go. + floor := c.now().Add(-recoveryFloorInterval).UnixMilli() + for sid, settled := range c.lastSettle { + if settled < floor { + delete(c.lastSettle, sid) + } + } + } c.mu.Unlock() if !pending { return diff --git a/go/internal/delivery/dispatch.go b/go/internal/delivery/dispatch.go index 2b72d6a13..b0bdf1bed 100644 --- a/go/internal/delivery/dispatch.go +++ b/go/internal/delivery/dispatch.go @@ -63,35 +63,45 @@ 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, or if the turn that posted it already settled, 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) + if !live || c.hold(ctx, authorSession, messageID, msg.GetAtUnixMs()) { + c.fanOutStored(ctx, messageID) + } +} + +// fanOutStored delivers a message now from its stored blocks: the author has no +// live turn, or the turn that posted it already settled. +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) { +// trace and tenant from ctx for fireHeld. It returns true without registering +// when the author already settled at or after atUnixMs, so the caller fires now. +// Clock skew between Server instances can still leave such a message held. +func (c *Consumer) hold(ctx context.Context, authorSession, messageID string, atUnixMs int64) (fireNow bool) { entry := heldEntry{messageID: messageID, traceparent: otelx.Traceparent(ctx)} if tenant, ok := store.TenantFromContext(ctx); ok { entry.tenant = tenant } c.mu.Lock() defer c.mu.Unlock() + if settled, ok := c.lastSettle[authorSession]; ok && settled >= atUnixMs { + return true + } c.held[authorSession] = append(c.held[authorSession], entry) + return false } // 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..0fbfb58c1 --- /dev/null +++ b/go/internal/delivery/held_gaps_test.go @@ -0,0 +1,154 @@ +//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 ( + "testing" + "time" + + compassv1 "github.com/RigelBuild/compass/go/gen/compass/v1" + "github.com/RigelBuild/compass/go/internal/store" +) + +// heldGapFixture wires a live agent author and a live subscribed recipient over +// a block-capturing dispatcher, with the consumer clock pinned to settleAt. +func heldGapFixture(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") + 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. +func TestHoldAfterSettleFiresAtOnce(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("m1", "settled body", settleAt.Add(-time.Second))) + 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 was held after the settle already happened") + } +} + +// 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 drops settle times older than the floor interval and keeps +// fresh ones, so the map stays bounded for authors that never settle again. +func TestRecoveryPrunesStaleSettleTimes(t *testing.T) { + c, _, _, _ := newTestConsumer(t) //nolint:dogsled // only the consumer's clock and settle map are exercised. + now := time.UnixMilli(2_000_000) + c.now = func() time.Time { return now } + + c.OnSessionSettled("sess-old", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) + now = now.Add(recoveryFloorInterval + time.Millisecond) + c.OnSessionSettled("sess-fresh", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) + if !c.hasLastSettle("sess-old") || !c.hasLastSettle("sess-fresh") { + t.Fatal("precondition: both settles should be recorded") + } + + c.requestRecovery() + c.drainRecovery(t.Context()) + if c.hasLastSettle("sess-old") { + t.Fatal("settle time older than the floor interval survived recovery") + } + if !c.hasLastSettle("sess-fresh") { + t.Fatal("fresh settle time was pruned") + } +} + +// 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) + } +} 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..603b98665 100644 --- a/go/internal/delivery/settle.go +++ b/go/internal/delivery/settle.go @@ -28,6 +28,9 @@ func (c *Consumer) OnSessionSettled(sessionID string, state compassv1.AgentSessi } c.mu.Lock() c.settleQueue = append(c.settleQueue, settleEvent{sessionID: sessionID, state: state}) + // 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). @@ -131,7 +134,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 held messages: 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) @@ -304,6 +309,7 @@ func (c *Consumer) OnSessionsReaped(sessionIDs []string) { defer c.mu.Unlock() for _, sid := range sessionIDs { delete(c.held, sid) + delete(c.lastSettle, sid) } } From 67ceb08fe774de9b6f302f5193e8bf0be6723a34 Mon Sep 17 00:00:00 2001 From: mintaka Date: Sat, 26 Sep 2026 05:37:32 -0400 Subject: [PATCH 2/3] fix(delivery): keep post order on fire-now; skip only live-held; prune by liveness (RIG-4031) Review fixes: - A hold that finds its settle already recorded now appends to c.held and queues a settle edge on the loop, instead of dispatching on the fabric goroutine. fireHeld sends the list in post order, so a late M2 no longer overtakes a held M1. - Sweeps skip only messages held for a live author session, so an entry stranded under a dead session is still delivered by the cursor sweep. - lastSettle is pruned by liveness, not age: settle >= at is correct at any age, and an age prune dropped the guard while a backlog replayed. Spec-impact: none. Refs RIG-4031 Co-authored-by: Matt Wilkinson --- go/internal/delivery/consumer.go | 46 +++--- go/internal/delivery/consumer_test.go | 8 ++ go/internal/delivery/dispatch.go | 46 ++++-- go/internal/delivery/held_gaps_test.go | 189 +++++++++++++++++++++---- go/internal/delivery/settle.go | 64 +++++---- 5 files changed, 250 insertions(+), 103 deletions(-) diff --git a/go/internal/delivery/consumer.go b/go/internal/delivery/consumer.go index 2c1c57908..6e75910a1 100644 --- a/go/internal/delivery/consumer.go +++ b/go/internal/delivery/consumer.go @@ -193,29 +193,16 @@ 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. - // Pruned by the recovery pass and dropped on reap. + // 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 @@ -396,19 +383,20 @@ func (c *Consumer) drainRecovery(ctx context.Context) { c.mu.Lock() pending := c.recoveryPending c.recoveryPending = false - if pending { - // A hold this late is not a same-turn race, so the entry can go. - floor := c.now().Add(-recoveryFloorInterval).UnixMilli() - for sid, settled := range c.lastSettle { - if settled < floor { - delete(c.lastSettle, sid) - } - } - } c.mu.Unlock() 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 b0bdf1bed..d6e128761 100644 --- a/go/internal/delivery/dispatch.go +++ b/go/internal/delivery/dispatch.go @@ -63,17 +63,19 @@ func (c *Consumer) onMessagePosted(ctx context.Context, msg *compassv1.Message) } // Agent-authored. If the author has a live session, HOLD until it settles; - // otherwise, or if the turn that posted it already settled, deliver now, - // re-reading the settled blocks from the store — 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 || c.hold(ctx, authorSession, messageID, msg.GetAtUnixMs()) { + if !live { 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, or the turn that posted it already settled. +// live turn. func (c *Consumer) fanOutStored(ctx context.Context, messageID string) { wire, channel, author, err := c.storeMessageToWire(ctx, messageID) if err != nil { @@ -87,21 +89,37 @@ func (c *Consumer) fanOutStored(ctx context.Context, messageID string) { // 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. It returns true without registering -// when the author already settled at or after atUnixMs, so the caller fires now. -// Clock skew between Server instances can still leave such a message held. -func (c *Consumer) hold(ctx context.Context, authorSession, messageID string, atUnixMs int64) (fireNow bool) { +// 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)} if tenant, ok := store.TenantFromContext(ctx); ok { entry.tenant = tenant } c.mu.Lock() - defer c.mu.Unlock() - if settled, ok := c.lastSettle[authorSession]; ok && settled >= atUnixMs { - return true - } c.held[authorSession] = append(c.held[authorSession], entry) - return false + settled, ok := c.lastSettle[authorSession] + fireNow := ok && settled >= atUnixMs + if fireNow { + // lastSettle stays as is: this is a replay of that settle, not a new one. + c.settleQueue = append(c.settleQueue, settleEvent{ + sessionID: authorSession, + state: compassv1.AgentSessionState_AGENT_SESSION_STATE_READY, + }) + } + 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 index 0fbfb58c1..de799d489 100644 --- a/go/internal/delivery/held_gaps_test.go +++ b/go/internal/delivery/held_gaps_test.go @@ -6,6 +6,8 @@ package delivery // sweep runs while an author streams. These cases pin both held-message gaps. import ( + "sync" + "sync/atomic" "testing" "time" @@ -13,9 +15,10 @@ import ( "github.com/RigelBuild/compass/go/internal/store" ) -// heldGapFixture wires a live agent author and a live subscribed recipient over -// a block-capturing dispatcher, with the consumer clock pinned to settleAt. -func heldGapFixture(t *testing.T, settleAt time.Time) (*Consumer, *blockCapturingDispatcher, *fakeReads) { +// 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() @@ -26,6 +29,13 @@ func heldGapFixture(t *testing.T, settleAt time.Time) (*Consumer, *blockCapturin 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 @@ -39,25 +49,93 @@ func messageAt(id, body string, at time.Time) store.Message { } // 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. +// 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) - c, disp, reads := heldGapFixture(t, settleAt) + 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", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) - c.waitSettleDrained(t) + c.OnSessionSettled("sess-author", tc.state) + c.waitSettleDrained(t) - postMessage(t, c, reads, messageAt("m1", "settled body", settleAt.Add(-time.Second))) - 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) + 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") + } + }) } - 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 was held after the settle already happened") +} + +// 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) + } + }) } } @@ -83,27 +161,28 @@ func TestHoldAfterSettleKeepsLaterMessageHeld(t *testing.T) { } } -// The recovery pass drops settle times older than the floor interval and keeps -// fresh ones, so the map stays bounded for authors that never settle again. -func TestRecoveryPrunesStaleSettleTimes(t *testing.T) { - c, _, _, _ := newTestConsumer(t) //nolint:dogsled // only the consumer's clock and settle map are exercised. +// 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-old", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) - now = now.Add(recoveryFloorInterval + time.Millisecond) - c.OnSessionSettled("sess-fresh", compassv1.AgentSessionState_AGENT_SESSION_STATE_READY) - if !c.hasLastSettle("sess-old") || !c.hasLastSettle("sess-fresh") { + 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-old") { - t.Fatal("settle time older than the floor interval survived recovery") + if !c.hasLastSettle("sess-live") { + t.Fatal("an old settle time of a live session was pruned") } - if !c.hasLastSettle("sess-fresh") { - t.Fatal("fresh settle time was pruned") + if c.hasLastSettle("sess-dead") { + t.Fatal("the settle time of a non-live session survived recovery") } } @@ -152,3 +231,55 @@ func TestStartSweepSkipsHeldMessage(t *testing.T) { 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/settle.go b/go/internal/delivery/settle.go index 603b98665..865d4ebe4 100644 --- a/go/internal/delivery/settle.go +++ b/go/internal/delivery/settle.go @@ -87,21 +87,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, in post order, from +// each message's CURRENT (settled) stored blocks (design.md:158-168), then +// clears the registry entry. An edge for a session with nothing held 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() @@ -134,8 +126,8 @@ func (c *Consumer) drainStarts(ctx context.Context) { ev := c.startQueue[0] c.startQueue = c.startQueue[1:] c.mu.Unlock() - // Skip held messages: fireHeld re-resolves recipients at settle, so this - // session still gets them, with their settled blocks. + // 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, @@ -295,15 +287,10 @@ 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() @@ -314,21 +301,26 @@ func (c *Consumer) OnSessionsReaped(sessionIDs []string) { } // 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. func (c *Consumer) heldIDs() map[string]struct{} { + live := c.liveSessionIDs() // resolver read stays outside c.mu c.mu.Lock() defer c.mu.Unlock() ids := make(map[string]struct{}) - for _, entries := range c.held { + for sid, entries := range c.held { + if _, ok := live[sid]; !ok { + continue + } for _, e := range entries { ids[e.messageID] = struct{}{} } @@ -336,6 +328,16 @@ func (c *Consumer) heldIDs() map[string]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. From a68a65f2d566f07bddf5408ec572bf1714bc78ac Mon Sep 17 00:00:00 2001 From: mintaka Date: Sat, 26 Sep 2026 05:51:22 -0400 Subject: [PATCH 3/3] fix(delivery): bound a replayed settle edge to its settle time (RIG-4031) The edge a late hold queues fired the whole held list, including a next-turn message still streaming. Held entries now carry their commit time, and the edge fires only entries at or before the settle it replays; a real settle still fires all. heldIDs reads liveness after the held snapshot, so a race leaves an entry skipped rather than sent partial. Spec-impact: none. Refs RIG-4031 Co-authored-by: Matt Wilkinson --- go/internal/delivery/consumer.go | 4 ++ go/internal/delivery/dispatch.go | 6 ++- go/internal/delivery/held_gaps_test.go | 60 ++++++++++++++++++++++++++ go/internal/delivery/settle.go | 56 ++++++++++++++++-------- 4 files changed, 106 insertions(+), 20 deletions(-) diff --git a/go/internal/delivery/consumer.go b/go/internal/delivery/consumer.go index 6e75910a1..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 diff --git a/go/internal/delivery/dispatch.go b/go/internal/delivery/dispatch.go index d6e128761..481eb05fa 100644 --- a/go/internal/delivery/dispatch.go +++ b/go/internal/delivery/dispatch.go @@ -98,7 +98,7 @@ func (c *Consumer) fanOutStored(ctx context.Context, messageID string) { // 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)} + entry := heldEntry{messageID: messageID, traceparent: otelx.Traceparent(ctx), atUnixMs: atUnixMs} if tenant, ok := store.TenantFromContext(ctx); ok { entry.tenant = tenant } @@ -107,10 +107,12 @@ func (c *Consumer) hold(ctx context.Context, authorSession, messageID string, at settled, ok := c.lastSettle[authorSession] fireNow := ok && settled >= atUnixMs if fireNow { - // lastSettle stays as is: this is a replay of that settle, not a new one. + // 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() diff --git a/go/internal/delivery/held_gaps_test.go b/go/internal/delivery/held_gaps_test.go index de799d489..b1a013910 100644 --- a/go/internal/delivery/held_gaps_test.go +++ b/go/internal/delivery/held_gaps_test.go @@ -139,6 +139,66 @@ func TestFireNowKeepsPostOrderBehindHeld(t *testing.T) { } } +// 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) { diff --git a/go/internal/delivery/settle.go b/go/internal/delivery/settle.go index 865d4ebe4..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,7 @@ 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() @@ -87,9 +88,9 @@ 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. An edge for a session with nothing held is a no-op. +// 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 @@ -104,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) } } @@ -258,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) @@ -312,17 +324,25 @@ func (c *Consumer) sweepAllLive(ctx context.Context) { // 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{} { - live := c.liveSessionIDs() // resolver read stays outside c.mu c.mu.Lock() - defer c.mu.Unlock() - ids := make(map[string]struct{}) + bySession := make(map[string][]string, len(c.held)) for sid, entries := range c.held { + for _, e := range entries { + 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 _, e := range entries { - ids[e.messageID] = struct{}{} + for _, id := range msgIDs { + ids[id] = struct{}{} } } return ids