jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180# /// script# dependencies = []# ///"""backfill-batch contract: a dispatch batch larger than the job queue muststill make durable progress.
This is the defect that ended the 27 July 2026 whole-network experiment, andthe 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 tosubmit every job in the batch before consuming a single result, so with awhole-network batch (100,000) against a 400-slot queue the submitting fiberblocked 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 thesubmit the fiber was blocked on never completed. Nothing was ever queued forcompletion: 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, and2. 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, andit passes on the broken code. The pinned simulator's repositories are ~4 KiB,so concurrent demand cannot exceed any budget that a single repository stillfits in -- measured, the whole 3,000-repository run peaks at ~50 KiB in flight.The deadlock needs per-repository memory large enough that concurrent workersexhaust a multi-GB budget while the submit loop stays blocked, which only realrepositories produce. This is the same fixture limit that keeps thedifferential-oracle row partial in docs/semantic-parity.md, in a second anddistinct way.
What it does pin is the weaker property that a dispatch batch many times largerthan the job queue converges, and that every repository handled becomes durablerather than merely counted. That is worth holding, but do not read a pass hereas evidence that the deadlock is fixed.
run: uv run tests/backfill_batch_contract.py (or `just backfill-batch-contract`)"""import osimport pathlibimport reimport shutilimport subprocessimport sysimport timeimport urllib.errorimport urllib.request
PORT = 6103BASE = 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 = 16REPOS = 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())