atproto relay in zig zlay.waow.tech
relay zig atproto
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751//! slurper — multi-host PDS crawl manager//!//! manages one Subscriber fiber per tracked PDS host. handles://! - loading known hosts from DB on startup//! - 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)//!//! all downstream components (Broadcaster, DiskPersist, Validator) are//! thread-safe for N concurrent producers, so this just orchestrates.
const std = @import("std");const Io = std.Io;const http = std.http;const broadcaster = @import("broadcaster.zig");const validator_mod = @import("validator.zig");const event_log_mod = @import("event_log.zig");const subscriber_mod = @import("subscriber.zig");const collection_index_mod = @import("collection_index/index.zig");const resync_mod = @import("collection_index/resync.zig");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 ws_ping_mod = @import("ws_ping.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);
pub const HostValidationError = atproto.HostValidationError;pub const validateHostname = atproto.validateHostname;const checkHost = atproto.checkHost;
pub const Options = struct { seed_host: []const u8 = "bsky.network", max_message_size: usize = 5 * 1024 * 1024, frame_workers: u16 = 16, frame_queue_capacity: u16 = 4096, /// max hosts to connect concurrently during initial startup ramp. /// prevents TLS handshake storm from starving the event loop. /// 0 = unlimited (legacy behavior). startup_batch_size: u16 = 50,};
const WorkerEntry = struct { future: Io.Future(void), subscriber: *subscriber_mod.Subscriber,};
pub const Slurper = struct { allocator: Allocator, bc: *broadcaster.Broadcaster, validator: *validator_mod.Validator, persist: *event_log_mod.DiskPersist, collection_index: ?*collection_index_mod.CollectionIndex = null, resyncer: ?*resync_mod.Resyncer = null, status_checker: ?*status_checker_mod.StatusChecker = null, host_ops: ?*host_ops_mod.HostOpsQueue = null, cursor_map: ?*host_ops_mod.CursorMap = null, db_queue: ?*event_log_mod.DbRequestQueue = null, ws_ping: ?*ws_ping_mod.WsPing = null, shutdown: *std.atomic.Value(bool), options: Options,
// frame processing pool — offloads heavy work from reader threads frame_pool: ?frame_worker_mod.FramePool = null,
// shared TLS CA bundle — loaded once, used by all subscriber connections ca_bundle: ?std.crypto.Certificate.Bundle = null,
// active subscriber threads, keyed by host_id workers: std.AutoHashMapUnmanaged(u64, WorkerEntry) = .empty, workers_mutex: Io.Mutex = Io.Mutex.init,
// crawl request queue crawl_queue: std.ArrayListUnmanaged([]const u8) = .empty, crawl_mutex: Io.Mutex = Io.Mutex.init, crawl_cond: Io.Condition = Io.Condition.init,
// background tasks startup_future: ?Io.Future(void) = null, crawl_future: ?Io.Future(void) = null,
io: Io, /// dedicated Threaded io for the frame worker pool — safe from plain OS threads pool_io: Io,
pub fn init( allocator: Allocator, bc: *broadcaster.Broadcaster, val: *validator_mod.Validator, persist: *event_log_mod.DiskPersist, shutdown: *std.atomic.Value(bool), options: Options, io: Io, pool_io: Io, ) Slurper { return .{ .allocator = allocator, .bc = bc, .validator = val, .persist = persist, .shutdown = shutdown, .options = options, .io = io, .pool_io = pool_io, }; }
/// start the slurper: bootstrap hosts from seed relay, load from DB, spawn workers. /// Go relay: pull-hosts bootstraps from bsky.network's listHosts, then crawls each PDS directly. pub fn start(self: *Slurper) !void { // load CA bundle once — shared by all subscriber TLS connections var bundle: std.crypto.Certificate.Bundle = .empty; try bundle.rescan(self.allocator, self.io, Io.Timestamp.now(self.io, .real)); self.ca_bundle = bundle; log.info("loaded shared CA bundle", .{});
// create frame processing pool. the pool's mutex/cond use the MAIN io, // not pool_io: submitters can be fibers, and a fiber that parks through // a Threaded io blocks its loop thread (contended lock = kernel futex // wait, full-queue backpressure = thread sleep), starving every other // fiber including the broadcaster. the main io's futex is safe from // both domains — fibers park in the sched table, the pool's own OS // worker threads fall back to the kernel futex, wakes reach both. // under Io.Threaded main io and pool_io behave identically here. self.frame_pool = try frame_worker_mod.FramePool.init(self.allocator, .{ .num_workers = self.options.frame_workers, .queue_capacity = self.options.frame_queue_capacity, .stack_size = @import("../main.zig").default_stack_size, }, self.io); log.info("frame pool started: {d} workers, queue capacity {d}", .{ self.options.frame_workers, self.options.frame_queue_capacity });
// spawn worker startup in background so HTTP server + probes come up immediately. // 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}); }
/// pull PDS host list from the seed relay's com.atproto.sync.listHosts endpoint. /// runs on its own std.Thread (pool_io) — DNS and outbound HTTP work. /// Go relay: cmd/relay/pull.go — one-time bootstrap, reads REST API, not firehose. pub fn pullHosts(self: *Slurper) void { var cursor: ?[]const u8 = null; var total: usize = 0; const limit = 500;
var client: http.Client = .{ .allocator = self.allocator, .io = self.pool_io }; defer client.deinit();
while (true) { if (self.shutdown.load(.acquire)) break;
// build URL with pagination var url_buf: [512]u8 = undefined; const url = if (cursor) |c| std.fmt.bufPrint(&url_buf, "https://{s}/xrpc/com.atproto.sync.listHosts?limit={d}&cursor={s}", .{ self.options.seed_host, limit, c }) catch break else std.fmt.bufPrint(&url_buf, "https://{s}/xrpc/com.atproto.sync.listHosts?limit={d}", .{ self.options.seed_host, limit }) catch break;
var aw: std.Io.Writer.Allocating = .init(self.allocator); defer aw.deinit();
const result = client.fetch(.{ .location = .{ .url = url }, .response_writer = &aw.writer, .method = .GET, }) catch |err| { log.warn("pullHosts: fetch failed: {s}", .{@errorName(err)}); break; };
if (result.status != .ok) { log.warn("pullHosts: got status {d}", .{@backingInt(result.status)}); break; }
const body = aw.written();
// parse JSON response: { "hosts": [{"hostname": "...", "status": "..."}, ...], "cursor": "..." } const parsed = std.json.parseFromSlice(ListHostsResponse, self.allocator, body, .{ .ignore_unknown_fields = true }) catch |err| { log.warn("pullHosts: JSON parse failed: {s}", .{@errorName(err)}); break; }; defer parsed.deinit();
const hosts = parsed.value.hosts orelse break; if (hosts.len == 0) break;
var added: usize = 0; for (hosts) |host| { // skip non-active hosts if (host.status) |s| { if (!std.mem.eql(u8, s, "active")) continue; } // validate hostname format (rejects IPs, localhost, etc.) const normalized = validateHostname(self.allocator, host.hostname) catch continue; defer self.allocator.free(normalized);
// skip banned domains (Threaded pool — runs on own std.Thread) if (self.persist.isDomainBanned(normalized)) continue;
// insert into DB (no describeServer check — the seed relay already vetted them) _ = self.persist.getOrCreateHost(normalized) catch continue; added += 1; } total += added; log.info("pullHosts: page fetched, {d} hosts added ({d} total)", .{ added, total });
// advance cursor if (parsed.value.cursor) |next_cursor| { // free previous cursor if we allocated one if (cursor) |prev| self.allocator.free(prev); cursor = self.allocator.dupe(u8, next_cursor) catch break; } else { break; // no more pages } }
// free final cursor if (cursor) |c| self.allocator.free(c); log.info("pullHosts: bootstrap complete, {d} hosts added from {s}", .{ total, self.options.seed_host }); }
const ListHostsResponse = struct { hosts: ?[]const ListHostEntry = null, cursor: ?[]const u8 = null, };
const ListHostEntry = struct { hostname: []const u8, status: ?[]const u8 = null, };
/// add a crawl request (from requestCrawl endpoint) pub fn addCrawlRequest(self: *Slurper, hostname: []const u8) !void { const duped = try self.allocator.dupe(u8, hostname); self.crawl_mutex.lockUncancelable(self.io); defer self.crawl_mutex.unlock(self.io); try self.crawl_queue.append(self.allocator, duped); self.crawl_cond.signal(self.io); }
/// validate and add a host: format check, domain ban, describeServer, then spawn. /// mirrors Go relay's requestCrawl → SubscribeToHost pipeline. /// uses phased approach: DB checks via DbRequestQueue, HTTP via temp thread. fn addHost(self: *Slurper, raw_hostname: []const u8) !void { const db_queue = self.db_queue orelse return error.DbQueueNotConfigured;
// step 1: validate and normalize hostname format // Go relay: host.go ParseHostname const hostname = validateHostname(self.allocator, raw_hostname) catch |err| { log.warn("host validation failed for '{s}': {s}", .{ raw_hostname, @errorName(err) }); return; }; defer self.allocator.free(hostname);
// phase 1: DB checks via DbRequestQueue const AddHostDbReq = struct { base: event_log_mod.DbRequest = .{ .callback = &execute }, hostname: []const u8, host_id: u64 = 0, last_seq: u64 = 0, rejected: bool = false,
fn execute(b: *event_log_mod.DbRequest, dp: *event_log_mod.DiskPersist) void { const s: *@This() = @fieldParentPtr("base", b); // domain ban check if (dp.isDomainBanned(s.hostname)) { s.rejected = true; return; } // host ban check if (dp.isHostBanned(s.hostname)) { s.rejected = true; return; } // get or create host const info = dp.getOrCreateHost(s.hostname) catch |e| { b.err = e; return; }; s.host_id = info.id; s.last_seq = info.last_seq; } }; var db_req: AddHostDbReq = .{ .hostname = hostname }; db_queue.push(&db_req.base); db_req.base.wait(self.io, self.shutdown); if (db_req.base.err != null) return error.DbRequestFailed; if (db_req.rejected) { log.warn("host {s}: banned/blocked, rejecting", .{hostname}); return; }
// phase 2: dedup — check if already tracked (Evented, local) { self.workers_mutex.lockUncancelable(self.io); defer self.workers_mutex.unlock(self.io); if (self.workers.contains(db_req.host_id)) { log.debug("host {s} already has a worker, skipping", .{hostname}); return; } }
// phase 3: describeServer liveness check (outbound HTTP via temp thread) const CheckHostReq = struct { allocator: Allocator, hostname_copy: []const u8, check_io: Io, check_err: ?HostValidationError = null, done: std.atomic.Value(bool) = .{ .raw = false },
fn run(s: *@This()) void { checkHost(s.allocator, s.hostname_copy, s.check_io) catch |e| { s.check_err = e; }; s.done.store(true, .release); } }; var check_req: CheckHostReq = .{ .allocator = self.allocator, .hostname_copy = hostname, .check_io = self.pool_io, }; const check_thread = std.Thread.spawn(.{}, CheckHostReq.run, .{&check_req}) catch { log.warn("host {s}: failed to spawn check thread", .{hostname}); return; }; // wait for check to complete while (!check_req.done.load(.acquire)) { if (self.shutdown.load(.acquire)) return; self.io.sleep(Io.Duration.fromMicroseconds(100), .awake) catch { while (!check_req.done.load(.acquire)) { if (self.shutdown.load(.acquire)) return; std.atomic.spinLoopHint(); } break; }; } check_thread.join(); if (check_req.check_err) |err| { log.warn("host {s}: describeServer check failed: {s}", .{ hostname, @errorName(err) }); return; }
// phase 4: reset status + failures via DbRequestQueue const ResetHostReq = struct { base: event_log_mod.DbRequest = .{ .callback = &execute }, host_id: u64,
fn execute(b: *event_log_mod.DbRequest, dp: *event_log_mod.DiskPersist) void { const s: *@This() = @fieldParentPtr("base", b); dp.updateHostStatus(s.host_id, "active") catch {}; dp.resetHostFailures(s.host_id) catch {}; } }; var reset_req: ResetHostReq = .{ .host_id = db_req.host_id }; db_queue.push(&reset_req.base); reset_req.base.wait(self.io, self.shutdown);
// phase 5: spawn worker (Evented) try self.spawnWorker(db_req.host_id, hostname, db_req.last_seq); log.info("added host {s} (id={d})", .{ hostname, db_req.host_id }); }
/// spawn a subscriber thread for a host fn spawnWorker(self: *Slurper, host_id: u64, hostname: []const u8, last_seq: u64) !void { const hostname_duped = try self.allocator.dupe(u8, hostname); errdefer self.allocator.free(hostname_duped);
const sub = try self.allocator.create(subscriber_mod.Subscriber); errdefer self.allocator.destroy(sub);
// get effective account count via DbRequestQueue var account_count: u64 = 0; if (self.db_queue) |db_queue| { const GetCountReq = struct { base: event_log_mod.DbRequest = .{ .callback = &execute }, hid: u64, count: u64 = 0,
fn execute(b: *event_log_mod.DbRequest, dp: *event_log_mod.DiskPersist) void { const s: *@This() = @fieldParentPtr("base", b); s.count = dp.getEffectiveAccountCount(s.hid); } }; var count_req: GetCountReq = .{ .hid = host_id }; db_queue.push(&count_req.base); count_req.base.wait(self.io, self.shutdown); account_count = count_req.count; }
sub.* = subscriber_mod.Subscriber.init( self.allocator, self.io, self.bc, self.validator, self.persist, self.shutdown, .{ .hostname = hostname_duped, .max_message_size = self.options.max_message_size, .host_id = host_id, .account_count = account_count, .ca_bundle = self.ca_bundle, }, ); // stamp before the worker fiber can run: a worker that wedges before // its first handshake must still age into reconnect candidacy. sub.last_progress_ms.store(milliTimestamp(self.io), .release); sub.collection_index = self.collection_index; sub.resyncer = self.resyncer; sub.status_checker = self.status_checker; if (self.frame_pool) |*fp| sub.pool = fp; sub.pool_io = self.pool_io; sub.host_ops = self.host_ops; sub.ws_ping = self.ws_ping; if (self.cursor_map) |cm| { sub.cursor_map = cm; sub.cursor_slot = cm.register(host_id, last_seq); } if (last_seq > 0) sub.last_upstream_seq = last_seq;
const future = try self.io.concurrent(runWorker, .{ self, host_id, sub });
self.workers_mutex.lockUncancelable(self.io); defer self.workers_mutex.unlock(self.io); try self.workers.put(self.allocator, host_id, .{ .future = future, .subscriber = sub, }); _ = self.bc.stats.connected_inbound.fetchAdd(1, .monotonic); }
/// worker thread wrapper — runs subscriber, cleans up on exit fn runWorker(self: *Slurper, host_id: u64, sub: *subscriber_mod.Subscriber) void { sub.run();
// return cursor slot to free list before destroying subscriber if (self.cursor_map) |cm| { if (sub.cursor_slot) |slot| { cm.unregister(slot); } }
// subscriber returned — remove from active workers self.workers_mutex.lockUncancelable(self.io); defer self.workers_mutex.unlock(self.io); _ = self.workers.remove(host_id); _ = self.bc.stats.connected_inbound.fetchSub(1, .monotonic);
log.info("worker for host_id={d} ({s}) exited", .{ host_id, sub.options.hostname });
self.allocator.free(sub.options.hostname); self.allocator.destroy(sub); }
/// background fiber: load hosts from DB and spawn all workers. /// runs in background so HTTP server + probes come up immediately. /// pullHosts runs on its own std.Thread in parallel (not gated). /// Go relay: ResubscribeAllHosts loops with 1ms sleep per host (goroutines). /// we batch-spawn with yields between batches to keep the event loop responsive /// for health checks and metrics during the initial TLS handshake ramp. fn spawnWorkers(self: *Slurper) void { const db_queue = self.db_queue orelse { log.err("spawnWorkers: db_queue not set", .{}); return; };
// load active hosts via DbRequestQueue const ListActiveHostsReq = 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.listActiveHosts(s.alloc) catch |e| { b.err = e; return; }; } }; var list_req: ListActiveHostsReq = .{ .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) { log.err("failed to load hosts: {s}", .{if (list_req.base.err) |e| @errorName(e) else "null result"}); return; }
const hosts = list_req.result.?; defer { for (hosts) |h| { self.allocator.free(h.hostname); self.allocator.free(h.status); } self.allocator.free(hosts); }
const batch: usize = if (self.options.startup_batch_size > 0) self.options.startup_batch_size else hosts.len; // 0 = unlimited
var spawned: usize = 0; for (hosts) |host| { if (self.shutdown.load(.acquire)) break; self.spawnWorker(host.id, host.hostname, host.last_seq) catch |err| { log.warn("failed to spawn worker for {s}: {s}", .{ host.hostname, @errorName(err) }); }; spawned += 1;
// yield between batches so the event loop can service health checks if (spawned % batch == 0 and spawned < hosts.len) { log.info("startup: spawned {d}/{d} hosts, yielding...", .{ spawned, hosts.len }); self.io.sleep(Io.Duration.fromMilliseconds(100), .awake) catch break; } }
log.info("startup complete: {d} host(s) spawned", .{hosts.len}); }
/// background thread: process crawl requests fn processCrawlQueue(self: *Slurper) void { while (!self.shutdown.load(.acquire)) { var hostname: ?[]const u8 = null; { self.crawl_mutex.lockUncancelable(self.io); defer self.crawl_mutex.unlock(self.io); while (self.crawl_queue.items.len == 0 and !self.shutdown.load(.acquire)) { self.crawl_cond.waitUncancelable(self.io, &self.crawl_mutex); } if (self.crawl_queue.items.len > 0) { hostname = self.crawl_queue.orderedRemove(0); } }
if (hostname) |h| { defer self.allocator.free(h); self.addHost(h) catch |err| { log.warn("crawl request failed for {s}: {s}", .{ h, @errorName(err) }); }; } } }
/// number of active workers pub fn workerCount(self: *Slurper) usize { self.workers_mutex.lockUncancelable(self.io); defer self.workers_mutex.unlock(self.io); return self.workers.count(); }
/// update rate limits for a running subscriber (called from admin API). /// if the host has a worker, recomputes and applies new limits immediately. pub fn updateHostLimits(self: *Slurper, host_id: u64, account_count: u64) void { self.workers_mutex.lockUncancelable(self.io); defer self.workers_mutex.unlock(self.io); if (self.workers.get(host_id)) |entry| { const trusted = subscriber_mod.isTrustedHost(entry.subscriber.options.hostname); const limits = subscriber_mod.computeLimits(trusted, account_count); entry.subscriber.rate_limiter.updateLimits(limits.sec, limits.hour, limits.day); log.info("updated rate limits for host_id={d}: sec={d} hour={d} day={d}", .{ host_id, limits.sec, limits.hour, limits.day, }); } }
/// 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 -- /// parked forever on a read that will never complete (see the 2026-08-18 /// incident). Neither existing path could do it: `addHost` dedupes on /// `workers.contains`, so requestCrawl is a no-op while the dead worker /// still exists, and admin block/unblock only writes DB status, leaving /// the host stranded as blocked without ever tearing the worker down. /// /// Cancelling is what actually reaches a wedged fiber: a flag cannot wake /// a fiber parked in netRead, but cancellation readies it and makes the /// read return `error.Canceled`. `host_shutdown` is set first so a /// *healthy* worker exits its reconnect loop cleanly instead of racing. /// `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. 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. /// /// `min_idle_ms` is the guard. A deaf host and a merely quiet host look /// identical from seq sampling, so a reconnect driven by silence alone /// would tear down healthy connections -- manufacturing exactly the /// connection churn that this class of bug feeds on (2026-08-18). 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_progress_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_progress_ms` and /// `util.timestamp` returns SECONDS, which is a live mixing hazard here. fn reconnectSkipAt(now_ms: i64, last_progress_ms: i64, min_idle_ms: i64) ?ReconnectOutcome { // 0 means the worker was never stamped. addHost stamps at spawn, so // this is defensive only; treat it as a reconnect candidate. if (last_progress_ms == 0) return null; const idle_ms = now_ms - last_progress_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 .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); @memcpy(hostname_buf[0..hostname_len], hn[0..hostname_len]); entry.subscriber.host_shutdown.store(true, .release); future = entry.future; }
// must not hold workers_mutex here: runWorker takes it on the way out. future.cancel(self.io);
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 .respawn_failed; }; log.info("forceReconnect: host_id={d} ({s}) worker torn down, respawn queued", .{ host_id, hostname }); return .reconnected; }
/// shutdown all workers and clean up pub fn deinit(self: *Slurper) void { // cancel background tasks if (self.startup_future) |*f| f.cancel(self.io); self.crawl_cond.signal(self.io); if (self.crawl_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; defer futures_to_cancel.deinit(self.allocator);
{ self.workers_mutex.lockUncancelable(self.io); defer self.workers_mutex.unlock(self.io); var it = self.workers.iterator(); while (it.next()) |entry| { futures_to_cancel.append(self.allocator, entry.value_ptr.future) catch {}; } }
// cancel all subscriber tasks FIRST (they stop submitting to pool) for (futures_to_cancel.items) |*f| f.cancel(self.io);
// then drain + join pool workers (processes remaining queued frames) if (self.frame_pool) |*fp| { fp.shutdown(); fp.deinit(); self.frame_pool = null; }
// clean up workers map self.workers.deinit(self.allocator);
// clean up crawl queue for (self.crawl_queue.items) |h| self.allocator.free(h); self.crawl_queue.deinit(self.allocator);
// free shared CA bundle 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));
// a worker wedged before its first handshake stamped only at spawn, so it // ages past the window and IS a candidate. This is the case the idle // guard must not hide: connectAndRead that never returns means // run()'s reconnect loop never iterates, so there is no backoff to // recover it either. try std.testing.expectEqual(@as(?S.ReconnectOutcome, null), S.reconnectSkipAt(now, now - 30 * 60 * std.time.ms_per_s, min_idle_ms));
// ...whereas an unreachable host looping in backoff re-stamps on every // attempt, so it stays inside the idle guard. try std.testing.expect(S.reconnectSkipAt(now, now - 5 * std.time.ms_per_s, min_idle_ms) != null);}