jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106import socket, ssl, base64, os, time, subprocess, re, json, statisticsip = 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+=cexcept socket.timeout: passs.close()report("getSegment Range request", d.split(b"\r\n\r\n")[0].decode(errors="replace")[:300])