From a17ef7287dbe317f1314e83ecbb6100a1a340048 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Tue, 21 Jul 2026 19:01:00 -0500 Subject: [PATCH] match upstream retry controls --- README.md | 18 +++ docs/configuration-parity.md | 3 +- docs/upstream-harness.md | 11 ++ justfile | 8 ++ src/internal/bootstrap/retry.zig | 107 ++++++++++++++---- src/main.zig | 34 +++++- src/write_sample.zig | 65 +++++++++++ tests/retry_config_contract.py | 188 +++++++++++++++++++++++++++++++ 8 files changed, 406 insertions(+), 28 deletions(-) create mode 100644 tests/retry_config_contract.py diff --git a/README.md b/README.md index 2dc237e..33668db 100644 --- a/README.md +++ b/README.md @@ -50,6 +50,7 @@ just archive-contract # archive XRPCs + resident manifest + pinned Go cl just plan-config-contract # non-default planner limits/threshold, real process just cursor-lookback-contract # v1/v2 replay-window behavior, real processes just compaction-config-contract # tombstone cap/rewrite workers, physical JSS +just retry-config-contract # global/per-host retry admission + backoff just status-contract # durable host rows + public HTTP view, offline just oracle # lifecycle crash + RocksDB fault recovery matrix just powerloss-image # one-time Linux NBD/ext4 tool-image bootstrap @@ -105,6 +106,23 @@ unlimited sentinel. Rewrite workers bound concurrent sealed-segment rewrites; zero selects `min(CPU count, 8)`. `just compaction-config-contract` proves both values reach the physical compactor rather than stopping at argument parsing. +Failed-repository healing exposes upstream's complete control set: + +``` +--failed-repo-retry-interval=4h +--failed-repo-retry-workers=16 +--failed-repo-retry-host-workers=4 +--failed-repo-retry-max-delay=168h +``` + +The interval and maximum delay accept non-negative Go-style durations; zero +interval disables the background loop. Zero worker values select upstream's +16-global/4-per-host defaults, and zero maximum delay selects seven days. +Transient failures use upstream's capped exponential delay plus uniform +`[0, delay/2)` jitter. `just retry-config-contract` holds real `getRepo` +requests to measure both admission limits and then reopens RocksDB to verify +the durable retry schedule. + The strict power-loss tier is deliberately separate from ordinary unit and process tests. It cross-builds the production ReleaseSafe Linux binary, runs RocksDB and JSS on ext4 over a kernel NBD device, kills the block backend before diff --git a/docs/configuration-parity.md b/docs/configuration-parity.md index 9349908..7e2cd06 100644 --- a/docs/configuration-parity.md +++ b/docs/configuration-parity.md @@ -10,8 +10,7 @@ the same runtime mechanism and an offline receipt observes that behavior. | relay, PLC, data directory | behavioral equivalent | `--upstream`, `--relay-http`, `--plc`, and `--data-dir`; public naming and `JETSTREAM_*` environment sources remain open | | max/selected backfill repos, backfill workers | behavioral equivalent | lifecycle and differential receipts; public max-repo flag naming remains open | | backfill batch size, async flush workers, skip merge discovery | open | Stream does not yet expose or implement these upstream controls | -| failed-repo retry interval | partial | real durable retry loop exists; Go-duration syntax and upstream flag name remain open | -| failed-repo retry global workers, per-host workers, max delay | open | mechanisms have defaults, but the complete public configuration path is not wired | +| failed-repo retry interval, global workers, per-host workers, max delay | closed | `just retry-config-contract` measures explicit and zero/default concurrency through held real HTTP requests, then reopens RocksDB to verify one durable failure and the configured jitter/cap window per row | | expensive repo-action rate limiter | closed | real HTTP/status receipt covers enabled and disabled modes | | cursor lookback | closed | `just cursor-lookback-contract` covers default 36h, widened, v1/v2 old-cursor behavior, and zero/pure-live | | segment cache max age | closed | real HTTP receipt covers Go duration, ceil-to-seconds, and cache headers | diff --git a/docs/upstream-harness.md b/docs/upstream-harness.md index d0b94da..ace56d0 100644 --- a/docs/upstream-harness.md +++ b/docs/upstream-harness.md @@ -182,6 +182,17 @@ chunk and report a two-worker group. Both public archives must contain only the two surviving marker rows at durable watermark 4. The receipt therefore checks the actual fold and rewrite path, not merely accepted CLI flags. +`just retry-config-contract` seeds durable failed-repository rows, serves real +`getRepo` failures from a threaded loopback HTTP server, and runs four +ReleaseSafe Stream processes. Held requests prove explicit global concurrency +3, explicit per-host concurrency 2, and upstream's zero-value defaults of 16 +global and 4 per host. Each process is stopped after exactly one completed +pass; the receipt reopens the real RocksDB store and requires every row to have +one additional attempt plus a next-attempt delay inside the configured +1–1.5-second exponential-jitter window. This covers scheduling, admission, +failure persistence, Go-duration parsing, zero defaults, and max-delay clamping +without network access. + ## pinned differential oracle (2026-07-20) `just differential-oracle` refuses to run unless the upstream checkout is at diff --git a/justfile b/justfile index 5153f53..f580d81 100644 --- a/justfile +++ b/justfile @@ -222,6 +222,14 @@ compaction-config-contract: run_case bounded 1 1 2 run_case unlimited 0 2 1 +# Exercise retry admission and persisted backoff through real HTTP, RocksDB, +# the production retry engine, and the ReleaseSafe daemon. +retry-config-contract: + #!/usr/bin/env bash + set -euo pipefail + zig build -Doptimize=ReleaseSafe + python3 tests/retry_config_contract.py + # Seed the real RocksDB repo/host schema, reopen it through the production # binary, and exercise the public host/account views without network access. status-contract: diff --git a/src/internal/bootstrap/retry.zig b/src/internal/bootstrap/retry.zig index 805ded7..bdb4e7d 100644 --- a/src/internal/bootstrap/retry.zig +++ b/src/internal/bootstrap/retry.zig @@ -30,8 +30,10 @@ pub const Stats = struct { deferred: u64 = 0, }; -const default_workers: usize = 16; -const default_host_workers: usize = 4; +pub const default_interval_ns: u64 = 4 * std.time.ns_per_hour; +pub const default_workers: usize = 16; +pub const default_host_workers: usize = 4; +pub const default_max_delay_ns: u64 = 7 * 24 * std.time.ns_per_hour; const HostGate = struct { mu: Io.Mutex = .init, @@ -60,15 +62,18 @@ pub const Retry = struct { meta: *meta_store.Store, data_dir: []const u8, relay_http: []const u8, - /// seconds between passes; 0 disables (run() returns immediately) - interval_s: u64 = 4 * 3600, + /// nanoseconds between passes; 0 disables (run() returns immediately) + interval_ns: u64 = default_interval_ns, /// backoff ceiling (upstream DefaultFailedRepoRetryMaxDelay: 7d) - max_delay_s: u64 = 7 * 24 * 3600, + max_delay_ns: u64 = default_max_delay_ns, workers: usize = default_workers, host_workers: usize = default_host_workers, stop_flag: std.atomic.Value(bool) = .init(false), + stop_event: Io.Event = .unset, stats: ?*metrics.Stats = null, emit_failure_logs: bool = true, + jitter_ctx: ?*anyopaque = null, + jitter: ?*const fn (?*anyopaque, u64) u64 = null, pub fn init(allocator: Allocator, io: Io, archive: *archive_mod.Archive, meta: *meta_store.Store, data_dir: []const u8, relay_http: []const u8) Retry { return .{ @@ -85,18 +90,29 @@ pub const Retry = struct { pub fn requestStop(self: *Retry) void { self.stop_flag.store(true, .release); + self.stop_event.set(self.io); } /// timer loop; spawn via io.concurrent, stop via requestStop. a /// failed pass logs and retries next interval — never tears down. pub fn run(self: *Retry) void { - if (self.interval_s == 0) return; - var since: u64 = 0; + if (self.interval_ns == 0) return; + const interval: Io.Clock.Duration = .{ + .raw = Io.Duration.fromNanoseconds(@intCast(self.interval_ns)), + .clock = .awake, + }; while (!self.stop_flag.load(.acquire)) { - self.io.sleep(Io.Duration.fromSeconds(1), .awake) catch break; - since += 1; - if (since < self.interval_s) continue; - since = 0; + const deadline = Io.Clock.Timestamp.fromNow(self.io, interval); + while (Io.Clock.Timestamp.now(self.io, .awake).compare(.lt, deadline)) { + self.stop_event.waitTimeout(self.io, .{ .deadline = deadline }) catch |err| switch (err) { + // Event waits may wake spuriously. Preserve the configured + // interval by checking the absolute monotonic deadline. + error.Timeout => continue, + error.Canceled => return, + }; + if (self.stop_flag.load(.acquire)) return; + } + if (self.stop_flag.load(.acquire)) return; const stats = self.runPass() catch |err| { log.warn("failed-repo retry pass failed: {s} (retrying next interval)", .{@errorName(err)}); continue; @@ -143,7 +159,7 @@ pub const Retry = struct { var work: RetryWork = .{ .retry = self, .store = &store, .dids = dids, .gates = &gates }; var group: Io.Group = .init; - const worker_count = @min(@max(self.workers, 1), dids.len); + const worker_count = @min(self.effectiveWorkers(), dids.len); for (0..worker_count) |_| try group.concurrent(self.io, RetryWork.runWorker, .{&work}); group.await(self.io) catch |err| return err; if (work.fatal_error) |err| return err; @@ -173,12 +189,37 @@ pub const Retry = struct { try self.meta.putDurable(key, &bytes); } - /// exponential from the interval, capped at max_delay + /// Upstream's exponential delay plus [0, delay/2) jitter, capped at + /// max_delay. retry_count is one-based after recording this failure. fn backoffDelayUs(self: *Retry, retry_count: u32) u64 { - const base = @max(self.interval_s, 1); - const shift = @min(retry_count -| 1, 20); - const delay_s = @min(base << @intCast(shift), self.max_delay_s); - return delay_s * std.time.us_per_s; + const attempt = retry_count -| 1; + const max_delay_ns = self.effectiveMaxDelayNs(); + var delay_ns = max_delay_ns; + if (self.interval_ns > 0 and attempt < @clz(self.interval_ns)) { + const shifted = self.interval_ns << @intCast(attempt); + if (shifted < max_delay_ns) delay_ns = shifted; + } + const half = delay_ns / 2; + if (half > 0) delay_ns +|= self.jitterBelow(half); + return @min(delay_ns, max_delay_ns) / std.time.ns_per_us; + } + + fn jitterBelow(self: *Retry, upper: u64) u64 { + if (self.jitter) |jitter| return @min(jitter(self.jitter_ctx, upper), upper - 1); + var source: std.Random.IoSource = .{ .io = self.io }; + return source.interface().uintLessThan(u64, upper); + } + + fn effectiveWorkers(self: *const Retry) usize { + return if (self.workers == 0) default_workers else self.workers; + } + + fn effectiveHostWorkers(self: *const Retry) usize { + return if (self.host_workers == 0) default_host_workers else self.host_workers; + } + + fn effectiveMaxDelayNs(self: *const Retry) u64 { + return if (self.max_delay_ns == 0) default_max_delay_ns else self.max_delay_ns; } }; @@ -220,7 +261,7 @@ const RetryWork = struct { return; } const gate = self.gates.get(candidate_host) orelse return error.MissingHostGate; - gate.acquire(self.retry.io, @max(self.retry.host_workers, 1)); + gate.acquire(self.retry.io, self.retry.effectiveHostWorkers()); defer gate.release(self.retry.io); // Recheck after waiting: a sibling on this host may have received a @@ -247,7 +288,7 @@ const RetryWork = struct { } const count = state.retry_count + 1; const retry_delay_us = if (err == error.RateLimited and details.retry_after_s != null) - @min(details.retry_after_s.? * std.time.us_per_s, self.retry.max_delay_s * std.time.us_per_s) + @min(details.retry_after_s.? *| std.time.ns_per_s, self.retry.effectiveMaxDelayNs()) / std.time.ns_per_us else self.retry.backoffDelayUs(count); const next_attempt = now_us + @as(i64, @intCast(retry_delay_us)); @@ -319,13 +360,39 @@ test "backoff delay doubles from the interval and caps" { .meta = undefined, .data_dir = "", .relay_http = "", - .interval_s = 4 * 3600, + .interval_ns = 4 * std.time.ns_per_hour, + .jitter = struct { + fn zero(_: ?*anyopaque, _: u64) u64 { + return 0; + } + }.zero, }; try testing.expectEqual(@as(u64, 4 * 3600 * std.time.us_per_s), r.backoffDelayUs(1)); try testing.expectEqual(@as(u64, 8 * 3600 * std.time.us_per_s), r.backoffDelayUs(2)); try testing.expectEqual(@as(u64, 7 * 24 * 3600 * std.time.us_per_s), r.backoffDelayUs(30)); } +test "backoff adds upstream half-window jitter and clamps to max delay" { + const max_jitter = struct { + fn call(_: ?*anyopaque, upper: u64) u64 { + return upper - 1; + } + }.call; + var r: Retry = .{ + .allocator = testing.allocator, + .io = undefined, + .archive = undefined, + .meta = undefined, + .data_dir = "", + .relay_http = "", + .interval_ns = std.time.ns_per_s, + .max_delay_ns = 1500 * std.time.ns_per_ms, + .jitter = max_jitter, + }; + try testing.expectEqual(@as(u64, 1_499_999), r.backoffDelayUs(1)); + try testing.expectEqual(@as(u64, 1_500_000), r.backoffDelayUs(2)); +} + test "runPass with no failed repos is a no-op" { var threaded: Io.Threaded = .init(testing.allocator, .{}); defer threaded.deinit(); diff --git a/src/main.zig b/src/main.zig index b54525e..a6ecbae 100644 --- a/src/main.zig +++ b/src/main.zig @@ -213,8 +213,10 @@ pub fn main(init: std.process.Init.Minimal) !void { var compaction_interval_s: u64 = 4 * 3600; var compaction_tombstone_cap: usize = compact_pass.default_tombstone_cap; var compaction_rewrite_workers: usize = 0; - // upstream DefaultFailedRepoRetryInterval: 4h; 0 disables the loop - var retry_interval_s: u64 = 4 * 3600; + var retry_interval_ns: u64 = retry_mod.default_interval_ns; + var retry_workers: usize = retry_mod.default_workers; + var retry_host_workers: usize = retry_mod.default_host_workers; + var retry_max_delay_ns: u64 = retry_mod.default_max_delay_ns; var timestamp_import_token: []const u8 = ""; var timestamp_import_dir_arg: ?[]const u8 = null; var repo_action_rate_limits = true; @@ -257,7 +259,19 @@ pub fn main(init: std.process.Init.Minimal) !void { } else if (std.mem.startsWith(u8, arg, "--compaction-rewrite-workers=")) { compaction_rewrite_workers = try std.fmt.parseInt(usize, arg["--compaction-rewrite-workers=".len..], 10); } else if (std.mem.startsWith(u8, arg, "--retry-interval=")) { - retry_interval_s = try std.fmt.parseInt(u64, arg["--retry-interval=".len..], 10); + const seconds = try std.fmt.parseInt(u64, arg["--retry-interval=".len..], 10); + retry_interval_ns = try std.math.mul(u64, seconds, std.time.ns_per_s); + } else if (std.mem.startsWith(u8, arg, "--failed-repo-retry-interval=")) { + retry_interval_ns = try parseDurationNanoseconds(arg["--failed-repo-retry-interval=".len..]); + } else if (std.mem.startsWith(u8, arg, "--failed-repo-retry-workers=")) { + const configured = try std.fmt.parseInt(usize, arg["--failed-repo-retry-workers=".len..], 10); + retry_workers = if (configured == 0) retry_mod.default_workers else configured; + } else if (std.mem.startsWith(u8, arg, "--failed-repo-retry-host-workers=")) { + const configured = try std.fmt.parseInt(usize, arg["--failed-repo-retry-host-workers=".len..], 10); + retry_host_workers = if (configured == 0) retry_mod.default_host_workers else configured; + } else if (std.mem.startsWith(u8, arg, "--failed-repo-retry-max-delay=")) { + const configured = try parseDurationNanoseconds(arg["--failed-repo-retry-max-delay=".len..]); + retry_max_delay_ns = if (configured == 0) retry_mod.default_max_delay_ns else configured; } else if (std.mem.startsWith(u8, arg, "--segment-cache-max-age=")) { segment_cache_max_age_s = try parseDurationCeilSeconds(arg["--segment-cache-max-age=".len..]); } else if (std.mem.startsWith(u8, arg, "--cursor-lookback=")) { @@ -347,7 +361,7 @@ pub fn main(init: std.process.Init.Minimal) !void { .stats = &stats, .bootstrap_enabled = do_backfill or backfill_max_repos > 0 or backfill_repos.len > 0, .compaction_enabled = compaction_interval_s > 0, - .retry_enabled = retry_interval_s > 0, + .retry_enabled = retry_interval_ns > 0, .plc_url = plc_url, .data_dir = data_dir, .repo_action_rate_limits = repo_action_rate_limits, @@ -538,9 +552,17 @@ pub fn main(init: std.process.Init.Minimal) !void { // archive writer mutex beside the live consumer. var retry = retry_mod.Retry.init(allocator, io, &archive, &meta, data_dir, relay_http); defer retry.deinit(); - retry.interval_s = retry_interval_s; + retry.interval_ns = retry_interval_ns; + retry.workers = retry_workers; + retry.host_workers = retry_host_workers; + retry.max_delay_ns = retry_max_delay_ns; retry.stats = &stats; - if (retry_interval_s != 0) log.info("failed-repo retry loop on (interval {d}s)", .{retry_interval_s}); + if (retry_interval_ns != 0) log.info("failed-repo retry loop on (interval {d}ns, workers {d}, host workers {d}, max delay {d}ns)", .{ + retry_interval_ns, + retry_workers, + retry_host_workers, + retry_max_delay_ns, + }); var retry_future = try io.concurrent(retry_mod.Retry.run, .{&retry}); defer { retry.requestStop(); diff --git a/src/write_sample.zig b/src/write_sample.zig index 9fd14f8..6e43c1a 100644 --- a/src/write_sample.zig +++ b/src/write_sample.zig @@ -101,6 +101,71 @@ pub fn main(init: std.process.Init.Minimal) !void { return; } + // --retry-status : seed durable failed rows + // for the production retry-control receipt. + if (std.mem.eql(u8, path, "--retry-status")) { + const dir = args.next() orelse return error.MissingPath; + const mode = args.next() orelse return error.MissingMode; + const count = try std.fmt.parseInt(usize, args.next() orelse return error.MissingCount, 10); + if (!std.mem.eql(u8, mode, "same") and !std.mem.eql(u8, mode, "unique")) return error.InvalidMode; + const meta_store = @import("internal/meta_store.zig"); + const repo_store = @import("internal/bootstrap/repo_store.zig"); + var meta = try meta_store.Store.open(allocator, dir); + defer meta.deinit(); + var store = try repo_store.Store.init(allocator, io, &meta); + defer store.deinit(); + for (0..count) |i| { + var did_buf: [96]u8 = undefined; + const did = try std.fmt.bufPrint(&did_buf, "did:plc:retryconfig{d}", .{i}); + var host_buf: [96]u8 = undefined; + const host = if (std.mem.eql(u8, mode, "same")) + "same.retry.invalid" + else + try std.fmt.bufPrint(&host_buf, "host{d}.retry.invalid", .{i}); + try store.put(did, .{ + .status = .failed, + .active = true, + .host = host, + .last_error = "seeded retry failure", + .error_class = .http_5xx, + .attempts = 1, + }); + } + return; + } + + // --retry-assert : reopen + // RocksDB after the process exits and prove one persisted retry failure + // per seed plus the configured jitter/cap window. + if (std.mem.eql(u8, path, "--retry-assert")) { + const dir = args.next() orelse return error.MissingPath; + const count = try std.fmt.parseInt(usize, args.next() orelse return error.MissingCount, 10); + const min_delay_us = try std.fmt.parseInt(i64, args.next() orelse return error.MissingDelay, 10); + const max_delay_us = try std.fmt.parseInt(i64, args.next() orelse return error.MissingDelay, 10); + const meta_store = @import("internal/meta_store.zig"); + const repo_store = @import("internal/bootstrap/repo_store.zig"); + var meta = try meta_store.Store.open(allocator, dir); + defer meta.deinit(); + var store = try repo_store.Store.init(allocator, io, &meta); + defer store.deinit(); + var saw_jitter = false; + for (0..count) |i| { + var did_buf: [96]u8 = undefined; + const did = try std.fmt.bufPrint(&did_buf, "did:plc:retryconfig{d}", .{i}); + const state = (try store.get(allocator, did)) orelse return error.MissingRetryState; + defer store.freeState(allocator, state); + if (state.status != .failed or state.retry_count != 1 or state.attempts != 2) + return error.InvalidRetryState; + const delay_us = state.next_attempt_us - state.last_attempt_us; + if (delay_us < min_delay_us or delay_us > max_delay_us) + return error.InvalidRetryDelay; + saw_jitter = saw_jitter or delay_us > min_delay_us; + } + if (count > 1 and !saw_jitter) return error.MissingRetryJitter; + std.debug.print("retry state: PASS ({d} rows, delay {d}..{d}us)\n", .{ count, min_delay_us, max_delay_us }); + return; + } + // --archive : drive the Archive instead (rotation + seal + active) if (std.mem.eql(u8, path, "--archive")) { const dir = args.next() orelse return error.MissingPath; diff --git a/tests/retry_config_contract.py b/tests/retry_config_contract.py new file mode 100644 index 0000000..3f53ab4 --- /dev/null +++ b/tests/retry_config_contract.py @@ -0,0 +1,188 @@ +"""Real-process retry worker, host gate, zero-default, and backoff receipt.""" + +import http.server +import os +import pathlib +import signal +import subprocess +import tempfile +import threading +import time +import urllib.error +import urllib.request + + +ROOT = pathlib.Path(__file__).resolve().parents[1] +RELAY_PORT = int(os.environ.get("STREAM_RETRY_CONFIG_RELAY_PORT", "6025")) +STREAM_PORT = int(os.environ.get("STREAM_RETRY_CONFIG_PORT", "6026")) + + +class State: + lock = threading.Lock() + active = 0 + peak = 0 + requests = 0 + first_at = None + + @classmethod + def reset(cls): + with cls.lock: + assert cls.active == 0 + cls.peak = 0 + cls.requests = 0 + cls.first_at = None + + +class Handler(http.server.BaseHTTPRequestHandler): + def do_GET(self): + if self.path.startswith("/xrpc/com.atproto.sync.getRepo?"): + with State.lock: + State.active += 1 + State.requests += 1 + if State.first_at is None: + State.first_at = time.monotonic() + State.peak = max(State.peak, State.active) + try: + time.sleep(0.25) + body = b'{"error":"ServiceUnavailable","message":"retry receipt"}' + self.send_response(503) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(body))) + self.end_headers() + self.wfile.write(body) + finally: + with State.lock: + State.active -= 1 + return + self.send_response(503) + self.send_header("Content-Length", "0") + self.end_headers() + + def log_message(self, *_): + pass + + +def run(*args): + subprocess.run(args, cwd=ROOT, check=True) + + +def run_case(name, mode, count, workers, host_workers, expected_peak): + State.reset() + with tempfile.TemporaryDirectory(prefix=f"stream-retry-{name}-") as raw: + data = pathlib.Path(raw) / "data" + log_path = pathlib.Path(raw) / "stream.log" + data.mkdir() + run("zig", "build", "write-sample", "--", "--retry-status", str(data), mode, str(count)) + args = [ + str(ROOT / "zig-out/bin/stream"), + f"--port={STREAM_PORT}", + f"--data-dir={data}", + f"--upstream=ws://127.0.0.1:{RELAY_PORT}", + f"--relay-http=http://127.0.0.1:{RELAY_PORT}", + f"--plc=http://127.0.0.1:{RELAY_PORT}", + "--compaction-interval=0", + "--no-verify", + "--failed-repo-retry-interval=1s", + f"--failed-repo-retry-workers={workers}", + f"--failed-repo-retry-host-workers={host_workers}", + "--failed-repo-retry-max-delay=1500ms", + ] + with log_path.open("wb") as log: + started_at = time.monotonic() + process = subprocess.Popen(args, cwd=ROOT, stdout=log, stderr=subprocess.STDOUT) + try: + deadline = time.monotonic() + 20 + while time.monotonic() < deadline: + if process.poll() is not None: + raise AssertionError(log_path.read_text(errors="replace")) + text = log_path.read_text(errors="replace") + if f"failed-repo retry pass: {count} candidates" in text: + break + time.sleep(0.05) + else: + raise AssertionError(log_path.read_text(errors="replace")) + with State.lock: + assert State.requests == count, (name, State.requests, text) + assert State.peak == expected_peak, (name, State.peak, text) + assert State.first_at is not None + assert State.first_at - started_at >= 0.9, (name, State.first_at - started_at) + finally: + if process.poll() is None: + process.send_signal(signal.SIGTERM) + process.wait(timeout=10) + run( + "zig", + "build", + "write-sample", + "--", + "--retry-assert", + str(data), + str(count), + "1000000", + "1500000", + ) + print(f"retry config: PASS ({name}, peak={expected_peak})") + + +def run_disabled(): + State.reset() + with tempfile.TemporaryDirectory(prefix="stream-retry-disabled-") as raw: + data = pathlib.Path(raw) / "data" + log_path = pathlib.Path(raw) / "stream.log" + data.mkdir() + run("zig", "build", "write-sample", "--", "--retry-status", str(data), "unique", "6") + args = [ + str(ROOT / "zig-out/bin/stream"), + f"--port={STREAM_PORT}", + f"--data-dir={data}", + f"--upstream=ws://127.0.0.1:{RELAY_PORT}", + f"--relay-http=http://127.0.0.1:{RELAY_PORT}", + f"--plc=http://127.0.0.1:{RELAY_PORT}", + "--compaction-interval=0", + "--no-verify", + "--failed-repo-retry-interval=0", + ] + with log_path.open("wb") as log: + process = subprocess.Popen(args, cwd=ROOT, stdout=log, stderr=subprocess.STDOUT) + try: + deadline = time.monotonic() + 10 + while time.monotonic() < deadline: + if process.poll() is not None: + raise AssertionError(log_path.read_text(errors="replace")) + try: + urllib.request.urlopen(f"http://127.0.0.1:{STREAM_PORT}/healthz", timeout=0.2) + break + except OSError: + time.sleep(0.05) + else: + raise AssertionError("disabled retry process never became healthy") + time.sleep(1.2) + with State.lock: + assert State.requests == 0, State.requests + finally: + if process.poll() is None: + process.send_signal(signal.SIGTERM) + process.wait(timeout=10) + print("retry config: PASS (disabled, no getRepo requests)") + + +server = http.server.ThreadingHTTPServer(("127.0.0.1", RELAY_PORT), Handler) +thread = threading.Thread(target=server.serve_forever, daemon=True) +thread.start() +try: + for _ in range(100): + try: + urllib.request.urlopen(f"http://127.0.0.1:{RELAY_PORT}/ready", timeout=0.1) + except urllib.error.HTTPError: + break + except OSError: + time.sleep(0.01) + run_case("explicit-global", "unique", 6, 3, 1, 3) + run_case("explicit-host", "same", 6, 4, 2, 2) + run_case("default-global", "unique", 20, 0, 1, 16) + run_case("default-host", "same", 8, 8, 0, 4) + run_disabled() +finally: + server.shutdown() + server.server_close() + thread.join(timeout=2) -- 2.51.2