#!/usr/bin/env python3 """Offline process receipt for pinned Jetstream V2 environment sources.""" from __future__ import annotations import json import os import pathlib import re import signal import socket import subprocess import tempfile import time import urllib.error import 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()