jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124"""Real WebSocket contract for Jetstream V2 cursor-lookback policy."""
import base64import hashlibimport jsonimport osimport socketimport 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")