jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635# /// script# dependencies = ["websockets"]# ///"""oracle: crash-matrix restart harness (docs/upstream-harness.md).
expects the simulator on :7777. for every lifecycle crashpoint, thisdriver runs the real binary until it aborts at the armed seam, restartsit clean, waits for the lifecycle to converge to steady_state, andasserts 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 = 6099BASE = 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()