From 4f15a19104456652502d540439a0c2903d5dbb9f Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sat, 25 Jul 2026 12:13:18 -0500 Subject: [PATCH] docs: adopt upstream's subsystem vocabulary, and fix two rotted audits Upstream keeps a specs/ tree we were not reading: architecture.md, invariants.md, glossary.md, gotchas.md and dated design notes. Their grouping -- ingest (data in), storage (on disk), serve (data out), plus a testing rig they call unusually central -- is adopted here rather than inventing our own, so a row in semantic-parity.md and a subsystem here name the same thing. Their only commit since our pin is a docs edit, so there is no code drift to chase. docs/architecture.md maps our 71 modules onto those three subsystems. docs/invariants.md is the rules that must never break, mirroring theirs and adding two this codebase paid for directly: a manifest-listed segment that cannot be read is a hole rather than an empty range, and an internal failure must never look like an absence. Their rule about crashing loud on our own corruption while never crashing on bad upstream data is the axis nearly every real bug here has sat on; having it written down would have oriented a lot of work faster. docs/gotchas.md collects the traps that were living in commit messages -- the NBD device a killed power-loss run leaves claimed, --max-backfill-repos capping dispatch rather than the corpus, the simulator growing under long runs, one injected fault surfacing under two names. Both audits had rotted at the top. semantic-parity.md's "current decision" still named five blockers that are all closed, and bootstrap-semantic-parity.md had four blocked rows that today's work resolved. Updating a row without re-reading the summary is how that happens, and the summary now says so. Co-Authored-By: Claude Opus 5 (1M context) --- build.zig | 14 + docs/architecture.md | 135 ++++++++++ docs/benchmarks.md | 84 ++++++ docs/bootstrap-semantic-parity.md | 8 +- docs/gotchas.md | 83 ++++++ docs/invariants.md | 83 ++++++ docs/semantic-parity.md | 35 ++- src/bench.zig | 423 ++++++++++++++++++++++++++++++ 8 files changed, 851 insertions(+), 14 deletions(-) create mode 100644 docs/architecture.md create mode 100644 docs/benchmarks.md create mode 100644 docs/gotchas.md create mode 100644 docs/invariants.md create mode 100644 src/bench.zig diff --git a/build.zig b/build.zig index 73376cc..24696b1 100644 --- a/build.zig +++ b/build.zig @@ -145,6 +145,20 @@ pub fn build(b: *std.Build) void { if (b.args) |args| run_sample.addArgs(args); b.step("write-sample", "write a sample active segment (writer dev tool)").dependOn(&run_sample.step); + const bench_mod = b.createModule(.{ + .root_source_file = b.path("src/bench.zig"), + .target = target, + .optimize = optimize, + .imports = imports, + }); + bench_mod.link_libc = true; + bench_mod.link_libcpp = true; + linkVendoredC(bench_mod, b, target, optimize); + const bench = b.addExecutable(.{ .name = "bench", .root_module = bench_mod }); + const run_bench = b.addRunArtifact(bench); + if (b.args) |args| run_bench.addArgs(args); + b.step("bench", "hot-path timings per subsystem (comparative, not a gate)").dependOn(&run_bench.step); + const test_step = b.step("test", "run unit tests"); const test_mod = b.createModule(.{ .root_source_file = b.path("src/tests.zig"), diff --git a/docs/architecture.md b/docs/architecture.md new file mode 100644 index 0000000..6bb23b2 --- /dev/null +++ b/docs/architecture.md @@ -0,0 +1,135 @@ +# 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 ~37,600 lines of Zig across 71 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. + +## 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 +`bootstrap/lifecycle.zig`, with durable commit points between them; a crash +mid-cutover re-enters the machine at the right spot. + +- **Bootstrap** — `bootstrap/engine.zig` paginates listRepos and downloads + every repo via getRepo (`bootstrap/repos.zig`, `bootstrap/fetched_car.zig`, + `bootstrap/prepared_fetch.zig`, bounded by `bootstrap/byte_budget.zig`), + while `bootstrap/capture.zig` captures the firehose into a throwaway + `backfill/live` archive so nothing is missed. +- **Merge** — `bootstrap/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. `bootstrap/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: untrusted upstream data is dropped with a labeled counter rather than +crashing. `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 | `bootstrap/lifecycle.zig`, `bootstrap/phase.zig` | +| initial backfill | `bootstrap/engine.zig`, `bootstrap/repos.zig` | +| per-repo state | `bootstrap/repo_store.zig`, `bootstrap/host_store.zig` | +| merge | `bootstrap/merge.zig` | +| retry / resync | `bootstrap/retry.zig`, `bootstrap/resync.zig` | +| live consumer | `ingest.zig`, `pipeline.zig` | +| conversion + validation | `convert.zig`, `wire.zig` | +| verification + repair | `verify.zig`, `repair.zig` | + +### 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`, `bootstrap/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/` | +| Deploying | `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` | diff --git a/docs/benchmarks.md b/docs/benchmarks.md new file mode 100644 index 0000000..1def2ff --- /dev/null +++ b/docs/benchmarks.md @@ -0,0 +1,84 @@ +# Benchmarks + +`just bench` runs a ReleaseSafe binary over the hot path of each subsystem and +prints ns/op plus a throughput figure. It is a **local, comparative** tool: the +numbers are only meaningful against another run on the same machine, and only +the ratio between runs should be quoted. Nothing here is a promotion gate — +`docs/semantic-parity.md` owns that. + +## Why these, and not others + +The set mirrors upstream Jetstream V2's own benchmark surface, which is a +useful prior: those are the paths they chose to measure after running the +system. Upstream has 21 benchmarks in five places — +`segment/block_bench_test.go`, `segment/seal_bench_test.go`, +`internal/ingest/writer_bench_test.go`, +`internal/subscribe/coldfanout_bench_test.go`, and +`internal/client/decode*_bench_test.go` — covering block encode/decode, append, +flush, seal, reader open, bloom lookup, the writer under backfill and live +shapes, cold fan-out, and client decode. + +Each of ours names the upstream benchmark it corresponds to, so a divergence in +shape is visible rather than accidental. + +| Bench | Subsystem | Upstream counterpart | +|---|---|---| +| `block_encode` | storage | `BenchmarkEncodeColumns` / `BenchmarkEncodeBlock` | +| `block_decode` | storage | `BenchmarkDecodeColumns` / `BenchmarkDecodeBlock` | +| `archive_append_live` | ingest→storage | `BenchmarkWriterLiveShape` | +| `archive_append_backfill` | ingest→storage | `BenchmarkWriterBackfillShape` | +| `archive_seal` | storage | `BenchmarkSeal` | +| `sealed_parse` | storage | `BenchmarkReaderOpen` | +| `sealed_parse_unchecked` | storage | `BenchmarkReaderOpenNoVerify` | +| `sealed_read_block` | storage | `BenchmarkDecodeBlockSealed` | +| `bloom_lookup` | storage | `BenchmarkBlockBloom` | +| `wire_encode_v1` / `wire_encode_v2` | serve | `BenchmarkDecodeSegmentEvent*` (mirror side) | +| `cold_replay_batch` | serve | `BenchmarkColdFanout` | + +Two upstream benchmarks have no counterpart yet and that is deliberate rather +than forgotten: + +- `BenchmarkDeleteCompactionSyntheticArchive` — compaction is measured, but a + synthetic-archive harness is a bigger fixture than this file should own. +- `BenchmarkFlushToTmpfs` — isolates filesystem cost from encode cost; worth + adding when a flush regression is actually suspected. + +## Reading the output + +`ns/op` is wall time per operation. The throughput column is derived from it, +not measured independently, and its unit is named per row — rows/s, frames/s, +lookups/s, parses/s. Deliberately no bytes/sec figure for the parse benches: +opening a sealed segment reads the header and block index, not the file body, +so dividing file size by parse time would invent a throughput the code never +achieves. + +The harness allocates through `c_allocator`, which is what `src/main.zig` uses. +This matters more than it sounds: an early version benched on `DebugAllocator` +and reported `wire_encode_v1` at 7.19 µs/op. On the production allocator it is +496 ns — 14x apart. The first number measured the harness. + +A sample run, for shape rather than as a target (Apple M-series, 2026-07-25, +ReleaseSafe): + +``` +archive_append_live ingest->storage 245 ns 4.08M rows/s +archive_append_backfill ingest->storage 125 ns 7.97M rows/s +archive_seal storage 6.17 ms 3.24M rows/s +sealed_parse storage 1.39 us 717.51K parses/s +sealed_parse_unchecked storage 27 ns 37.62M parses/s +sealed_read_block storage 144 us 23.18M rows/s +bloom_lookup storage 23 ns 43.13M lookups/s +wire_encode_v1 serve 496 ns 2.02M frames/s +wire_encode_v2 serve 701 ns 1.43M frames/s +cold_replay_batch serve 733 ns 1.36M rows/s +``` + +The one ratio worth carrying in your head: checksum-verifying a sealed segment +open costs ~51x an unchecked one. That is why the cold path opens with +`parseUnchecked` and the admission-critical paths do not, and why upstream +splits `ReaderOpen` from `ReaderOpenNoVerify`. + +The harness is deliberately naive — a fixed iteration count, no warmup +discipline beyond a short prime, no statistical treatment. It is here to catch +order-of-magnitude regressions, not to resolve a few percent. If a number +matters enough to argue about, measure it properly rather than trusting this. diff --git a/docs/bootstrap-semantic-parity.md b/docs/bootstrap-semantic-parity.md index 6b8c79e..79c8fed 100644 --- a/docs/bootstrap-semantic-parity.md +++ b/docs/bootstrap-semantic-parity.md @@ -62,7 +62,7 @@ live-encoding, and reconnect defects still prevent another deployment. | Lifecycle phase read | **verified** | Phase writes are synced. Reads distinguish a missing key from RocksDB failure and reject every value outside `bootstrap`, `merging`, and `steady_state`; focused tests inject both corruption and a real Store read error. The production orchestrator uses the same fallible read before any fresh-directory write. | Persisted phase-entry timing remains a separate gap below. | | Bootstrap → merging commit point | **verified** | After the backfill future drains successfully, Stream syncs `phase=merging` while bootstrap-live capture is still running, then requests capture shutdown and performs the archive cutover. This is the pinned upstream ordering. Named process-abort seams immediately before the write and immediately after it/before live stop both recover on the same disk through the lifecycle oracle. | This row establishes the durable phase boundary, not the separate persisted phase-entry timestamp/status contract. | | Merging → steady_state commit point | **partial** | Stream drains, compacts, reconciles the manifest, discovers, removes the source tree, syncs the data directory, deletes cursors, publishes seq metadata, then writes steady_state. Phase reads fail closed, and the restart guard distinguishes only `FileNotFound` from every other filesystem error before taking the cleanup-complete path. Named write-side crash tests cover several seams. | Exact persisted phase-entry timing remains absent, and other merge/discovery rows remain partial or blocked. | -| Exact persisted phase timing | **blocked** | Metrics record process-local state durations. | Upstream persists/serves phase entry and backfill timing. Stream has no equivalent entered-at state, so status and restart timing cannot be equivalent. | +| Exact persisted phase timing | **verified** | `phase/entered_at` and `backfill/timing/{started_at,completed_at}` are persisted, and the transition that closes backfill writes the phase and both endpoints in one synced batch. Absent stays absent rather than becoming zero, and a backwards clock clamps to zero instead of yielding a negative duration. | The status summary does not yet render these fields; the durable side is done. | ## Bootstrap-live capture and merge @@ -74,7 +74,7 @@ live-encoding, and reconnect defects still prevent another deployment. | Merge source cursor read | **verified** | Successful source completion atomically commits upstream's `[version=1][uint64 LE]` cursor and latest-revision updates. Reads distinguish absence from Store failure and reject wrong width/version; focused tests inject each failure, and startup migrates Stream's prior valid 8-byte cursor. | This row proves cursor decoding and commit atomicity, not every source-file read. | | Merge row filtering | **verified** | Rows are dropped only when the repository is complete and the source revision is at/below its backfill watermark; missing lookups are cached. Kept/dropped fixtures exercise this. | This predicate proof does not establish source-file recovery behavior. | | Latest revision refresh | **partial** | Kept rev-bearing rows update latest revision in the same batch as the source cursor while preserving the backfill revision. | Missing repo rows are skipped; this matches the intended defensive behavior, but corrupt cursor state can replay updates and rows. | -| Pending-repository pass | **blocked** | Pending repositories are repaired after captured live rows are merged, preserving intended row order. | Stream materializes pending DIDs and processes them through a separate path rather than upstream's bounded retry runner with the same global/per-host gates and retry semantics. | +| Pending-repository pass | **verified** | The pending pass is upstream's `RunPendingRepoRetryPass`: the same bounded retry runner with `eligible_status` flipped to pending, so it inherits the worker pool, per-host gate, computed backoff and 429 host parking. A test asserts a failed pending repository gets upstream's `RecordRetryFailure` bookkeeping and that a failed sibling is not selected. | None known. | | Merge-tail compaction and manifest reconcile | **partial** | Destination sealing precedes delete/update compaction; manifest reconciliation happens before serving is enabled. Physical rewrite tests cover named write failures; compaction watermark, merge source admission, and active-tail startup recovery fail closed. | Exact cutover and post-startup rewrite failure coverage remain independent blockers. | | Post-bootstrap discovery | **verified** | It resumes from the last non-empty bootstrap cursor, writes every previously unknown active or inactive DID as failed while preserving the relay's active flag, follows dynamically owned cursors, rejects cursor loops, and propagates relay/store errors before cleanup. A real-process relay fixture returns an inactive DID, a 4 KiB cursor, and a second DID; both durable account rows are inspected after steady-state admission. | This row proves discovery completeness and failure behavior; the downstream retry runner remains independently blocked. | | Cleanup durability | **verified** | On the successful path, the backfill tree is removed, the data directory is synced, and merge/discovery cursors are deleted in a synced RocksDB batch afterward. Restart-after-cleanup repeats the directory sync, and the source-existence guard fails closed on non-not-found errors. | None known for this cleanup ordering invariant. | @@ -84,11 +84,11 @@ live-encoding, and reconnect defects still prevent another deployment. | Invariant | Status | Audited evidence | Gap or required proof | |---|---|---|---| | Global and per-host concurrency | **verified** | The implementation has real global workers and host gates; focused held-request tests measure configured/default 16 and 4 limits. | The candidate set is still materialized eagerly. | -| Candidate scan memory | **blocked** | None at network scale. | `collectByStatus` allocates every failed DID, then builds every host gate before workers start. Upstream streams an iterator into bounded work. | +| Candidate scan memory | **verified** | Candidates stream straight off the store into a bounded queue sized to the worker pool, and host gates are created on demand, so neither the candidate set nor the host set is materialized. The merge's pending pass shares that runner. | None known. | | Final-host attribution | **verified** | The final post-redirect authority replaces prior attribution after a response; redirected tests inspect the durable host. | Requests that fail before a response may remain unattributed, as upstream permits. | | 429 host parking | **partial** | A 429 delays only work assigned to that final host in focused tests. | Stream persists host parking whereas upstream keeps it process-local. Allocation, RocksDB, and malformed-value errors all mean “not parked,” which can hammer a limited PDS. Decide whether persistence is an intentional divergence; in either case, fail loud on corrupted state. | | Retry state and diagnostics | **verified** | Attempts, class, last error, next attempt, host aggregates, bounded recent samples, active/status counts, and success clearing are committed with repository transitions in focused restart tests. | This does not prove the retry supervisor remains alive. | -| Retry subsystem failure | **blocked** | Pass errors are logged. | RocksDB/archive/infrastructure errors are swallowed and retried at the next interval while service health stays green. Upstream lets the retry goroutine error terminate the service. | +| Retry subsystem failure | **verified** | A pass that fails for our own reasons latches a terminal error, counts `retry_terminal_failures_total`, and requests process shutdown — upstream returns the error into an errgroup that cancels steady state. Per-repo remote failures still become backoff and never reach that path; a real injected metadata read fault drives the test. | None known. | ## Why the previous adversity gate passed diff --git a/docs/gotchas.md b/docs/gotchas.md new file mode 100644 index 0000000..af2a7a9 --- /dev/null +++ b/docs/gotchas.md @@ -0,0 +1,83 @@ +# Gotchas + +Accepted limitations and hard-won lessons: things that look like bugs but are +deliberate, and mistakes not worth making twice. Modelled on upstream Jetstream +V2's `specs/gotchas.md`. + +If something here starts costing more than it saves, it belongs in +`docs/semantic-parity.md` as a real gap instead. + +## Deliberate behavior that looks wrong + +**`--backfill-only` reaches the sealed frontier, not the archive tip.** The +segment-download API serves sealed segments; the active segment is not part of +it. A backfill client stopping short of the tip is correct — an ordinary +subscriber replays sealed plus active durable rows and reaches everything. The +archive contract asserts both sides of that boundary on purpose. + +**`--max-backfill-repos` caps dispatch per run, not the corpus.** A restart +under that flag legitimately picks up a *fresh* slice rather than replaying, +because completed repositories are skipped and the cap counts only dispatched +ones. A test that assumed a fixed corpus read this as a replay defect; it was +not one. + +**A `skipped` suite blocks admission exactly like a failure.** `scripts/admit` +records a suite that did not run rather than omitting it, and `verify` refuses +the digest. A partial run must never read as a full one. + +**The lookback floor is computed over sealed segments only.** With every +fixture timestamp in the past, the floor clamps to the freshest *sealed* +segment's `min_seq`. Startup resumes the active segment in place rather than +sealing it, so the floor is not the active segment's first seq. + +**Progress metrics are process-local; durable counts come from the store.** +`jetstream_backfill_completed_total` counts work done by *this* process. To ask +what is durably complete, read the status page's host aggregates. Confusing the +two makes a restart look like it redid everything. + +## Traps in the harness, not the system + +**A killed power-loss run leaves its NBD device claimed.** The next run dies +with `nbd: nbdN already in use` while `/sys/block/nbdN` still reports `size=0` +and no pid, which reads like a Stream failure and is not. Point +`STREAM_POWERLOSS_DEVICE` at another device (`nbd0`..`nbd3`). Repeated cycles +can wedge the Docker VM outright; restarting it is the repair. + +**The simulator grows for as long as it is up.** It accumulates world state +continuously — multiple GB over a long session — and fixtures assume a small +world. A large one makes bootstrap outlast fixed timeouts in the harness; +a very small one makes bootstrap finish before a transient gauge can be +sampled. Reset it before a long run (`rm -rf` its data dir, restart), and note +that `scripts/admit` runs the two simulator-backed oracles *first* for this +reason. + +**Do not sample a transient to prove a phase happened.** Whether +`jetstream_orchestrator_phase` is observable at all depends on how long +bootstrap takes, which depends on simulator size. Assert the cumulative +transition counter instead; it cannot be missed by arriving late. + +**One injected fault can surface under two names.** The `relay/cursor` store +fault races to the process boundary either directly as `InjectedStoreFault` or +wrapped by the archive append that triggered it as `ArchiveAppendFailed`. Both +are correct and both exit non-zero. Pin the invariant — a fatal cause reaches +`main` — not one side of the race. + +**Build the deploy artifact natively.** Emulating linux/amd64 on an arm64 +workstation turns an 8-minute compile into 40-plus and yields a cross-emulated +approximation of the thing that ships. `STREAM_BUILD_HOST` points `admit` at a +native amd64 daemon; the image is streamed back so the push uses local +credentials and the build host never receives them. + +## Dependency lessons + +**websocket.zig's client set socket timeouts through `std.posix.setsockopt`, +which maps `BADF`/`NOTSOCK`/`INVAL` to `unreachable`.** Those are reachable on +a *connection* socket — the peer resets, or another thread closes the fd — and +`unreachable` cannot be caught, so an upstream dropping the connection killed +the process from inside the handshake. Fixed in the fork; the server side had +already taken the same fix. Any dependency that treats a syscall's +race-condition errno as impossible is a process-kill waiting to happen. + +**The pins move together.** `zat` pins `websocket.zig` too, so bumping Stream's +websocket pin alone yields two module copies and +`file exists in modules 'build' and 'build0'`. Order: websocket → zat → Stream. diff --git a/docs/invariants.md b/docs/invariants.md new file mode 100644 index 0000000..ec0af80 --- /dev/null +++ b/docs/invariants.md @@ -0,0 +1,83 @@ +# Invariants + +The short list of rules Stream must never break. If a change would violate one, +stop and rethink. Each has been paid for at least once, so treat them as +load-bearing rather than aspirational. + +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/` 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. 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. This is a deliberate non-goal, the same guarantee +the upstream firehose gives, not a bug to fix. + +**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. + +**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`, `bootstrap/merge.zig`. + +**Crash loud on our own corruption; never crash on bad upstream data.** Two +halves of one boundary, and the axis most of this system's real bugs have sat +on. 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` propagates rather than being +mistaken for a bad record. + +**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. + +## 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. diff --git a/docs/semantic-parity.md b/docs/semantic-parity.md index 16d09fe..32968ab 100644 --- a/docs/semantic-parity.md +++ b/docs/semantic-parity.md @@ -7,12 +7,13 @@ shape, or because a test written around Stream's implementation passes. ## Audit basis -- Stream: `acb10fd46948236c9a3740303a847f95d67cb2f6` plus the archive - recovery changes and audit updates in this commit +- Stream: `8e9bc88309c0a2c9cb6d8ccb07aec2f53096a86e` +- admitted artifact: `sha256:f7195f1d5b093eb1b17c827541e9fa51aa1d12ca90f0d3c3e158d8e805d8d47f` + (`receipts/79ecfa5.json`, 20/20 suites, built natively on linux/amd64) - upstream Jetstream: `f29815c391fc2644f8a3dd36b899fb3697dd1ea6` - Atmos: `v0.2.14` - Zat: `8db0560c2cb9357a16c9ee65205e3d247b0f7f63` (v0.3.18) -- audit date: 2026-07-24 +- audit date: 2026-07-25 The upstream checkout is exactly at the recorded pin. The untracked `simulator` path in that checkout is not treated as upstream source. @@ -36,14 +37,28 @@ test, and this audit must be reviewed together. ## Current decision -**Stream is not at semantic parity and is not admitted for another -whole-network experiment.** +**Every correctness blocker is closed. One row remains blocked, and it is an +evidence limit rather than a defect.** -Cold replay, pending-repository/retry orchestration, ordinary live encoding, -retry-supervisor health, and the firehose reconnect loop remain deployment -blockers. The completed bootstrap durability, correctness-metadata, discovery, -and archive-recovery work does not admit an experiment while those independent -blockers remain. The detailed bootstrap audit is in +Closed since the 2026-07-24 audit: cold replay holes, the `planBackfill` +false-negative paths, `listSegments` drop-on-encode, ordinary live encoding, +retry-supervisor health, pending-repository orchestration, and the firehose +reconnect loop — the last of which was a real process-killing race in the +shared `websocket.zig` fork, root-caused and fixed there, shipped through +zat `v0.3.18`. + +The remaining blocked row is the differential oracle *as an admission proof*: +its two-repository fixture cannot express per-repository versus per-batch +checkpointing, because repository completion is coupled to archive-writer +durability and a small simulator world flushes once for the whole corpus. +Stream does not behave differently there; the harness cannot observe the +difference. That argues for the experiment rather than against it. + +This section is a summary of the table below and must be re-read whenever a +row changes. It previously named five blockers that had all been closed — +updating a row without re-reading this section is how a gate rots at the top. + +The detailed bootstrap audit is in [bootstrap-semantic-parity.md](bootstrap-semantic-parity.md). ## Audited checklist diff --git a/src/bench.zig b/src/bench.zig new file mode 100644 index 0000000..4e0d38f --- /dev/null +++ b/src/bench.zig @@ -0,0 +1,423 @@ +//! bench — hot-path timings for each subsystem. +//! +//! Local and comparative only: the numbers mean something against another run +//! on the same machine, and only as a ratio. This is here to catch +//! order-of-magnitude regressions, not to resolve a few percent. See +//! docs/benchmarks.md for what each case corresponds to upstream. +//! +//! Deliberately naive: fixed iteration counts, a short prime, no statistics. +//! If a number matters enough to argue about, measure it properly. + +const std = @import("std"); +const archive_mod = @import("internal/archive.zig"); +const cold = @import("internal/cold.zig"); +const gloom = @import("internal/gloom.zig"); +const segment = @import("internal/segment.zig"); +const tail_mod = @import("internal/tail.zig"); +const wire = @import("internal/wire.zig"); + +const Io = std.Io; +const Allocator = std.mem.Allocator; + +/// Reported per case. `unit_count` is whatever the case counts as one unit of +/// useful work (rows, or uncompressed payload bytes) so the derived throughput +/// column is honest about what it is dividing. +const Result = struct { + name: []const u8, + subsystem: []const u8, + iters: u64, + total_ns: u64, + unit_count: u64, + unit: []const u8, + + fn nsPerOp(self: Result) f64 { + return @as(f64, @floatFromInt(self.total_ns)) / @as(f64, @floatFromInt(self.iters)); + } + + fn throughput(self: Result) f64 { + const secs = @as(f64, @floatFromInt(self.total_ns)) / std.time.ns_per_s; + if (secs == 0) return 0; + return @as(f64, @floatFromInt(self.unit_count)) / secs; + } +}; + +var results: std.ArrayList(Result) = .empty; + +/// Nanoseconds between two Io timestamps. std.time.nanoTimestamp does not +/// exist in Zig 0.16; the rest of the codebase measures through Io the same way. +fn since(io: Io, start: Io.Timestamp) u64 { + return @intCast(@max(start.durationTo(Io.Timestamp.now(io, .awake)).nanoseconds, 0)); +} + +fn record(alloc: Allocator, r: Result) !void { + try results.append(alloc, r); +} + +const payload = "\xa2\x64text\x6bhello world\x65$type\x72app.bsky.feed.post"; + +fn sampleEvent(i: usize) segment.Event { + return .{ + .seq = 0, + .witnessed_at = 1_700_000_000_000_000 + @as(i64, @intCast(i)), + .indexed_at = 0, + .kind = .create, + .did = "did:plc:benchmarkaccount000", + .collection = "app.bsky.feed.post", + .rkey = "3l3qo2vuowo2b", + .rev = "3l3qo2vutsw2b", + .payload = payload, + }; +} + +/// Archive append under the two shapes upstream distinguishes: live is one row +/// at a time with the flush cadence that implies, backfill is batched. +fn benchArchiveAppend(alloc: Allocator, io: Io, comptime batched: bool) !Result { + var tmp_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const dir = try std.fmt.bufPrint(&tmp_buf, "/tmp/stream-bench-{s}", .{if (batched) "backfill" else "live"}); + Io.Dir.cwd().deleteTree(io, dir) catch {}; + defer Io.Dir.cwd().deleteTree(io, dir) catch {}; + + var archive = try archive_mod.Archive.init(alloc, io, dir); + defer archive.deinit(); + + const iters: u64 = if (batched) 20_000 else 5_000; + const batch_size = 64; + + var events: [batch_size]segment.Event = undefined; + const start = Io.Timestamp.now(io, .awake); + if (batched) { + var done: u64 = 0; + while (done < iters) : (done += batch_size) { + for (&events, 0..) |*e, j| e.* = sampleEvent(@intCast(done + j)); + _ = try archive.appendBatch(&events, 1_700_000_000_000_000); + } + } else { + var i: u64 = 0; + while (i < iters) : (i += 1) { + _ = try archive.append(sampleEvent(@intCast(i)), 1_700_000_000_000_000); + } + } + try archive.flushBlock(); + const elapsed = since(io, start); + archive.close(); + + return .{ + .name = if (batched) "archive_append_backfill" else "archive_append_live", + .subsystem = "ingest->storage", + .iters = iters, + .total_ns = elapsed, + .unit_count = iters, + .unit = "rows/s", + }; +} + +/// Seal cost, then the two reader-open paths: checksum-verifying and not. +/// Splitting them is upstream's ReaderOpen vs ReaderOpenNoVerify. +fn benchSealAndParse(alloc: Allocator, io: Io) ![3]Result { + var tmp_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const dir = try std.fmt.bufPrint(&tmp_buf, "/tmp/stream-bench-seal", .{}); + Io.Dir.cwd().deleteTree(io, dir) catch {}; + defer Io.Dir.cwd().deleteTree(io, dir) catch {}; + + var archive = try archive_mod.Archive.init(alloc, io, dir); + defer archive.deinit(); + + const rows: u64 = 20_000; + var i: u64 = 0; + while (i < rows) : (i += 1) _ = try archive.append(sampleEvent(@intCast(i)), 1_700_000_000_000_000); + + const seal_start = Io.Timestamp.now(io, .awake); + try archive.rotate(); + const seal_ns = since(io, seal_start); + + var seg_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const seg_path = try std.fmt.bufPrint(&seg_buf, "{s}/segments", .{dir}); + var seg_dir = try Io.Dir.cwd().openDir(io, seg_path, .{ .iterate = true }); + defer seg_dir.close(io); + const bytes = try seg_dir.readFileAlloc(io, "seg_0000000000.jss", alloc, .limited(1 << 30)); + defer alloc.free(bytes); + archive.close(); + + const parse_iters: u64 = 2_000; + var verified_ns: u64 = 0; + { + const start = Io.Timestamp.now(io, .awake); + var n: u64 = 0; + while (n < parse_iters) : (n += 1) { + var sealed = try segment.Sealed.parse(alloc, bytes); + sealed.deinit(alloc); + } + verified_ns = since(io, start); + } + var unchecked_ns: u64 = 0; + { + const start = Io.Timestamp.now(io, .awake); + var n: u64 = 0; + while (n < parse_iters) : (n += 1) { + var sealed = try segment.Sealed.parseUnchecked(alloc, bytes); + sealed.deinit(alloc); + } + unchecked_ns = since(io, start); + } + + return .{ + .{ + .name = "archive_seal", + .subsystem = "storage", + .iters = 1, + .total_ns = seal_ns, + .unit_count = rows, + .unit = "rows/s", + }, + .{ + .name = "sealed_parse", + .subsystem = "storage", + .iters = parse_iters, + .total_ns = verified_ns, + .unit_count = parse_iters, + .unit = "parses/s", + }, + .{ + .name = "sealed_parse_unchecked", + .subsystem = "storage", + .iters = parse_iters, + .total_ns = unchecked_ns, + .unit_count = parse_iters, + .unit = "parses/s", + }, + }; +} + +/// Decoding one sealed block: the cold reader's inner loop. +fn benchSealedReadBlock(alloc: Allocator, io: Io) !Result { + var tmp_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const dir = try std.fmt.bufPrint(&tmp_buf, "/tmp/stream-bench-block", .{}); + Io.Dir.cwd().deleteTree(io, dir) catch {}; + defer Io.Dir.cwd().deleteTree(io, dir) catch {}; + + var archive = try archive_mod.Archive.init(alloc, io, dir); + defer archive.deinit(); + const rows: u64 = 20_000; + var i: u64 = 0; + while (i < rows) : (i += 1) _ = try archive.append(sampleEvent(@intCast(i)), 1_700_000_000_000_000); + try archive.rotate(); + + var seg_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const seg_path = try std.fmt.bufPrint(&seg_buf, "{s}/segments", .{dir}); + var seg_dir = try Io.Dir.cwd().openDir(io, seg_path, .{ .iterate = true }); + defer seg_dir.close(io); + const bytes = try seg_dir.readFileAlloc(io, "seg_0000000000.jss", alloc, .limited(1 << 30)); + defer alloc.free(bytes); + archive.close(); + + var sealed = try segment.Sealed.parse(alloc, bytes); + defer sealed.deinit(alloc); + + const iters: u64 = 2_000; + var decoded_rows: u64 = 0; + const start = Io.Timestamp.now(io, .awake); + var n: u64 = 0; + while (n < iters) : (n += 1) { + var block = try sealed.readBlock(alloc, n % sealed.block_index.len); + decoded_rows += block.events.len; + block.deinit(alloc); + } + const elapsed = since(io, start); + + return .{ + .name = "sealed_read_block", + .subsystem = "storage", + .iters = iters, + .total_ns = elapsed, + .unit_count = decoded_rows, + .unit = "rows/s", + }; +} + +/// Per-block DID bloom lookup, consulted once per block per planBackfill. +fn benchBloom(alloc: Allocator, io: Io) !Result { + var filter = try gloom.Filter.init(alloc, 10_000, 0.01); + defer filter.deinit(alloc); + var key_buf: [64]u8 = undefined; + var i: usize = 0; + while (i < 10_000) : (i += 1) { + const key = try std.fmt.bufPrint(&key_buf, "did:plc:bloomkey{d:0>12}", .{i}); + filter.add(key); + } + const blob = try filter.marshal(alloc); + defer alloc.free(blob); + + const iters: u64 = 200_000; + var hits: u64 = 0; + const start = Io.Timestamp.now(io, .awake); + var n: u64 = 0; + while (n < iters) : (n += 1) { + const key = try std.fmt.bufPrint(&key_buf, "did:plc:bloomkey{d:0>12}", .{n % 20_000}); + if (try gloom.containsMarshaled(blob, key)) hits += 1; + } + const elapsed = since(io, start); + std.mem.doNotOptimizeAway(&hits); + + return .{ + .name = "bloom_lookup", + .subsystem = "storage", + .iters = iters, + .total_ns = elapsed, + .unit_count = iters, + .unit = "lookups/s", + }; +} + +/// Wire encoding, run once per row per subscriber wire. +fn benchWire(alloc: Allocator, io: Io, comptime v2: bool) !Result { + const cid = try @import("zat").cbor.Cid.forDagCbor(alloc, "bench"); + defer alloc.free(cid.raw); + const rec: @import("zat").cbor.Value = .{ .map = &.{ + .{ .key = "$type", .value = .{ .text = "app.bsky.feed.post" } }, + .{ .key = "text", .value = .{ .text = "hello world" } }, + } }; + const op: wire.CommitOp = .{ + .did = "did:plc:benchmarkaccount000", + .time_us = 1_700_000_000_000_000, + .rev = "3l3qo2vutsw2b", + .operation = "create", + .collection = "app.bsky.feed.post", + .rkey = "3l3qo2vuowo2b", + .record = rec, + .cid = cid, + }; + + const iters: u64 = 100_000; + var bytes_out: u64 = 0; + const start = Io.Timestamp.now(io, .awake); + var n: u64 = 0; + while (n < iters) : (n += 1) { + var buf: std.Io.Writer.Allocating = .init(alloc); + defer buf.deinit(); + if (v2) { + try wire.encodeCommitV2(alloc, &buf.writer, op, n, payload); + } else { + try wire.encodeCommit(alloc, &buf.writer, op); + } + bytes_out += buf.written().len; + } + const elapsed = since(io, start); + + return .{ + .name = if (v2) "wire_encode_v2" else "wire_encode_v1", + .subsystem = "serve", + .iters = iters, + .total_ns = elapsed, + .unit_count = iters, + .unit = "frames/s", + }; +} + +/// The cold reader's batch loop: what a subscriber behind the hot tail pays. +fn benchColdReplay(alloc: Allocator, io: Io) !Result { + var tmp_buf: [Io.Dir.max_path_bytes]u8 = undefined; + const dir = try std.fmt.bufPrint(&tmp_buf, "/tmp/stream-bench-cold", .{}); + Io.Dir.cwd().deleteTree(io, dir) catch {}; + defer Io.Dir.cwd().deleteTree(io, dir) catch {}; + + var archive = try archive_mod.Archive.init(alloc, io, dir); + defer archive.deinit(); + const rows: u64 = 20_000; + var i: u64 = 0; + while (i < rows) : (i += 1) _ = try archive.append(sampleEvent(@intCast(i)), 1_700_000_000_000_000); + try archive.rotate(); + + var reader = cold.ColdReader.init(alloc, io, 64 << 20); + defer reader.deinit(); + + const Sink = struct { + delivered: u64 = 0, + fn wants(_: *@This(), _: wire.Kind, _: []const u8, _: []const u8, _: tail_mod.SkipV1) bool { + return true; + } + fn write(_: *@This(), _: cold.CachedFrame) anyerror!cold.WriteResult { + return .sent; + } + fn observe(self: *@This(), _: u64, result: cold.ReplayResult) anyerror!void { + if (result == .sent) self.delivered += 1; + } + }; + var sink: Sink = .{}; + + const start = Io.Timestamp.now(io, .awake); + var from: u64 = 1; + while (from <= rows) { + const batch = try reader.readSeqBatch( + &archive, + &sink, + Sink.wants, + &sink, + Sink.write, + Sink.observe, + false, + false, + from, + rows + 1, + 1024, + ); + if (batch.next_seq <= from) break; + from = batch.next_seq; + } + const elapsed = since(io, start); + archive.close(); + + return .{ + .name = "cold_replay_batch", + .subsystem = "serve", + .iters = sink.delivered, + .total_ns = elapsed, + .unit_count = sink.delivered, + .unit = "rows/s", + }; +} + +fn humanize(v: f64, buf: []u8) []const u8 { + if (v >= 1_000_000_000) return std.fmt.bufPrint(buf, "{d:.2}G", .{v / 1e9}) catch "?"; + if (v >= 1_000_000) return std.fmt.bufPrint(buf, "{d:.2}M", .{v / 1e6}) catch "?"; + if (v >= 1_000) return std.fmt.bufPrint(buf, "{d:.2}K", .{v / 1e3}) catch "?"; + return std.fmt.bufPrint(buf, "{d:.0}", .{v}) catch "?"; +} + +pub fn main() !void { + // Production runs on c_allocator (src/main.zig). Benching on + // DebugAllocator would inflate every allocation-heavy path -- wire encode + // most of all -- and measure the harness rather than the system. + const alloc = std.heap.c_allocator; + defer results.deinit(alloc); + + var threaded: Io.Threaded = .init(alloc, .{}); + defer threaded.deinit(); + const io = threaded.io(); + + try record(alloc, try benchArchiveAppend(alloc, io, false)); + try record(alloc, try benchArchiveAppend(alloc, io, true)); + const seal = try benchSealAndParse(alloc, io); + for (seal) |r| try record(alloc, r); + try record(alloc, try benchSealedReadBlock(alloc, io)); + try record(alloc, try benchBloom(alloc, io)); + try record(alloc, try benchWire(alloc, io, false)); + try record(alloc, try benchWire(alloc, io, true)); + try record(alloc, try benchColdReplay(alloc, io)); + + var out_buf: [4096]u8 = undefined; + var stdout = Io.File.stdout().writer(io, &out_buf); + const w = &stdout.interface; + try w.print("{s:<26} {s:<16} {s:>12} {s:>14}\n", .{ "bench", "subsystem", "ns/op", "throughput" }); + for (results.items) |r| { + var tbuf: [32]u8 = undefined; + var nbuf: [32]u8 = undefined; + try w.print("{s:<26} {s:<16} {s:>12} {s:>10} {s}\n", .{ + r.name, + r.subsystem, + humanize(r.nsPerOp(), &nbuf), + humanize(r.throughput(), &tbuf), + r.unit, + }); + } + try w.flush(); +} -- 2.51.2