Skip to content

StreamInbox

Defined in: packages/client/src/circuits/stream-inbox.ts:35

The staging buffer between K durable-streams subscriptions and the applier — the offset-keyed successor to the LSN-keyed ShapeInbox (ADR-0056).

Three things it deliberately no longer has, each because the reason for it was Electric-shaped:

No commit floor. ADR-0031’s floor existed to work around a cached catch-up body asserting a stale global_last_seen_lsn: a quiet shape’s cacheable response could claim a watermark captured before a busy sibling’s writes, holding delivered changes buffered until that shape’s first live poll. Durable-streams carries up-to-date as a response header on a live request, and returns 204 with it set on every long-poll timeout, so a quiet shape re-asserts freshness each cycle. The failure mode the floor compensated for cannot occur, so the floor is deleted rather than ported.

No cross-shape position comparison. Offsets are per-stream and, per the protocol, comparable only within a stream. So the commit gate is not a min over positions but a predicate over reports: commit when every shape’s most recent response asserted up-to-date. That is the same happens-before argument ADR-0056 makes for alignment, applied per commit — with the caveat StreamInbox.isGroupUpToDate states: it is a complete argument at alignment, and in the live steady state it bounds cross-shape tearing to two responses’ inter-arrival rather than excluding it (docs/backlog/0014-cross-stream-commit-fence.md).

No snapshot-acceptance flag. It existed because a re-snapshot’s rows floored to LSN 0 while the frontier might already sit high. A reset here rewinds the applied offset to nothing, and the re-snapshot arrives at real offsets above it, so the problem does not arise.

new StreamInbox(shapeNames): StreamInbox

Defined in: packages/client/src/circuits/stream-inbox.ts:54

Iterable<string>

StreamInbox

ackAll(epochsAtPeek): void

Defined in: packages/client/src/circuits/stream-inbox.ts:198

Drop everything peeked and advance each shape’s applied offset — call only after the commit that consumed a matching peekAll succeeded.

A shape whose epoch changed since the peek is skipped entirely: a must-refetch landing mid-commit replaced its buffer with post-reset content that this commit never saw, and acking it would discard the rebuild.

Map<string, number>

void


appliedOffsetFor(shapeName): string | null

Defined in: packages/client/src/circuits/stream-inbox.ts:209

The offset a shape has applied and persisted, or null before its first commit.

string

string | null


epochFor(shapeName): number

Defined in: packages/client/src/circuits/stream-inbox.ts:146

string

number


everyShapeReported(): boolean

Defined in: packages/client/src/circuits/stream-inbox.ts:107

Whether every shape has reported up-to-date at least once — boot alignment’s precondition.

boolean


hasBufferedBatches(): boolean

Defined in: packages/client/src/circuits/stream-inbox.ts:139

Whether any shape holds a batch at all, including the empty position-carrying ones.

boolean


hasBufferedChanges(): boolean

Defined in: packages/client/src/circuits/stream-inbox.ts:129

boolean


ingestBatch(shapeName, changes, offset, upToDate): void

Defined in: packages/client/src/circuits/stream-inbox.ts:72

Buffer one delivery.

A batch at or below the applied offset is already applied and dropped. The comparison is lexicographic, which the protocol explicitly sanctions within a stream (offsets are opaque but lexicographically sortable and strictly increasing) and equally explicitly does not sanction across streams — which is why nothing here ever compares two shapes’ offsets.

string

ChangeLike[]

string

boolean

void


isGroupUpToDate(): boolean

Defined in: packages/client/src/circuits/stream-inbox.ts:99

The commit gate: every shape’s most recent response asserted up-to-date.

This is the whole steady-state condition, and it is stronger than it looks in one direction and weaker in another. Stronger: a shape mid-catch-up reports false, so a group never commits while any member is still draining backfill, and at alignment the whole held batch commits together. Weaker: once caught up, every shape’s flag is LATCHED from its last completed response while its next long poll sits parked — a 204 timeout carries up-to-date too — so this predicate is effectively always true and each delivery commits at once. Two streams carrying halves of one server transaction are answered by two separate polls, so the client can apply one half before the other arrives; the window is their inter-arrival. Nothing here can close that — the engine’s per-stream appends carry no cross-stream fence — so it is recorded rather than papered over (docs/backlog/0014-cross-stream-commit-fence.md).

boolean


markAligned(): void

Defined in: packages/client/src/circuits/stream-inbox.ts:125

Record that the barrier was satisfied and the group aligned. One-time per reset generation.

void


needsAlignment(): boolean

Defined in: packages/client/src/circuits/stream-inbox.ts:120

Whether the one-time boot alignment still has to run. Alignment additionally requires the engine barrier (ADR-0056 decision 3) — pendingFlips > 0 with every stream up-to-date is a computed revocation the engine has not yet delivered, and a boot that claimed consistency there would present a store missing an eviction.

boolean


peekAll(): Map<string, ChangeLike[]>

Defined in: packages/client/src/circuits/stream-inbox.ts:161

Peek — without removing — everything held, per shape, in stream order.

There is no target position to peek up to. The gate is “every shape has drained”, so what is held IS the complete unit: a partial peek could only ever split a batch the server sent whole.

Map<string, ChangeLike[]>


pendingOffsets(): Map<string, string>

Defined in: packages/client/src/circuits/stream-inbox.ts:180

The offset each shape would resume from if the currently-held batches were committed — the last buffered batch’s offset, or the already-applied one where nothing is held.

Persisted in the SAME transaction as the rows it acknowledges. Persisted ahead, a crash loses the envelopes in between; persisted behind, they are re-delivered and re-applied. Only the second is survivable.

Map<string, string>


resetShape(shapeName): void

Defined in: packages/client/src/circuits/stream-inbox.ts:217

Reset a shape on must-refetch: drop its buffer, rewind it to the start of its stream, and re-arm the group’s alignment so the barrier is consulted again once every shape has re-reported.

string

void


snapshotEpochs(): Map<string, number>

Defined in: packages/client/src/circuits/stream-inbox.ts:151

Snapshot every shape’s epoch at peek time, to be handed back to ackAll.

Map<string, number>