"""Real WebSocket contract for Jetstream V2 cursor-lookback policy.""" import base64 import hashlib import json import os import socket import urllib.request HOST = "127.0.0.1" PORT = int(os.environ["STREAM_CURSOR_LOOKBACK_PORT"]) MODE = os.environ["STREAM_CURSOR_LOOKBACK_MODE"] def handshake(path): connection = socket.create_connection((HOST, PORT), timeout=2) connection.settimeout(2) key = base64.b64encode(os.urandom(16)).decode() 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) lines = head.decode().split("\r\n") headers = {} for line in lines[1:]: name, value = line.split(":", 1) headers[name.lower()] = value.strip() if lines[0].startswith("HTTP/1.1 101"): expected = base64.b64encode( hashlib.sha1((key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11").encode()).digest() ).decode() assert headers["sec-websocket-accept"] == expected return connection, lines[0], headers, buffered def read_exact(connection, buffered, length): while len(buffered) < length: buffered += connection.recv(length - len(buffered)) return buffered[:length], buffered[length:] def read_frame(connection, buffered): header, buffered = read_exact(connection, buffered, 2) opcode = header[0] & 0x0F length = header[1] & 0x7F assert header[1] & 0x80 == 0, "server frames must not be masked" if length == 126: raw, buffered = read_exact(connection, buffered, 2) length = int.from_bytes(raw, "big") elif length == 127: raw, buffered = read_exact(connection, buffered, 8) length = int.from_bytes(raw, "big") payload, buffered = read_exact(connection, buffered, length) assert opcode in (1, 2) return payload, buffered # The sample archive's events are stamped in 2023, so every segment is older # than now-36h and the floor takes upstream's clamp branch: the MinSeq of the # freshest SEALED segment. The sample's third segment is left active (unsealed) # and startup resumes it in place, as upstream's ingest.Open does, so the floor # is segment 1's min_seq (2001) -- not the active segment's first seq (4001). if MODE == "default": connection, status, headers, buffered = handshake("/xrpc/network.bsky.jetstream.subscribeEvents?cursor=1") try: assert status.startswith("HTTP/1.1 400") length = int(headers["content-length"]) body, _ = read_exact(connection, buffered, length) text = body.decode() # v2 rejections are XRPC JSON error envelopes with a structured name err = json.loads(text) assert err["error"] == "CursorTooOld", err assert "cursor 1 is below the lookback floor 2001" in err["message"], err finally: connection.close() connection, status, _, buffered = handshake("/subscribe?cursor=1") try: assert status.startswith("HTTP/1.1 101") payload, _ = read_frame(connection, buffered) event = json.loads(payload) assert event["time_us"] == 1_700_000_002_001_000 finally: connection.close() elif MODE == "wide": connection, status, _, buffered = handshake("/xrpc/network.bsky.jetstream.subscribeEvents?cursor=1") try: assert status.startswith("HTTP/1.1 101") payload, _ = read_frame(connection, buffered) frame = json.loads(payload) assert frame["$type"] == "message", frame assert frame["payload"]["seq"] == 1, frame finally: connection.close() elif MODE == "disabled": connection, status, _, buffered = handshake("/xrpc/network.bsky.jetstream.subscribeEvents?cursor=1") try: assert status.startswith("HTTP/1.1 101") connection.settimeout(0.2) try: unexpected = buffered or connection.recv(1) except TimeoutError: unexpected = b"" assert unexpected == b"", "disabled cursor unexpectedly replayed archive data" metrics = urllib.request.urlopen(f"http://{HOST}:{PORT}/metrics").read().decode() assert 'stream_subscribe_cursor_requests_total{mode="disabled"} 1' in metrics finally: connection.close() else: raise AssertionError(f"unknown mode: {MODE}") print(f"cursor lookback {MODE}: PASS")