diff --git a/src/broadcaster.zig b/src/broadcaster.zig index a14ed46..70c1417 100644 --- a/src/broadcaster.zig +++ b/src/broadcaster.zig @@ -65,6 +65,15 @@ pub const Stats = struct { host_authority_reject_bad_url: std.atomic.Value(u64) = .{ .raw = 0 }, host_authority_reject_unknown_host: std.atomic.Value(u64) = .{ .raw = 0 }, host_authority_reject_host_mismatch: std.atomic.Value(u64) = .{ .raw = 0 }, + // subscriber reconnect-cycle disconnect breakdown. only handshake-or-earlier + // failures accumulate to failed_attempts (reset_failures fires after + // ws_handshake succeeds, before read_loop). 15 consecutive of these → + // host moves to 'exhausted'. read_loop disconnects do not directly + // exhaust unless the next handshake also fails. + subscriber_disconnect_dns_connect: std.atomic.Value(u64) = .{ .raw = 0 }, + subscriber_disconnect_tls_init: std.atomic.Value(u64) = .{ .raw = 0 }, + subscriber_disconnect_ws_handshake: std.atomic.Value(u64) = .{ .raw = 0 }, + subscriber_disconnect_read_loop: std.atomic.Value(u64) = .{ .raw = 0 }, // frame pool memory pressure pool_queued_bytes: std.atomic.Value(u64) = .{ .raw = 0 }, // persist/broadcast pipeline contention @@ -1094,6 +1103,13 @@ pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, attrib \\relay_host_authority_reject{{branch="unknown_host"}} {d} \\relay_host_authority_reject{{branch="host_mismatch"}} {d} \\ + \\# TYPE relay_subscriber_disconnect_total counter + \\# HELP relay_subscriber_disconnect_total subscriber reconnect-cycle disconnects by phase + \\relay_subscriber_disconnect_total{{phase="dns_connect"}} {d} + \\relay_subscriber_disconnect_total{{phase="tls_init"}} {d} + \\relay_subscriber_disconnect_total{{phase="ws_handshake"}} {d} + \\relay_subscriber_disconnect_total{{phase="read_loop"}} {d} + \\ , .{ stats.failed_bad_did.load(.acquire), stats.failed_bad_rev.load(.acquire), @@ -1108,6 +1124,10 @@ pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, attrib stats.host_authority_reject_bad_url.load(.acquire), stats.host_authority_reject_unknown_host.load(.acquire), stats.host_authority_reject_host_mismatch.load(.acquire), + stats.subscriber_disconnect_dns_connect.load(.acquire), + stats.subscriber_disconnect_tls_init.load(.acquire), + stats.subscriber_disconnect_ws_handshake.load(.acquire), + stats.subscriber_disconnect_read_loop.load(.acquire), }) catch return w.buffered(); // memory attribution — internal capacities help identify what's consuming RSS diff --git a/src/subscriber.zig b/src/subscriber.zig index 3dc904d..36b7a7b 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -322,17 +322,28 @@ pub const Subscriber = struct { // DNS + TCP connect through pool_io (Threaded — has working netLookup). // The resulting fd is used by the Evented io for TLS + WebSocket I/O. + // each phase labels its catch site so operators can see which layer + // is failing across the fleet via relay_subscriber_disconnect_total. const dns_io = self.pool_io orelse self.io; - const host_name = try Io.net.HostName.init(self.options.hostname); - const net_stream = try host_name.connect(dns_io, 443, .{ .mode = .stream }); + const host_name = Io.net.HostName.init(self.options.hostname) catch |e| { + _ = self.bc.stats.subscriber_disconnect_dns_connect.fetchAdd(1, .monotonic); + return e; + }; + const net_stream = host_name.connect(dns_io, 443, .{ .mode = .stream }) catch |e| { + _ = self.bc.stats.subscriber_disconnect_dns_connect.fetchAdd(1, .monotonic); + return e; + }; - var client = try websocket.Client.initWithStream(self.io, self.allocator, net_stream, .{ + var client = websocket.Client.initWithStream(self.io, self.allocator, net_stream, .{ .host = self.options.hostname, .port = 443, .tls = true, .max_size = self.options.max_message_size, .ca_bundle = self.options.ca_bundle, - }); + }) catch |e| { + _ = self.bc.stats.subscriber_disconnect_tls_init.fetchAdd(1, .monotonic); + return e; + }; defer client.deinit(); var host_header_buf: [256]u8 = undefined; @@ -342,7 +353,10 @@ pub const Subscriber = struct { .{self.options.hostname}, ) catch self.options.hostname; - try client.handshake(path, .{ .headers = host_header }); + client.handshake(path, .{ .headers = host_header }) catch |e| { + _ = self.bc.stats.subscriber_disconnect_ws_handshake.fetchAdd(1, .monotonic); + return e; + }; log.info("host {s}: connected", .{self.options.hostname}); // reset failures on successful connect (via host_ops queue) @@ -363,10 +377,16 @@ pub const Subscriber = struct { // heartbeat read loop: merged keepalive + frame reading in a single task. // SO_RCVTIMEO fires after interval_ms, triggering a ping. closes after // max_failures consecutive idle intervals with no frames received. - try client.readLoopWithHeartbeat(&handler, .{ + // read_loop disconnects only count toward `failed_attempts` if the + // *next* handshake also fails — reset_failures was already pushed + // above on successful handshake. + client.readLoopWithHeartbeat(&handler, .{ .interval_ms = ping_interval_sec * 1000, .max_failures = max_ping_failures, - }); + }) catch |e| { + _ = self.bc.stats.subscriber_disconnect_read_loop.fetchAdd(1, .monotonic); + return e; + }; } };