From ab467d1f595b7be28db099d9f7b09ec02902ab64 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Fri, 7 Aug 2026 16:18:01 -0500 Subject: [PATCH] =?UTF-8?q?handoff:=20zat.ArchiveBackfill=20spec=20?= =?UTF-8?q?=E2=80=94=20consume=20the=20jetstream=20archive=20backfill=20AP?= =?UTF-8?q?I=20(pollz=20first)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- HANDOFF-archive-backfill.md | 110 ++++++++++++++++++++++++++++++++++++ 1 file changed, 110 insertions(+) create mode 100644 HANDOFF-archive-backfill.md diff --git a/HANDOFF-archive-backfill.md b/HANDOFF-archive-backfill.md new file mode 100644 index 0000000..0181f4b --- /dev/null +++ b/HANDOFF-archive-backfill.md @@ -0,0 +1,110 @@ +# handoff: zat.ArchiveBackfill — consume a jetstream archive's backfill API + +Written 2026-08-07 by the stream operator session. Context: stream.waow.tech +serves the full network (~23B events, sealed segments back to seq 1) behind +upstream jetstream V2's archive API. Consumers want full-history rebuilds — +pollz is the first committed adopter — and the decode machinery belongs in +zat so every zig consumer gets it once, rather than each app vendoring zstd +and a jss parser. (The Bluesky team validated the consumption pattern on +2026-08-07: whole-network backfill via this API into ClickHouse in ~17h.) + +## the goal, from the consumer's side + +```zig +var client = zat.JetstreamClient.init(io, allocator, .{ + .hosts = &.{ "stream.waow.tech", ... }, + .wanted_collections = &.{ "tech.waow.pollz.poll", "tech.waow.pollz.vote" }, +}); +// NEW: replay the archive through the same handler the live tail uses, +// then hand off to subscribe at the seq the plan covered. +try zat.ArchiveBackfill.run(io, allocator, .{ + .host = "https://stream.waow.tech", + .collections = &.{ "tech.waow.pollz.poll", "tech.waow.pollz.vote" }, +}, &handler); // same onEvent as JetstreamClient +client.subscribe(&handler); +``` + +The contract that makes it composable: ArchiveBackfill delivers +`zat.JetstreamEvent`s (same shape as the live client — commit events with +parsed record) so an app's existing handler works unmodified. pollz's entire +diff should be ~15 lines. + +## the API being consumed (all verified live on stream.waow.tech) + +1. `POST /xrpc/network.bsky.jetstream.planBackfill` body + `{"collections":[...]}` (also accepts `dids`) → `{plannedThroughSeq, + sealedTipSeq, segments:[{name, index, checksum, minSeq, maxSeq, mode, + blocks?:[{first,last}]}]}`. `mode` is `"blocks"` (fetch listed ranges) or + whole-segment. +2. `GET .../getBlock?segment=&blockIndex=` → one plain zstd frame + (magic `28 b5 2f fd`, content checksum on, **no dictionary** — the + getZstdDictionary endpoint is for subscribe-v2 wire compression only). +3. `GET .../getSegment?name=` → whole sealed jss file (~277MB; + supports Range/206 + ETag = segment checksum). Parse header offsets, + walk the block index, decompress each frame. +4. Resume recipe: after the archive pass, connect the live client with a + cursor at/before `plannedThroughSeq` (subscribe-v2 seq cursor is exact; + v1 time_us also works). Overlap is fine for idempotent consumers; + document at-least-once. + +## jss v1 format essentials + +Authoritative spec: `tangled.org/zat.dev/stream` → `docs/jss-format-v1.md`. +Short version: + +- header 256B: magic "jss0"; `block_count u32 @14`; `block_index_offset u64 + @90`; block index entries are 52B: `offset u64, compressed_size u32, + uncompressed_size u32, event_count u32, min_seq u64, max_seq u64, + min_witnessed_at i64, max_witnessed_at i64`. Block frame at `offset` is + `len u64` then `len` bytes of zstd. +- decompressed block is columnar: + `event_count u32; seq[]u64; witnessed_at[]i64; indexed_at[]i64; kind[]u8; + collection_len[]u8; did_len[]u16; rkey_len[]u8; rev_len[]u8; + event_len[]u32;` then blobs `collections|dids|rkeys|revs|payloads` + (prefix-sum sliced; trailing bytes are an error). payloads are raw + DAG-CBOR records. +- kind enum and the mapping to jetstream commit operations: see the spec. + Filter rows by collection client-side even in blocks mode (a planned + block can contain other collections' rows). + +A known-good reference implementation of plan→fetch→decode (python, +~150 lines) is `examples/backfill-collection-sqlite.py` in the stream repo. + +## what zat needs that it doesn't have + +- **zstd decompression** (decode-only is enough). Nothing in zat links zstd + today; stream vendors facebook/zstd v1.5.7 — same pin recommended. +- **DAG-CBOR record → the same parsed-record representation the live + jetstream client hands to onEvent** (zat already decodes DAG-CBOR on the + firehose path; reuse that). +- The plan/fetch client is plain XRPC + HTTP GETs (HttpTransport exists). + +## sizing / behavior notes from operating the producer + +- Collection-filtered plans can still be tens of GB when the collection + touches many blocks (streamplace chat: 138k blocks ≈ 40GB fetched for + ~170k rows). Stream the work: fetch → decode → deliver → drop; never + accumulate. Whole-segment mode holds one ~277MB body; 2GB RAM machines + are fine. +- Throughput observed: ~50MB/s serial from adjacent-DC; block fetches are + small — a little concurrency (4-8 in flight) goes a long way and should + be an option, not a default (be a polite client). +- Verify segment checksum only when the caller asks (Range-resumable + getSegment makes this cheap: etag == checksum). + +## also queued for zat, same session's findings (separate, smaller) + +1. **primary failback**: JetstreamClient never returns to hosts[0] after + failover; consumers stick to fallbacks until their own process restarts. +2. **idle-timeout knob** in the read loop: a host that holds the socket + open but goes silent never triggers failover today. + +## first consumer, ready to go + +pollz (`~/tangled.org/zzstoatzz.io/pollz`): live-only today (no cursor → +every restart is a permanent tally hole). With ArchiveBackfill its boot +becomes: archive pass through the existing `Handler`/`processCommit` path +(insertPoll/insertVote are idempotent upserts), then `subscribe` as now +with stream.waow.tech first in hosts. Collections `tech.waow.pollz.poll`, +`tech.waow.pollz.vote` — tiny corpus, seconds to replay, provably-complete +tallies forever after. -- 2.51.2