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(.{