# /// script # dependencies = [] # /// """backfill-batch contract: a dispatch batch larger than the job queue must still make durable progress. This is the defect that ended the 27 July 2026 whole-network experiment, and the one that made the 23 July run look like a dispatch-boundary metric problem. The job queue is bounded at `worker_count * 2`. `finishWorkBatch` used to submit every job in the batch before consuming a single result, so with a whole-network batch (100,000) against a 400-slot queue the submitting fiber blocked on a full queue while finished repositories accumulated in `results`, each still owning its CAR against the in-flight budget. The budget exhausted, workers could no longer allocate, they stopped draining the queue, and the submit the fiber was blocked on never completed. Nothing was ever queued for completion: 12,686 repositories emitted, zero durable, every row still `not_started`. Two conditions are required, which is why nothing caught it: 1. more jobs in one batch than the queue holds -- every other test corpus is smaller than the queue, and 2. an in-flight budget small enough that retained results exhaust it -- the default 8 GiB is never reached by a small fixture. WHAT THIS TEST DOES NOT DO: it does not reproduce the production deadlock, and it passes on the broken code. The pinned simulator's repositories are ~4 KiB, so concurrent demand cannot exceed any budget that a single repository still fits in -- measured, the whole 3,000-repository run peaks at ~50 KiB in flight. The deadlock needs per-repository memory large enough that concurrent workers exhaust a multi-GB budget while the submit loop stays blocked, which only real repositories produce. This is the same fixture limit that keeps the differential-oracle row partial in docs/semantic-parity.md, in a second and distinct way. What it does pin is the weaker property that a dispatch batch many times larger than the job queue converges, and that every repository handled becomes durable rather than merely counted. That is worth holding, but do not read a pass here as evidence that the deadlock is fixed. run: uv run tests/backfill_batch_contract.py (or `just backfill-batch-contract`) """ import os import pathlib import re import shutil import subprocess import sys import time import urllib.error import urllib.request PORT = 6103 BASE = f"http://127.0.0.1:{PORT}" DATA = pathlib.Path("data-backfill-batch") BIN = "./zig-out/bin/stream" SIM = os.environ.get("STREAM_SIM_URL", "http://127.0.0.1:7777") # The deadlock needs concurrent demand to exceed the budget, so it needs many # workers holding memory at once -- one worker can never exhaust it, which is # why a single-worker version of this test passes on the broken code. WORKERS = 16 REPOS = int(os.environ.get("STREAM_BATCH_REPOS", "3000")) # Generous relative to the work: 64 tiny simulator repositories converge in # seconds. Exceeding this means no progress is being made, not slow progress. TIMEOUT_S = 180 SAMPLE_RE = re.compile(r"^([a-zA-Z_:][a-zA-Z0-9_:]*)(?:\{([^}]*)\})?\s+([^\s]+)$") LABEL_RE = re.compile(r'(\w+)="([^"]*)"') def metrics(): with urllib.request.urlopen(BASE + "/metrics", timeout=5) as r: body = r.read().decode() out = {} for line in body.splitlines(): if line.startswith("#"): continue m = SAMPLE_RE.match(line.strip()) if not m: continue name, labels, value = m.groups() key = name if labels: pairs = dict(LABEL_RE.findall(labels)) if "status" in pairs: key = f'{name}{{status="{pairs["status"]}"}}' elif "result" in pairs: key = f'{name}{{result="{pairs["result"]}"}}' try: out[key] = float(value) except ValueError: pass return out def main(): shutil.rmtree(DATA, ignore_errors=True) log_path = pathlib.Path("/tmp/backfill-batch-contract.log") log = open(log_path, "w") proc = subprocess.Popen( [ BIN, f"--port={PORT}", f"--data-dir={DATA}", f"--relay-url={SIM}", f"--plc-url={SIM}", f"--max-backfill-repos={REPOS}", f"--backfill-workers={WORKERS}", "--compaction-interval=0", "--retry-interval=0", ], stdout=log, stderr=log, ) deadline = time.time() + TIMEOUT_S observed = {} try: while time.time() < deadline: if proc.poll() is not None: raise SystemExit( f"stream exited early with {proc.returncode}; see {log_path}" ) try: observed = metrics() except (urllib.error.URLError, ConnectionError, TimeoutError): time.sleep(0.5) continue complete = observed.get('stream_backfill_repos_durable{status="complete"}', 0) handled_now = observed.get( "jetstream_backfill_handle_repo_duration_seconds_count", 0 ) if handled_now > WORKERS * 2 * 8 and complete >= handled_now: break time.sleep(0.5) else: queued = observed.get("jetstream_backfill_completion_queued_total", 0) handled = observed.get( "jetstream_backfill_handle_repo_duration_seconds_count", 0 ) raise SystemExit( "batch never reached durable completion within " f"{TIMEOUT_S}s: handled={handled:.0f} queued_for_completion={queued:.0f} " f"durable_complete={observed.get('stream_backfill_repos_durable{status=\"complete\"}', 0):.0f} " "A full queue with unconsumed results is the batch deadlock." ) # Anti-vacuity: the batch must actually have exceeded the queue, or # this test proves nothing about the condition it exists to pin. handled = observed.get("jetstream_backfill_handle_repo_duration_seconds_count", 0) assert handled > WORKERS * 2 * 8, f"only {handled:.0f} repositories were handled" queue_capacity = WORKERS * 2 assert REPOS > queue_capacity * 4, ( f"{REPOS} jobs vs {queue_capacity}-slot queue is not an overflow" ) complete = observed['stream_backfill_repos_durable{status="complete"}'] # Not every discovered repository is eligible, so the bar is that every # repository actually handled became durable -- not an exact count. assert complete >= handled, f"durable complete {complete:.0f} < handled {handled:.0f}" print( f"ok: {complete:.0f}/{REPOS} durably complete with a " f"{queue_capacity}-slot queue" ) finally: proc.terminate() try: proc.wait(timeout=30) except subprocess.TimeoutExpired: proc.kill() log.close() if __name__ == "__main__": sys.exit(main())