#!/usr/bin/env python3 """Drive Stream's real archive through the pinned upstream Go client. The target must contain the deterministic 5,000-row archive produced by: zig build write-sample -- --archive /tmp/stream-xrpc-data No network resolution is permitted. The script rejects any upstream checkout whose HEAD is not the semantic-parity pin and invokes the exact Go 1.26.5 toolchain from its module-cache path when the ambient `go` is older. """ import json import os import pathlib import platform import re import subprocess import urllib.request UPSTREAM_PIN = "289b0328c2e1a0ccf8c870cb45de0b2397de19fb" UPSTREAM = pathlib.Path( os.environ.get( "STREAM_UPSTREAM_REPO", pathlib.Path.home() / "github.com/bluesky-social/jetstream", ) ) HOST = os.environ.get("STREAM_BASE_URL", "http://localhost:6018").rstrip("/") def output(*cmd, cwd=None, env=None): return subprocess.run( cmd, cwd=cwd, env=env, check=True, text=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE, ).stdout.strip() def go_binary(): explicit = os.environ.get("STREAM_GO") if explicit: return explicit ambient = output("go", "version") if "go1.26.5 " in ambient: return "go" modcache = pathlib.Path(output("go", "env", "GOMODCACHE")) goos = {"Darwin": "darwin", "Linux": "linux"}.get(platform.system()) machine = {"arm64": "arm64", "aarch64": "arm64", "x86_64": "amd64"}.get( platform.machine() ) if goos is None or machine is None: raise SystemExit(f"unsupported cached Go toolchain platform: {platform.platform()}") cached = ( modcache / f"golang.org/toolchain@v0.0.1-go1.26.5.{goos}-{machine}" / "bin/go" ) if not cached.is_file(): raise SystemExit( "Go 1.26.5 is not available offline; install or cache it before this contract" ) return str(cached) head = output("git", "rev-parse", "HEAD", cwd=UPSTREAM) if head != UPSTREAM_PIN: raise SystemExit(f"upstream HEAD {head} != required pin {UPSTREAM_PIN}") GO = go_binary() GO_ENV = os.environ.copy() GO_ENV.update({"GOTOOLCHAIN": "local", "GOPROXY": "off", "GOSUMDB": "off"}) def client(*args): proc = subprocess.run( [GO, "run", "./cmd/client", "--host", HOST, *args], cwd=UPSTREAM, env=GO_ENV, check=True, text=True, stdout=subprocess.PIPE, stderr=subprocess.PIPE, ) if proc.stderr: raise AssertionError(f"official client wrote stderr: {proc.stderr}") return proc.stdout def replay(name, expected_count, expected_first, expected_last, *filters): raw = client("--after-seq", "0", "--backfill-only", *filters, "--print") rows = [json.loads(line) for line in raw.splitlines()] cursors = [row["cursor"] for row in rows] got = (len(rows), cursors[0], cursors[-1]) want = (expected_count, expected_first, expected_last) assert got == want, f"{name}: got count/range {got}, want {want}" assert all(a < b for a, b in zip(cursors, cursors[1:])), name assert len(set(cursors)) == len(cursors), f"{name}: duplicate cursor" print(f"{name}: PASS ({len(rows)} rows, {cursors[0]}..{cursors[-1]})") # --backfill-only reaches the SEALED frontier, not the archive tip. The sample # leaves its third segment active (unsealed) and startup resumes it in place as # upstream's ingest.Open does, so these stop at seq 4000. The full cutover run # below has no such bound and must still reach all 5,000 events -- that gap is # the boundary, not a shortfall. replay("whole-segment", 4000, 1, 4000) replay( "did-block-plan", 2000, 2, 4000, "--did", "did:plc:streamsampleaccount", ) replay( "sentinel-only-block-plan", 40, 100, 4000, "--collection", "com.example.nothing", ) bounded = client( "--after-seq", "1000", "--before-seq", "1234", "--backfill-only", "--print", ) bounded_rows = [json.loads(line) for line in bounded.splitlines()] bounded_cursors = [row["cursor"] for row in bounded_rows] assert bounded_cursors == list(range(1001, 1235)), "bounded (after,before] replay" print("bounded-range: PASS (234 rows, 1001..1234)") cutover = client( "--after-seq", "0", "--duration", "1s", "--report-interval", "1s", ) final = next(line for line in reversed(cutover.splitlines()) if line.startswith("final ")) assert re.search(r"events=5,000\b", final), final assert re.search(r"last_cursor=5000\b", final), final print("backfill-live-cutover: PASS (5,000 unique rows, cursor 5000)") # The official client's archive-to-live cutover must traverse the same # all-resident manifest as upstream, not rescan block indexes from disk. segments = json.load( urllib.request.urlopen( HOST + "/xrpc/network.bsky.jetstream.listSegments", timeout=5 ) )["segments"] metrics = urllib.request.urlopen(HOST + "/metrics", timeout=5).read().decode() def scalar(name): prefix = name + " " line = next((line for line in metrics.splitlines() if line.startswith(prefix)), None) assert line is not None, f"missing metric {name}" return float(line[len(prefix) :]) assert scalar("jetstream_manifest_segments_loaded") == len(segments) assert scalar("jetstream_manifest_block_index_load_seconds_count") == len(segments) assert scalar("jetstream_manifest_block_index_cache_hits_total") > 0 assert scalar("jetstream_manifest_block_index_cache_misses_total") == 0 print( "resident-manifest: PASS " f"({len(segments)} segments, {int(scalar('jetstream_manifest_block_index_cache_hits_total'))} block-index hits)" ) print(f"OFFICIAL CLIENT ARCHIVE CONTRACT PASS ({UPSTREAM_PIN[:7]}, {GO})")