//! websocket keepalive pinger — per-connection liveness for the zio backend. //! //! the 2026-08-18 and 2026-08-27 wedges were mass half-open TCP: the network //! path died without delivering FIN/RST, the kernel kept every connection //! ESTABLISHED, and level-triggered epoll stayed silent forever. capture //! forensics measured the precondition directly: ~2,900 subscriber sockets //! with no keepalive, no retransmits in flight, and no read deadline — //! nothing in the stack could distinguish a quiet peer from a dead one. //! //! websocket.zig's readLoopWithHeartbeat closes that gap under Io.Threaded, //! where SO_RCVTIMEO makes read() return after interval_ms. under zio it is //! structurally inert: netRead parks the fiber in epoll and the socket //! timeout never fires. this module is the zio-side replacement — the //! subscriber fiber runs a ping loop (subscriber.pingLoop) against the //! decision function here, while the read loop runs as a cancelable task. //! //! silence alone is NEVER grounds for a close: of ~3,000 connected hosts //! only ~120 deliver anything in a given 30s, and a silence-judging sweep //! tore production down on 2026-08-19 (9730b66). the discriminator is the //! pong: a quiet-but-alive peer answers pings, a half-open one cannot. //! sending the ping also arms the kernel's own detector — unacknowledged //! data hits the tcp_retries2 retransmission limit in ~15-25 minutes and //! errors the socket, a backstop that pure reading never provides. //! //! every knob is an atomic so /admin/ws-ping can change it at runtime //! without a restart. nothing is persisted: a restart returns to the env //! defaults. `{"enabled":false}` is the kill switch — pingers keep looping //! but never ping and never close, restoring today's (pre-fix) behavior. const std = @import("std"); const util = @import("util/util.zig"); pub const WsPing = struct { enabled: std.atomic.Value(bool), interval_sec: std.atomic.Value(u64), max_failures: std.atomic.Value(u64), pub fn init(enabled: bool, interval_sec: u64, max_failures: u64) WsPing { return .{ .enabled = .{ .raw = enabled }, .interval_sec = .{ .raw = interval_sec }, .max_failures = .{ .raw = max_failures }, }; } /// env defaults match indigo's slurper keepalive (30s ping, 4 missed /// intervals) — the reference relay's operationally proven cadence. pub fn fromEnv() WsPing { const enabled = if (util.getenv("RELAY_WS_PING_ENABLED")) |v| parseBool(v) orelse true else true; return init( enabled, util.parseEnvInt(u64, "RELAY_WS_PING_INTERVAL_SEC", 30), util.parseEnvInt(u64, "RELAY_WS_PING_MAX_FAILURES", 4), ); } pub fn formatJson(self: *const WsPing, buf: []u8) []const u8 { return std.fmt.bufPrint(buf, "{{\"enabled\":{},\"interval_sec\":{d},\"max_failures\":{d},\"persisted\":false}}", .{ self.enabled.load(.acquire), self.interval_sec.load(.acquire), self.max_failures.load(.acquire), }) catch "{\"error\":\"Internal\"}"; } }; pub const Decision = enum { /// connection proved life this interval (or the pinger is disabled, or /// the reader is blocked in frame processing) — reset the failure count. idle, /// silent past the interval with pings outstanding but below the limit: /// send another ping. ping, /// max_failures pings went unanswered — declare the connection half-open /// and tear it down. kill, }; /// one pinger-tick verdict. pure so the policy is testable without sockets. /// /// `in_handler` guards the rate-limit/backpressure case: a reader blocked /// inside serverMessage (waitForAllow can hold a day-limited host for hours, /// by design) is not reading pongs, so unanswered pings prove nothing. this /// mirrors indigo/SO_RCVTIMEO semantics, where the deadline only ticks while /// actually blocked in a read. pub fn decide( cfg: *const WsPing, now_ms: i64, activity_ms: i64, in_handler: bool, pending: u64, ) Decision { if (!cfg.enabled.load(.acquire)) return .idle; if (in_handler) return .idle; const interval_ms: i64 = @intCast(cfg.interval_sec.load(.acquire) * std.time.ms_per_s); if (now_ms - activity_ms < interval_ms) return .idle; if (pending >= cfg.max_failures.load(.acquire)) return .kill; return .ping; } fn parseBool(v: []const u8) ?bool { if (std.mem.eql(u8, v, "1") or std.ascii.eqlIgnoreCase(v, "true")) return true; if (std.mem.eql(u8, v, "0") or std.ascii.eqlIgnoreCase(v, "false")) return false; return null; } const ms = std.time.ms_per_s; test "disabled pinger never pings or kills" { const cfg = WsPing.init(false, 30, 4); try std.testing.expectEqual(Decision.idle, decide(&cfg, 1000 * ms, 0, false, 0)); try std.testing.expectEqual(Decision.idle, decide(&cfg, 1000 * ms, 0, false, 99)); } test "recent activity is idle and resets" { const cfg = WsPing.init(true, 30, 4); // activity 10s ago, interval 30s — alive try std.testing.expectEqual(Decision.idle, decide(&cfg, 100 * ms, 90 * ms, false, 3)); } test "silence past interval pings until max_failures then kills" { const cfg = WsPing.init(true, 30, 4); const now = 1000 * ms; const stale = now - 31 * ms; try std.testing.expectEqual(Decision.ping, decide(&cfg, now, stale, false, 0)); try std.testing.expectEqual(Decision.ping, decide(&cfg, now, stale, false, 3)); try std.testing.expectEqual(Decision.kill, decide(&cfg, now, stale, false, 4)); } test "reader blocked in handler is never judged" { const cfg = WsPing.init(true, 30, 4); const now = 10_000 * ms; // hours of silence with pings outstanding — but the reader is inside // serverMessage (rate-limit block), so pongs cannot be read. idle. try std.testing.expectEqual(Decision.idle, decide(&cfg, now, 0, true, 4)); } test "a pong between pings resets via idle" { const cfg = WsPing.init(true, 30, 4); // pending is 3, but activity is fresh: the caller sees idle and resets. try std.testing.expectEqual(Decision.idle, decide(&cfg, 200 * ms, 195 * ms, false, 3)); } test "runtime knob changes take effect immediately" { var cfg = WsPing.init(true, 30, 4); const now = 1000 * ms; const stale = now - 31 * ms; try std.testing.expectEqual(Decision.ping, decide(&cfg, now, stale, false, 0)); cfg.interval_sec.store(60, .release); try std.testing.expectEqual(Decision.idle, decide(&cfg, now, stale, false, 0)); cfg.interval_sec.store(30, .release); cfg.max_failures.store(1, .release); try std.testing.expectEqual(Decision.kill, decide(&cfg, now, stale, false, 1)); } test "formatJson carries the knobs" { const cfg = WsPing.init(true, 30, 4); var buf: [256]u8 = undefined; try std.testing.expectEqualStrings( "{\"enabled\":true,\"interval_sec\":30,\"max_failures\":4,\"persisted\":false}", cfg.formatJson(&buf), ); }