atproto pds in zig pds.zat.dev
pds atproto
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116#!/usr/bin/env python3"""Measure process launch to a successful HTTP health response; never use a live DB."""import argparsefrom contextlib import closingimport hashlibimport jsonimport osfrom pathlib import Pathimport platformimport reimport socketimport sqlite3import statisticsimport subprocessimport tempfileimport timeimport urllib.request
def run(binary, db, directory, timeout): with socket.socket() as sock: sock.bind(('127.0.0.1', 0)) port = sock.getsockname()[1] origin = f'http://127.0.0.1:{port}' env = {key: value for key, value in os.environ.items() if not key.startswith('ZDS_')} command = [str(binary), '--host', '127.0.0.1', '--port', str(port), '--db', str(db), '--blobstore-path', str(directory / 'blobs'), '--public-url', origin, '--server-did', 'did:web:localhost', '--jwt-secret', 'disposable-startup-benchmark', '--crawlers', ''] opener = urllib.request.build_opener(urllib.request.ProxyHandler({})) with tempfile.TemporaryFile(mode='w+b') as log: started = time.perf_counter() process = subprocess.Popen(command, env=env, stdout=log, stderr=log) try: while True: if process.poll() is not None: raise RuntimeError(f'server exited with status {process.returncode}') if time.perf_counter() - started > timeout: raise TimeoutError('server did not become healthy') try: with opener.open(origin + '/xrpc/_health', timeout=.2) as response: if response.status == 200 and json.load(response).get('status') == 'ok': elapsed = (time.perf_counter() - started) * 1000 break except (OSError, ValueError): pass time.sleep(.005) finally: process.terminate() try: process.wait(timeout=5) except subprocess.TimeoutExpired: process.kill() process.wait() log.seek(0) # Only stage timings are exported; a snapshot can contain private account data. stages = {key: int(value) for key, value in re.findall( r'startup (\w+_ms)=(\d+)', log.read().decode(errors='replace'))} return {'ready_ms': round(elapsed, 2), **stages}
def seed_history(db, commits): # Storage-shaped fixture only: these are not signed ATProto repo blocks. # An inactive account avoids account-announcement work; NULL revisions # model historical tree blocks that match neither a record nor a commit. with closing(sqlite3.connect(db)) as conn, conn: conn.execute("INSERT INTO accounts(did,handle,email,password_hash,account_status) VALUES ('did:plc:bench','bench.test','bench@example.test','unused','deactivated')") conn.executemany("INSERT INTO commits(seq,did,cid,rev) VALUES (?,'did:plc:bench',?,?)", ((i+1, f'commit-{i}', f'rev-{i}') for i in range(commits))) conn.executemany("INSERT INTO seq_events(seq,did,commit_cid,evt) VALUES (?,'did:plc:bench',?,zeroblob(4096))", ((i+1, f'commit-{i}') for i in range(commits))) conn.executemany("INSERT INTO repo_blocks(did,cid,data,repo_rev) VALUES ('did:plc:bench',?,zeroblob(4096),?)", ((f'commit-{i}', f'rev-{i}') for i in range(commits))) conn.executemany("INSERT INTO repo_blocks(did,cid,data,repo_rev) VALUES ('did:plc:bench',?,zeroblob(4096),?)", ((f'tree-{i}', None if i < commits//5 else 'rev-0') for i in range(commits*5)))
def main(): parser = argparse.ArgumentParser(description=__doc__) parser.add_argument('--binary', type=Path, default=Path('zig-out/bin/zds')) parser.add_argument('--runs', type=int, default=5) parser.add_argument('--timeout', type=float, default=30) parser.add_argument('--history', type=int, default=10000, help='synthetic commits for the history restart lane (0 disables)') parser.add_argument('--snapshot', type=Path, help='SQLite backup source; opened read-only and copied') args = parser.parse_args() if args.runs < 1 or args.timeout <= 0 or args.history < 0: parser.error('runs and timeout must be positive; history must be nonnegative') binary = args.binary.resolve(strict=True) result = {'platform': platform.platform(), 'binary': str(binary), 'binary_sha256': hashlib.sha256(binary.read_bytes()).hexdigest(), 'poll_interval_ms': 5, 'history_commits': args.history, 'block_bytes': 4096, 'event_bytes': 4096, 'note': 'warm filesystem caches; excludes build, VM boot and proxy routing', 'scenarios': {}} with tempfile.TemporaryDirectory(prefix='zds-startup-') as root: directory = Path(root) for scenario in ['fresh', 'restart'] + (['history'] if args.history else []) + (['snapshot'] if args.snapshot else []): samples = [] db = directory / f'{scenario}.sqlite3' if scenario in ('restart', 'history'): run(binary, db, directory, args.timeout) if scenario == 'history': seed_history(db, args.history) for index in range(args.runs): if scenario == 'fresh': db = directory / f'fresh-{index}.sqlite3' if scenario == 'snapshot': with closing(sqlite3.connect(args.snapshot.resolve().as_uri() + '?mode=ro', uri=True)) as source: with closing(sqlite3.connect(db)) as dest: source.backup(dest) samples.append(run(binary, db, directory, args.timeout)) values = [sample['ready_ms'] for sample in samples] result['scenarios'][scenario] = {'median_ms': statistics.median(values), 'min_ms': min(values), 'max_ms': max(values), 'samples': samples} print(json.dumps(result, indent=2))
if __name__ == '__main__': main()