jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182#!/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 checkoutwhose HEAD is not the semantic-parity pin and invokes the exact Go 1.26.5toolchain from its module-cache path when the ambient `go` is older."""
import jsonimport osimport pathlibimport platformimport reimport subprocessimport 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), finalassert re.search(r"last_cursor=5000\b", final), finalprint("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") > 0assert scalar("jetstream_manifest_block_index_cache_misses_total") == 0print( "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})")