jetstream v2 in zig stream.waow.tech
stream docs upstream-bootstrap-spec.md
6.6 kB
Markdown

upstream bootstrap + lifecycle spec (condensed from agent study 2026-07-09) #

Current review pin: bluesky-social/jetstream f29815c391fc2644f8a3dd36b899fb3697dd1ea6, Atmos v0.2.14 (the Go atproto library upstream runs). Refresh this comparison against upstream before provisioning any full-network experiment.

This describes upstream's behaviour, not Stream's. A rule here is what Stream is measured against; see the audit note at the end, and bootstrap-semantic-parity.md for what Stream actually proves.

This file and upstream-compaction-spec.md condense a 2026-07-09 agent study that is not committed in this repository. For anything not covered here, consult upstream source at the pin above.

pebble below is upstream's embedded key-value store (Go). Stream uses RocksDB in the same role, so "pebble phase" means the durable phase record, whatever the engine.

keystones for the zig port:

lifecycle #

phases: bootstrap → merging → steady_state, durable (upstream: pebble phase). two commit points: WritePhase(merging) AFTER backfill drains BEFORE teardown; WritePhase(steady_state) after merge fully completes. Run() dispatches by phase with fallthrough; crash matrix per phase (bootstrap: re-run, completed repos skip; merging: seal-guard tolerates unsealed source tail; mid-merge: per-source cursor, at-least-once; post-cleanup-pre-phase-write: guard detects missing backfill/ dir). fsync data before pebble commit, always; fsync dir before deleting durable state.

bootstrap #

  • listRepos pages of 1000 → batches of 100k, shuffled, dispatched to 100 workers. The shuffle is within a batch, so batch N still draws from listRepos position N×100k. That distinction matters when reading throughput: repos/s stays roughly flat inside one batch and shifts between batches, because it is the batch's place in listRepos order — not the shuffle — that decides the repository sizes in it.
  • cursors: relay/list_repos_cursor (checkpointed only after durability drain; drains to "" on completion) + bootstrap/last_listrepos_cursor (for merge discovery), saved atomically.
  • getRepo VIA THE RELAY (302 → PDS), NO signature verification on this path, xrpc retries off. The engine owns two independent budgets: 3 ordinary transient retries (4 attempts total; 1s base, 30s cap) and 20 additional 429 attempts (Retry-After honored up to 330s). Each download attempt has a 5m deadline; hitting it fails that repo for this run without an in-loop retry. RepoNotFound completes terminally with no rows; RepoDeactivated, RepoSuspended, and RepoTakendown become terminal unavailable rows.
  • CAR → rows: walk MST, one KindCreate per record, Rev = head commit.Rev for EVERY row (whole-repo watermark — merge filter depends on it), one witnessed_at per repo, batched 1024/append. gates: bad commit rev → repo fails (retryable); bad NSID/rkey path → drop record only; field-too-long → drop w/ own counter; missing block → repo fails.
  • per-DID status rows: not_started|pending|complete|failed|unavailable + Backfill.Rev watermark. Selected backfill resolves the DID document and stores declared Handle, PDS endpoint, and normalized Host; handle index changes commit atomically with the repo row. Completion sets both immutable Backfill.Rev and top-level latest Rev. Completion is committed after segment data. interrupted-bootstrap rows (pre-existing not_started, not discovered this run) defer to pending — merge repairs after live drain (re-downloading now would put low seqs under stale tombstones).
  • --max-backfill-repos: debug truncation of listRepos, no durable cursor advance, auto-skip merge discovery.
  • --backfill-repos: debug explicit DID set, mutually exclusive with max-repos; process the caller's order serially, bypass listRepos, delete stale merge discovery state, resolve Handle/PDS metadata even for already-complete rows, maintain the local handle index, and auto-skip merge discovery. Ordinary whole-network listRepos backfill does not perform per-DID identity resolution.

live capture during bootstrap #

second live consumer → backfill/live_segments/ with THROWAWAY seq space (live_segments/seq/next; merge discards + re-assigns) but SHARED relay/cursor. starts concurrently with backfill; no tombstone machinery; not user-visible.

merge #

open dst on segments/ (real seq space); drain sources in index order with merge/next_source_idx cursor; keep/drop per row: drop iff DID complete AND row.Rev <= Backfill.Rev (lexicographic TID); re-stamp witnessed_at = now; the Backfill.Rev lookup is cached once per DID for the complete merge run, including missing rows; per-source atomic commit (cursor + repo Rev refresh). then: pending-repo retry pass (HandleRepoResync: KindSync tombstone + KindCreateResync rows, pre-validated), dst seal, merge-tail delete/update compaction, discovery from bootstrap/last_listrepos_cursor (unknown DIDs → failed rows for steady retry), cleanup (RemoveAll backfill/, fsync dir, delete cursors).

cutover + serving #

steady state: ONE live consumer owns the writer (never two on one active segment). /subscribe + /subscribe-v2 + XRPC return 503 "bootstrap in progress" until steady_state. failed-repo retry loop: 4h interval, 16 workers, 4/host, 7d max backoff, 429 parks the host.

Stream audit note #

This file condenses the pinned upstream behavior; it does not describe Stream's implementation status.

A 2026-07-24 source audit replaced Stream's earlier dispatch-result list, which waited for the full dispatch batch and could replay completed repositories. Completion tracking records a per-repository final-sequence watermark and stages covered completions from the ordinary archive durability hook. It also preserves upstream's empty-CAR, RepoNotFound, duplicate-DID, captured-timestamp, and forced-drain rules. Focused tests hold one sibling incomplete, cross the writer boundary, and fully reopen the store without replaying the completed sibling.

Phase, relay-cursor, and compaction-watermark reads fail closed on Store errors and malformed or unknown encodings, with injected read-fault tests. The merge source cursor uses the same strict versioned encoding, and the restart-after-cleanup guard treats only FileNotFound as proof of absence. Current Stream status and the evidence needed to close each remaining gap (post-bootstrap discovery, retry candidate scan bounds, retry-supervisor health) are maintained in bootstrap-semantic-parity.md and semantic-parity.md.