diff --git a/README.md b/README.md index 9581702..ac9f1fc 100644 --- a/README.md +++ b/README.md @@ -132,8 +132,11 @@ delay the others, drains cooperative subscribers until the client budget, and shuts down public and debug HTTP traffic from the same instant while the accepted ingest pipeline reaches its durable boundary. `just shutdown-contract` exercises the real ReleaseSafe process offline with both cooperative and silent -peers; the transport suite also saturates one peer's real socket send buffer and -proves another peer still receives code 1001 before the shared deadline. +peers. It also terminates bootstrap while live capture is connected and +`listRepos` is hung, then reopens the same data directory and repeats; active +backfill downloads are canceled while the archive flushes its recoverable +prefix. The transport suite also saturates one peer's real socket send buffer +and proves another peer still receives code 1001 before the shared deadline. As upstream does, a non-positive HTTP shutdown timeout selects 30 seconds, whereas a non-positive client-drain timeout is already expired: teardown is immediate and delivery of the best-effort close frame is not guaranteed. diff --git a/docs/bootstrap-semantic-parity.md b/docs/bootstrap-semantic-parity.md index 0e1e17f..fc9535d 100644 --- a/docs/bootstrap-semantic-parity.md +++ b/docs/bootstrap-semantic-parity.md @@ -27,7 +27,7 @@ every experiment. | Steady retry host semantics | persist actual post-redirect PDS host; 16 global / 4 per-host workers; 429 parks only that host | redirected `Request.uri` is normalized and persisted; worker gates and durable host-specific parks match upstream limits; a real redirected 429 fixture proves the final host replaces the predicted host and only the final host is parked | closed | | Retry diagnostics | persist attempts, last error/class, host, retry count, next attempt, and durable per-host aggregates | RocksDB row v5 persists every repo field and backward-decodes v2/v3/v4; retry, terminal resync, active-flip, and crash-recovery RMW paths preserve identity/latest fields, while success clears stale repo failure state. The same synced transition batch maintains normalized `host/` current counts plus cumulative error-class counts and five bounded recent samples; host moves remove the repo from the old bucket, successful retry retains diagnostic history, and a one-time migration rebuilds current aggregates without inventing history. The public host view reopens those durable rows through the production binary and exercises both upstream sort modes plus restart, text/HTML, filtering, and escaping offline. | closed | | Deadline coverage | getRepo streaming read + CAR load/completeness check, excluding handler/durability | same `prepared_fetch.zig` future covers both and returns a prepared repo | closed | -| Oracle/adversity suite | truncation, hanging server, terminal errors, crash/restart, active flips, and power loss at canonical lifecycle boundaries | focused offline fixtures cover truncation, hanging fetches, terminal errors, retry/429 interaction, final-host parking, and hostile embedded DIDs; the ordinary crash matrix survived all 8 local lifecycle seams and rejects shutdown OOM/fatal-handler evidence. The strict Linux oracle runs production RocksDB/JSS on ext4 over calibrated NBD storage, kills storage before the process, reconstructs only FLUSH/FUA-acknowledged bytes, and recovered all 8 canonical lifecycle schedules plus the two-crash restart-after-cleanup guard through public archive and V2 replay. | closed; rerun both oracles against the exact candidate before deployment | +| Oracle/adversity suite | truncation, hanging server, terminal errors, crash/restart, active flips, and power loss at canonical lifecycle boundaries | focused offline fixtures cover truncation, hanging fetches, terminal errors, retry/429 interaction, final-host parking, and hostile embedded DIDs. The shutdown receipt now drives the production ReleaseSafe process into bootstrap with connected live capture and a hung `listRepos`, sends SIGTERM, requires prompt exit without SIGKILL, and reopens the same data directory to repeat the proof. The ordinary crash matrix survived all 8 local lifecycle seams and rejects shutdown OOM/fatal-handler evidence. The strict Linux oracle runs production RocksDB/JSS on ext4 over calibrated NBD storage, kills storage before the process, reconstructs only FLUSH/FUA-acknowledged bytes, and recovered all 8 canonical lifecycle schedules plus the two-crash restart-after-cleanup guard through public archive and V2 replay. | closed; rerun both oracles against the exact candidate before deployment | The deleted July experiment is not evidence for closing these rows. A fresh experiment may be provisioned only after every blocker is closed and the gate diff --git a/docs/configuration-parity.md b/docs/configuration-parity.md index 7ec1f54..6225cfa 100644 --- a/docs/configuration-parity.md +++ b/docs/configuration-parity.md @@ -26,7 +26,7 @@ the same runtime mechanism and an offline receipt observes that behavior. | compaction tombstone cap and rewrite workers | closed | `just compaction-config-contract` observes chunk count, actual rewrite group size, physical survivors, and durable watermark; zero sentinels match upstream | | timestamp import token and directory | closed | real HTTP/import, restart, and strict power-loss receipts | | public/debug bind addresses | closed | `serve --addr` binds the protocol-only public listener, an empty `--debug-addr` binds nothing, and a configured debug listener exclusively serves health/readiness/metrics; `just listener-contract` proves route isolation, public WebSocket upgrade, rejected debug upgrade, disabled-debug refusal, listener readiness, and the historical combined-port compatibility mode through real offline processes | -| shutdown/client drain timeouts | closed | canonical signed Go-duration flags, 5s/10s defaults, net/http's non-positive 30s sentinel, the client drain's already-expired non-positive sentinel (immediate teardown with a best-effort close frame), listener-admission gating, concurrent public/debug/pipeline shutdown, close code 1001 for positive budgets, cooperative early exit, and one shared silent-client deadline are covered by unit and real ReleaseSafe process receipts. Per-client close hooks now fan out concurrently with borrowed-handler lifetime protection; a real-socket transport regression saturates the first peer's send buffer, proves a later peer receives 1001 before the deadline, and proves the blocked writer is interrupted at the shared deadline. | +| shutdown/client drain timeouts | closed | canonical signed Go-duration flags, 5s/10s defaults, net/http's non-positive 30s sentinel, the client drain's already-expired non-positive sentinel (immediate teardown with a best-effort close frame), listener-admission gating, concurrent public/debug/pipeline shutdown, close code 1001 for positive budgets, cooperative early exit, and one shared silent-client deadline are covered by unit and real ReleaseSafe process receipts. Bootstrap, merge, and discovery run behind a cancellable lifecycle boundary: a real process establishes live capture, hangs `listRepos`, exits promptly on SIGTERM, then reopens the same durable directory and proves the contract again. Per-client close hooks fan out concurrently with borrowed-handler lifetime protection; a real-socket transport regression saturates the first peer's send buffer, proves a later peer receives 1001 before the deadline, and proves the blocked writer is interrupted at the shared deadline. | | log level/format | closed | `--log-level` retains debug calls in ReleaseSafe and applies the upstream case-insensitive aliases plus info default to every `std.log` producer; `--log-format` defaults to newline-delimited JSON and selects canonical text fields. `just logging-contract` proves default JSON, warn filtering, text/debug output, persistent root-flag placement, and fail-closed invalid values through real offline processes. | | OpenTelemetry | closed | provider bootstrap is real and empirically covered: endpoint activation/precedence, HTTP/protobuf path rules, headers, gzip, insecure mode, request timeout, batch queue/delay/size/export timeout, upstream retry/backoff/throttle behavior, the six SDK samplers, process plus service resource attributes, global provider installation, W3C TraceContext+Baggage propagation (including no-exporter mode), bounded shutdown export, custom CA trust, and client certificate authentication all run through pinned `otel-zig`. Every pinned production span matches upstream names, scopes, parentage, attributes, and success/error finalization, including `xrpcapi/getBlock`; real offline OTLP collectors exercise the production sync/async archive, success/failure XRPC boundaries, native custom-CA HTTPS, rejection without a client certificate, and a successful verified mTLS request. The production image carries the libcurl runtime selected only by a complete client certificate/key pair. See `opentelemetry.md`. | | `JETSTREAM_*` environment sources and unknown-variable rejection | closed | the exact pinned flag-source map applies root, serve, and `inspect-all` values with CLI-over-environment precedence; overridden invalid values are never parsed. Unknown names are sorted/deduplicated and rejected before command execution while upstream's Kubernetes `JETSTREAM_APP_*` and its `JETSTREAM_ORACLE_*`/`JETSTREAM_SIM_*` test namespaces are exempt. `just environment-contract` proves real ReleaseSafe root/inspect/serve processes, representative string behavior, every typed source, timestamp-import authentication, and precedence offline; the multiseed oracle proves its inherited upstream seed reaches the crash campaign without leaking into Stream configuration. | diff --git a/src/internal/bootstrap/engine.zig b/src/internal/bootstrap/engine.zig index a104a6f..962be0d 100644 --- a/src/internal/bootstrap/engine.zig +++ b/src/internal/bootstrap/engine.zig @@ -276,9 +276,12 @@ const Work = struct { try self.group.concurrent(self.io, workerLoop, .{self}); } - fn shutdown(self: *Work) void { + fn shutdown(self: *Work, completed: bool) void { self.queue.close(self.io); - self.group.await(self.io) catch {}; + if (completed) + self.group.await(self.io) catch {} + else + self.group.cancel(self.io); self.publishGauges(); } @@ -554,7 +557,8 @@ pub fn runTraced( var work = try Work.init(allocator, io, work_config, archive, trace.context()); defer work.deinit(); try work.start(); - defer work.shutdown(); + var work_completed = false; + defer work.shutdown(work_completed); var stats: Stats = .{}; if (config.metrics) |m| m.backfill_progress_completed.store(0, .monotonic); @@ -695,6 +699,7 @@ pub fn runTraced( try store.flush(); } work.publishGauges(); + work_completed = true; return stats; } diff --git a/src/main.zig b/src/main.zig index e258d31..c7ca0a1 100644 --- a/src/main.zig +++ b/src/main.zig @@ -53,6 +53,20 @@ const default_stack_size: usize = 8 * 1024 * 1024; var shutdown_requested: std.atomic.Value(bool) = .init(false); +const LifecycleRun = struct { + allocator: std.mem.Allocator, + io: Io, + options: lifecycle.Options, + archive: *archive_mod.Archive, + meta: *meta_store.Store, + cursor_store: *cursor_mod.Store, + verifier: ?*verify.Verifier, + stats: *metrics.Stats, + trace_context: observability.Context, + done: std.atomic.Value(bool) = .init(false), + result: ?anyerror = null, +}; + fn signalHandler(_: std.posix.SIG) callconv(.c) void { // Atomic storage is the only work performed in signal context. The main // task observes it and owns cancellation, pipeline drain, durable flush, @@ -1121,23 +1135,59 @@ pub fn main(init: std.process.Init.Minimal) !void { // or --max-backfill-repos (debug) start one; serving is 503-gated // until the lifecycle commits steady_state. if (mid_lifecycle or bootstrap_requested) { - lifecycle.runTraced(allocator, io, .{ - .data_dir = data_dir, - .upstream = upstream, - .relay_http = relay_http, - .plc_url = plc_url, - .max_repos = backfill_max_repos, - .selected_repos = backfill_repos, - .workers = backfill_workers, - .batch_size = backfill_batch_size, - .async_flush_workers = backfill_async_flush_workers, - .max_inflight_bytes = backfill_max_inflight_bytes, - .skip_merge_discovery = skip_merge_discovery, - .compaction_enabled = compaction_interval_ns != 0, - }, &archive, &meta, &cursor_store, if (verifier) |*v| v else null, &stats, orchestrator_trace.context()) catch |err| { + var lifecycle_run: LifecycleRun = .{ + .allocator = allocator, + .io = io, + .options = .{ + .data_dir = data_dir, + .upstream = upstream, + .relay_http = relay_http, + .plc_url = plc_url, + .max_repos = backfill_max_repos, + .selected_repos = backfill_repos, + .workers = backfill_workers, + .batch_size = backfill_batch_size, + .async_flush_workers = backfill_async_flush_workers, + .max_inflight_bytes = backfill_max_inflight_bytes, + .skip_merge_discovery = skip_merge_discovery, + .compaction_enabled = compaction_interval_ns != 0, + }, + .archive = &archive, + .meta = &meta, + .cursor_store = &cursor_store, + .verifier = if (verifier) |*v| v else null, + .stats = &stats, + .trace_context = orchestrator_trace.context(), + }; + var lifecycle_future = try io.concurrent(runLifecycle, .{&lifecycle_run}); + while (!shutdown_requested.load(.acquire) and !lifecycle_run.done.load(.acquire)) + Io.sleep(io, Io.Duration.fromMilliseconds(25), .awake) catch |err| switch (err) { + error.Canceled => {}, + }; + if (shutdown_requested.load(.acquire)) { + log.info("shutdown requested during bootstrap; canceling lifecycle and flushing its durable prefix", .{}); + hub.draining.store(true, .release); + hub.listeners_ready.store(false, .release); + + var listener_shutdowns: Io.Group = .init; + errdefer listener_shutdowns.cancel(io); + try listener_shutdowns.concurrent(io, cancelWsServer, .{ &server_future, io }); + if (debug_future) |*future| + try listener_shutdowns.concurrent(io, cancelWsServer, .{ future, io }); + + lifecycle_future.cancel(io); + archive.close(); + archive_closed = true; + try listener_shutdowns.await(io); + orchestrator_trace.succeed(); + try shutdownTracing(&tracing, shutdown_timeout); + return; + } + lifecycle_future.await(io); + if (lifecycle_run.result) |err| { logPersistenceFailure(data_dir, err); return @errorCast(err); - }; + } hub.serving.store(true, .release); } if (cursor == null) { @@ -1343,6 +1393,23 @@ fn runConsumer(consumer: *ingest.Consumer, done: *std.atomic.Value(bool)) !void try consumer.run(); } +fn runLifecycle(run: *LifecycleRun) void { + defer run.done.store(true, .release); + lifecycle.runTraced( + run.allocator, + run.io, + run.options, + run.archive, + run.meta, + run.cursor_store, + run.verifier, + run.stats, + run.trace_context, + ) catch |err| { + run.result = err; + }; +} + fn runWsServer(server: *websocket.Server(server_mod.Handler), listener: *Io.net.Server, context: *server_mod.ListenerContext) void { server.runIo(listener, context); } diff --git a/tests/shutdown_contract.py b/tests/shutdown_contract.py index e619ba6..60c9f7b 100644 --- a/tests/shutdown_contract.py +++ b/tests/shutdown_contract.py @@ -5,12 +5,15 @@ from __future__ import annotations import atexit import base64 +import hashlib +import http.server import os import signal import socket import struct import subprocess import tempfile +import threading import time import urllib.request @@ -33,16 +36,18 @@ def terminate(process: subprocess.Popen[str]) -> None: process.wait(timeout=3) -def start(root: str, port: int, *timeouts: str) -> subprocess.Popen[str]: +def start( + root: str, port: int, *timeouts: str, upstream_port: int = 17997 +) -> subprocess.Popen[str]: process = subprocess.Popen( [ BIN, "serve", f"--addr=127.0.0.1:{port}", f"--data-dir={root}", - "--upstream=ws://127.0.0.1:17997", - "--relay-http=http://127.0.0.1:17997", - "--plc-url=http://127.0.0.1:17997", + f"--upstream=ws://127.0.0.1:{upstream_port}", + f"--relay-http=http://127.0.0.1:{upstream_port}", + f"--plc-url=http://127.0.0.1:{upstream_port}", "--compaction-interval=0", "--retry-interval=0", "--no-verify", @@ -56,6 +61,69 @@ def start(root: str, port: int, *timeouts: str) -> subprocess.Popen[str]: return process +class BootstrapFixture: + def __init__(self) -> None: + self.capture_connected = threading.Event() + self.backfill_started = threading.Event() + self.stop = threading.Event() + self.requests: list[str] = [] + fixture = self + + class Handler(http.server.BaseHTTPRequestHandler): + protocol_version = "HTTP/1.1" + + def do_GET(self) -> None: + fixture.requests.append(self.path) + if self.path.startswith("/xrpc/com.atproto.sync.subscribeRepos"): + key = self.headers["Sec-WebSocket-Key"] + assert key is not None + accept = base64.b64encode( + hashlib.sha1( + (key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11").encode() + ).digest() + ).decode() + self.send_response(101) + self.send_header("Upgrade", "websocket") + self.send_header("Connection", "Upgrade") + self.send_header("Sec-WebSocket-Accept", accept) + self.end_headers() + fixture.capture_connected.set() + self.connection.settimeout(0.1) + while not fixture.stop.is_set(): + try: + if not self.connection.recv(4096): + break + except TimeoutError: + continue + except OSError: + break + return + if self.path.startswith("/xrpc/com.atproto.sync.listRepos"): + fixture.backfill_started.set() + while not fixture.stop.wait(0.1): + pass + return + self.send_error(404) + + def log_message(self, *_: object) -> None: + pass + + self.server = http.server.ThreadingHTTPServer(("127.0.0.1", 0), Handler) + self.server.daemon_threads = True + self.thread = threading.Thread(target=self.server.serve_forever, daemon=True) + self.thread.start() + + @property + def port(self) -> int: + return self.server.server_address[1] + + def close(self) -> None: + self.stop.set() + self.server.shutdown() + self.server.server_close() + self.thread.join(timeout=2) + + def wait_ready(process: subprocess.Popen[str], port: int) -> None: deadline = time.monotonic() + 10 while time.monotonic() < deadline: @@ -227,16 +295,55 @@ def nonpositive_case(root: str) -> float: return elapsed +def bootstrap_case(root: str) -> tuple[float, float]: + fixture = BootstrapFixture() + try: + elapsed: list[float] = [] + for attempt in range(2): + port = unused_port() + process = start( + root, + port, + "--shutdown-timeout=2s", + "--client-drain-timeout=200ms", + "--backfill", + "--backfill-workers=1", + "--backfill-async-flush-workers=0", + upstream_port=fixture.port, + ) + wait_ready(process, port) + assert fixture.capture_connected.wait(3), "bootstrap capture did not connect" + assert fixture.backfill_started.wait(3), ( + f"bootstrap listRepos did not start; requests={fixture.requests!r}" + ) + started = time.monotonic() + process.send_signal(signal.SIGTERM) + wait_clean(process, 3) + duration = time.monotonic() - started + assert duration < 1.5, ( + f"bootstrap attempt {attempt + 1} waited for the hanging " + f"upstream request: {duration:.3f}s" + ) + elapsed.append(duration) + fixture.capture_connected.clear() + fixture.backfill_started.clear() + return elapsed[0], elapsed[1] + finally: + fixture.close() + + def main() -> None: with tempfile.TemporaryDirectory() as root: cooperative = cooperative_case(f"{root}/cooperative") silent = silent_case(f"{root}/silent") nonpositive = nonpositive_case(f"{root}/nonpositive") + bootstrap, bootstrap_restart = bootstrap_case(f"{root}/bootstrap") print( "shutdown contract: " f"cooperative={cooperative:.3f}s silent={silent:.3f}s " f"nonpositive={nonpositive:.3f}s positive-close=1001 " - "nonpositive-close=best-effort shared-deadline=200ms" + f"bootstrap={bootstrap:.3f}s bootstrap-restart={bootstrap_restart:.3f}s " + "nonpositive-close=best-effort shared-deadline=200ms durable-restart=ok" )