From 143605c83addff39a0d89b5801cb4d986a72b6e9 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Fri, 10 Apr 2026 23:42:12 -0500 Subject: [PATCH] persistent host reconnect: never give up on known hosts MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit subscribers now retry with exponential backoff capped at 30 min (was 60s cap with hard kill at 15 failures). on successful connect, backoff resets to 1s and host flips back to active. hosts that fail 15+ consecutive times are marked dormant (observable) but the subscriber keeps retrying. a reconciliation loop every 5 min respawns any active/dormant host missing from the worker map. this eliminates dependence on the external reconnect cron for host retention — it can be reduced to discovery-only. Co-Authored-By: Claude Opus 4.6 (1M context) --- docs/design.md | 15 +++++++++ src/api/xrpc.zig | 2 +- src/event_log.zig | 16 +++++++++- src/host_ops.zig | 9 ++++-- src/slurper.zig | 77 +++++++++++++++++++++++++++++++++++++++++++++- src/subscriber.zig | 21 +++++++------ 6 files changed, 125 insertions(+), 15 deletions(-) diff --git a/docs/design.md b/docs/design.md index 612db8f..6f43252 100644 --- a/docs/design.md +++ b/docs/design.md @@ -54,6 +54,7 @@ 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/dormant hosts missing from worker map | | metrics server | 1 | 8 MB | HTTP on internal port, prometheus scrape | | main thread | 1 | default | signal handling, shutdown coordination | @@ -127,6 +128,7 @@ 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, subscriber keeps retrying) → `active` (on successful reconnect). `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 @@ -222,6 +224,19 @@ 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 never gives up on a known host — subscribers retry with exponential +backoff capped at 30 minutes, and backoff resets to 1 second on any +successful connection. after sustained failure (15+ consecutive +disconnects), the host is marked `dormant` for observability but the +subscriber keeps retrying. a reconciliation loop (every 5 min) respawns +any active/dormant host missing from the worker map, catching hosts lost +to any gap. this eliminates the need for an external reconnect cron — +the cron can be reduced to discovery-only (new hosts not yet in the DB). + ### 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 c6db9ce..2debef1 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")) + else if (std.mem.eql(u8, raw_status, "exhausted") or std.mem.eql(u8, raw_status, "dormant")) "offline" else raw_status; diff --git a/src/event_log.zig b/src/event_log.zig index 9f9b6c1..f4c05db 100644 --- a/src/event_log.zig +++ b/src/event_log.zig @@ -702,6 +702,20 @@ pub const DiskPersist = struct { } fn listActiveHostsImpl(allocator: Allocator, db: *pg.Pool) ![]Host { + return listHostsByStatusImpl(allocator, db, "status = 'active'"); + } + + /// list hosts eligible for reconnection: active + dormant (Threaded pool). + /// used by the reconciliation loop to respawn missing workers. + pub fn listReconnectableHosts(self: *DiskPersist, allocator: Allocator) ![]Host { + return listReconnectableHostsImpl(allocator, self.db); + } + + fn listReconnectableHostsImpl(allocator: Allocator, db: *pg.Pool) ![]Host { + return listHostsByStatusImpl(allocator, db, "status IN ('active', 'dormant')"); + } + + fn listHostsByStatusImpl(allocator: Allocator, db: *pg.Pool, comptime where: []const u8) ![]Host { var hosts: std.ArrayListUnmanaged(Host) = .empty; errdefer { for (hosts.items) |host| { @@ -712,7 +726,7 @@ pub const DiskPersist = struct { } var result = try db.query( - "SELECT id, hostname, status, last_seq, failed_attempts, account_limit FROM host WHERE status = 'active' ORDER BY id ASC", + "SELECT id, hostname, status, last_seq, failed_attempts, account_limit FROM host WHERE " ++ where ++ " ORDER BY id ASC", .{}, ); defer result.deinit(); diff --git a/src/host_ops.zig b/src/host_ops.zig index 0dd1448..663e916 100644 --- a/src/host_ops.zig +++ b/src/host_ops.zig @@ -281,15 +281,18 @@ pub const HostOpsQueue = struct { .increment_failures => { const failures = self.persist.incrementHostFailures(op.host_id) catch 0; if (failures >= self.max_consecutive_failures) { - 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); + // mark dormant for observability — subscriber keeps retrying. + // reconciliation loop will also periodically re-check dormant hosts. + log.warn("host_ops: host_id={d} marked dormant after {d} failures", .{ op.host_id, failures }); + self.persist.updateHostStatus(op.host_id, "dormant") catch {}; } }, .reset_failures => { 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 9c3e99f..091a060 100644 --- a/src/slurper.zig +++ b/src/slurper.zig @@ -5,7 +5,8 @@ //! - spawning/stopping subscriber workers //! - processing crawl requests (adding new hosts) //! - host validation (format, domain ban, describeServer, relay loop detection) -//! - tracking host lifecycle (active → exhausted → blocked) +//! - tracking host lifecycle (active → dormant, with persistent retry) +//! - periodic reconciliation of DB host table vs running workers //! //! all downstream components (Broadcaster, DiskPersist, Validator) are //! thread-safe for N concurrent producers, so this just orchestrates. @@ -253,6 +254,7 @@ 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 @@ -302,6 +304,7 @@ 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. @@ -709,6 +712,77 @@ pub const Slurper = struct { } } + /// background fiber: periodically reconcile DB host table with running workers. + /// respawns any active/dormant host that doesn't have a running subscriber. + /// catches hosts lost to prior exhaustion, bugs, or any other gap. + fn reconcileHosts(self: *Slurper) void { + const interval = Io.Duration.fromSeconds(5 * 60); // 5 minutes + + // wait for initial startup to finish before first reconciliation + if (self.startup_future) |*f| { + f.wait(self.io) catch {}; + } + + while (!self.shutdown.load(.acquire)) { + self.io.sleep(interval, .awake) catch return; + if (self.shutdown.load(.acquire)) return; + + const db_queue = self.db_queue orelse continue; + + // query reconnectable hosts (active + dormant) via DbRequestQueue + const ListReconnectableReq = 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.listReconnectableHosts(s.alloc) catch |e| { + b.err = e; + return; + }; + } + }; + var list_req: ListReconnectableReq = .{ .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) 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); @@ -737,6 +811,7 @@ 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 3dc904d..5fe9526 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -219,6 +219,7 @@ pub const Subscriber = struct { shutdown: *std.atomic.Value(bool), last_upstream_seq: ?u64 = null, last_cursor_flush: i64 = 0, + backoff: u64 = 1, rate_limiter: RateLimiter = .{}, // per-host shutdown (e.g. FutureCursor — stops only this subscriber) @@ -256,11 +257,12 @@ pub const Subscriber = struct { return self.shutdown.load(.acquire) or self.host_shutdown.load(.acquire); } - /// run the subscriber loop. reconnects with exponential backoff. - /// blocks until shutdown flag is set or host is exhausted. + /// run the subscriber loop. reconnects with exponential backoff (cap 30 min). + /// blocks until shutdown flag is set. never gives up on a known host — + /// host_ops marks hosts dormant after sustained failure for observability, + /// but the subscriber keeps retrying. pub fn run(self: *Subscriber) void { - var backoff: u64 = 1; - const max_backoff: u64 = 60; + const max_backoff: u64 = 1800; // 30 minutes // cursor is set at spawn time by slurper if (self.last_upstream_seq) |seq| { @@ -272,7 +274,7 @@ pub const Subscriber = struct { self.connectAndRead() catch |err| { if (self.shouldStop()) return; - log.err("host {s}: error: {s}, reconnecting in {d}s...", .{ self.options.hostname, @errorName(err), backoff }); + log.err("host {s}: error: {s}, reconnecting in {d}s...", .{ self.options.hostname, @errorName(err), self.backoff }); }; if (self.shouldStop()) return; @@ -283,19 +285,19 @@ pub const Subscriber = struct { hq.push(.{ .host_id = self.options.host_id, .kind = .increment_failures, - .payload = .{ .host_shutdown = &self.host_shutdown }, + .payload = .{ .none = {} }, }); } } // backoff sleep in small increments so we can check shutdown - var remaining: u64 = backoff; + var remaining: u64 = self.backoff; while (remaining > 0 and !self.shouldStop()) { const chunk = @min(remaining, 1); self.io.sleep(Io.Duration.fromSeconds(@intCast(chunk)), .awake) catch {}; remaining -= chunk; } - backoff = @min(backoff * 2, max_backoff); + self.backoff = @min(self.backoff * 2, max_backoff); } } @@ -345,7 +347,8 @@ pub const Subscriber = struct { try client.handshake(path, .{ .headers = host_header }); log.info("host {s}: connected", .{self.options.hostname}); - // reset failures on successful connect (via host_ops queue) + // reset failures + backoff on successful connect + self.backoff = 1; if (self.options.host_id > 0) { if (self.host_ops) |hq| { hq.push(.{ -- 2.51.2