jetstream v2 in zig stream.waow.tech
stream docs architecture.md
8.9 kB
Markdown

Architecture overview #

A fast orientation map: what the big pieces are, how they fit, and where to read more. Deliberately high-level so it doesn't rot when a function moves. Each module's //! header is authoritative for its own internals; when this file disagrees with one, the module is right and this file should be fixed.

The grouping below is deliberately the same one upstream Jetstream V2 uses in its specs/architecture.md — ingest, storage, serve, plus a testing rig. Sharing the vocabulary makes the parity audit legible: a row in semantic-parity.md names a surface, and that surface lands in exactly one subsystem here.

Stream is ~38,800 lines of Zig across 72 files. It is one static binary on one machine: it ingests every record from every repo on a relay, stores them in the JSS columnar format, and lets clients replay that history and cut over to the live firehose.

Reading the paths below. Every foo.zig / dir/foo.zig in this file is relative to src/internal/, which is grouped one directory per subsystem:

src/internal/{ingest,storage,serve,compact,runtime,timestamp}/

So ingest/backfill/merge.zig is src/internal/ingest/backfill/merge.zig, and a bare pipeline.zig lives under one of those six. If a name here does not resolve, the tree moved and this file is stale — find src -name pipeline.zig settles it in one command.

The three subsystems #

Everything is ingest (data in), storage (data on disk), or serve (data out).

Ingest — getting data onto disk #

A lifecycle state machine with three phases, owned by ingest/backfill/lifecycle.zig, with durable commit points between them; a crash mid-cutover re-enters the machine at the right spot.

  • Bootstrap — ingest/backfill/engine.zig paginates listRepos and downloads every repo via getRepo (ingest/backfill/repos.zig, ingest/backfill/fetched_car.zig, ingest/backfill/prepared_fetch.zig; worker count is the only concurrency bound, as it is upstream), while ingest/backfill/capture.zig captures the firehose into a throwaway backfill/live archive so nothing is missed.
  • Merge — ingest/backfill/merge.zig drains captured live segments into the permanent archive, dropping events already covered by each repo's backfilled head, then a compaction pass runs before serving ungates.
  • Steady state — ingest.zig pumps the firehose through pipeline.zig (32 workers, FIFO per DID) into the archive. ingest/backfill/retry.zig re-downloads repos that failed; the merge's pending pass is the same runner with eligible_status flipped, which is how it inherits backoff and host parking.

convert.zig turns firehose events into wire frames and is the validation gate: invalid upstream data is dropped and counted with a labeled metric. verify.zig does Sync 1.1 signature/MST verification and repair.zig performs authenticated whole-repo repair. cursor.zig holds the durable relay cursor.

lifecycle / cutover ingest/backfill/lifecycle.zig, ingest/backfill/phase.zig
initial backfill ingest/backfill/engine.zig, ingest/backfill/repos.zig
per-repo state ingest/backfill/repo_store.zig, ingest/backfill/host_store.zig
merge ingest/backfill/merge.zig
retry / resync ingest/backfill/retry.zig, ingest/backfill/resync.zig
live consumer ingest.zig, pipeline.zig
conversion + validation convert.zig, wire.zig
verification + repair verify.zig, repair.zig

Four worker pools, four different counts #

"32 workers" and "100 workers" describe different pools:

pool workers bounded queue = workers * 2 source
initial backfill crawl 100 200 ingest/backfill/engine.zig default_workers
live scheduler (steady state) 32 64 ingest/pipeline.zig worker_count
live repair 32 64 ingest/repair.zig worker_count
failed-repo retry 16 (4 per host) — ingest/backfill/retry.zig default_workers

Only the first is operator-tunable in practice (--backfill-workers); the live counts match upstream Atmos defaults.

workers * 2 is one expression serving three pools. It bounds backfill dispatch, the live scheduler's per-DID pending capacity, and the repair job queue — so stream_backfill_queue_depth reads 200 while live-scheduler.md and live-repair.md both describe a "64", and all three are the same rule. It is also the bound in the batch submit/consume ordering rule (invariants.md, "a bounded producer must not outrun its consumer"), so raising --backfill-workers moves a deadlock-relevant bound, not just concurrency.

Storage — the segment format and the metadata store #

Two places hold state, and the durability ordering between them is the invariant that keeps a crash safe: segment fsync first, metadata commit second.

  • Segment files — the columnar, zstd-compressed, append-only JSS logs. An active segment is a file state machine (append → flush → fsync → seal). segment_writer.zig writes blocks, segment_footer.zig finalizes the footer and header, segment.zig reads sealed files, and archive.zig owns rotation, sealing, and startup recovery. manifest.zig holds resident sealed-segment metadata (a directory scan plus self-describing headers, not a database).
  • The metadata store — meta_store.zig, RocksDB, holds everything not cheaply re-derivable from segments: relay cursor, lifecycle phase, per-DID backfill status, host aggregates, compaction watermark.

compact/ performs delete/update compaction (pass.zig chunked passes, steady.zig the long-running compactor, watermark.zig the durable watermark), with tombstone.zig deciding which rows may be dropped and segment_rewrite.zig doing the physical rewrite. gloom.zig (blocked bloom), zstd.zig, and xxhash128.zig are the format's primitives.

Serve — getting data out #

  • /subscribe websocket (server.zig) — pull-based fan-out. Every subscriber is served from wherever its cursor points: tail.zig (the in-memory hot tail) for recent events, or cold.zig + cold_cache.zig (a bounded disk walk through a shared block cache) for older cursors. There is no per-client outbound queue, so a slow reader cannot blow up memory.
  • Archive download (xrpcapi.zig) — planBackfill → getSegment / getBlock, the paginated path clients use to pull sealed history.
  • Operator surfaces — status_page.zig (hosts, accounts, segments, collections), homepage.zig, repo_export.zig.

filter.zig implements wantedCollections/wantedDids.

Testing rig #

Unusually central, and worth knowing even if you are not touching tests.

  • crashpoint.zig and segment_io.zig are deterministic fault seams: named crash points at durable commit boundaries, and injected segment I/O faults.
  • tests/oracle.py boots the real binary against the pinned upstream simulator and aborts it at every lifecycle crashpoint, verifying restart convergence.
  • tests/powerloss_oracle.py runs a real Linux binary on ext4 over kernel NBD and cuts power at named write/fsync/rename boundaries.
  • tests/differential_oracle.py compares Stream against the pinned upstream Go implementation, reader, and client.
  • scripts/admit binds a passing run of all of the above to one image digest.

Cross-cutting #

metrics.zig and observability.zig (OpenTelemetry), logging.zig, process_metrics.zig, disk_space.zig, environment.zig, operational.zig / inspect_all.zig (offline operator commands), and timestamp/ (the operator timestamp-import pipeline).

Where to look #

I want to understand… Start here
What is proven and what is not docs/semantic-parity.md
The on-disk segment format docs/jss-format-v1.md, segment.zig
Sealing rules docs/jss-seal-spec.md, segment_footer.zig
The ingest lifecycle / cutover docs/upstream-bootstrap-spec.md, ingest/backfill/lifecycle.zig
Compaction docs/upstream-compaction-spec.md, compact/
The live scheduler docs/live-scheduler.md, pipeline.zig
Repair docs/live-repair.md, repair.zig
The subscribe protocol docs/upstream-subscribe-protocol.md, server.zig
Timestamp import docs/timestamp-import.md, timestamp/
Running it — flags and env vars docs/configuration.md, runtime/cli.zig
The offline inspection subcommands docs/cli.md
Deploying docs/deployment-runbook.md, deploy/README.md, scripts/admit
Performance docs/benchmarks.md, just bench
Accepted limitations and traps docs/gotchas.md
The rules that must not break docs/invariants.md
Why a thing is the way it is docs/lessons-from-zlay.md