jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788# /// script# requires-python = ">=3.12"# dependencies = ["requests", "zstandard"]# ///"""Bounded metadata-only storage investigation. No block payloads or Microcosm calls."""from pathlib import Pathimport collectionsimport jsonimport structimport subprocessimport timeimport requestsimport zstandard
ROOT = Path(__file__).parentkey = 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()session = requests.Session()session.headers['Authorization'] = 'Bearer '+keyBASE = 'https://stream.waow.tech/xrpc/network.bsky.jetstream.'requests_made=0bytes_read=0
def get(endpoint, params, byte_range=None): global requests_made,bytes_read if requests_made>=20: raise ValueError('Request budget exhausted') requests_made+=1 headers={'Range':f'bytes={byte_range}'} if byte_range else {} cap=4*1024*1024 with session.get(BASE+endpoint,params=params,headers=headers,timeout=30,stream=True) as r: r.raise_for_status() if byte_range and r.status_code!=206: raise ValueError('Range request was not honored') body=r.raw.read(cap+1) if len(body)>cap: raise ValueError('Metadata exceeds 4 MiB read limit') bytes_read+=len(body) return body
segments=[]cursor=Nonefor page_no in range(10): params={'limit':1000} if cursor is not None: params['cursor']=cursor page=json.loads(get('listSegments',params));segments.extend(page['segments']) cursor=page.get('cursor') if not cursor: breakelse: raise ValueError('Listing exceeded page limit')
probes=[]for seg in [segments[0],segments[len(segments)//2],segments[-1]]: params={'name':seg['name']} header=get('getSegment',params,'0-255') if len(header)!=256 or header[:4]!=b'jss0' or struct.unpack_from('<H',header,12)[0]!=1: raise ValueError('Invalid header') checksum=struct.unpack_from('<Q',header,4)[0] if f'{checksum:016x}'!=seg['checksum']: raise ValueError('Segment changed since listing') blocks=struct.unpack_from('<I',header,14)[0] footer_offset=struct.unpack_from('<Q',header,58)[0] collection_offset=struct.unpack_from('<Q',header,82)[0] footer_size=seg['sizeBytes']-footer_offset if blocks*52>4*1024*1024 or seg['sizeBytes']-collection_offset>4*1024*1024: raise ValueError('Index exceeds read budget') block_index=get('getSegment',params,f'{footer_offset}-{footer_offset+blocks*52-1}') ix=get('getSegment',params,f'{collection_offset}-{seg["sizeBytes"]-1}') if len(block_index)!=blocks*52 or len(ix)!=seg['sizeBytes']-collection_offset: raise ValueError('Truncated index') if get('getSegment',params,'0-255')!=header: raise ValueError('Segment changed during probe') entries=[struct.unpack_from('<QIIIQQqq',block_index,i*52) for i in range(blocks)] compressed=sum(e[1] for e in entries) uncompressed=sum(e[2] for e in entries) framing=blocks*8 if compressed+framing+256+footer_size!=seg['sizeBytes']: raise ValueError('Byte accounting mismatch') count,block_count,mask_len,decoded_size=struct.unpack_from('<IIII',ix) if block_count!=blocks or mask_len!=(count+7)//8 or decoded_size>32*1024*1024: raise ValueError('Invalid collection index') decoded=zstandard.ZstdDecompressor().decompress(ix[16:],max_output_size=32*1024*1024) if len(decoded)!=decoded_size: raise ValueError('Invalid decoded size') names=[];position=0 for i in range(count): length=decoded[position]; names.append(decoded[position+5:position+5+length].decode());position+=5+length if position+blocks*mask_len!=len(decoded): raise ValueError('Invalid bitmask length') exclusive=collections.Counter();mixed_bytes=0;exclusive_bytes=0;empty_bytes=0;mixed_blocks=0 for i,entry in enumerate(entries): mask=decoded[position+i*mask_len:position+(i+1)*mask_len] groups={'.'.join(name.split('.')[:2]) if not name.startswith('$') else name for j,name in enumerate(names) if mask[j//8]&(1<<(j%8))} if len(groups)==1: exclusive[next(iter(groups))]+=entry[1];exclusive_bytes+=entry[1] elif len(groups)>1: mixed_bytes+=entry[1];mixed_blocks+=1 else: empty_bytes+=entry[1] probes.append({'name':seg['name'],'index':seg['index'],'checksum':seg['checksum'],'fileBytes':seg['sizeBytes'],'compressedFrameBytes':compressed,'uncompressedBlockBytes':uncompressed,'framingBytes':framing,'headerBytes':256,'footerBytes':footer_size,'blocks':blocks,'collections':count,'mixedNamespaceBlocks':mixed_blocks,'mixedNamespaceFrameBytes':mixed_bytes,'singleNamespaceFrameBytes':exclusive_bytes,'emptyFrameBytes':empty_bytes,'largestSingleNamespaceFrameBytes':exclusive.most_common(5)})report={'measuredAt':int(time.time()*1000),'source':'https://stream.waow.tech','segmentCount':len(segments),'sealedFileBytes':sum(s['sizeBytes'] for s in segments),'smallestFileBytes':min(s['sizeBytes'] for s in segments),'largestFileBytes':max(s['sizeBytes'] for s in segments),'requests':requests_made,'metadataBytesRead':bytes_read,'probes':probes,'scope':'Listing of sealed files; first, middle, and last listing entries probed. Not a representative sample or a per-namespace disk total.'}(ROOT/'storage-report.json').write_text(json.dumps(report,indent=2)+'\n')print(json.dumps(report,indent=2))