Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 31 additions & 22 deletions go/internal/delivery/consumer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -281,13 +279,15 @@ 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,
newFloorTicker: func() (<-chan time.Time, func()) {
t := time.NewTicker(recoveryFloorInterval)
return t.C, t.Stop
},
now: time.Now,
}
}

Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
}
Expand Down
8 changes: 8 additions & 0 deletions go/internal/delivery/consumer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ package delivery

import (
"context"
"slices"
"sync"
"testing"
"time"
Expand Down Expand Up @@ -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) {
Expand Down
62 changes: 46 additions & 16 deletions go/internal/delivery/dispatch.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading