From 39a0886acbacf49c7f621f709d8f5def1ad33467 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sun, 12 Apr 2026 11:25:53 -0500 Subject: [PATCH] revert host retention changes (143605c..e5a2a14) reverts all four host retention commits that caused production restart loops. the interaction between reconciliation, dormant logic, startup jitter, and cold-start ramp with ~2,800 hosts was not testable via unit tests and each deploy regressed relay-eval coverage (alternating 0%/97% from kubelet kills). back to 80eca78 behavior: exhausted hosts stop after 15 failures, cron handles re-discovery. stable baseline for 24+ hours at 97-99%. the feature needs a local test harness that validates startup ramp behavior against a realistic host table before any production deploy. Co-Authored-By: Claude Opus 4.6 (1M context) --- docs/design.md | 15 -------- src/api/xrpc.zig | 2 +- src/host_ops.zig | 8 +--- src/slurper.zig | 94 +++------------------------------------------- src/subscriber.zig | 25 ++++-------- 5 files changed, 16 insertions(+), 128 deletions(-) diff --git a/docs/design.md b/docs/design.md index 5ab30c1..612db8f 100644 --- a/docs/design.md +++ b/docs/design.md @@ -54,7 +54,6 @@ additionally: | flush thread | 1 | 8 MB | batched fsync of event log (100ms or 400 events) | | GC thread | 1 | 8 MB | event log file cleanup every 10 minutes | | crawl queue thread | 1 | 8 MB | process `requestCrawl` — validate hostname, describeServer, spawn worker | -| reconcile thread | 1 | 8 MB | every 5 min: respawn workers for active hosts missing from worker map | | metrics server | 1 | 8 MB | HTTP on internal port, prometheus scrape | | main thread | 1 | default | signal handling, shutdown coordination | @@ -128,7 +127,6 @@ tables: - `account` — uid, did, status, upstream_status, host_id - `account_repo` — uid, rev, commit_data_cid (latest repo state) - `host` — id, hostname, status, last_seq, failed_attempts - - status lifecycle: `active` → `dormant` (after sustained failures, thread exits, DB row persists) → `active` (when discovery re-activates via requestCrawl). `blocked`/`banned` via admin. `idle` on FutureCursor. - `log_file_refs` — seq→file mapping for cursor binary search - `domain_ban` — banned domain suffixes - `backfill_progress` — collection backfill cursor tracking @@ -224,19 +222,6 @@ zlay trades a brief trust window for throughput. the window is bounded: - signature failures trigger key eviction + re-resolution (sync spec guidance) - next commit from the same DID hits the refreshed cache -### host retention model - -indigo gives up on a host after 15 consecutive dial failures (~3 min) -and marks it `offline`, relying on external `requestCrawl` to recover. -zlay retries with exponential backoff capped at 30 minutes (vs indigo's -60s cap). backoff resets to 1 second on any successful connection. after -15 consecutive failures, the host is marked `dormant` — the worker -thread exits (freeing resources) but the DB row persists. a -reconciliation loop (every 5 min) respawns any active host that lost -its worker to a bug or restart. dormant hosts are re-activated by the -discovery cron (requestCrawl), which can now be lightweight and frequent -since it only needs to handle genuinely new hosts + dormant recovery. - ### consumer buffer sizing zlay uses an 8K-entry per-consumer ring buffer (vs indigo's 16K-entry channel). diff --git a/src/api/xrpc.zig b/src/api/xrpc.zig index 2debef1..c6db9ce 100644 --- a/src/api/xrpc.zig +++ b/src/api/xrpc.zig @@ -582,7 +582,7 @@ const GetHostStatusReq = struct { const status = if (std.mem.eql(u8, raw_status, "blocked")) "banned" - else if (std.mem.eql(u8, raw_status, "exhausted") or std.mem.eql(u8, raw_status, "dormant")) + else if (std.mem.eql(u8, raw_status, "exhausted")) "offline" else raw_status; diff --git a/src/host_ops.zig b/src/host_ops.zig index 92b4d0f..0dd1448 100644 --- a/src/host_ops.zig +++ b/src/host_ops.zig @@ -281,10 +281,8 @@ pub const HostOpsQueue = struct { .increment_failures => { const failures = self.persist.incrementHostFailures(op.host_id) catch 0; if (failures >= self.max_consecutive_failures) { - // mark dormant and stop the subscriber thread — free resources. - // the DB row persists so discovery can re-activate the host later. - log.warn("host_ops: host_id={d} dormant after {d} failures, stopping worker", .{ op.host_id, failures }); - self.persist.updateHostStatus(op.host_id, "dormant") catch {}; + log.warn("host_ops: host_id={d} exhausted after {d} failures", .{ op.host_id, failures }); + self.persist.updateHostStatus(op.host_id, "exhausted") catch {}; op.payload.host_shutdown.store(true, .release); } }, @@ -292,8 +290,6 @@ pub const HostOpsQueue = struct { self.persist.resetHostFailures(op.host_id) catch |err| { log.debug("host_ops: reset failures failed for host_id={d}: {s}", .{ op.host_id, @errorName(err) }); }; - // if host was dormant, flip back to active on successful reconnect - self.persist.updateHostStatus(op.host_id, "active") catch {}; }, .update_status => { self.persist.updateHostStatus(op.host_id, op.payload.status.slice()) catch |err| { diff --git a/src/slurper.zig b/src/slurper.zig index f09a74e..9c3e99f 100644 --- a/src/slurper.zig +++ b/src/slurper.zig @@ -5,8 +5,7 @@ //! - spawning/stopping subscriber workers //! - processing crawl requests (adding new hosts) //! - host validation (format, domain ban, describeServer, relay loop detection) -//! - tracking host lifecycle (active → dormant, with persistent retry) -//! - periodic reconciliation of DB host table vs running workers +//! - tracking host lifecycle (active → exhausted → blocked) //! //! all downstream components (Broadcaster, DiskPersist, Validator) are //! thread-safe for N concurrent producers, so this just orchestrates. @@ -254,7 +253,6 @@ pub const Slurper = struct { // background tasks startup_future: ?Io.Future(void) = null, crawl_future: ?Io.Future(void) = null, - reconcile_future: ?Io.Future(void) = null, io: Io, /// dedicated Threaded io for the frame worker pool — safe from plain OS threads @@ -304,7 +302,6 @@ pub const Slurper = struct { // pullHosts + listActiveHosts + spawnWorker all happen in the background thread. self.startup_future = try self.io.concurrent(spawnWorkers, .{self}); self.crawl_future = try self.io.concurrent(processCrawlQueue, .{self}); - self.reconcile_future = try self.io.concurrent(reconcileHosts, .{self}); } /// pull PDS host list from the seed relay's com.atproto.sync.listHosts endpoint. @@ -530,13 +527,12 @@ pub const Slurper = struct { reset_req.base.wait(self.io, self.shutdown); // phase 5: spawn worker (Evented) - try self.spawnWorker(db_req.host_id, hostname, db_req.last_seq, 0); + try self.spawnWorker(db_req.host_id, hostname, db_req.last_seq); log.info("added host {s} (id={d})", .{ hostname, db_req.host_id }); } - /// spawn a subscriber thread for a host. - /// startup_jitter_ms: random delay before first connect (0 = no jitter). - fn spawnWorker(self: *Slurper, host_id: u64, hostname: []const u8, last_seq: u64, startup_jitter_ms: u32) !void { + /// spawn a subscriber thread for a host + fn spawnWorker(self: *Slurper, host_id: u64, hostname: []const u8, last_seq: u64) !void { const hostname_duped = try self.allocator.dupe(u8, hostname); errdefer self.allocator.free(hostname_duped); @@ -587,7 +583,6 @@ pub const Slurper = struct { sub.cursor_slot = cm.register(host_id, last_seq); } if (last_seq > 0) sub.last_upstream_seq = last_seq; - sub.startup_jitter_ms = startup_jitter_ms; const future = try self.io.concurrent(runWorker, .{ self, host_id, sub }); @@ -672,16 +667,10 @@ pub const Slurper = struct { else hosts.len; // 0 = unlimited - // spread initial connections across a 30-second window to avoid - // DNS/TLS handshake storm that starves health probes during startup. - // each host gets a deterministic jitter based on its ID. - const jitter_window_ms: u32 = 30_000; - var spawned: usize = 0; for (hosts) |host| { if (self.shutdown.load(.acquire)) break; - const jitter: u32 = @intCast((host.id *% 2654435761) % jitter_window_ms); - self.spawnWorker(host.id, host.hostname, host.last_seq, jitter) catch |err| { + self.spawnWorker(host.id, host.hostname, host.last_seq) catch |err| { log.warn("failed to spawn worker for {s}: {s}", .{ host.hostname, @errorName(err) }); }; spawned += 1; @@ -720,78 +709,6 @@ pub const Slurper = struct { } } - /// background fiber: periodically reconcile DB host table with running workers. - /// respawns any active host that doesn't have a running subscriber — catches - /// hosts lost to bugs, restarts, or other gaps. dormant hosts are excluded; - /// they wait for discovery (requestCrawl) to flip them back to active. - fn reconcileHosts(self: *Slurper) void { - const interval_secs: u64 = 5 * 60; // 5 minutes - - while (!self.shutdown.load(.acquire)) { - // sleep in 1-second increments so we can check shutdown - var remaining: u64 = interval_secs; - while (remaining > 0 and !self.shutdown.load(.acquire)) { - self.io.sleep(Io.Duration.fromSeconds(1), .awake) catch return; - remaining -= 1; - } - if (self.shutdown.load(.acquire)) return; - - const db_queue = self.db_queue orelse continue; - - // query active hosts only — dormant hosts wait for discovery - const ListActiveReq = struct { - base: event_log_mod.DbRequest = .{ .callback = &execute }, - alloc: Allocator, - result: ?[]event_log_mod.DiskPersist.Host = null, - - fn execute(b: *event_log_mod.DbRequest, dp: *event_log_mod.DiskPersist) void { - const s: *@This() = @fieldParentPtr("base", b); - s.result = dp.listActiveHosts(s.alloc) catch |e| { - b.err = e; - return; - }; - } - }; - var list_req: ListActiveReq = .{ .alloc = self.allocator }; - db_queue.push(&list_req.base); - list_req.base.wait(self.io, self.shutdown); - - if (list_req.base.err != null or list_req.result == null) continue; - - const hosts = list_req.result.?; - defer { - for (hosts) |h| { - self.allocator.free(h.hostname); - self.allocator.free(h.status); - } - self.allocator.free(hosts); - } - - var respawned: usize = 0; - for (hosts) |host| { - if (self.shutdown.load(.acquire)) break; - - const missing = blk: { - self.workers_mutex.lockUncancelable(self.io); - defer self.workers_mutex.unlock(self.io); - break :blk !self.workers.contains(host.id); - }; - - if (missing) { - self.spawnWorker(host.id, host.hostname, host.last_seq, 0) catch |err| { - log.warn("reconcile: failed to respawn {s}: {s}", .{ host.hostname, @errorName(err) }); - continue; - }; - respawned += 1; - } - } - - if (respawned > 0) { - log.info("reconcile: respawned {d} host(s)", .{respawned}); - } - } - } - /// number of active workers pub fn workerCount(self: *Slurper) usize { self.workers_mutex.lockUncancelable(self.io); @@ -820,7 +737,6 @@ pub const Slurper = struct { if (self.startup_future) |*f| f.cancel(self.io); self.crawl_cond.signal(self.io); if (self.crawl_future) |*f| f.cancel(self.io); - if (self.reconcile_future) |*f| f.cancel(self.io); // collect futures to cancel (can't cancel while holding workers_mutex) var futures_to_cancel: std.ArrayListUnmanaged(Io.Future(void)) = .empty; diff --git a/src/subscriber.zig b/src/subscriber.zig index 6f90998..3dc904d 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -219,9 +219,6 @@ pub const Subscriber = struct { shutdown: *std.atomic.Value(bool), last_upstream_seq: ?u64 = null, last_cursor_flush: i64 = 0, - backoff: u64 = 1, - /// one-time random delay (ms) before first connect — spreads TLS storm at startup - startup_jitter_ms: u32 = 0, rate_limiter: RateLimiter = .{}, // per-host shutdown (e.g. FutureCursor — stops only this subscriber) @@ -259,28 +256,23 @@ pub const Subscriber = struct { return self.shutdown.load(.acquire) or self.host_shutdown.load(.acquire); } - /// run the subscriber loop. reconnects with exponential backoff (cap 30 min). - /// blocks until shutdown or host marked dormant. on dormancy the thread exits - /// and frees resources — the DB row persists so discovery can re-activate later. + /// run the subscriber loop. reconnects with exponential backoff. + /// blocks until shutdown flag is set or host is exhausted. pub fn run(self: *Subscriber) void { - const max_backoff: u64 = 1800; // 30 minutes + var backoff: u64 = 1; + const max_backoff: u64 = 60; // cursor is set at spawn time by slurper if (self.last_upstream_seq) |seq| { log.info("host {s}: resuming from cursor {d}", .{ self.options.hostname, seq }); } - // one-time startup jitter — spread DNS/TLS storm across a window - if (self.startup_jitter_ms > 0 and !self.shouldStop()) { - self.io.sleep(Io.Duration.fromMilliseconds(self.startup_jitter_ms), .awake) catch {}; - } - while (!self.shouldStop()) { log.info("host {s}: connecting...", .{self.options.hostname}); self.connectAndRead() catch |err| { if (self.shouldStop()) return; - log.err("host {s}: error: {s}, reconnecting in {d}s...", .{ self.options.hostname, @errorName(err), self.backoff }); + log.err("host {s}: error: {s}, reconnecting in {d}s...", .{ self.options.hostname, @errorName(err), backoff }); }; if (self.shouldStop()) return; @@ -297,13 +289,13 @@ pub const Subscriber = struct { } // backoff sleep in small increments so we can check shutdown - var remaining: u64 = self.backoff; + var remaining: u64 = backoff; while (remaining > 0 and !self.shouldStop()) { const chunk = @min(remaining, 1); self.io.sleep(Io.Duration.fromSeconds(@intCast(chunk)), .awake) catch {}; remaining -= chunk; } - self.backoff = @min(self.backoff * 2, max_backoff); + backoff = @min(backoff * 2, max_backoff); } } @@ -353,8 +345,7 @@ pub const Subscriber = struct { try client.handshake(path, .{ .headers = host_header }); log.info("host {s}: connected", .{self.options.hostname}); - // reset failures + backoff on successful connect - self.backoff = 1; + // reset failures on successful connect (via host_ops queue) if (self.options.host_id > 0) { if (self.host_ops) |hq| { hq.push(.{ -- 2.51.2