diff --git a/src/main.zig b/src/main.zig index 4eadc51..2663e11 100644 --- a/src/main.zig +++ b/src/main.zig @@ -203,6 +203,7 @@ pub fn main() !void { const frame_workers = parseEnvInt(u16, "FRAME_WORKERS", 16); const frame_queue_capacity = parseEnvInt(u16, "FRAME_QUEUE_CAPACITY", 4096); const startup_batch_size = parseEnvInt(u16, "STARTUP_BATCH_SIZE", 50); + const max_concurrent_connects = parseEnvInt(u32, "MAX_CONCURRENT_CONNECTS", 50); const db_pool_size = parseEnvInt(u16, "DB_POOL_SIZE", 20); // install signal handlers (including SIGPIPE ignore) @@ -310,6 +311,7 @@ pub fn main() !void { .frame_workers = frame_workers, .frame_queue_capacity = frame_queue_capacity, .startup_batch_size = startup_batch_size, + .max_concurrent_connects = max_concurrent_connects, }, io, pool_io, diff --git a/src/slurper.zig b/src/slurper.zig index 9c3e99f..4a036fb 100644 --- a/src/slurper.zig +++ b/src/slurper.zig @@ -34,6 +34,9 @@ pub const Options = struct { /// prevents TLS handshake storm from starving the event loop. /// 0 = unlimited (legacy behavior). startup_batch_size: u16 = 50, + /// max concurrent DNS/TLS connection attempts across all subscriber threads. + /// bounds both the initial startup ramp and reconnect storms. + max_concurrent_connects: u32 = 50, }; // --- host validation --- @@ -241,6 +244,9 @@ pub const Slurper = struct { // shared TLS CA bundle — loaded once, used by all subscriber connections ca_bundle: ?std.crypto.Certificate.Bundle = null, + // connect gate — limits concurrent DNS/TLS handshakes across all subscribers + connect_gate: subscriber_mod.ConnectGate = .{ .limit = 50 }, + // active subscriber threads, keyed by host_id workers: std.AutoHashMapUnmanaged(u64, WorkerEntry) = .empty, workers_mutex: Io.Mutex = Io.Mutex.init, @@ -277,6 +283,7 @@ pub const Slurper = struct { .options = options, .io = io, .pool_io = pool_io, + .connect_gate = .{ .limit = options.max_concurrent_connects }, }; } @@ -578,6 +585,7 @@ pub const Slurper = struct { if (self.frame_pool) |*fp| sub.pool = fp; sub.pool_io = self.pool_io; sub.host_ops = self.host_ops; + sub.connect_gate = &self.connect_gate; if (self.cursor_map) |cm| { sub.cursor_map = cm; sub.cursor_slot = cm.register(host_id, last_seq); diff --git a/src/subscriber.zig b/src/subscriber.zig index 3dc904d..d78b35e 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -199,6 +199,13 @@ const RateLimiter = struct { } }; +/// semaphore that limits concurrent connection attempts across all subscriber threads. +/// bounds the DNS/TLS storm during startup and reconnect storms. +pub const ConnectGate = struct { + active: std.atomic.Value(u32) = .{ .raw = 0 }, + limit: u32, +}; + pub const Subscriber = struct { allocator: Allocator, io: Io, @@ -216,6 +223,8 @@ pub const Subscriber = struct { /// coalescing cursor map — subscriber writes latest seq atomically, worker sweeps every 5s cursor_map: ?*host_ops_mod.CursorMap = null, cursor_slot: ?u32 = null, + /// connect gate — limits concurrent DNS/TLS handshakes across all subscribers + connect_gate: ?*ConnectGate = null, shutdown: *std.atomic.Value(bool), last_upstream_seq: ?u64 = null, last_cursor_flush: i64 = 0, @@ -268,6 +277,18 @@ pub const Subscriber = struct { } while (!self.shouldStop()) { + // wait for connect permit — limits concurrent DNS/TLS handshakes + if (self.connect_gate) |gate| { + while (gate.active.load(.acquire) >= gate.limit) { + if (self.shouldStop()) return; + self.io.sleep(Io.Duration.fromMilliseconds(100), .awake) catch {}; + } + _ = gate.active.fetchAdd(1, .monotonic); + } + defer if (self.connect_gate) |gate| { + _ = gate.active.fetchSub(1, .monotonic); + }; + log.info("host {s}: connecting...", .{self.options.hostname}); self.connectAndRead() catch |err| {