Skip to content
Merged
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
11 changes: 10 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -1539,7 +1539,16 @@ bus.listen(
```

Retryable failures (e.g. transient `NotFound`) are nacked for redelivery; the runner
never silently acks a handler error.
never silently acks a handler error. NATS JetStream NAKs back off by delivery count
(50 ms doubling to 5 s by default) so a repeatedly failing message cannot monopolize
a consumer; retries stay unlimited. A deterministic, durably recorded rejection is a
permanent error, not a retryable one.

A service can place independent route bundles in named delivery lanes
(`Service::lane(name, routes)`). Each lane keeps broker order and runs
concurrently with the others, and a delivery is settled only after every lane
finished it. See [consumer delivery lanes](docs/consumer-delivery-lanes.md) for
the contract and when a route may leave the default lane.

### Transport boundaries (producer vs consumer)

Expand Down
2 changes: 1 addition & 1 deletion distributed_cli/src/contracts/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1259,7 +1259,7 @@ fn migration_inventory_is_deterministic_and_preserves_runtime_order() {
.iter()
.map(|migration| migration.version)
.collect::<Vec<_>>();
assert_eq!(versions, vec![1, 2, 3, 4, 5, 6, 7, 8]);
assert_eq!(versions, vec![1, 2, 3, 4, 5, 6, 7, 8, 9]);
assert_eq!(
inventory.canonical_bytes().expect("canonical inventory"),
inventory
Expand Down
1 change: 1 addition & 0 deletions distributed_cli/tests/cli_lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -286,6 +286,7 @@ exit 1
fs::write(path, serde_json::to_vec_pretty(&config).unwrap()).unwrap();
}

#[track_caller]
fn wait_until(timeout: Duration, predicate: impl Fn() -> bool) {
let started = std::time::Instant::now();
while !predicate() {
Expand Down
99 changes: 99 additions & 0 deletions docs/consumer-delivery-lanes.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,99 @@
# Consumer delivery lanes and retry backoff

One durable consumer (`listen`/`subscribe` group) delivers messages in broker
order and, by default, runs every registered route for a message before it
receives the next one. That is the simplest ordering contract, but it couples
latency across unrelated work: a process-manager policy that only turns a fact
into the next idempotent command waits behind every projection and slow
external effect registered for earlier messages. It also lets one message
that keeps failing monopolize the consumer.

## Retry backoff for retryable failures

A retryable failure NAKs the delivery so the broker redelivers it later. On
NATS JetStream the NAK carries a delay derived from the delivery count:
`base * 2^(delivered - 1)`, capped (defaults 50 ms and 5 s;
`NatsBus::with_nack_backoff` / `NatsJetStreamSource::with_nack_backoff`).
The first retry is still prompt. A message that fails every time stops
being redelivered ahead of newer messages in a hot loop. Retries remain
unlimited and the message is still retained; only the redelivery time moves.
Setting the base to zero restores an immediate NAK.

This is not a substitute for classifying failures correctly. A command
rejection that is durably recorded under a deterministic command identity
replays the same rejection on every retry. The handler must return a
permanent error, which the configured failure policy (dead-letter by default)
settles and records in transport metrics. Transient infrastructure failures
stay retryable.

## Lanes

A router may partition its routes into named **lanes**. `Service::lane(name,
routes)` registers a route bundle in a lane. Plain `Service::routes` uses the
default lane. Lanes are opt-in. A service that registers no named lane runs
exactly as before.

When a router has more than one lane and the source settles each delivery
independently (`MessageSource::settles_independently`, true for NATS
JetStream), the runner:

1. receives messages in broker order and hands each one to every lane that
has a route for it;
2. runs each lane's messages strictly in receive order, one at a time, with
the lane's routes in registration order. This is the same order a
single-lane consumer gives those routes;
3. runs different lanes concurrently, so a slow route in one lane cannot delay
another lane;
4. settles a delivery only after every lane that received it has finished it.
Any retryable lane failure NAKs the delivery, and redelivery runs all of its
lanes again. A permanent failure applies the failure policy once;
5. bounds the number of received but unsettled deliveries
(`RunOptions::with_lane_window`, default 16). When the window is full the
runner stops receiving until a delivery settles;
6. stops a lane at its first stop-class failure (`FailurePolicy::Stop` or a
retain-and-stop error such as `ApplicationReloading`). That delivery is
settled as in sequential mode. Later deliveries queued for the halted lane
are NAKed, not skipped. Other in-flight deliveries finish and settle, then
the run returns the error.

The default lane keeps its existing behavior, including route-order
dependencies between its routes within one delivery.

Sources whose acknowledgement is positional (Kafka offset commits) or
lease-based (SQL table rows) report `settles_independently() == false`. For
them the runner keeps the sequential loop and dispatches every lane's routes
in order for each message. Lanes then change nothing about delivery.

### What a lane must satisfy

Put a route bundle in its own lane only when **no route in another lane
depends on its effects within the same delivery, and it depends on none of
theirs**. Typical examples are process-manager policies that read only the
delivered fact and send idempotent commands. A route that reads a projection
written by an earlier route for the same message must stay in that route's
lane.

Lanes also reorder effects **across deliveries**: one lane may finish later
messages while another lane is still working on earlier ones. A route whose
correctness (or whose downstream readers' correctness) needs another lane's
effects from an *earlier* message — for example, a process policy that
completes a workflow the UI then reads through a projection maintained in
another lane — must share that lane, or its readers must tolerate the
projection arriving later. Forge hit this: provisioning completed before the
read-access projection existed, producing transient 404s, until the
derivation moved into the process lane.

Lanes do not weaken existing guarantees:

- **Ordering:** each lane observes deliveries in broker order. A redelivery
(after NAK or ack-wait expiry) can arrive after later messages. That is
already true for a single-lane consumer.
- **At-least-once:** a delivery is acknowledged only after all of its lanes
succeeded. A crash or failure before that redelivers it to every lane, so
every route must stay idempotent, as today.
- **Failure policy and DLQ:** unchanged, applied once per delivery after all
lanes finish.
- **Ack deadline:** a delivery waiting in a busy lane still counts against the
broker's ack wait. Keep the window small enough that a lane's backlog clears
well within it. An expired delivery is redelivered, which is safe under the
idempotency requirement.
78 changes: 78 additions & 0 deletions docs/external-facts.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
# Authenticated external facts

`DomainEventOccurrence::capture_external` represents an authenticated fact from
an external ledger, webhook or other durable source. It does not create an
aggregate, command receipt or business authorization. The adapter must verify
the source before calling the typed constructor. Deserializing an occurrence
alone does not authenticate its origin.

The source identity consists of producer, stream, numeric position and member
key. Members of one external transaction share its position, with distinct keys.
The logical occurrence ID depends only on that identity. Changing its descriptor,
body, timestamp or metadata preserves the ID and must fail the existing durable
input fingerprint fence. Retries must reproduce the same canonical bytes. When
the source has no timestamp, an adapter may explicitly use the Unix epoch as an
unknown-time sentinel; it must not present adapter receipt time as source time.

An `external_snapshot` projection applies full-row source snapshots. It orders
each external stream/member independently of actual broker delivery positions;
it does not invent an aggregate sequence. Existing aggregate projection and
command APIs retain their authority and fencing rules. Derived external facts
retain their source provenance, but are not accepted as direct source snapshots.

NATS publishes external facts with a content-bound broker dedup key and a
separate reserved logical occurrence ID header. Thus identical retries within
the broker dedup window collapse, but altered bytes reach durable validation.
Different logical source facts with identical bodies never collapse. Retained
archive reads verify these headers and preserve external occurrences for replay.
Malformed or ambiguous identity headers fail permanently before dispatch and
remain unacknowledged. The supervisor sees the error; the adapter does not
silently terminate the message or advance progress past it.

An identical logical input redelivered at a new broker cursor advances only the
delivery checkpoint within its execution generation. It does not repeat row
changes, observations or business effects. A new execution generation applies
its first delivery normally. Original and alias cursor bindings remain immutable,
and one canonical topology-wide message identity arbitrates concurrent partition
claims. Migration 0009 adds delivery aliases without deleting canonical identity
or failure records. SQL mutation continues to use the existing partition lock.

Permanent input identities and their aliases must be retained for the replay
and source-conflict horizon. A bounded broker dedup window is not such a fence.
An ingress that acknowledges an external cursor must first establish its own
required durable qualification (for example, a history projection atomically
committed with the protocol fingerprint), and retain source replay until then.
Publishing is not approval, and waiting for all UI consumers is unnecessary.

Offline snapshot rebuild preserves these same identities. It distinguishes
original aggregate positions, original external stream/position/member keys,
and derived occurrence IDs; unrelated external facts cannot collide at empty
aggregate fields. Duplicate logical IDs must retain identical canonical bytes.
Aggregate projections still require a complete original aggregate sequence
prefix. External positions may be sparse or shared by different members, so
the source adapter must certify the complete retained source range instead of
inventing a contiguous aggregate history. Rebuild still rejects omitted stored
source versions, changed content at one source position, a different source
claiming an existing row, and derived facts used as snapshot authority. It does
not run business handlers, republish events, or reset delivery cursors.

Some original aggregate records are intentionally private and therefore absent
from the public archive. Public-archive-only rebuild remains fail-closed at such
a gap. An offline adapter can instead supply `AggregateRebuildCoverage`, built
from a complete original quiescent event-store stream, its independently read
head, an explicit authored private-contract inventory, and the retained public
occurrences. Missing or reordered records, unknown private contracts, missing
publications, changed public identities, and cross-stream substitution fail.
Every witnessed public occurrence must remain byte-identical in the rebuild
history. No private record is fabricated as a public event.
This constructor deliberately supports only authored single-publication streams
(one ordinal-zero occurrence per non-private record). The adapter must establish
that emitter contract; arbitrary one-to-many publication completeness cannot be
inferred from archive absence and is not covered by this API.

This is an operator trust boundary, not cryptographic authentication of an
arbitrary export: retain the original source provenance and review the artifact
digest. The digest prevents substitution of the reviewed artifact. It does not
prove arbitrary private payloads reconstruct a public JSON state. Authored
event-to-public contract mappings and typed projection body validation remain
required; an absent broker message never makes a source record private.
2 changes: 2 additions & 0 deletions docs/gateway/live-sharing.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,8 @@ If a slow consumer cannot preserve all evidence within its queue, it receives
`LIVE_RESET_REQUIRED`; a blocked socket is closed so it must reconnect. No
latest-value replacement silently discards confirmation proof. Group deadline,
upstream loss and incomplete/invalid origin envelopes also require recovery.
An origin failure frame (an error-only live execution result) carries no evidence
and is not fanned out; the group's consumers receive `LIVE_RESET_REQUIRED`.

Dropping a consumer releases only its lease. Last leave aborts the actual upstream
stream, socket and origin change-feed receiver, including a pending origin SQL
Expand Down
59 changes: 59 additions & 0 deletions docs/live-query-delivery.md
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,47 @@ Every live response declares `extensions.distributed.live.mode`:
- `resumable`: a result with matching, nonempty index and cursor vectors. Existing
resume validation, reset, replay and causal reconciliation rules apply.

A failed execution has no result and therefore no `live` metadata. It is sent
as the terminal failure frame described below.

## Failed live executions

A live execution that fails (for example a storage error or statement timeout)
has no result. The server sends exactly one **failure frame** and then completes
the operation. The failure ends that subscription: the server sends no more
frames for it and never attaches a later execution's metadata to the failure
frame. A failure frame has:

- `data` that is `null` or absent;
- a nonempty `errors` array, where each entry is an object with a string
`message`;
- a valid base `extensions.distributed` envelope (`protocolVersion`,
`schemaHash`, `authorizationGeneration`, `cacheScope`, `operation`, and
`trustedPresets` when the surface has any), with no `snapshot`, `live` or
`command`.

The client checks the envelope's binding, schema, operation and authorization
generation in the same way as for any other response. The failure frame then
works like an HTTP error response: it admits nothing. It writes no data or
membership, advances no cursor or operation generation, takes no ownership and
confirms no command. The client shows the original GraphQL errors for that
operation and keeps any previously admitted data readable. It closes the
failed stream. After a bounded backoff (1s, doubling to at most 30s), it opens
a fresh subscription for the operation's current watches. That subscription
resumes from the last admitted cursors, if there are any. The backoff resets
when an admitted frame arrives, and that frame also replaces the errors. This is
how a live query recovers after its storage comes back, without a page reload.
Disposing the last watch or ending the authorization generation cancels a
pending reopen.

Only an error-only frame is a failure frame. A live frame that carries non-null
`data` (including partial data with errors) without both `snapshot` and `live`
is invalid, as are frames with `errors` that are absent, empty or malformed and
lack live metadata, and frames with a `snapshot` or `live` but not both. A
response without the `extensions.distributed` envelope is invalid as before.
Intermediaries that cannot relay a failure frame, such as shared gateway live
fan-out, end the consumer with `LIVE_RESET_REQUIRED` instead.

Snapshot delivery does not relax read permissions or invent causal evidence.
Changes affecting only denied rows must not produce activity frames. A row
leaving the authorized result disappears from that operation's membership;
Expand All @@ -30,6 +71,24 @@ empty result acquire rows without waiting for another page load. It does not
override an independently active live stream or query ownership acquired after
the subscription started; those results have no safe cross-stream ordering.

Independent streams often share an index without disagreeing about it. A
layout and a page may both select
`repositories(where: { id: { _eq: $id } }, limit: 1)` and then different
relationships below that row. A snapshot frame whose shared indexes have the
same membership as the independent owner's is admitted. Each such index must
contain the same normalized record keys in the same order and the same null
value, and the owner's index must be complete and not stale. The owner keeps
those shared indexes, and the frame writes only its other indexes and its
records. The resulting graph is exactly the frame's server result, so nothing
is fabricated and no cross-stream order is assumed.

The whole frame is rejected as before if any shared membership differs or
cannot be established. That includes GraphQL errors, operation-local embedded
rows and rows without record evidence. The subscription is then reopened after
another operation next writes one of those shared indexes. The fresh result is
checked against the owner's new membership, so a stream whose frame arrived
before the owner caught up does not wait for its own next server change.

When the last watch for a live operation is disposed, its local ownership is
retired at a monotonically increasing local boundary. A later live subscription
may take over an incomparable shared index only when that subscription started
Expand Down
21 changes: 21 additions & 0 deletions docs/live-retired-owner-handoff.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
# Snapshot stream handoff after owner disposal

A snapshot live stream can begin while a previous route still owns a shared
relationship index. Its buffered frames must not acquire that graph merely
because the old route later disposes: their start fence predates disposal.
However, retaining that fence forever leaves the new route permanently pending.

When an incoming snapshot is blocked by a retired owner's index, reopen that
contending operation at a fresh local start fence. Reject the triggering old
frame and all subsequent callbacks from its retired transport. Only a new
server-authorized initial frame may take over the graph. A still-active owner
continues to block; normal same-scope hydration must not restart retained layout
subscriptions. No application polling, forced refresh, or optimistic-state
suppression is part of this recovery.

Runtime evidence: Forge's creation page received completed/ready rows, but the
prior NewRepository operation owned a shared topology relationship at revision
24, retired at 27, after the new subscription had already started. Entity clocks
advanced while the complete-empty root index remained fenced. The existing
pre-disposal protocol test covers safety; a fresh-receiver continuation covers
liveness after the same boundary.
Loading
Loading