diff --git a/src/subscriber.zig b/src/subscriber.zig index 38eae7a..c4138d9 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -268,6 +268,11 @@ pub const Subscriber = struct { return self.shutdown.load(.acquire) or self.host_shutdown.load(.acquire); } + /// release connect gate permit (idempotent-safe: only decrements if gate exists) + fn releaseConnectGate(self: *Subscriber) void { + if (self.connect_gate) |gate| _ = gate.active.fetchSub(1, .monotonic); + } + /// run the subscriber loop. reconnects with exponential backoff. /// blocks until shutdown or host marked dormant. pub fn run(self: *Subscriber) void { @@ -279,7 +284,9 @@ pub const Subscriber = struct { } while (!self.shouldStop()) { - // wait for connect permit — limits concurrent DNS/TLS handshakes + // acquire connect permit — limits concurrent DNS/TLS handshakes. + // released immediately after connectAndRead returns (success or failure) + // so the permit is not held during backoff sleep. if (self.connect_gate) |gate| { while (gate.active.load(.acquire) >= gate.limit) { if (self.shouldStop()) return; @@ -287,16 +294,20 @@ pub const Subscriber = struct { } _ = 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| { - if (self.shouldStop()) return; + if (self.shouldStop()) { + self.releaseConnectGate(); + return; + } log.err("host {s}: error: {s}, reconnecting in {d}s...", .{ self.options.hostname, @errorName(err), self.backoff }); }; + // release permit after connectAndRead completes (success or failure) — + // connected subscribers reading frames don't count against the limit, + // and failed subscribers shouldn't hold permits during backoff sleep + self.releaseConnectGate(); if (self.shouldStop()) return;