jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118# /// script# requires-python = ">=3.12"# dependencies = ["requests", "zstandard", "cbor2"]# ///"""Loopback-only Strata exploration. Archive key never reaches the browser."""from http.server import ThreadingHTTPServer, SimpleHTTPRequestHandlerfrom pathlib import Pathfrom urllib.parse import urlparse, parse_qsimport base64import importlib.utilimport jsonimport subprocessimport threadingimport timeimport requestsimport cbor2import zstandard
ROOT = Path(__file__).parentSTREAM = 'https://stream.waow.tech'STRATA = 'https://strata.zat.dev'# Reuse the instance's documented columnar decoder rather than invent a format.spec = importlib.util.spec_from_file_location('archive_example', ROOT.parent / 'examples/backfill-collection-sqlite.py')archive_example = importlib.util.module_from_spec(spec)spec.loader.exec_module(archive_example)KEY = subprocess.check_output(['sops', '--decrypt', '--extract', '["prefect"]["blocks"]["stream-archive-key-strata"]', str(Path.home() / 'tangled.org/zzstoatzz.io/secrets/prod.yaml')], stderr=subprocess.DEVNULL).decode().strip()HEADERS = {'Authorization': 'Bearer ' + KEY}CACHE = {}LOCK = threading.Lock()READ_LOCK = threading.Semaphore(2)
def fetch_json(url, **kwargs): response = requests.get(url, timeout=30, **kwargs) response.raise_for_status() return response.json()
def overview(): with LOCK: if CACHE and time.time() - CACHE['fetchedAt']/1000 < 300: return CACHE progress = fetch_json(STRATA + '/api/progress') collections = fetch_json(STRATA + '/api/collections')['collections'] heat = fetch_json(STRATA + '/api/heat?step=100&top=1') status = requests.get(STREAM + '/status', timeout=15).text CACHE.update(progress=progress, collections=collections, buckets=heat['buckets'], status=status, fetchedAt=int(time.time()*1000), source=STREAM, summarySource=STRATA) return CACHE
def json_cbor(value): if isinstance(value, bytes): return {'base64': base64.b64encode(value).decode()} if isinstance(value, cbor2.CBORTag): return {'tag':value.tag,'value':json_cbor(value.value)} if isinstance(value, dict): return {str(k):json_cbor(v) for k,v in value.items()} if isinstance(value, (list,tuple)): return [json_cbor(v) for v in value] return value
def sample(collection, after, before): if collection not in {c['nsid'] for c in overview()['collections']}: raise ValueError('Unknown collection') if after < 0 or before <= after: raise ValueError('Invalid archive range') started = time.monotonic() response = requests.post(STREAM+'/xrpc/network.bsky.jetstream.planSnapshot',headers=HEADERS,json={'collections':[collection], 'afterSeq':after, 'beforeSeq':before},timeout=30) response.raise_for_status() plan = response.json() rows=[]; blocks=0; downloaded=0 # Inspect at most four planned blocks; never fetch an entire segment. for segment in plan.get('segments',[]): groups = segment.get('blocks') or [{'first':0,'last':3}] for group in groups: for index in range(group['first'],group['last']+1): if blocks >= 4 or len(rows) >= 8: break block = requests.get(STREAM+'/xrpc/network.bsky.jetstream.getBlock', params={'segment':segment['name'],'blockIndex':index},headers=HEADERS,timeout=20,stream=True) if block.status_code == 404: block.close() break block.raise_for_status() with block: compressed = block.raw.read(4*1024*1024+1) if len(compressed)>4*1024*1024: raise ValueError('Block exceeds prototype read limit') raw=zstandard.ZstdDecompressor().decompress(compressed,max_output_size=32*1024*1024) downloaded+=len(compressed);blocks+=1 for seq,witnessed,did,rkey,payload in archive_example.decode_block(raw,collection): if not after < seq <= before: continue rows.append({'seq':seq,'witnessedAt':witnessed,'uri':f'at://{did}/{collection}/{rkey}','record':json_cbor(cbor2.loads(payload)) if payload else None}) if len(rows)>=8: break if blocks>=4 or len(rows)>=8: break if blocks>=4 or len(rows)>=8: break return {'collection':collection,'records':rows,'blocksRead':blocks,'bytesRead':downloaded,'seconds':round(time.monotonic()-started,2),'afterSeq':after,'beforeSeq':before,'plannedThroughSeq':plan.get('plannedThroughSeq'),'sealedTipSeq':plan.get('sealedTipSeq'),'boundedSample':True}
class Handler(SimpleHTTPRequestHandler): def __init__(self,*args,**kwargs): super().__init__(*args,directory=str(ROOT),**kwargs) def do_GET(self): parsed=urlparse(self.path) if self.headers.get('Host') not in ('127.0.0.1:4179','localhost:4179') or self.headers.get('Origin') not in (None,'http://127.0.0.1:4179','http://localhost:4179'): self.send_error(403); return if parsed.path.endswith('.py'): self.send_error(404); return if not parsed.path.startswith('/api/'): return super().do_GET() try: if parsed.path=='/api/overview': result=overview() elif parsed.path=='/api/sample': q=parse_qs(parsed.query) with READ_LOCK: result=sample(q['collection'][0],int(q.get('after',['0'])[0]),int(q['before'][0])) elif parsed.path=='/api/heat': ns=parse_qs(parsed.query)['ns'][0] if ns not in {c['ns'] for c in overview()['collections']}: raise ValueError('Unknown namespace') result=fetch_json(STRATA+'/api/heat',params={'step':100,'ns':ns}) else: self.send_error(404); return body=json.dumps(result,ensure_ascii=False).encode() self.send_response(200) except (requests.RequestException, ValueError, KeyError) as error: body=json.dumps({'error':'Archive request failed. Try again or choose a different slice.','kind':type(error).__name__}).encode() self.send_response(502) self.send_header('Content-Type','application/json');self.send_header('Cache-Control','no-store');self.send_header('Content-Length',str(len(body)));self.end_headers();self.wfile.write(body) def log_message(self,format,*args): pass
if __name__=='__main__': print('Strata archive prototype: http://127.0.0.1:4179',flush=True) ThreadingHTTPServer(('127.0.0.1',4179),Handler).serve_forever()