jetstream v2 in zig stream.waow.tech
stream docs invariants.md
6.4 kB
Markdown

Invariants #

The short list of rules Stream must never break. If a change would violate one, stop and rethink. Each is backed by at least one incident or test.

These mirror upstream Jetstream V2's specs/invariants.md, deliberately: where a rule is theirs, we hold it because clients depend on it, and stating it in their words keeps the parity audit legible. Where ours differs (RocksDB rather than pebble, for instance) the difference is noted.

docs/semantic-parity.md says what is proven; this file says what must be true.

The rules #

Sealed segments are immutable. Once a segment file is sealed its bytes never change. Compaction, merge, and timestamp import produce new files and swap them in with a tmp-write-then-rename; they never edit a sealed file in place. Anything caching segment contents can trust a sealed file's bytes are stable for the life of that file. See docs/jss-seal-spec.md, segment_rewrite.zig, segment_patch.zig.

fsync the segment before you commit metadata. For every block: append and fsync into the active segment file first, then commit the synced RocksDB batch that advances relay/cursor and the per-DID repo/<did> bookkeeping. Never the other way around. Because the cursor is written after the data, a crash between the two leaves the cursor at or before the last durable event, so restart replays a few events instead of losing them. The power-loss oracle exists to hold this line — tests/powerloss_oracle.py, which runs the real Linux binary on ext4 over a kernel NBD device and cuts power at named write/fsync/rename boundaries, then checks what survived. "Oracle" throughout this project means a test that runs the real binary and judges its behaviour from the outside, as opposed to a unit test calling a function. See archive.zig, meta_store.zig.

The cursor is inclusive, and seq 0 means "nothing yet." ?cursor=N replays starting at seq N. Sequence numbers start at 1; seq 0 is a reserved "before the beginning" sentinel, so ?cursor=0 replays everything. Seqs are assigned at ingestion and are instance-local.

At-least-once delivery; clients must be idempotent. The same event can be delivered more than once — a resuming client re-receives its last-seen event, a re-merge can re-emit a row. Exactly-once delivery is a non-goal; this matches the upstream firehose guarantee.

Per-DID order is preserved, always. Events for one DID replay in ingestion order across every phase and every segment generation. The live scheduler is FIFO per DID for this reason (pipeline.zig), and the property survives compaction and merge. AppView correctness depends on it — an AppView is a downstream atproto service that rebuilds queryable state by replaying the firehose, so if two events for one repository arrive out of order it can persist the older one last and serve stale data indefinitely.

Segment files sort in creation order, and that order is time order. Zero-padded names mean a lexicographic sort is creation order, and every event in seg_N was witnessed before every event in seg_N+1. Block topology is self-describing in each sealed footer and does not depend on an external index staying in sync.

A manifest-listed segment that cannot be read is a hole, not an empty range. Cold replay and merge propagate the error rather than skipping the file; skipping would advance a cursor past durable events nothing revisits. This one was a live bug — both cold walkers skipped, and merge's SourceIndexGap guard was untested. See cold.zig, ingest/backfill/merge.zig.

Crash loud on our own corruption; never crash on bad upstream data. Two halves of one boundary. Invalid internal state — metadata corruption, fsync failure, impossible segment structure, a broken durability invariant — should stop the process rather than limp on and risk corrupting the archive; that is why a terminal retry-pass failure requests shutdown instead of sleeping until the next interval. Invalid upstream data — a malformed frame, an over-limit field, a record no encoder can render — must never crash, stop, or exit the server: drop it, count it with a labeled metric, and keep running. Resource exhaustion is ours, not the remote's: OutOfMemory must propagate, not be classified as a bad record. Bootstrap retires a repository permanently on a structural verdict like InvalidCar, so an allocation failure reported under that name discards good data. An allocation-failure sweep over prepare/emit pins this (ingest/backfill/repos.zig).

An internal failure must never look like an absence. A read error is not a missing key; a planner failure is not "no matching data"; an encode failure is not an empty result set. Each of these was a real false-negative path. Fail closed, or answer 5xx — never return a smaller truthful-looking answer, because a client cannot tell the difference and will never come back for the gap.

A bounded producer must not outrun its consumer. Bootstrap dispatch, the bounded job queue and the completion batcher form a cycle: the queue holds worker_count * 2 while a batch is 100,000 by default (--backfill-batch-size), so submitting a whole batch before consuming from it blocks the submitting fiber while finished repositories pile up unconsumed. This deadlocked a full-network run at zero completions. A now-removed in-flight byte budget (upstream has none) contributed. The ordering rule stands on its own: submit and consume together. See ingest/backfill/engine.zig.

This survived the test suite because the pinned simulator's repositories are ~4 KiB, so the queue never filled offline. It reproduced in 90 seconds on a small cloud box crawling the real network. docs/gotchas.md has that recipe under "the pinned simulator's repositories are ~4 KiB".

Do not add a bound upstream does not have without recording the decision. Both halves of that deadlock were Stream-only mechanisms introduced without justification in their commit messages. Diverging is allowed when it is a deliberate, recorded choice.

See also #

  • docs/architecture.md — how the subsystems these rules govern fit together.
  • docs/gotchas.md — accepted limitations and traps; the things that are deliberately not invariants.
  • docs/semantic-parity.md — what is proven, row by row, against upstream.