jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372#!/usr/bin/env python3"""Offline process receipt for pinned Jetstream V2 environment sources."""
from __future__ import annotations
import jsonimport osimport pathlibimport reimport signalimport socketimport subprocessimport tempfileimport timeimport urllib.errorimport urllib.request
BIN = "./zig-out/bin/stream"PIN = "289b0328c2e1a0ccf8c870cb45de0b2397de19fb"ROOT = pathlib.Path(__file__).resolve().parents[1]# STREAM_UPSTREAM_REPO is the name the justfile, scripts/admit, and the other# oracles already use; this file was the only one reading JETSTREAM_UPSTREAM,# so on any box that does not lay the upstream out as ~/{forge}/{org}/{repo}# it looked for a checkout that was never there. JETSTREAM_UPSTREAM still wins# when set, for anyone who has it in their shell.UPSTREAM = pathlib.Path( os.environ.get("JETSTREAM_UPSTREAM") or os.environ.get("STREAM_UPSTREAM_REPO") or (pathlib.Path.home() / "github.com/bluesky-social/jetstream")).resolve()
def clean_env() -> dict[str, str]: return {key: value for key, value in os.environ.items() if not key.startswith("JETSTREAM_")}
def exact_pinned_map_case() -> None: revision = subprocess.run( ["git", "rev-parse", "HEAD"], cwd=UPSTREAM, text=True, capture_output=True, check=True ).stdout.strip() assert revision == PIN, f"upstream checkout is {revision}, expected {PIN}"
upstream: dict[str, str] = {} for source in (UPSTREAM / "cmd/jetstream/main.go", UPSTREAM / "cmd/jetstream/inspect_all.go"): text = source.read_text() for block in re.findall(r"&cli\.[A-Za-z0-9]+Flag\{(.*?)\n\s*\}", text, re.DOTALL): name = re.search(r'Name:\s*"([^"]+)"', block) env = re.search(r'Sources:\s*cli\.EnvVars\("(JETSTREAM_[^"]+)"\)', block) if name and env: prior = upstream.setdefault(env.group(1), name.group(1)) assert prior == name.group(1)
stream_text = (ROOT / "src/internal/runtime/environment.zig").read_text() stream = { env: flag for env, flag in re.findall( r'\.env = "(JETSTREAM_[^"]+)", \.flag = "([^"]+)"', stream_text ) } # Stream-only variables: each entry is an intentional divergence recorded # in docs/configuration-parity.md. Anything not listed here must match # upstream exactly. stream_only = { "JETSTREAM_REBLOOM_SWEEP": "rebloom-sweep", # Archive bearer gate: upstream's OSS server has no archive auth (the # hosted Bluesky instances gate at their edge); stream carries the # gate in-process. Empty = open. "JETSTREAM_ARCHIVE_API_KEY": "archive-api-key", # Named, metered, file-revocable keys beside the fleet key — the # in-process approximation of the hosted gateway (semantic-parity.md) "JETSTREAM_ARCHIVE_API_KEYS_FILE": "archive-api-keys-file", "JETSTREAM_MAX_SUBSCRIBERS": "max-subscribers", "JETSTREAM_MAX_COLD_READERS": "max-cold-readers", # Upstream's atmos client reconnects only on a closed or errored # socket; a relay trickling frames holds it indefinitely. 0 = off. "JETSTREAM_UPSTREAM_SLOW_MIN_RATE": "upstream-slow-min-rate", } # Upstream variables Stream knowingly has NOT ported: the PDS-direct # fleet bootstrap (upstream specs/notes/2026-08-03-pds-direct-backfill- # design.md, landed by d4dd2f0). Stream's bootstrap remains relay-getRepo; # declaring these flags without the behavior would mislead operators. # Each entry is asserted ABSENT from Stream's map, so porting the feature # forces this list to shrink rather than rot. upstream_unported = { "JETSTREAM_BACKFILL_GLOBAL_DOWNLOADS", "JETSTREAM_BACKFILL_HOST_WORKERS_MAX", "JETSTREAM_BACKFILL_MAX_ACTIVE_HOSTS", "JETSTREAM_BACKFILL_MAX_HOSTS", } assert not (stream_only.keys() & upstream.keys()), ( f"stream_only entries collide with upstream: {sorted(stream_only.keys() & upstream.keys())}" ) assert upstream_unported <= upstream.keys(), ( f"upstream_unported names not declared upstream anymore: {sorted(upstream_unported - upstream.keys())}" ) assert not (upstream_unported & stream.keys()), ( f"unported variables now present in stream; remove them from upstream_unported: {sorted(upstream_unported & stream.keys())}" ) expected = {k: v for k, v in upstream.items() if k not in upstream_unported} | stream_only assert stream == expected, ( f"environment map drift\n" f"missing/wrong: {sorted(expected.items() - stream.items())}\n" f"extra/wrong: {sorted(stream.items() - expected.items())}" )
def unused_port() -> int: with socket.socket() as sock: sock.bind(("127.0.0.1", 0)) return sock.getsockname()[1]
def wait_ready(process: subprocess.Popen[str], port: int) -> None: deadline = time.monotonic() + 10 while time.monotonic() < deadline: if process.poll() is not None: raise AssertionError(process.stderr.read() if process.stderr else "stream exited") try: # listener-first startup: wait past the startup gate, not just # for an answering socket — xrpc paths 503 as "starting" until # storage is published, and this contract asserts on the auth # layer behind that gate body = urllib.request.urlopen(f"http://127.0.0.1:{port}/status", timeout=0.2).read() if b"state starting" not in body: return time.sleep(0.025) except OSError: time.sleep(0.025) raise AssertionError(f"listener {port} did not become ready")
def stop(process: subprocess.Popen[str]) -> list[dict[str, object]]: process.send_signal(signal.SIGTERM) process.wait(timeout=5) stderr = process.stderr.read() if process.stderr else "" assert process.returncode == 0, stderr return [json.loads(line) for line in stderr.splitlines() if line]
def http_status(request: urllib.request.Request) -> int: try: with urllib.request.urlopen(request, timeout=2) as response: response.read() return response.status except urllib.error.HTTPError as error: error.read() return error.code
def unknown_case() -> None: env = clean_env() env.update( { "JETSTREAM_ZZZ": "1", "JETSTREAM_AAA": "value=with=equals", "JETSTREAM_ADDR": "127.0.0.1:0", "JETSTREAM_APP_INSTANCE": "stream-0", "JETSTREAM_ORACLE_SEED": "7001", "JETSTREAM_SIM_DATA_DIR": "./sim", } ) result = subprocess.run( [BIN, "version"], env=env, text=True, capture_output=True, timeout=3, check=False ) assert result.returncode != 0 assert "unrecognized JETSTREAM_ environment variables: JETSTREAM_AAA, JETSTREAM_ZZZ" in result.stderr assert "JETSTREAM_ADDR " not in result.stderr assert "JETSTREAM_APP_INSTANCE" not in result.stderr assert "JETSTREAM_ORACLE_SEED" not in result.stderr assert "JETSTREAM_SIM_DATA_DIR" not in result.stderr
def root_and_inspect_case(root: pathlib.Path) -> None: env = clean_env() env["JETSTREAM_LOG_LEVEL"] = "loud" failed = subprocess.run( [BIN, "version"], env=env, text=True, capture_output=True, timeout=3, check=False ) assert failed.returncode != 0 and "InvalidLogLevel" in failed.stderr
overridden = subprocess.run( [BIN, "--log-level=info", "version"], env=env, text=True, capture_output=True, timeout=3, check=False, ) assert overridden.returncode == 0, overridden.stderr assert overridden.stdout.startswith("jetstream version ")
inspect_root = root / "inspect-from-env" inspect_root.mkdir() env = clean_env() env["JETSTREAM_DATA_DIR"] = str(inspect_root) inspected = subprocess.run( [BIN, "inspect-all"], env=env, text=True, capture_output=True, timeout=3, check=False ) assert inspected.returncode == 0, inspected.stderr assert f"data-dir: {inspect_root}" in inspected.stdout assert f"[0] {inspect_root}/segments" in inspected.stdout
def serve_sources_case(root: pathlib.Path) -> None: port = unused_port() data = root / "serve-from-env" imports = root / "imports-from-env" imports.mkdir() env = clean_env() env.update( { "JETSTREAM_LOG_LEVEL": "debug", "JETSTREAM_LOG_FORMAT": "json", "JETSTREAM_ADDR": f"127.0.0.1:{port}", "JETSTREAM_DEBUG_ADDR": "", "JETSTREAM_RELAY_URL": "http://127.0.0.1:17990", "JETSTREAM_PLC_URL": "http://127.0.0.1:17990", "JETSTREAM_DATA_DIR": str(data), "JETSTREAM_SHUTDOWN_TIMEOUT": "2s", "JETSTREAM_CLIENT_DRAIN_TIMEOUT": "3s", "JETSTREAM_SKIP_MERGE_DISCOVERY": "false", "JETSTREAM_FAILED_REPO_RETRY_INTERVAL": "0", "JETSTREAM_DISABLE_REPO_ACTION_RATE_LIMITS": "false", "JETSTREAM_CURSOR_LOOKBACK": "17m", "JETSTREAM_SEGMENT_CACHE_MAX_AGE": "37s", "JETSTREAM_PLAN_MAX_DIDS": "13", "JETSTREAM_PLAN_MAX_COLLECTIONS": "14", "JETSTREAM_PLAN_MAX_ENTRIES": "15", "JETSTREAM_PLAN_WHOLE_SEGMENT_THRESHOLD": "0.75", "JETSTREAM_SUBSCRIBE_READ_LOG_RETENTION_BYTES": "12345", "JETSTREAM_SUBSCRIBE_BLOCK_CACHE_BYTES": "23456", "JETSTREAM_SUBSCRIBE_READ_BATCH": "7", "JETSTREAM_SUBSCRIBE_SLOW_WINDOW": "8s", "JETSTREAM_SUBSCRIBE_SLOW_MIN_RATE": "2.5", "JETSTREAM_CURSOR_BLOCK_INDEX_CACHE_SIZE": "-9", "JETSTREAM_COMPACTION_INTERVAL": "0", "JETSTREAM_COMPACTION_TOMBSTONE_CAP": "123", "JETSTREAM_COMPACTION_REWRITE_WORKERS": "2", "JETSTREAM_TIMESTAMP_IMPORT_TOKEN": "env-secret", "JETSTREAM_TIMESTAMP_IMPORT_DIR": str(imports), } ) process = subprocess.Popen( [BIN, "serve"], env=env, text=True, stdout=subprocess.DEVNULL, stderr=subprocess.PIPE ) wait_ready(process, port)
endpoint = f"http://127.0.0.1:{port}/xrpc/network.bsky.jetstream.importTimestamps" unauthenticated = urllib.request.Request(endpoint, data=b"{}", method="POST") assert http_status(unauthenticated) == 401 authenticated = urllib.request.Request( endpoint, data=b"{}", method="POST", headers={"Authorization": "Bearer env-secret"} ) assert http_status(authenticated) != 401
records = stop(process) messages = [str(record["msg"]) for record in records] assert any("logging configured: level=debug format=json" in message for message in messages) assert any(f"archive at {data}/segments" in message for message in messages) assert any( "subscribe controls: read-log-retention=12345 bytes block-cache=23456 bytes " "read-batch=7 slow-window=8000000000ns slow-min-rate=2.5" in message for message in messages ) assert any( f"public listener 127.0.0.1:{port}; debug listener disabled" in message for message in messages )
def cli_precedence_case(root: pathlib.Path) -> None: port = unused_port() data = root / "cli-wins" env = clean_env() env.update( { "JETSTREAM_LOG_LEVEL": "loud", "JETSTREAM_ADDR": "not-an-address", "JETSTREAM_DATA_DIR": "/dev/null/not-a-directory", "JETSTREAM_RELAY_URL": "not-a-url", "JETSTREAM_SUBSCRIBE_READ_BATCH": "not-an-int", "JETSTREAM_DISABLE_REPO_ACTION_RATE_LIMITS": "not-a-bool", "JETSTREAM_FAILED_REPO_RETRY_INTERVAL": "not-a-duration", "JETSTREAM_COMPACTION_INTERVAL": "not-a-duration", } ) process = subprocess.Popen( [ BIN, "serve", "--log-level=info", f"--addr=127.0.0.1:{port}", f"--data-dir={data}", "--relay-url=http://127.0.0.1:17990", "--subscribe-read-batch=7", "--disable-repo-action-rate-limits=false", "--failed-repo-retry-interval=0", "--compaction-interval=0", ], env=env, text=True, stdout=subprocess.DEVNULL, stderr=subprocess.PIPE, ) wait_ready(process, port) records = stop(process) assert any("read-batch=7" in str(record["msg"]) for record in records)
def invalid_typed_sources_case(root: pathlib.Path) -> None: keys = ( "JETSTREAM_SHUTDOWN_TIMEOUT", "JETSTREAM_CLIENT_DRAIN_TIMEOUT", "JETSTREAM_MAX_BACKFILL_REPOS", "JETSTREAM_BACKFILL_WORKERS", "JETSTREAM_BACKFILL_BATCH_SIZE", "JETSTREAM_BACKFILL_ASYNC_FLUSH_WORKERS", "JETSTREAM_BACKFILL_REPOS", "JETSTREAM_SKIP_MERGE_DISCOVERY", "JETSTREAM_FAILED_REPO_RETRY_INTERVAL", "JETSTREAM_FAILED_REPO_RETRY_WORKERS", "JETSTREAM_FAILED_REPO_RETRY_HOST_WORKERS", "JETSTREAM_FAILED_REPO_RETRY_MAX_DELAY", "JETSTREAM_DISABLE_REPO_ACTION_RATE_LIMITS", "JETSTREAM_CURSOR_LOOKBACK", "JETSTREAM_SEGMENT_CACHE_MAX_AGE", "JETSTREAM_PLAN_MAX_DIDS", "JETSTREAM_PLAN_MAX_COLLECTIONS", "JETSTREAM_PLAN_MAX_ENTRIES", "JETSTREAM_PLAN_WHOLE_SEGMENT_THRESHOLD", "JETSTREAM_SUBSCRIBE_READ_LOG_RETENTION_BYTES", "JETSTREAM_SUBSCRIBE_BLOCK_CACHE_BYTES", "JETSTREAM_SUBSCRIBE_READ_BATCH", "JETSTREAM_SUBSCRIBE_SLOW_WINDOW", "JETSTREAM_SUBSCRIBE_SLOW_MIN_RATE", "JETSTREAM_CURSOR_BLOCK_INDEX_CACHE_SIZE", "JETSTREAM_COMPACTION_INTERVAL", "JETSTREAM_COMPACTION_TOMBSTONE_CAP", "JETSTREAM_COMPACTION_REWRITE_WORKERS", ) for index, key in enumerate(keys): env = clean_env() env.update( { "JETSTREAM_ADDR": "127.0.0.1:0", "JETSTREAM_DATA_DIR": str(root / f"invalid-{index}"), "JETSTREAM_RELAY_URL": "http://127.0.0.1:17990", key: "definitely-invalid", } ) result = subprocess.run( [BIN, "serve"], env=env, text=True, capture_output=True, timeout=3, check=False ) assert result.returncode != 0, f"{key} was accepted or ignored"
def main() -> None: exact_pinned_map_case() unknown_case() with tempfile.TemporaryDirectory(prefix="stream-env-") as raw: root = pathlib.Path(raw) root_and_inspect_case(root) serve_sources_case(root) cli_precedence_case(root) invalid_typed_sources_case(root) print("environment contract: unknowns sources precedence inspect typed-validation passed")
if __name__ == "__main__": main()