diff --git a/HANDOFF-archive-backfill.md b/HANDOFF-archive-backfill.md deleted file mode 100644 index 0181f4b..0000000 --- a/HANDOFF-archive-backfill.md +++ /dev/null @@ -1,110 +0,0 @@ -# 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.