From 5d3a7b54aba8256a900e71421e12b362a47712a8 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sat, 11 Apr 2026 01:59:11 -0500 Subject: [PATCH] fix host retention: dormant hosts stop their thread, not just relabel MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit the previous change kept dormant subscriber threads running forever, meaning thread count could only go up. dormant now correctly stops the worker thread (freeing resources) while preserving the DB row for discovery to re-activate later. reconciliation loop queries only active hosts — dormant hosts wait for requestCrawl. separated "don't forget the host" (DB row persists) from "don't stop the thread" (thread exits on dormancy). removed unused listReconnectableHosts. Co-Authored-By: Claude Opus 4.6 (1M context) --- docs/design.md | 20 ++++++++++---------- src/event_log.zig | 16 +--------------- src/host_ops.zig | 7 ++++--- src/slurper.zig | 13 +++++++------ src/subscriber.zig | 7 +++---- 5 files changed, 25 insertions(+), 38 deletions(-) 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 }, }); } } -- 2.51.2