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.zigpaginates 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), whileingest/backfill/capture.zigcaptures the firehose into a throwawaybackfill/livearchive so nothing is missed. - Merge —
ingest/backfill/merge.zigdrains 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.zigpumps the firehose throughpipeline.zig(32 workers, FIFO per DID) into the archive.ingest/backfill/retry.zigre-downloads repos that failed; the merge's pending pass is the same runner witheligible_statusflipped, 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.zigwrites blocks,segment_footer.zigfinalizes the footer and header,segment.zigreads sealed files, andarchive.zigowns rotation, sealing, and startup recovery.manifest.zigholds 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 #
/subscribewebsocket (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, orcold.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.zigandsegment_io.zigare deterministic fault seams: named crash points at durable commit boundaries, and injected segment I/O faults.tests/oracle.pyboots the real binary against the pinned upstream simulator and aborts it at every lifecycle crashpoint, verifying restart convergence.tests/powerloss_oracle.pyruns a real Linux binary on ext4 over kernel NBD and cuts power at named write/fsync/rename boundaries.tests/differential_oracle.pycompares Stream against the pinned upstream Go implementation, reader, and client.scripts/admitbinds 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 |