jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322# /// script# dependencies = ["zstandard"]# ///"""backfill api checks. seed + run: zig build write-sample -- --archive ./data-e2e ./zig-out/bin/stream --port=6018 --data-dir=./data-e2e & UV_OFFLINE=1 uv run tests/backfill_api.py
Set STREAM_BASE_URL and STREAM_DATA_DIR to exercise another loopbackcandidate without copying or mutating its archive."""import datetimeimport emailimport email.utilsimport json, urllib.request, zstandard, hashlib, os, pathlib, sys
SERVER = os.environ.get("STREAM_BASE_URL", "http://localhost:6018").rstrip("/")DATA_DIR = pathlib.Path(os.environ.get("STREAM_DATA_DIR", "data-e2e"))BASE = SERVER + "/xrpc/network.bsky.jetstream."
ls = json.load(urllib.request.urlopen(BASE + "listSegments"))segs = ls["segments"]assert len(segs) >= 1, "need sealed segments"s0 = segs[0]assert len(s0["checksum"]) == 16 and s0["eventCount"] > 0 and s0["maxSeq"] >= s0["minSeq"]assert [x["index"] for x in segs] == sorted(x["index"] for x in segs)
# getSegment: bytes identical to the on-disk file; ETag == checksumresp = urllib.request.urlopen(BASE + f"getSegment?name={s0['name']}")body = resp.read()assert resp.headers["ETag"] == f'"{s0["checksum"]}"', resp.headers["ETag"]disk = (DATA_DIR / "segments" / s0["name"]).read_bytes()disk_stat = (DATA_DIR / "segments" / s0["name"]).stat()last_modified = email.utils.formatdate(int(disk_stat.st_mtime), usegmt=True)assert resp.headers["Last-Modified"] == last_modified, resp.headersassert hashlib.sha256(body).digest() == hashlib.sha256(disk).digest(), "served bytes != disk bytes"assert len(body) == s0["sizeBytes"]
# getBlock: raw zstd frame that decompresses to a columnar blockresp = urllib.request.urlopen(BASE + f"getBlock?segment={s0['name']}&blockIndex=0")frame = resp.read()assert resp.headers["Last-Modified"] == last_modified, resp.headersblock = zstandard.ZstdDecompressor().decompress(frame, max_output_size=1 << 30)n = int.from_bytes(block[:4], "little")assert n > 0, "block event_count"
# errors: unknown segment -> SegmentNotFound; oob block -> BlockNotFoundimport urllib.errorfor url, want in [ (BASE + "getSegment?name=seg_zzzzzzzzzz.jss", "SegmentNotFound"), (BASE + f"getBlock?segment={s0['name']}&blockIndex=9999", "BlockNotFound"),]: try: urllib.request.urlopen(url) sys.exit(f"expected error for {url}") except urllib.error.HTTPError as e: assert json.load(e)["error"] == want
# xrpcserver's public middleware owns unknown NSIDs: 501 with the canonical# MethodNotImplemented envelope, not an ordinary InvalidRequest-shaped 400.try: urllib.request.urlopen(BASE + "doesNotExist") sys.exit("expected unknown XRPC method to return 501")except urllib.error.HTTPError as e: assert e.code == 501 assert json.load(e)["error"] == "MethodNotImplemented"
print(f"backfill api: PASS ({len(segs)} segments, seg0 {s0['eventCount']} events, block0 {n} events)")
# --- planSnapshot ---def post_plan(body): req = urllib.request.Request(BASE + "planSnapshot", data=json.dumps(body).encode(), headers={"Content-Type": "application/json"}, method="POST") return json.load(urllib.request.urlopen(req))
# Planning must use the metadata retained at process startup. Hide every JSS# data file for one request: the old implementation reopened and reread the# entire archive here, while the upstream contract performs no segment I/O.hidden = []for path in sorted((DATA_DIR / "segments").glob("seg_*.jss")): moved = path.with_suffix(".jss.hidden") path.rename(moved) hidden.append((path, moved))try: resident_plan = post_plan({}) assert resident_plan["segments"], "resident plan unexpectedly empty"finally: for path, moved in hidden: moved.rename(path)
# unfiltered: whole archive planned as whole segments (density 1.0)plan = post_plan({})assert plan["sealedTipSeq"] > 0assert plan["plannedThroughSeq"] == plan["sealedTipSeq"], "one page should cover all"assert all(p["mode"] == "segment" for p in plan["segments"])assert plan["stats"]["segmentsMatched"] == len(plan["segments"]) > 0
# collection filter matching nothing still selects DID-marker blocks# ($identity/$account/$sync sentinels ride along by contract — a folding# consumer must see account-level deletions even when collection-filtered)nothing = post_plan({"collections": ["com.example.nothing"]})assert nothing["stats"]["blocksMatched"] <= plan["stats"]["blocksMatched"]
# wildcard matches archived collectionswild = post_plan({"collections": ["app.bsky.feed.*"]})assert len(wild["segments"]) > 0
# DID filter: known DID matches; unknown DID prunes everything via bloomsknown = post_plan({"dids": ["did:plc:streamsampleaccount"]})assert len(known["segments"]) > 0, "bloom must not false-negative a real DID"unknown = post_plan({"dids": ["did:plc:neverwritten0000000000"]})assert unknown["stats"]["blocksMatched"] == 0, "blooms should prune unknown DIDs"
# seq window: afterSeq at the tip -> nothing leftdone = post_plan({"afterSeq": plan["sealedTipSeq"]})assert done["segments"] == []
print(f"planSnapshot: PASS (tip={plan['sealedTipSeq']}, {plan['stats']['segmentsMatched']} segs, unknown-did blocks={unknown['stats']['blocksMatched']})")
# --- Range on getSegment ---req = urllib.request.Request(BASE + f"getSegment?name={s0['name']}", headers={"Range": "bytes=0-255"})resp = urllib.request.urlopen(req)assert resp.status == 206 and resp.headers["Content-Range"] == f"bytes 0-255/{s0['sizeBytes']}"part = resp.read()assert part == disk[:256] and part[:4] == b"jss0"req = urllib.request.Request(BASE + f"getSegment?name={s0['name']}", headers={"Range": "bytes=-100"})assert urllib.request.urlopen(req).read() == disk[-100:]try: urllib.request.urlopen(urllib.request.Request(BASE + f"getSegment?name={s0['name']}", headers={"Range": f"bytes={s0['sizeBytes']}-"})) sys.exit("expected 416")except urllib.error.HTTPError as e: assert e.code == 416 assert e.headers["Content-Range"] == f"bytes */{len(disk)}" assert e.read() == b"invalid range: failed to overlap\n"print("range: PASS")
def multipart_ranges(url, representation, specs): raw = ", ".join(f"{start}-{end}" for start, end in specs) response = urllib.request.urlopen( urllib.request.Request(url, headers={"Range": "bytes=" + raw}) ) payload = response.read() content_type = response.headers["Content-Type"] assert response.status == 206 assert content_type.startswith("multipart/byteranges; boundary=") assert response.headers.get("Content-Range") is None assert int(response.headers["Content-Length"]) == len(payload) message = email.message_from_bytes( b"Content-Type: " + content_type.encode() + b"\r\nMIME-Version: 1.0\r\n\r\n" + payload ) parts = message.get_payload() assert len(parts) == len(specs) for part, (start, end) in zip(parts, specs): assert part.get_content_type() == "application/octet-stream" assert part["Content-Range"] == f"bytes {start}-{end}/{len(representation)}" assert part.get_payload(decode=True) == representation[start:end + 1]
# Pinned upstream delegates both archive byte endpoints to http.ServeContent.# Multiple satisfiable ranges therefore produce a real multipart/byteranges# response rather than 416; exercise the physical file and raw zstd frame.multipart_ranges( BASE + f"getSegment?name={s0['name']}", disk, [(0, 15), (len(disk) - 16, len(disk) - 1)],)multipart_ranges( BASE + f"getBlock?segment={s0['name']}&blockIndex=0", frame, [(0, 7), (len(frame) - 8, len(frame) - 1)],)
# ServeContent discards non-overlapping members when another member overlaps,# and ignores the entire Range header when duplicate ranges ask it to transfer# more bytes than the representation itself.mixed = urllib.request.urlopen( urllib.request.Request( BASE + f"getSegment?name={s0['name']}", headers={"Range": f"bytes=0-0, {len(disk) + 10}-"}, ))assert mixed.status == 206 and mixed.read() == disk[:1]duplicate = urllib.request.urlopen( urllib.request.Request( BASE + f"getSegment?name={s0['name']}", headers={"Range": f"bytes=0-{len(disk) - 1},0-0"}, ))assert duplicate.status == 200 and duplicate.read() == disk
# If-Range uses strong comparison: a wildcard is not a validator. If-None-# Match uses weak comparison and therefore accepts the weak form of the ETag.wildcard = urllib.request.urlopen( urllib.request.Request( BASE + f"getSegment?name={s0['name']}", headers={"Range": "bytes=0-0", "If-Range": "*"}, ))assert wildcard.status == 200 and wildcard.read() == disklater_match = urllib.request.urlopen( urllib.request.Request( BASE + f"getSegment?name={s0['name']}", headers={ "Range": "bytes=0-0", "If-Range": f'"other", "{s0["checksum"]}"', }, ))assert later_match.status == 200 and later_match.read() == diskfirst_match = urllib.request.urlopen( urllib.request.Request( BASE + f"getSegment?name={s0['name']}", headers={ "Range": "bytes=0-0", "If-Range": f'"{s0["checksum"]}", "other"', }, ))assert first_match.status == 206 and first_match.read() == disk[:1]try: urllib.request.urlopen( urllib.request.Request( BASE + f"getSegment?name={s0['name']}", headers={"If-None-Match": f'W/"{s0["checksum"]}"'}, ) ) sys.exit("expected weak If-None-Match to return 304")except urllib.error.HTTPError as e: assert e.code == 304 assert e.headers.get("Last-Modified") is None
def expect_error(url, code, headers): try: urllib.request.urlopen(urllib.request.Request(url, headers=headers)) sys.exit(f"expected HTTP {code} for {url}") except urllib.error.HTTPError as exc: assert exc.code == code, (exc.code, exc.read()) return exc
segment_url = BASE + f"getSegment?name={s0['name']}"block_url = BASE + f"getBlock?segment={s0['name']}&blockIndex=0"etag = f'"{s0["checksum"]}"'old_date = email.utils.formatdate(int(disk_stat.st_mtime) - 60, usegmt=True)future_date = email.utils.formatdate(int(disk_stat.st_mtime) + 60, usegmt=True)modified_dt = datetime.datetime.fromtimestamp(int(disk_stat.st_mtime), datetime.UTC)rfc850_date = modified_dt.strftime("%A, %d-%b-%y %H:%M:%S GMT")asctime_date = modified_dt.strftime("%a %b %e %H:%M:%S %Y")
# Pinned Go ServeContent evaluates representation preconditions in RFC order.# If-Match suppresses If-Unmodified-Since, while any present If-None-Match# suppresses If-Modified-Since even when its validator does not match.for url in (segment_url, block_url): assert urllib.request.urlopen(urllib.request.Request(url, headers={"If-Match": "*"})).status == 200 assert urllib.request.urlopen(urllib.request.Request(url, headers={"If-Match": etag if url == segment_url else f'"{s0["checksum"]}:0"'})).status == 200 failed_match = expect_error(url, 412, {"If-Match": '"other"'}) assert failed_match.headers["Last-Modified"] == last_modified assert failed_match.headers["ETag"] is not None assert urllib.request.urlopen(urllib.request.Request(url, headers={ "If-Match": "*", "If-Unmodified-Since": old_date, })).status == 200 expect_error(url, 412, {"If-Unmodified-Since": old_date}) assert urllib.request.urlopen(urllib.request.Request(url, headers={"If-Unmodified-Since": future_date})).status == 200 expect_error(url, 304, {"If-Modified-Since": last_modified}) assert urllib.request.urlopen(urllib.request.Request(url, headers={"If-Modified-Since": old_date})).status == 200 assert urllib.request.urlopen(urllib.request.Request(url, headers={ "If-None-Match": '"other"', "If-Modified-Since": future_date, })).status == 200
# Go's http.ParseTime accepts its two obsolete HTTP layouts and time.Parse's# fractional-second extension even though generated Last-Modified is IMF-date.for compatible_date in (rfc850_date, asctime_date, last_modified.replace(" GMT", ".5 GMT")): expect_error(segment_url, 304, {"If-Modified-Since": compatible_date})
# If-Range accepts either the first strong ETag or a date exactly equal to the# representation mtime. A stale date makes ServeContent ignore Range.dated_range = urllib.request.urlopen(urllib.request.Request( segment_url, headers={"Range": "bytes=0-0", "If-Range": last_modified},))assert dated_range.status == 206 and dated_range.read() == disk[:1]stale_range = urllib.request.urlopen(urllib.request.Request( segment_url, headers={"Range": "bytes=0-0", "If-Range": old_date},))assert stale_range.status == 200 and stale_range.read() == disk
# ServeContent distinguishes malformed syntax from a valid range that simply# does not overlap. Both are 416, but only the latter carries Content-Range.malformed = expect_error(segment_url, 416, {"Range": "items=0-1"})assert malformed.headers.get("Content-Range") is Noneassert malformed.read() == b"invalid range\n"
print("servecontent conditional/range semantics: PASS")
# Operator archive views. These render from resident manifest metadata only,# so they must agree exactly with what listSegments reports rather than being# a second, drifting source of truth.seg_text = urllib.request.urlopen(SERVER + "/status?tab=segments").read().decode()assert "stream sealed segments" in seg_text, seg_text[:200]for seg in segs: assert f"{seg['index']}\t{seg['eventCount']}" in seg_text, (seg, seg_text[:400])total_events = sum(s["eventCount"] for s in segs)assert f"{len(segs)} sealed segments, {total_events} events" in seg_text, seg_text[-200:]
col_text = urllib.request.urlopen(SERVER + "/status?tab=collections").read().decode()assert "stream archived collections" in col_text, col_text[:200]assert "app.bsky.feed.post" in col_text, col_text[:400]# DID-level marker sentinels are not collections and must not be listed.assert "$" not in col_text, col_text[:400]
req = urllib.request.Request( SERVER + "/status?tab=segments", headers={"Accept": "text/html"})seg_html = urllib.request.urlopen(req).read().decode()assert "Sealed segments" in seg_html# The new views must join the existing nav, not replace it.for label in ("Summary", "Hosts", "Accounts", "Segments", "Collections"): assert f">{label}</a>" in seg_html, label
print(f"archive status views: PASS ({len(segs)} segments, {total_events} events)")