# /// script # requires-python = ">=3.12" # dependencies = ["zstandard", "requests"] # /// """backfill one collection from a stream instance's archive into sqlite. the consumer-side recipe for network.bsky.jetstream.planSnapshot: plan (collection filter) -> fetch planned blocks/segments -> decompress (plain zstd frames, no dictionary) -> columnar decode -> keep matching rows -> insert. prints a timing/throughput report at the end. block fetches run concurrently (--workers, default 12): sparse collections touch one block per ~4096 events, so a serial fetch is pure round-trip latency. decode/insert stays on the main thread; sqlite needs no locking. usage: uv run backfill-collection-sqlite.py \ --host http://127.0.0.1:8080 \ --collection place.stream.chat.message \ --db streamplace-chat.db """ from __future__ import annotations import argparse import json import sqlite3 import struct import time from concurrent.futures import ThreadPoolExecutor import requests import zstandard def decode_block(raw: bytes, wanted: str) -> list[tuple]: n = struct.unpack_from(" None: ap = argparse.ArgumentParser() ap.add_argument("--host", default="http://127.0.0.1:8080") ap.add_argument("--collection", required=True) ap.add_argument("--db", required=True) ap.add_argument("--workers", type=int, default=12) args = ap.parse_args() db = sqlite3.connect(args.db) db.execute("pragma journal_mode=wal") db.execute( "create table if not exists records(" "seq integer primary key, witnessed_at integer, did text, rkey text, record blob)" ) s = requests.Session() dctx = zstandard.ZstdDecompressor() t0 = time.monotonic() plan = s.post( f"{args.host}/xrpc/network.bsky.jetstream.planSnapshot", json={"collections": [args.collection]}, timeout=60, ).json() segs = plan["segments"] t_plan = time.monotonic() - t0 def fetch_block(job): name, bi = job resp = s.get( f"{args.host}/xrpc/network.bsky.jetstream.getBlock", params={"segment": name, "blockIndex": bi}, timeout=120, ) resp.raise_for_status() return resp.content pool = ThreadPoolExecutor(max_workers=args.workers) rows = fetched_bytes = blocks_read = 0 for i, seg in enumerate(segs): if seg["mode"] == "blocks": jobs = [ (seg["name"], bi) for r in seg["blocks"] for bi in range(r["first"], r["last"] + 1) ] for content in pool.map(fetch_block, jobs): fetched_bytes += len(content) blocks_read += 1 batch = decode_block(dctx.decompress(content), args.collection) db.executemany("insert or replace into records values(?,?,?,?,?)", batch) rows += len(batch) else: # whole segment resp = s.get( f"{args.host}/xrpc/network.bsky.jetstream.getSegment", params={"name": seg["name"]}, timeout=600, ) resp.raise_for_status() body = resp.content fetched_bytes += len(body) block_index_offset = struct.unpack_from("