jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970"""Offline production-WebSocket receipt for raw subscriber read batching."""
import base64import hashlibimport osimport socketimport timeimport urllib.request
HOST = "127.0.0.1"PORT = int(os.environ["STREAM_SUBSCRIBE_BATCH_PORT"])EXPECTED_COLD_BATCHES = (5000 + 17 - 1) // 17
def metric(name): body = urllib.request.urlopen(f"http://{HOST}:{PORT}/metrics", timeout=2).read().decode() prefix = name + " " line = next(line for line in body.splitlines() if line.startswith(prefix)) return int(float(line.split()[1]))
connection = socket.create_connection((HOST, PORT), timeout=2)connection.settimeout(0.2)key = base64.b64encode(os.urandom(16)).decode()path = "/xrpc/network.bsky.jetstream.subscribeEvents?cursor=1&dids=did%3Aplc%3Aabsentbatchfixture"connection.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: response += connection.recv(4096)head, buffered = response.split(b"\r\n\r\n", 1)assert head.startswith(b"HTTP/1.1 101"), headaccept = base64.b64encode( hashlib.sha1((key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11").encode()).digest())assert b"sec-websocket-accept: " + accept.lower() in head.lower()
deadline = time.monotonic() + 5while metric("jetstream_subscribe_cold_reads_total") < EXPECTED_COLD_BATCHES: assert time.monotonic() < deadline, "cold replay did not finish" time.sleep(0.01)
first = metric("jetstream_subscribe_cold_reads_total")time.sleep(0.2)second = metric("jetstream_subscribe_cold_reads_total")assert first == EXPECTED_COLD_BATCHES, (first, EXPECTED_COLD_BATCHES)assert second == first, "cold reader polled after reaching the empty hot handoff"assert metric("jetstream_subscribe_hot_reads_total") == 0assert buffered == b"", "fully filtered replay unexpectedly emitted a frame"try: unexpected = connection.recv(1)except TimeoutError: unexpected = b""assert unexpected == b"", "fully filtered replay unexpectedly emitted a frame"connection.close()
print( "subscribe read batch: PASS " f"(5000 raw filtered rows / 17 = {EXPECTED_COLD_BATCHES} cold pulls; stable hot handoff)")