import socket, ssl, base64, os, time, subprocess, re, json, statistics ip = subprocess.check_output(["docker","inspect","-f","{{range .NetworkSettings.Networks}}{{.IPAddress}}{{end}}","stream-experiment-stream-1"],text=True).strip() def http(path, method="GET", body=None, accept=None, timeout=10): s = socket.create_connection((ip,8080),timeout=timeout) h = f"{method} {path} HTTP/1.1\r\nHost: x\r\nConnection: close\r\n" if accept: h += f"Accept: {accept}\r\n" if body is not None: h += f"Content-Type: application/json\r\nContent-Length: {len(body)}\r\n" h += "\r\n" s.sendall(h.encode() + (body.encode() if body else b"")) data=b"" try: while True: c = s.recv(65536) if not c: break data += c if len(data) > 300000: break except socket.timeout: pass s.close() head = data.split(b"\r\n",1)[0].decode(errors="replace") bodyb = data.split(b"\r\n\r\n",1)[1] if b"\r\n\r\n" in data else b"" return head, bodyb[:400] def ws(path, secs=4): s = socket.create_connection((ip,8080),timeout=8) key = base64.b64encode(os.urandom(16)).decode() s.sendall((f"GET {path} HTTP/1.1\r\nHost: x\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Key: {key}\r\nSec-WebSocket-Version: 13\r\n\r\n").encode()) buf=b"" try: while b"\r\n\r\n" not in buf: buf += s.recv(4096) except Exception as e: return f"ERR {e}", b"" head, data = buf.split(b"\r\n\r\n",1) s.settimeout(secs); t0=time.time() try: while time.time()-t0 < secs: data += s.recv(1<<20) except socket.timeout: pass s.close() return head.decode(errors="replace").splitlines()[0], data def report(name, result): print(f"### {name}\n{result}\n") now_us = int(time.time()*1e6) # --- subscribe cursor edges --- h,d = ws("/subscribe?cursor=0", 4) report("cursor=0 (oldest retained expected)", f"{h}; bytes={len(d)}; first_time_us={re.search(rb'time_us.:(\\d+)',d).group(1).decode() if re.search(rb'time_us.:(\\d+)',d) else None}") h,d = ws(f"/subscribe?cursor={now_us + 3600_000000}", 4) report("future cursor (+1h) — expect silent live", f"{h}; bytes={len(d)}") h,d = ws("/subscribe?cursor=999999999999999999999999", 3) report("absurd numeric cursor", f"{h}; bytes={len(d)}") h,d = ws("/subscribe?cursor=banana", 3) report("non-numeric cursor", f"{h}; bytes={len(d)}") h,d = ws("/subscribe?cursor=-5", 3) report("negative cursor", f"{h}; bytes={len(d)}") h,d = ws("/subscribe?wantedCollections=app.bsky.feed.post", 5) kinds = set(re.findall(rb'"collection":"([^"]+)"', d)) report("wantedCollections filter (coral's exact query)", f"{h}; bytes={len(d)}; collections seen={sorted(k.decode() for k in kinds)[:5]}") h,d = ws("/subscribe?wantedCollections=bogus.collection.name", 4) report("filter matching nothing — expect connected+quiet", f"{h}; bytes={len(d)}") h,d = ws("/subscribe?wantedDids=did:plc:doesnotexist123", 4) report("wantedDids no match", f"{h}; bytes={len(d)}") # --- v2 --- h,d = ws("/subscribe-v2", 4) report("subscribe-v2 no cursor", f"{h}; bytes={len(d)}; sample={d[:80]!r}") # --- xrpc --- for name, path in [ ("listSegments default", "/xrpc/network.bsky.jetstream.listSegments"), ("listSegments limit=1", "/xrpc/network.bsky.jetstream.listSegments?limit=1"), ("listSegments limit=0", "/xrpc/network.bsky.jetstream.listSegments?limit=0"), ("listSegments limit=-1", "/xrpc/network.bsky.jetstream.listSegments?limit=-1"), ("listSegments limit=99999", "/xrpc/network.bsky.jetstream.listSegments?limit=99999"), ("listSegments cursor=zzz", "/xrpc/network.bsky.jetstream.listSegments?cursor=zzz"), ("getSegment missing name", "/xrpc/network.bsky.jetstream.getSegment"), ("getSegment bad name", "/xrpc/network.bsky.jetstream.getSegment?name=../../etc/passwd"), ("getSegment nonexistent", "/xrpc/network.bsky.jetstream.getSegment?name=seg_zzzzzzzzzz.jss"), ("getBlock no params", "/xrpc/network.bsky.jetstream.getBlock"), ("getBlock bad index", "/xrpc/network.bsky.jetstream.getBlock?segment=seg_0000000000.jss&blockIndex=999999"), ("getBlock valid", "/xrpc/network.bsky.jetstream.getBlock?segment=seg_0000000000.jss&blockIndex=0"), ("getZstdDictionary", "/xrpc/network.bsky.jetstream.getZstdDictionary"), ("unknown xrpc", "/xrpc/com.atproto.sync.getRepo"), ]: h,b = http(path) report(name, f"{h}; body[:120]={b[:120]!r}") h,b = http("/xrpc/network.bsky.jetstream.planBackfill", "POST", json.dumps({"fromSeq":1,"toSeq":1000000})) report("planBackfill basic", f"{h}; body={b[:200]!r}") h,b = http("/xrpc/network.bsky.jetstream.planBackfill", "POST", "not json") report("planBackfill garbage body", f"{h}; body={b[:120]!r}") # getSegment Range support (matters for resumable downloads) s = socket.create_connection((ip,8080),timeout=10) s.sendall(b"GET /xrpc/network.bsky.jetstream.getSegment?name=seg_0000000000.jss HTTP/1.1\r\nHost: x\r\nRange: bytes=0-99\r\nConnection: close\r\n\r\n") d=b"" try: while len(d)<600: c=s.recv(4096) if not c: break d+=c except socket.timeout: pass s.close() report("getSegment Range request", d.split(b"\r\n\r\n")[0].decode(errors="replace")[:300])