# /// 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 loopback candidate without copying or mutating its archive. """ import datetime import email import email.utils import 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 == checksum resp = 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.headers assert 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 block resp = urllib.request.urlopen(BASE + f"getBlock?segment={s0['name']}&blockIndex=0") frame = resp.read() assert resp.headers["Last-Modified"] == last_modified, resp.headers block = 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 -> BlockNotFound import urllib.error for 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"] > 0 assert 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 collections wild = post_plan({"collections": ["app.bsky.feed.*"]}) assert len(wild["segments"]) > 0 # DID filter: known DID matches; unknown DID prunes everything via blooms known = 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 left done = 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() == disk later_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() == disk first_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 None assert 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}" in seg_html, label print(f"archive status views: PASS ({len(segs)} segments, {total_events} events)")