jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216# /// script# dependencies = ["websockets", "zstandard"]# ///"""e2e checks for the /subscribe endpoint.
expects the simulator on :7777 and stream on :6008 (see `just e2e`)."""import asyncio, json, websockets
URL = "ws://localhost:6008/subscribe"
def replay_url(base=URL, *params): return base + "?cursor=1" + "".join(f"&{param}" for param in params)
async def recv_json(ws, timeout=10): return json.loads(await asyncio.wait_for(ws.recv(), timeout))
async def recv_commit(ws, timeout=10): # identity/account events bypass wantedCollections by contract while True: e = await recv_json(ws, timeout) if e["kind"] == "commit": return e
async def read_n(url, n): out = [] async with websockets.connect(url) as ws: for _ in range(n): out.append(json.loads(await ws.recv())) return out
async def basics(): evs = await read_n(replay_url(), 5) assert all("did" in e and "time_us" in e and "kind" in e for e in evs)
evs = await read_n(replay_url(URL, "wantedCollections=app.bsky.feed.like"), 5) cols = {e["commit"]["collection"] for e in evs if e["kind"] == "commit"} assert cols <= {"app.bsky.feed.like"}, cols
evs = await read_n(replay_url(URL, "wantedCollections=app.bsky.graph.*"), 5) cols = {e["commit"]["collection"] for e in evs if e["kind"] == "commit"} assert all(c.startswith("app.bsky.graph.") for c in cols), cols
a, b = await asyncio.gather( read_n(replay_url(URL, "wantedCollections=app.bsky.feed.post"), 3), read_n(replay_url(), 3), ) assert len(a) == 3 and len(b) == 3
try: await read_n(replay_url(URL, "wantedDids=bob"), 1) raise AssertionError("bad did accepted") except websockets.InvalidStatus: pass print("basics: PASS")
async def check_v2(): # proposal-0015 message frames: one self-describing JSON object per frame seqs = [] async with websockets.connect(replay_url("ws://localhost:6008/xrpc/network.bsky.jetstream.subscribeEvents"), subprotocols=["xrpc.v1.json"]) as ws: assert ws.subprotocol == "xrpc.v1.json", f"negotiation echo missing: {ws.subprotocol}" for _ in range(10): frame = json.loads(await asyncio.wait_for(ws.recv(), 10)) assert frame["$type"] == "message", frame p = frame["payload"] assert p["$type"].startswith("network.bsky.jetstream.subscribeEvents#"), p["$type"] assert p["seq"] > 0 and p["time"].endswith("Z") and "." in p["time"] seqs.append(p["seq"]) if p["$type"].endswith("#commit") and p["operation"] != "delete": assert "record" in p and "record_cbor" not in p assert seqs == sorted(seqs), "v2 seqs must be monotonic"
# inclusive seq-cursor resume: reconnect from an observed seq replays it async with websockets.connect(f"ws://localhost:6008/xrpc/network.bsky.jetstream.subscribeEvents?cursor={seqs[0]}") as ws: frame = json.loads(await asyncio.wait_for(ws.recv(), 10)) assert frame["payload"]["seq"] >= seqs[0]
# the legacy v1 parameter names are rejected with the replacement named try: async with websockets.connect(f"ws://localhost:6008/xrpc/network.bsky.jetstream.subscribeEvents?wantedDids=did%3Aplc%3Ax") as ws: raise AssertionError("wantedDids must be rejected on subscribeEvents") except websockets.InvalidStatus as e: assert e.response.status_code == 400 print(f"subscribeEvents: PASS (seqs {seqs[0]}..{seqs[-1]})")
async def check_zstd(): import zstandard, urllib.request dict_bytes = open("src/assets/zstd_dictionary", "rb").read() dctx = zstandard.ZstdDecompressor(dict_data=zstandard.ZstdCompressionDict(dict_bytes)) # Jetstream's application-level zstd mode is mutually exclusive with # WebSocket permessage-deflate. Newer websockets clients offer the latter # by default, so suppress it explicitly for this protocol check. async with websockets.connect(replay_url(URL, "compress=true"), compression=None) as ws: for _ in range(5): frame = await asyncio.wait_for(ws.recv(), 10) assert isinstance(frame, bytes), "zstd mode must send binary frames" e = json.loads(dctx.decompress(frame, max_output_size=1 << 20)) assert "did" in e and "kind" in e print("zstd compress: PASS")
def check_metrics(): import urllib.request body = urllib.request.urlopen("http://localhost:6008/metrics").read().decode() assert 'stream_build_info{git_sha="' in body def val(name): for line in body.splitlines(): if line.startswith(name + " ") or line.startswith(name + "{"): pass for line in body.splitlines(): if line.split(" ")[0] == name: return float(line.split(" ")[1]) raise AssertionError(f"metric {name} missing") assert val("stream_events_total") > 0 assert val("stream_frames_total") > 0 # a verifier that silently passes everything unverified is # indistinguishable from a working one without this valid = [l for l in body.splitlines() if l.startswith('stream_verify_total{result="valid"}')] assert valid and float(valid[0].split(" ")[1]) > 0, "verification never verified anything" invalid = [l for l in body.splitlines() if l.startswith('stream_verify_total{result="invalid"}')] assert float(invalid[0].split(" ")[1]) == 0, "honest traffic produced invalid signatures" def labeled(name, label): prefix = f'{name}{{{label}}} ' line = next((line for line in body.splitlines() if line.startswith(prefix)), None) assert line is not None, f"metric {name}{{{label}}} missing" return float(line.split(" ")[1]) assert labeled("jetstream_subscribe_events_sent_total", 'compression="none"') > 0 assert labeled("jetstream_subscribe_events_sent_total", 'compression="deflate"') > 0 assert labeled("jetstream_subscribe_events_sent_total", 'compression="zstd"') > 0 assert labeled("jetstream_subscribe_bytes_encoded_total", 'compression="zstd"') > 0 assert labeled("jetstream_subscribe_bytes_sent_total", 'compression="zstd"') > 0 assert val("jetstream_subscribe_events_filtered_total") > 0 assert val("jetstream_subscribe_events_oversize_total") > 0 health = urllib.request.urlopen("http://localhost:6008/healthz").read().decode() assert health.strip() == "ok" print("metrics: PASS")
async def main(): await basics() # options_update swaps the filter atomically async with websockets.connect(replay_url(URL, "wantedCollections=app.bsky.feed.like")) as ws: e = await recv_commit(ws) assert e["commit"]["collection"] == "app.bsky.feed.like" await ws.send(json.dumps({"type": "options_update", "payload": {"wantedCollections": ["app.bsky.feed.post"]}})) # drain until only posts arrive (a few likes may already be in flight) seen_post = False for _ in range(50): e = await recv_commit(ws) if e["commit"]["collection"] == "app.bsky.feed.post": seen_post = True break assert seen_post print("options_update: PASS")
# unknown type ignored, connection stays up async with websockets.connect(replay_url()) as ws: await ws.send(json.dumps({"type": "wat", "payload": {}})) await recv_json(ws) print("unknown type ignored: PASS")
# malformed envelope -> close 1007 async with websockets.connect(replay_url()) as ws: await ws.send("not json") try: while True: await asyncio.wait_for(ws.recv(), 5) except websockets.ConnectionClosed as e: assert e.rcvd.code == 1007, e.rcvd print("bad envelope 1007: PASS")
# invalid payload -> close 1008 async with websockets.connect(replay_url()) as ws: await ws.send(json.dumps({"type": "options_update", "payload": {"wantedDids": ["bob"]}})) try: while True: await asyncio.wait_for(ws.recv(), 5) except websockets.ConnectionClosed as e: assert e.rcvd.code == 1008, e.rcvd print("invalid payload 1008: PASS")
# requireHello: silence until hello, then events flow async with websockets.connect(replay_url(URL, "requireHello=true")) as ws: try: await asyncio.wait_for(ws.recv(), 2) raise AssertionError("got event before hello") except asyncio.TimeoutError: pass await ws.send(json.dumps({"type": "options_update", "payload": {"wantedCollections": ["app.bsky.feed.like"]}})) e = await recv_commit(ws) assert e["commit"]["collection"] == "app.bsky.feed.like" print("requireHello: PASS")
# requireHello quirk: only exact literal "true" async with websockets.connect(replay_url(URL, "requireHello=True")) as ws: await recv_json(ws) # events flow immediately print("requireHello=True is false: PASS")
# maxMessageSizeBytes: tiny cap suppresses all events (garbage cap doesn't) async with websockets.connect(replay_url(URL, "maxMessageSizeBytes=10")) as ws: try: await asyncio.wait_for(ws.recv(), 2) raise AssertionError("oversize event delivered") except asyncio.TimeoutError: print("maxMessageSizeBytes cap: PASS") async with websockets.connect(replay_url(URL, "maxMessageSizeBytes=garbage")) as ws: await recv_json(ws) print("maxMessageSizeBytes garbage->0: PASS")
await check_v2() await check_zstd() check_metrics() print("ALL PASS")
asyncio.run(main())