diff --git a/docs/design.md b/docs/design.md index 6f43252..5ab30c1 100644 --- a/docs/design.md +++ b/docs/design.md @@ -54,7 +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 | +| 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 +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. + - 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 @@ -228,14 +228,14 @@ zlay trades a brief trust window for throughput. the window is bounded: 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). +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 diff --git a/src/event_log.zig b/src/event_log.zig index f4c05db..9f9b6c1 100644 --- a/src/event_log.zig +++ b/src/event_log.zig @@ -702,20 +702,6 @@ 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| { @@ -726,7 +712,7 @@ pub const DiskPersist = struct { } var result = try db.query( - "SELECT id, hostname, status, last_seq, failed_attempts, account_limit FROM host WHERE " ++ where ++ " ORDER BY id ASC", + "SELECT id, hostname, status, last_seq, failed_attempts, account_limit FROM host WHERE status = 'active' ORDER BY id ASC", .{}, ); defer result.deinit(); diff --git a/src/host_ops.zig b/src/host_ops.zig index 663e916..92b4d0f 100644 --- a/src/host_ops.zig +++ b/src/host_ops.zig @@ -281,10 +281,11 @@ pub const HostOpsQueue = struct { .increment_failures => { const failures = self.persist.incrementHostFailures(op.host_id) catch 0; if (failures >= self.max_consecutive_failures) { - // 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 }); + // 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 {}; + op.payload.host_shutdown.store(true, .release); } }, .reset_failures => { diff --git a/src/slurper.zig b/src/slurper.zig index 00979ff..a30e913 100644 --- a/src/slurper.zig +++ b/src/slurper.zig @@ -713,8 +713,9 @@ 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. + /// 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 @@ -729,21 +730,21 @@ pub const Slurper = struct { const db_queue = self.db_queue orelse continue; - // query reconnectable hosts (active + dormant) via DbRequestQueue - const ListReconnectableReq = struct { + // 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.listReconnectableHosts(s.alloc) catch |e| { + s.result = dp.listActiveHosts(s.alloc) catch |e| { b.err = e; return; }; } }; - var list_req: ListReconnectableReq = .{ .alloc = self.allocator }; + var list_req: ListActiveReq = .{ .alloc = self.allocator }; db_queue.push(&list_req.base); list_req.base.wait(self.io, self.shutdown); diff --git a/src/subscriber.zig b/src/subscriber.zig index 5fe9526..3a266aa 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -258,9 +258,8 @@ pub const Subscriber = struct { } /// 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. + /// 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. pub fn run(self: *Subscriber) void { const max_backoff: u64 = 1800; // 30 minutes @@ -285,7 +284,7 @@ pub const Subscriber = struct { hq.push(.{ .host_id = self.options.host_id, .kind = .increment_failures, - .payload = .{ .none = {} }, + .payload = .{ .host_shutdown = &self.host_shutdown }, }); } }