# /// script # dependencies = ["websockets"] # /// """oracle: crash-matrix restart harness (docs/upstream-harness.md). expects the simulator on :7777. for every lifecycle crashpoint, this driver runs the real binary until it aborts at the armed seam, restarts it clean, waits for the lifecycle to converge to steady_state, and asserts serving invariants: - /healthz stays reachable; serving ungates (no 503) - listSegments returns sealed segments with events - /subscribe replay: seqs strictly increasing, events well-formed - serving proves the durable phase reached steady_state; backfill/ is gone - the unified RocksDB metadata directory exists and no legacy stores appear run: uv run tests/oracle.py (or `just oracle`) """ import asyncio, json, pathlib, re, shutil, signal, subprocess, sys, time, urllib.request import websockets PORT = 6099 BASE = f"http://127.0.0.1:{PORT}" DATA = pathlib.Path("data-oracle") BIN = "./zig-out/bin/stream" SIM = "http://127.0.0.1:7777" CRASHPOINTS = [ "mid_backfill_download", "after_backfill_before_phase_merging", "after_phase_merging_before_live_stop", "after_live_close_before_seal", "after_merge_drain_before_pending", "after_pending_pass", "after_seal_before_discovery", "after_discovery_before_cleanup", "after_cleanup_before_phase_write", ] ARGS = [ BIN, f"--port={PORT}", f"--data-dir={DATA}", f"--relay-url={SIM}", f"--plc-url={SIM}", "--max-backfill-repos=25", "--compaction-interval=0", "--retry-interval=0", ] SAMPLE_RE = re.compile(r"^([a-zA-Z_:][a-zA-Z0-9_:]*)(?:\{([^}]*)\})?\s+([^\s]+)$") LABEL_RE = re.compile(r'(\w+)="([^"]*)"') def http_status(path): try: with urllib.request.urlopen(BASE + path, timeout=5) as r: return r.status, r.read() except urllib.error.HTTPError as e: return e.code, b"" except OSError: return None, b"" def wait_expected_exit(p, timeout=300): """Reap an expected process exit; on timeout, kill and reap before raising. A leaked process squats the port via SO_REUSEPORT and the kernel then round-robins later cases' requests into its 503s — the f0f84f1 gate leaked exactly one such squatter and every following run coin-flipped until it was found and killed.""" try: return p.wait(timeout=timeout) except subprocess.TimeoutExpired: p.kill() p.wait() raise def wait_serving(proc, timeout=120): """poll until the xrpc surface ungates (lifecycle hit steady_state)""" deadline = time.time() + timeout while time.time() < deadline: if proc.poll() is not None: raise AssertionError(f"stream exited early: rc={proc.returncode}") st, _ = http_status("/xrpc/network.bsky.jetstream.listSegments") if st == 200: return time.sleep(0.5) raise AssertionError("timed out waiting for steady_state serving") def metrics_snapshot(): st, body = http_status("/metrics") assert st == 200, st samples = {} for line in body.decode().splitlines(): if line.startswith("#"): continue match = SAMPLE_RE.match(line) if not match: continue name, labels, value = match.groups() samples[(name, tuple(sorted(LABEL_RE.findall(labels or ""))))] = float(value) return samples def metric(samples, name, **labels): return samples.get((name, tuple(sorted(labels.items()))), 0.0) async def replay_invariants(n): """seqs strictly increasing from a cold cursor; frames well-formed""" url = f"ws://127.0.0.1:{PORT}/xrpc/network.bsky.jetstream.subscribeEvents?cursor=0" prev = 0 got = 0 async with websockets.connect(url) as ws: while got < n: raw = await asyncio.wait_for(ws.recv(), timeout=15) frame = json.loads(raw) assert frame.get("$type") == "message", frame e = frame["payload"] assert e.get("$type", "").startswith("network.bsky.jetstream.subscribeEvents#"), e seq = e.get("seq") assert isinstance(seq, int) and seq > prev, (prev, e) # A non-string did means the server served a malformed frame — seen # as `'list' object has no attribute 'startswith'` (2026-08-09, # intermittent, under live-capture fault injection). Report the # frame itself: an AttributeError here hides the evidence. did = e.get("did", "") assert isinstance(did, str) and did.startswith("did:"), ( f"malformed did on subscribeEvents: type={type(did).__name__} " f"value={did!r}\nraw frame: {raw[:600]!r}" ) # The encoder substitutes U+FFFD for invalid UTF-8 so a bad byte # can never change a field's type (wire.zig writeString). That is # a wire-safety net, NOT a licence to serve corrupt rows: if a # replacement character reaches a subscriber, fault recovery # admitted a torn row and that is the bug to chase. assert "�" not in did, ( f"corrupt did served after recovery: {did!r}\n" f"raw frame: {raw[:600]!r}" ) prev = seq got += 1 return got def run_case(point): shutil.rmtree(DATA, ignore_errors=True) log = open(f"/tmp/oracle-{point}.log", "w") # phase 1: run armed, expect the abort at the seam p = subprocess.Popen(ARGS + [f"--crashpoint={point}"], stdout=log, stderr=log) rc = wait_expected_exit(p) assert rc != 0, f"{point}: armed run exited cleanly (seam never hit?)" marker = f"crashpoint {point}: aborting" logged = pathlib.Path(f"/tmp/oracle-{point}.log").read_text() assert marker in logged, f"{point}: abort marker missing (anti-vacuity)" # phase 2: restart clean, converge, verify p = subprocess.Popen(ARGS, stdout=log, stderr=log) try: wait_serving(p) st, body = http_status("/xrpc/network.bsky.jetstream.listSegments") assert st == 200, st segs = json.loads(body)["segments"] assert segs and sum(s["eventCount"] for s in segs) > 0, segs st, _ = http_status("/healthz") assert st == 200, st metrics = metrics_snapshot() assert metric(metrics, "jetstream_orchestrator_phase") == 3 assert metric( metrics, "jetstream_orchestrator_phase_transitions_total", **{"from": "merging", "to": "steady_state"}, ) == 1 for state in ("merge", "write_phase_steady"): assert metric( metrics, "jetstream_orchestrator_state_duration_seconds_count", state=state, ) == 1 merge_counters = ( "jetstream_orchestrator_merge_events_kept_total", "jetstream_orchestrator_merge_events_dropped_total", "jetstream_orchestrator_merge_segments_consumed_total", "jetstream_orchestrator_merge_did_lookups_total", "jetstream_orchestrator_merge_repo_revs_updated_total", "jetstream_orchestrator_merge_dids_discovered_post_bootstrap_total", ) for name in merge_counters: assert (name, ()) in metrics, f"missing registered orchestrator counter {name}" if point in { "mid_backfill_download", "after_backfill_before_phase_merging", }: assert metric( metrics, "jetstream_orchestrator_phase_transitions_total", **{"from": "bootstrap", "to": "merging"}, ) == 1 for state in ( "drain_bootstrap", "seal_bootstrap", "close_backfill", "write_phase_merging", ): assert metric( metrics, "jetstream_orchestrator_state_duration_seconds_count", state=state, ) == 1 if point in { "after_phase_merging_before_live_stop", "after_live_close_before_seal", }: # The killed process already synced phase=merging. Recovery must # enter merge directly rather than replaying bootstrap or claiming # a fresh bootstrap -> merging transition in this process. assert ( metric( metrics, "jetstream_orchestrator_phase_transitions_total", **{"from": "bootstrap", "to": "merging"}, ) == 0 ) assert ( metric( metrics, "jetstream_orchestrator_state_duration_seconds_count", state="write_phase_merging", ) == 0 ) if point == "mid_backfill_download": assert metric(metrics, "jetstream_orchestrator_merge_segments_consumed_total") > 0 assert metric(metrics, "jetstream_orchestrator_merge_did_lookups_total") > 0 assert ( metric(metrics, "jetstream_orchestrator_merge_events_kept_total") + metric(metrics, "jetstream_orchestrator_merge_events_dropped_total") ) > 0 # A fresh simulator has small repos; older simulator state can have # thousands of rows. Verify the complete replay when it is small and a # bounded 200-event prefix when it is large. Waiting for an arbitrary # 200 events when only (say) 168 exist turns success into a timeout. replay_count = min(200, sum(s["eventCount"] for s in segs)) assert replay_count > 0 n = asyncio.run(replay_invariants(replay_count)) assert n == replay_count logged = pathlib.Path(f"/tmp/oracle-{point}.log").read_text() assert "OutOfMemory" not in logged, "capture/archive shutdown raced" assert "fatal event handling failure" not in logged assert (DATA / "meta.rocksdb").is_dir(), "unified metadata store missing" assert not (DATA / "phase").exists(), "legacy phase store reappeared" assert not (DATA / "repos.log").exists(), "legacy repo store reappeared" assert not (DATA / "backfill").exists(), "backfill tree survived cleanup" if point == "mid_backfill_download": # anti-vacuity: the crash left a durable dispatch batch, so the # restart's merge MUST have resynced at least one row through the # pending pass. Parallel bootstrap persists the page before fanout, # so more than one interrupted row is expected. logged = pathlib.Path(f"/tmp/oracle-{point}.log").read_text() match = re.search( r"pending-repo pass: (\d+) candidates, (\d+) resynced, (\d+) deferred", logged ) assert match, "pending pass did not report a receipt" candidates, resynced, deferred = (int(g) for g in match.groups()) assert resynced > 0, "pending pass never repaired interrupted repos" # Anti-vacuity in both directions: every candidate must be accounted # for, so a pass that silently drops rows cannot read as a clean one. assert candidates == resynced + deferred, (candidates, resynced, deferred) # A process born directly into durable steady_state must restore # phase=3 without inventing transitions that occurred in an older # process lifetime. p.send_signal(signal.SIGTERM) p.wait(timeout=15) p = subprocess.Popen(ARGS, stdout=log, stderr=log) wait_serving(p) restarted_metrics = metrics_snapshot() assert metric(restarted_metrics, "jetstream_orchestrator_phase") == 3 assert metric( restarted_metrics, "jetstream_orchestrator_phase_transitions_total", **{"from": "merging", "to": "steady_state"}, ) == 0 finally: p.send_signal(signal.SIGTERM) try: p.wait(timeout=15) except subprocess.TimeoutExpired: p.kill() log.close() def run_store_fault_case(): """A merge-cursor RocksDB failure must abort loudly, then recover.""" shutil.rmtree(DATA, ignore_errors=True) log_path = pathlib.Path("/tmp/oracle-store-fault.log") log = open(log_path, "w") p = subprocess.Popen(ARGS + [ "--store-fault-prefix=merge/next_source_idx", "--store-fault-ordinal=1", ], stdout=log, stderr=log) rc = wait_expected_exit(p) assert rc != 0, "armed store-fault run completed instead of failing loud" log.flush() logged = log_path.read_text() marker = "store fault: injected I/O failure for prefix merge/next_source_idx at ordinal 1" assert marker in logged, "store fault never fired (anti-vacuity)" assert "InjectedStoreFault" in logged, "injected store failure did not reach the process boundary" # The failed batch did not advance the source cursor. Reopen the same real # RocksDB/JSS state without a fault and require full lifecycle convergence. p = subprocess.Popen(ARGS, stdout=log, stderr=log) try: wait_serving(p) st, body = http_status("/xrpc/network.bsky.jetstream.listSegments") assert st == 200, st segs = json.loads(body)["segments"] assert segs and sum(s["eventCount"] for s in segs) > 0, segs # A bounded prefix is enough here: the ordinary crash matrix already # sweeps long replay, while this case is specifically the store-fault # fail-loud/reopen contract. replay_count = min(50, sum(s["eventCount"] for s in segs)) assert replay_count > 0 assert asyncio.run(replay_invariants(replay_count)) == replay_count assert not (DATA / "backfill").exists(), "recovery did not finish cleanup" finally: p.send_signal(signal.SIGTERM) try: p.wait(timeout=15) except subprocess.TimeoutExpired: p.kill() log.close() def run_live_capture_fault_case(): """A fatal bootstrap-live cursor write must abort the concurrent crawl.""" shutil.rmtree(DATA, ignore_errors=True) log_path = pathlib.Path("/tmp/oracle-live-capture-fault.log") log = open(log_path, "w") p = subprocess.Popen(ARGS + [ "--max-backfill-repos=100", "--backfill-workers=1", "--store-fault-prefix=relay/cursor", "--store-fault-ordinal=1", ], stdout=log, stderr=log) rc = wait_expected_exit(p) assert rc != 0, "fatal live-capture fault did not stop bootstrap" log.flush() logged = log_path.read_text() marker = "store fault: injected I/O failure for prefix relay/cursor at ordinal 1" assert marker in logged, "live cursor store fault never fired (anti-vacuity)" # The injected relay/cursor fault reaches the process boundary by two racing # paths: directly as InjectedStoreFault, or wrapped by the archive append # that triggered the cursor write as ArchiveAppendFailed. Both name the same # injected fault and both are correctly fatal; which one wins is timing. # The invariant is that a fatal cause surfaces at main, not which spelling # of it does -- pinning one made this assertion a coin flip. boundary = [ f'"msg":"{cause}","component":"main"' for cause in ("InjectedStoreFault", "ArchiveAppendFailed") ] assert any(marker in logged for marker in boundary), \ "live pipeline fatal cause did not reach the process boundary" assert '"msg":"LiveCaptureEnded","component":"main"' not in logged, \ "bootstrap replaced the pipeline's original fatal cause" # The same RocksDB/JSS state must remain restartable. A clean process # resumes bootstrap from durable boundaries and converges without relying # on any frame which the failed callback had not acknowledged. p = subprocess.Popen(ARGS, stdout=log, stderr=log) try: wait_serving(p) st, body = http_status("/xrpc/network.bsky.jetstream.listSegments") assert st == 200, st segs = json.loads(body)["segments"] replay_count = min(50, sum(s["eventCount"] for s in segs)) assert replay_count > 0, segs assert asyncio.run(replay_invariants(replay_count)) == replay_count assert not (DATA / "backfill").exists(), \ "live-capture fault recovery did not finish cleanup" finally: p.send_signal(signal.SIGTERM) try: p.wait(timeout=15) except subprocess.TimeoutExpired: p.kill() log.close() def run_segment_fault_cases(): """Real JSS write/fsync failures must fail loud and recover cleanly.""" cases = [ ("write:1:enospc", ("NoSpaceLeft",), True), ("sync:2:eio", ("InputOutput",), False), # This ordinal lands inside a per-repository archive append. The # repository boundary intentionally maps every local persistence # failure to the fatal ArchiveAppendFailed class, while the concurrent # live archive propagates ShortWrite directly. Either archive may win # the process-wide ordinal; the marker below proves the exact leaf # fault and both accepted boundaries are fatal. ("write:3:shortwrite", ("ShortWrite", "ArchiveAppendFailed"), False), ] for spec, error_names, wants_disk_message in cases: shutil.rmtree(DATA, ignore_errors=True) label = spec.replace(":", "-") log_path = pathlib.Path(f"/tmp/oracle-segment-fault-{label}.log") log = open(log_path, "w") p = subprocess.Popen(ARGS + [f"--segment-fault={spec}"], stdout=log, stderr=log) rc = wait_expected_exit(p) assert rc != 0, f"{spec}: armed segment-fault run completed instead of failing loud" log.flush() logged = log_path.read_text() op, ordinal, kind = spec.split(":") marker = f"segment fault: injected {kind} on {op} at ordinal {ordinal}" assert marker in logged, f"{spec}: segment fault never fired (anti-vacuity)" assert any(name in logged for name in error_names), \ f"{spec}: injected segment failure did not reach process boundary" if wants_disk_message: assert "fatal persistence error: disk full" in logged assert str(DATA) in logged assert "restart stream" in logged p = subprocess.Popen(ARGS, stdout=log, stderr=log) try: wait_serving(p) st, body = http_status("/xrpc/network.bsky.jetstream.listSegments") assert st == 200, st segs = json.loads(body)["segments"] replay_count = min(50, sum(s["eventCount"] for s in segs)) assert replay_count > 0, segs assert asyncio.run(replay_invariants(replay_count)) == replay_count assert not (DATA / "backfill").exists(), f"{spec}: recovery did not finish cleanup" assert not list((DATA / "segments").glob("*.tmp")), f"{spec}: stale segment tmp survived" finally: p.send_signal(signal.SIGTERM) try: p.wait(timeout=15) except subprocess.TimeoutExpired: p.kill() log.close() def run_backfill_metrics_receipt(): """A clean production run must exercise, not merely register, backfill metrics.""" shutil.rmtree(DATA, ignore_errors=True) log = open("/tmp/oracle-backfill-metrics.log", "w") p = subprocess.Popen(ARGS, stdout=log, stderr=log) try: wait_serving(p) samples = metrics_snapshot() completed = metric(samples, "jetstream_backfill_completed_total") queued = metric(samples, "jetstream_backfill_completion_queued_total") durable_repos = metric(samples, "jetstream_backfill_completion_durable_repos_total") assert metric(samples, "jetstream_backfill_discovered_total") > 0 assert completed > 0 assert queued == completed assert durable_repos == completed assert metric(samples, "jetstream_backfill_progress_completed") == completed assert metric(samples, "jetstream_backfill_completion_queue_depth") == 0 assert metric(samples, "jetstream_backfill_completion_durable_batches_total") > 0 assert metric(samples, "jetstream_backfill_completion_queue_wait_seconds_count") == completed assert metric(samples, "jetstream_backfill_handle_repo_duration_seconds_count") > 0 assert metric(samples, "jetstream_backfill_forced_checkpoint_flushes_total") > 0 assert metric(samples, "jetstream_backfill_completion_stage_errors_total") == 0 assert metric(samples, "jetstream_backfill_on_fail_store_errors_total") == 0 # Disabled retry still has registered canonical families; it must not # fabricate work to make the dashboard look active. # # The merge's pending pass is upstream's RunPendingRepoRetryPass: the # same runner with eligibleStatus flipped, so it increments this shared # family exactly once even when the steady loop is disabled. Upstream # behaves identically -- RunPendingRepoRetryPass calls runPass, which # calls incRetryPasses. So one pass is correct here; more than one would # mean the disabled steady loop ran, and any nonzero work counter would # mean work was invented. assert ( "jetstream_backfill_failed_repo_retry_passes_total", (), ) in samples, "missing registered retry passes family" assert metric(samples, "jetstream_backfill_failed_repo_retry_passes_total") == 1 for name in ( "jetstream_backfill_failed_repo_retry_candidates_total", "jetstream_backfill_failed_repo_retry_attempts_total", "jetstream_backfill_failed_repo_retry_succeeded_total", "jetstream_backfill_failed_repo_retry_failed_total", "jetstream_backfill_failed_repo_retry_skipped_host_parked_total", ): assert (name, ()) in samples, f"missing registered retry family {name}" assert metric(samples, name) == 0 finally: p.send_signal(signal.SIGTERM) try: p.wait(timeout=15) except subprocess.TimeoutExpired: p.kill() log.close() shutil.rmtree(DATA, ignore_errors=True) def run_bootstrap_pipeline_metrics_receipt(): """Pipeline scheduler metrics must exist and move before serving ungates.""" shutil.rmtree(DATA, ignore_errors=True) log = open("/tmp/oracle-bootstrap-pipeline-metrics.log", "w") p = subprocess.Popen(ARGS + [ "--max-backfill-repos=100", "--backfill-workers=1", ], stdout=log, stderr=log) try: # Do not wait to catch phase==bootstrap live. Whether that gauge is # observable at all depends on how long bootstrap happens to take: on a # small simulator the whole lifecycle finishes in under four seconds and # the window closes before the first scrape lands, while a large one # holds it open for minutes. Sampling a transient made this fixture a # function of simulator size rather than of Stream's behavior. # # Sample whenever the pipeline has demonstrably moved, then prove # bootstrap really happened from the cumulative transition counter, # which cannot be missed by arriving late. deadline = time.time() + 120 receipt = None while time.time() < deadline: if p.poll() is not None: raise AssertionError(f"stream exited before bootstrap metrics receipt: rc={p.returncode}") try: samples = metrics_snapshot() except AssertionError: time.sleep(0.02) continue if metric(samples, "stream_pipeline_ticket", stage="submitted") > 0: receipt = samples break time.sleep(0.02) assert receipt is not None, "pipeline never scheduled any work" # Anti-vacuity: the run really did pass through bootstrap, and reached # merging, before any of the counters below were read. assert ( metric( receipt, "jetstream_orchestrator_phase_transitions_total", **{"from": "bootstrap", "to": "merging"}, ) > 0 or metric(receipt, "jetstream_orchestrator_phase") == 1 ), "run never entered or left bootstrap" submitted = metric(receipt, "stream_pipeline_ticket", stage="submitted") claimed = metric(receipt, "stream_pipeline_ticket", stage="claimed") emitted = metric(receipt, "stream_pipeline_ticket", stage="emitted") assert submitted >= claimed assert submitted >= emitted queued_gap = metric( receipt, "stream_pipeline_stage_gap", boundary="submitted_to_claimed" ) writer_gap = metric( receipt, "stream_pipeline_stage_gap", boundary="claimed_to_emitted" ) total_gap = metric( receipt, "stream_pipeline_stage_gap", boundary="submitted_to_emitted" ) assert queued_gap >= 0 and writer_gap >= 0 assert total_gap == submitted - emitted assert queued_gap + writer_gap == total_gap assert metric(receipt, "stream_pipeline_fatal") == 0 for stage in ("scheduler", "emit", "repair", "cursor"): key = ("stream_pipeline_failures_total", (("stage", stage),)) assert key in receipt, f"missing registered pipeline failure stage {stage}" assert metric(receipt, "stream_pipeline_failures_total", stage=stage) == 0 finally: p.send_signal(signal.SIGTERM) try: p.wait(timeout=15) except subprocess.TimeoutExpired: p.kill() log.close() shutil.rmtree(DATA, ignore_errors=True) def main(): try: with urllib.request.urlopen(SIM + "/xrpc/com.atproto.sync.listRepos?limit=1", timeout=5) as r: assert r.status == 200 except OSError: sys.exit("simulator not reachable on :7777 — run `just simulator` first") t0 = time.time() run_bootstrap_pipeline_metrics_receipt() print(f"ok bootstrap pipeline metrics ({time.time() - t0:.1f}s)") t0 = time.time() run_backfill_metrics_receipt() print(f"ok backfill metrics ({time.time() - t0:.1f}s)") t0 = time.time() run_live_capture_fault_case() print(f"ok live capture fault recovery ({time.time() - t0:.1f}s)") t0 = time.time() run_store_fault_case() print(f"ok store fault recovery ({time.time() - t0:.1f}s)") t0 = time.time() run_segment_fault_cases() print(f"ok segment fault recovery ({time.time() - t0:.1f}s)") for point in CRASHPOINTS: t0 = time.time() run_case(point) print(f"ok {point} ({time.time() - t0:.1f}s)") shutil.rmtree(DATA, ignore_errors=True) print(f"oracle: {len(CRASHPOINTS)} crashpoints survived") if __name__ == "__main__": main()