diff --git a/src/internal/api/admin.zig b/src/internal/api/admin.zig index 94bbc39..372bf65 100644 --- a/src/internal/api/admin.zig +++ b/src/internal/api/admin.zig @@ -18,6 +18,11 @@ const resync_mod = @import("../collection_index/resync.zig"); const util = @import("../util/util.zig"); const log = std.log.scoped(.relay); + +/// default idle window for /admin/hosts/reconnect. Long enough that an +/// ordinarily quiet PDS is not mistaken for a deaf one, short enough to be +/// useful during an incident. Overridable per call. +const default_reconnect_min_idle_s: i64 = 60; const getenv = util.getenv; const formatTimestamp = util.formatTimestamp; @@ -302,9 +307,31 @@ pub fn handleAdminUnblockHost(conn: *h.Conn, body: []const u8, headers: *const w /// In-band recovery for a worker that is alive but deaf. Unlike /// block+unblock this never changes the host's status, so a host cannot be /// left stranded as blocked if the operator's second call does not land. -pub fn handleAdminReconnectHost(conn: *h.Conn, body: []const u8, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { +/// +/// POST /admin/hosts/reconnect?min_idle_seconds=60&dry_run=true +/// body: {"hostname": "..."} +/// +/// `min_idle_seconds` guards against tearing down a healthy worker: a deaf +/// host and a merely quiet host are indistinguishable from seq sampling, so +/// a bulk sweep would otherwise manufacture the connection churn it is +/// trying to recover from. It is a query parameter, not a constant, so the +/// window can be widened mid-incident without a redeploy. +pub fn handleAdminReconnectHost(conn: *h.Conn, query: []const u8, body: []const u8, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { if (!checkAdmin(conn, headers)) return; + const min_idle_s = blk: { + const raw = h.queryParam(query, "min_idle_seconds") orelse break :blk default_reconnect_min_idle_s; + break :blk std.fmt.parseInt(i64, raw, 10) catch { + h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"min_idle_seconds must be an integer\"}"); + return; + }; + }; + if (min_idle_s < 0) { + h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"min_idle_seconds must not be negative\"}"); + return; + } + const dry_run = if (h.queryParam(query, "dry_run")) |v| std.mem.eql(u8, v, "true") else false; + const parsed = std.json.parseFromSlice(struct { hostname: []const u8 }, ctx.persist.allocator, body, .{ .ignore_unknown_fields = true }) catch { h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"invalid JSON\"}"); return; @@ -349,13 +376,38 @@ pub fn handleAdminReconnectHost(conn: *h.Conn, body: []const u8, headers: *const return; } - if (!ctx.slurper.forceReconnect(req.host_id)) { - h.respondJson(conn, .not_found, "{\"error\":\"NoWorker\",\"message\":\"host has no active worker\"}"); + const min_idle_ms = min_idle_s * std.time.ms_per_s; + var buf: [256]u8 = undefined; + + if (dry_run) { + // report the decision, and for a skip the idle time behind it -- the + // skips are what tell an operator whether their silent-set is real. + const outcome = ctx.slurper.inspectReconnect(req.host_id, min_idle_ms); + const payload = switch (outcome) { + .reconnected => std.fmt.bufPrint(&buf, "{{\"dry_run\":true,\"would\":\"reconnect\"}}", .{}), + .skipped_recent_frame => |idle_ms| std.fmt.bufPrint(&buf, "{{\"dry_run\":true,\"would\":\"skip\",\"reason\":\"recent_frame\",\"idle_seconds\":{d},\"min_idle_seconds\":{d}}}", .{ @divFloor(idle_ms, std.time.ms_per_s), min_idle_s }), + .no_worker => std.fmt.bufPrint(&buf, "{{\"dry_run\":true,\"would\":\"skip\",\"reason\":\"no_worker\"}}", .{}), + .respawn_failed => std.fmt.bufPrint(&buf, "{{\"dry_run\":true,\"would\":\"reconnect\"}}", .{}), + } catch { + h.respondJson(conn, .internal_server_error, "{\"error\":\"Internal\"}"); + return; + }; + h.respondJson(conn, .ok, payload); return; } - log.info("admin: forced reconnect for host {s} (id={d})", .{ parsed.value.hostname, req.host_id }); - h.respondJson(conn, .ok, "{\"success\":true}"); + switch (ctx.slurper.forceReconnect(req.host_id, min_idle_ms)) { + .reconnected => { + log.info("admin: forced reconnect for host {s} (id={d})", .{ parsed.value.hostname, req.host_id }); + h.respondJson(conn, .ok, "{\"success\":true,\"reconnected\":true}"); + }, + .skipped_recent_frame => |idle_ms| { + const payload = std.fmt.bufPrint(&buf, "{{\"success\":true,\"reconnected\":false,\"reason\":\"recent_frame\",\"idle_seconds\":{d},\"min_idle_seconds\":{d}}}", .{ @divFloor(idle_ms, std.time.ms_per_s), min_idle_s }) catch "{\"success\":true,\"reconnected\":false,\"reason\":\"recent_frame\"}"; + h.respondJson(conn, .ok, payload); + }, + .no_worker => h.respondJson(conn, .not_found, "{\"error\":\"NoWorker\",\"message\":\"host has no active worker\"}"), + .respawn_failed => h.respondJson(conn, .internal_server_error, "{\"error\":\"RespawnFailed\",\"message\":\"worker torn down but respawn could not be queued\"}"), + } } /// set or clear the account_limit override for a host. diff --git a/src/internal/api/router.zig b/src/internal/api/router.zig index 7a6c79d..9ca8b06 100644 --- a/src/internal/api/router.zig +++ b/src/internal/api/router.zig @@ -135,7 +135,7 @@ fn handlePost(conn: *websocket.Conn, path: []const u8, query: []const u8, body: } else if (std.mem.eql(u8, path, "/admin/hosts/unblock")) { admin.handleAdminUnblockHost(conn, body, headers, ctx); } else if (std.mem.eql(u8, path, "/admin/hosts/reconnect")) { - admin.handleAdminReconnectHost(conn, body, headers, ctx); + admin.handleAdminReconnectHost(conn, query, body, headers, ctx); } else if (std.mem.eql(u8, path, "/admin/hosts/changeLimits")) { admin.handleAdminChangeLimits(conn, body, headers, ctx); } else if (std.mem.eql(u8, path, "/admin/backfill-collections")) { diff --git a/src/internal/slurper.zig b/src/internal/slurper.zig index 454a546..ab0ee91 100644 --- a/src/internal/slurper.zig +++ b/src/internal/slurper.zig @@ -23,6 +23,7 @@ const status_checker_mod = @import("status_checker.zig"); const frame_worker_mod = @import("frame_worker.zig"); const host_ops_mod = @import("host_ops.zig"); const atproto = @import("atproto/main.zig"); +const milliTimestamp = @import("util/util.zig").milliTimestamp; const Allocator = std.mem.Allocator; const log = std.log.scoped(.relay); @@ -564,6 +565,16 @@ pub const Slurper = struct { } } + /// outcome of a forceReconnect attempt, so callers can distinguish + /// "nothing to do" from "declined on purpose". + pub const ReconnectOutcome = union(enum) { + reconnected, + /// the worker showed signs of life this recently — left alone. + skipped_recent_frame: i64, + no_worker, + respawn_failed, + }; + /// force a host's worker to drop and re-establish its connection. /// /// This is the in-band recovery for a worker that is alive but deaf -- @@ -580,17 +591,56 @@ pub const Slurper = struct { /// `future.cancel` awaits completion, so by the time it returns runWorker /// has already freed the cursor slot, dropped the map entry and /// decremented `connected_inbound` -- which is why the respawn below no - /// longer trips the dedup. + /// longer trips the dedup. Cancelling an already-finished future is the + /// designed path, not a use-after-free: zio's Task is refcounted so it + /// outlives its fiber until awaited or cancelled exactly once. /// - /// Returns false if the host has no worker. - pub fn forceReconnect(self: *Slurper, host_id: u64) bool { + /// `min_idle_ms` is the guard. A deaf host and a merely quiet host look + /// identical from seq sampling, so a bulk sweep driven by a silent-host + /// probe would otherwise tear down healthy connections -- manufacturing + /// exactly the connection churn that this class of bug feeds on. Any + /// worker that has seen a frame (or completed a handshake) more recently + /// than this is left alone. + /// non-null when this worker should be left alone. Shared by + /// forceReconnect and inspectReconnect so the two cannot drift; callers + /// must hold workers_mutex. + fn reconnectSkip(self: *Slurper, entry: WorkerEntry, min_idle_ms: i64) ?ReconnectOutcome { + return reconnectSkipAt( + milliTimestamp(self.io), + entry.subscriber.last_frame_ms.load(.acquire), + min_idle_ms, + ); + } + + /// the guard decision, free of clock and lock so it can be tested. all + /// three arguments are milliseconds -- the field is `last_frame_ms` and + /// `util.timestamp` returns SECONDS, which is a live mixing hazard here. + fn reconnectSkipAt(now_ms: i64, last_frame_ms: i64, min_idle_ms: i64) ?ReconnectOutcome { + // 0 = never saw a frame or a handshake, so it is a reconnect candidate. + if (last_frame_ms == 0) return null; + const idle_ms = now_ms - last_frame_ms; + if (idle_ms < min_idle_ms) return .{ .skipped_recent_frame = idle_ms }; + return null; + } + + /// what forceReconnect would decide, without tearing anything down. + pub fn inspectReconnect(self: *Slurper, host_id: u64, min_idle_ms: i64) ReconnectOutcome { + self.workers_mutex.lockUncancelable(self.io); + defer self.workers_mutex.unlock(self.io); + const entry = self.workers.get(host_id) orelse return .no_worker; + return self.reconnectSkip(entry, min_idle_ms) orelse .reconnected; + } + + pub fn forceReconnect(self: *Slurper, host_id: u64, min_idle_ms: i64) ReconnectOutcome { var hostname_buf: [256]u8 = undefined; var hostname_len: usize = 0; var future: Io.Future(void) = undefined; { self.workers_mutex.lockUncancelable(self.io); defer self.workers_mutex.unlock(self.io); - const entry = self.workers.get(host_id) orelse return false; + const entry = self.workers.get(host_id) orelse return .no_worker; + if (self.reconnectSkip(entry, min_idle_ms)) |skip| return skip; + // copy the hostname now: runWorker frees it once the future ends. const hn = entry.subscriber.options.hostname; hostname_len = @min(hn.len, hostname_buf.len); @@ -605,10 +655,10 @@ pub const Slurper = struct { const hostname = hostname_buf[0..hostname_len]; self.addCrawlRequest(hostname) catch |e| { log.err("forceReconnect: host_id={d} ({s}) torn down but respawn enqueue failed: {s}", .{ host_id, hostname, @errorName(e) }); - return false; + return .respawn_failed; }; log.info("forceReconnect: host_id={d} ({s}) worker torn down, respawn queued", .{ host_id, hostname }); - return true; + return .reconnected; } /// shutdown all workers and clean up @@ -652,3 +702,31 @@ pub const Slurper = struct { if (self.ca_bundle) |*b| b.deinit(self.allocator); } }; + +test "forceReconnect guard: only workers idle past the window are torn down" { + const S = Slurper; + const min_idle_ms: i64 = 60 * std.time.ms_per_s; + const now: i64 = 1_000_000; + + // a worker that delivered a frame one second ago is healthy -- a + // silent-host probe cannot tell it apart from a deaf one, so the endpoint + // has to. Tearing this down would manufacture the connection churn the + // 2026-08-18 bug feeds on. + const recent = S.reconnectSkipAt(now, now - 1_000, min_idle_ms); + try std.testing.expect(recent != null); + try std.testing.expectEqual(@as(i64, 1_000), recent.?.skipped_recent_frame); + + // idle well past the window: a reconnect candidate. + try std.testing.expectEqual(@as(?S.ReconnectOutcome, null), S.reconnectSkipAt(now, now - 120 * std.time.ms_per_s, min_idle_ms)); + + // never saw a frame or a handshake: also a candidate. + try std.testing.expectEqual(@as(?S.ReconnectOutcome, null), S.reconnectSkipAt(now, 0, min_idle_ms)); + + // boundary: exactly at the window is idle enough (the guard skips only + // strictly-more-recent activity). + try std.testing.expectEqual(@as(?S.ReconnectOutcome, null), S.reconnectSkipAt(now, now - min_idle_ms, min_idle_ms)); + + // min_idle_seconds=0 disables the guard entirely, which is what an + // operator wants when they are certain a host is deaf. + try std.testing.expectEqual(@as(?S.ReconnectOutcome, null), S.reconnectSkipAt(now, now, 0)); +} diff --git a/src/internal/subscriber.zig b/src/internal/subscriber.zig index 1a70101..048eda7 100644 --- a/src/internal/subscriber.zig +++ b/src/internal/subscriber.zig @@ -211,6 +211,13 @@ pub const Subscriber = struct { shutdown: *std.atomic.Value(bool), last_upstream_seq: ?u64 = null, last_cursor_flush: i64 = 0, + /// ms timestamp of the last sign of life from upstream: a frame, or the + /// handshake that opened the connection. Written by the subscriber fiber, + /// read by the admin fiber, hence atomic. This is the liveness signal + /// `Slurper.forceReconnect` uses so a bulk recovery sweep cannot tear + /// down workers that are merely quiet (see the 2026-08-18 incident: a + /// deaf host and a quiet host look identical from seq sampling alone). + last_frame_ms: std.atomic.Value(i64) = .{ .raw = 0 }, rate_limiter: RateLimiter = .{}, // per-host shutdown (e.g. FutureCursor — stops only this subscriber) @@ -386,6 +393,9 @@ pub const Subscriber = struct { return e; }; log.info("host {s}: connected", .{self.options.hostname}); + // a fresh connection counts as liveness: a worker that just completed + // its handshake must not look idle to a recovery sweep. + self.last_frame_ms.store(milliTimestamp(self.io), .release); // reset failures on successful connect (via host_ops queue) if (self.options.host_id > 0) { @@ -473,13 +483,15 @@ const FrameHandler = struct { // count every successfully decoded event (matches Go relay's events_received_counter) _ = sub.bc.stats.frames_in.fetchAdd(1, .monotonic); + const now_ms = milliTimestamp(io); + sub.last_frame_ms.store(now_ms, .release); // extract seq for cursor tracking (deferred until after pool accepts) const upstream_seq = payload.getUint("seq"); // time-based cursor flush (Go relay: every 4 seconds) { - const now = timestamp(io); + const now = @divFloor(now_ms, std.time.ms_per_s); if (now - sub.last_cursor_flush >= cursor_flush_interval_sec) { sub.flushCursor(); sub.last_cursor_flush = now;