jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224#!/usr/bin/env python3"""Offline production-binary contract for upstream HTTP/getBlock metrics."""
from __future__ import annotations
import base64import hashlibimport jsonimport osimport pathlibimport reimport socketimport structimport timeimport urllib.errorimport urllib.parseimport urllib.request
BASE = os.environ.get("STREAM_BASE_URL", "http://127.0.0.1:6021").rstrip("/")DATA_DIR = pathlib.Path(os.environ["STREAM_DATA_DIR"])XRPC = BASE + "/xrpc/network.bsky.jetstream."SAMPLE_RE = re.compile(r"^([a-zA-Z_:][a-zA-Z0-9_:]*)(?:\{([^}]*)\})?\s+([^\s]+)$")LABEL_RE = re.compile(r'(\w+)="([^"]*)"')
def request(path: str, *, method: str = "GET", headers: dict[str, str] | None = None): req = urllib.request.Request(BASE + path, method=method, headers=headers or {}) try: with urllib.request.urlopen(req, timeout=5) as response: return response.status, response.read(), {key.lower(): value for key, value in response.headers.items()} except urllib.error.HTTPError as error: return error.code, error.read(), {key.lower(): value for key, value in error.headers.items()}
def scrape() -> dict[tuple[str, tuple[tuple[str, str], ...]], float]: status, raw, _ = request("/metrics") assert status == 200 samples: dict[tuple[str, tuple[tuple[str, str], ...]], float] = {} for line in raw.decode().splitlines(): if line.startswith("#"): continue match = SAMPLE_RE.match(line) if not match: continue name, labels_raw, value = match.groups() labels = tuple(sorted(LABEL_RE.findall(labels_raw or ""))) samples[(name, labels)] = float(value) return samples
def value(samples, name: str, **labels: str) -> float: return samples.get((name, tuple(sorted(labels.items()))), 0.0)
def delta(before, after, name: str, **labels: str) -> float: return value(after, name, **labels) - value(before, name, **labels)
def websocket_round_trip(path: str) -> None: parsed = urllib.parse.urlsplit(BASE) host = parsed.hostname or "127.0.0.1" port = parsed.port or 80 key = base64.b64encode(os.urandom(16)).decode() expected_accept = base64.b64encode( hashlib.sha1((key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11").encode()).digest() ).decode() conn = socket.create_connection((host, port), timeout=5) try: conn.sendall( ( f"GET {path} HTTP/1.1\r\n" f"Host: {host}:{port}\r\n" "Upgrade: websocket\r\n" "Connection: Upgrade\r\n" f"Sec-WebSocket-Key: {key}\r\n" "Sec-WebSocket-Version: 13\r\n\r\n" ).encode() ) response = b"" while b"\r\n\r\n" not in response: chunk = conn.recv(4096) assert chunk, "websocket handshake closed before headers" response += chunk headers = response.split(b"\r\n\r\n", 1)[0].decode().split("\r\n") assert headers[0].startswith("HTTP/1.1 101 "), headers[0] parsed_headers = { key.strip().lower(): value.strip() for key, value in (line.split(":", 1) for line in headers[1:]) } assert parsed_headers.get("sec-websocket-accept") == expected_accept
# A real masked client close frame; the server must finish the inline # subscription handler before its HTTP middleware observation appears. payload = struct.pack("!H", 1000) mask = os.urandom(4) masked = bytes(byte ^ mask[i % 4] for i, byte in enumerate(payload)) conn.sendall(bytes((0x88, 0x80 | len(payload))) + mask + masked) conn.settimeout(2) try: while conn.recv(4096): pass except (TimeoutError, socket.timeout): pass finally: conn.close()
def main() -> None: # The archive is real and sealed. Startup may still be completing its # offline live-source attempt, so wait for the XRPC readiness gate. listing = None for _ in range(100): status, body, _ = request("/xrpc/network.bsky.jetstream.listSegments") if status == 200: listing = json.loads(body) break time.sleep(0.05) assert listing is not None, "archive XRPC never became ready" segment = listing["segments"][0] name = urllib.parse.quote(segment["name"])
status, sealed, headers = request( f"/xrpc/network.bsky.jetstream.getSegment?name={name}" ) assert status == 200 and sealed assert headers["cache-control"] == "public, max-age=2"
before = scrape()
block_path = f"/xrpc/network.bsky.jetstream.getBlock?segment={name}&blockIndex=0" status, frame, headers = request(block_path) assert status == 200 and frame assert headers["cache-control"] == "public, max-age=2" etag = headers["etag"] status, _, headers = request(block_path, headers={"If-None-Match": etag}) assert status == 304 and headers["cache-control"] == "public, max-age=2" status, _, headers = request(block_path, headers={"If-Match": '"other"'}) assert status == 412 and headers["cache-control"] == "public, max-age=2" range_status, partial, _ = request(block_path, headers={"Range": "bytes=0-0"}) assert range_status == 206 and partial == frame[:1] status, _, headers = request(block_path, headers={"Range": "bytes=999999999-"}) assert status == 416 and "cache-control" not in headers assert request("/xrpc/network.bsky.jetstream.getBlock")[0] == 400 assert request(f"/xrpc/network.bsky.jetstream.getBlock?segment={name}&blockIndex=999999")[0] == 404 assert request("/xrpc/network.bsky.jetstream.notImplemented")[0] == 501 assert request(block_path, method="POST")[0] == 405 assert request("/")[0] == 200 assert request("/status", method="HEAD")[0] == 200 assert request("/healthz")[0] == 200 assert request("/not-found")[0] == 404 websocket_round_trip("/subscribe")
after = None for _ in range(100): candidate = scrape() if delta(before, candidate, "jetstream_http_request_duration_seconds_count", handler="subscribe", method="GET", code="200") == 1: after = candidate break time.sleep(0.05) assert after is not None, "subscription HTTP lifetime was not observed"
expected_http = { ("xrpc/", "GET", "200"): 1, ("xrpc/", "GET", "304"): 1, ("xrpc/", "GET", "206"): 1, ("xrpc/", "GET", "416"): 1, ("xrpc/", "GET", "412"): 1, ("xrpc/", "GET", "400"): 1, ("xrpc/", "GET", "404"): 1, ("xrpc/", "GET", "501"): 1, ("xrpc/", "POST", "405"): 1, ("root", "GET", "200"): 1, ("status", "HEAD", "200"): 1, ("subscribe", "GET", "200"): 1, } for (handler, method, code), expected in expected_http.items(): got = delta( before, after, "jetstream_http_request_duration_seconds_count", handler=handler, method=method, code=code, ) assert got == expected, ((handler, method, code), got, expected)
total_http_delta = sum( value(after, name, **dict(labels)) - value(before, name, **dict(labels)) for name, labels in set(before) | set(after) if name == "jetstream_http_request_duration_seconds_count" ) assert total_http_delta == sum(expected_http.values()), total_http_delta
expected_results = {"ok": 3, "bad_request": 1, "not_found": 1, "error": 2} for result, expected in expected_results.items(): assert delta(before, after, "jetstream_getblock_requests_total", result=result) == expected assert delta(before, after, "jetstream_getblock_served_bytes_total") == len(frame) assert delta(before, after, "jetstream_getblock_duration_seconds_count") == 7
# Match upstream's manifest/file race semantics through the real server: # a segment which remains known to the resident manifest but disappears # from disk is an internal storage failure, not SegmentNotFound. (DATA_DIR / "segments" / segment["name"]).unlink() segment_path = f"/xrpc/network.bsky.jetstream.getSegment?name={name}" segment_failure = request(segment_path) block_failure = request(block_path) assert segment_failure[0] == 500 assert block_failure[0] == 500 assert json.loads(segment_failure[1])["error"] == "InternalServerError" assert json.loads(block_failure[1])["error"] == "InternalServerError" failed = scrape() assert delta(after, failed, "jetstream_http_request_duration_seconds_count", handler="xrpc/", method="GET", code="500") == 2 assert delta(after, failed, "jetstream_getblock_requests_total", result="error") == 1
print( "HTTP/getBlock metrics PASS: " f"http={int(total_http_delta)} getBlock=8 full_bytes={len(frame)} websocket=1 storage_race=500" )
if __name__ == "__main__": main()