diff --git a/src/slurper.zig b/src/slurper.zig index a30e913..f09a74e 100644 --- a/src/slurper.zig +++ b/src/slurper.zig @@ -530,12 +530,13 @@ pub const Slurper = struct { 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); + try self.spawnWorker(db_req.host_id, hostname, db_req.last_seq, 0); 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 { + /// spawn a subscriber thread for a host. + /// startup_jitter_ms: random delay before first connect (0 = no jitter). + fn spawnWorker(self: *Slurper, host_id: u64, hostname: []const u8, last_seq: u64, startup_jitter_ms: u32) !void { const hostname_duped = try self.allocator.dupe(u8, hostname); errdefer self.allocator.free(hostname_duped); @@ -586,6 +587,7 @@ pub const Slurper = struct { sub.cursor_slot = cm.register(host_id, last_seq); } if (last_seq > 0) sub.last_upstream_seq = last_seq; + sub.startup_jitter_ms = startup_jitter_ms; const future = try self.io.concurrent(runWorker, .{ self, host_id, sub }); @@ -670,10 +672,16 @@ pub const Slurper = struct { else hosts.len; // 0 = unlimited + // spread initial connections across a 30-second window to avoid + // DNS/TLS handshake storm that starves health probes during startup. + // each host gets a deterministic jitter based on its ID. + const jitter_window_ms: u32 = 30_000; + var spawned: usize = 0; for (hosts) |host| { if (self.shutdown.load(.acquire)) break; - self.spawnWorker(host.id, host.hostname, host.last_seq) catch |err| { + const jitter: u32 = @intCast((host.id *% 2654435761) % jitter_window_ms); + self.spawnWorker(host.id, host.hostname, host.last_seq, jitter) catch |err| { log.warn("failed to spawn worker for {s}: {s}", .{ host.hostname, @errorName(err) }); }; spawned += 1; @@ -770,7 +778,7 @@ pub const Slurper = struct { }; if (missing) { - self.spawnWorker(host.id, host.hostname, host.last_seq) catch |err| { + self.spawnWorker(host.id, host.hostname, host.last_seq, 0) catch |err| { log.warn("reconcile: failed to respawn {s}: {s}", .{ host.hostname, @errorName(err) }); continue; }; diff --git a/src/subscriber.zig b/src/subscriber.zig index 3a266aa..6f90998 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -220,6 +220,8 @@ pub const Subscriber = struct { last_upstream_seq: ?u64 = null, last_cursor_flush: i64 = 0, backoff: u64 = 1, + /// one-time random delay (ms) before first connect — spreads TLS storm at startup + startup_jitter_ms: u32 = 0, rate_limiter: RateLimiter = .{}, // per-host shutdown (e.g. FutureCursor — stops only this subscriber) @@ -268,6 +270,11 @@ pub const Subscriber = struct { log.info("host {s}: resuming from cursor {d}", .{ self.options.hostname, seq }); } + // one-time startup jitter — spread DNS/TLS storm across a window + if (self.startup_jitter_ms > 0 and !self.shouldStop()) { + self.io.sleep(Io.Duration.fromMilliseconds(self.startup_jitter_ms), .awake) catch {}; + } + while (!self.shouldStop()) { log.info("host {s}: connecting...", .{self.options.hostname});