# /// script # requires-python = ">=3.12" # dependencies = ["requests", "zstandard", "cbor2"] # /// """Loopback-only Strata exploration. Archive key never reaches the browser.""" from http.server import ThreadingHTTPServer, SimpleHTTPRequestHandler from pathlib import Path from urllib.parse import urlparse, parse_qs import base64 import importlib.util import json import subprocess import threading import time import requests import cbor2 import zstandard ROOT = Path(__file__).parent STREAM = '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()