diff --git a/docs/deployment.md b/docs/deployment.md index 31a58de..dd940cd 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -71,6 +71,31 @@ change the defaults, set the `RELAY_INGEST_STALL_*` env vars in `deploy/zlay-values.yaml`. `relay_ingest_stalled`, `relay_ingest_rate_fps`, `relay_ingest_below_threshold_seconds` and `relay_ingest_stall_enabled` are on `/metrics`; when `stalled` is 1 the 503 body carries the rate and window. + +**the keepalive pinger detects half-open connections.** the 2026-08-18 and +2026-08-27 wedges were mass half-open TCP: the network path died without +delivering FIN/RST and nothing in the stack could tell a quiet peer from a +dead one. each subscriber connection now pings its upstream after +`RELAY_WS_PING_INTERVAL_SEC` (30s) of silence and reconnects after +`RELAY_WS_PING_MAX_FAILURES` (4) unanswered pings — any received frame +(data or pong) proves life, so quiet-but-alive hosts are never touched. +under zio this runs as a per-connection pinger with the read loop as a +cancelable task (SO_RCVTIMEO, which drives the Threaded path's heartbeat, +never fires on an epoll-parked fiber). knobs are runtime-settable, not +persisted, same contract as the stall detector: + +```bash +curl -H "Authorization: Bearer $RELAY_ADMIN_PASSWORD" https://zlay.waow.tech/admin/ws-ping +# kill switch (reverts to pre-fix no-keepalive behavior): +curl -X POST -H "Authorization: Bearer $RELAY_ADMIN_PASSWORD" \ + https://zlay.waow.tech/admin/ws-ping -d '{"enabled":false}' +# retune: any subset of {"interval_sec":N,"max_failures":N} +``` + +on `/metrics`: `relay_ws_pings_sent_total` (baseline traffic — quiet hosts +make this tick steadily), `relay_ws_ping_timeout_closes_total` (each is a +connection declared half-open; a mass spike is the 08-18 signature being +caught in the act) and `relay_ws_ping_write_failures_total`. - `-Dtarget=x86_64-linux-gnu` — production target, glibc. glibc malloc (per-thread arenas + `madvise` page-return + `malloc_trim`) is load-bearing for RSS at ~2,800 threads. **NB:** the old "musl breaks RocksDB / illegal instructions" claim was a zig 0.15 artifact and is **retired** — on 0.16 musl builds, links, and runs SIGILL-free at thread scale (canary, 2026-06). musl stays off prod only because the one static-musl allocator we validated as buildable — mimalloc — leaks RSS under our thread model (v2/v3/purge-forced/even under glibc). full matrix in [musl-investigation.md](musl-investigation.md). - `-Dcpu=baseline` — required when building inside Docker/QEMU (not needed for `zlay-publish-remote` since it builds natively). - `-Doptimize=ReleaseSafe` — safety checks on, optimizations on. production default since 2026-03-05. previously caused OOM (see [incident-2026-03-04.md](incident-2026-03-04.md)) — resolved by the frame pool moving heavy work off reader threads. diff --git a/src/internal/api/admin.zig b/src/internal/api/admin.zig index de053b2..7608c0d 100644 --- a/src/internal/api/admin.zig +++ b/src/internal/api/admin.zig @@ -682,3 +682,62 @@ pub fn handleIngestStallUpdate(conn: *h.Conn, body: []const u8, headers: *const var buf: [512]u8 = undefined; h.respondJson(conn, .ok, il.formatJson(&buf, util.milliTimestamp(ctx.io))); } + +/// GET /admin/ws-ping — the websocket keepalive pinger's knobs. +pub fn handleWsPingStatus(conn: *h.Conn, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { + if (!checkAdmin(conn, headers)) return; + var buf: [256]u8 = undefined; + h.respondJson(conn, .ok, ctx.ws_ping.formatJson(&buf)); +} + +/// POST /admin/ws-ping +/// body: any subset of {"enabled":bool,"interval_sec":N,"max_failures":N} +/// +/// takes effect on each pinger's next tick and is NOT persisted: a restart +/// returns to the RELAY_WS_PING_* env defaults. `{"enabled":false}` is the +/// kill switch — connections revert to no-keepalive behavior. +pub fn handleWsPingUpdate(conn: *h.Conn, body: []const u8, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { + if (!checkAdmin(conn, headers)) return; + + const Update = struct { + enabled: ?bool = null, + interval_sec: ?u64 = null, + max_failures: ?u64 = null, + }; + const parsed = std.json.parseFromSlice(Update, ctx.persist.allocator, body, .{ .ignore_unknown_fields = true }) catch { + h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"expected JSON with any of enabled, interval_sec, max_failures\"}"); + return; + }; + defer parsed.deinit(); + const u = parsed.value; + + if (u.enabled == null and u.interval_sec == null and u.max_failures == null) { + h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"no fields to update\"}"); + return; + } + if (u.interval_sec) |v| if (v == 0) { + h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"interval_sec must be positive\"}"); + return; + }; + if (u.max_failures) |v| if (v == 0) { + h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"max_failures must be positive\"}"); + return; + }; + + const wp = ctx.ws_ping; + if (u.enabled) |v| { + wp.enabled.store(v, .release); + log.warn("admin: ws keepalive pinger {s} (not persisted; env default returns on restart)", .{if (v) "enabled" else "DISABLED"}); + } + if (u.interval_sec) |v| { + wp.interval_sec.store(v, .release); + log.warn("admin: ws keepalive interval set to {d}s (not persisted)", .{v}); + } + if (u.max_failures) |v| { + wp.max_failures.store(v, .release); + log.warn("admin: ws keepalive max_failures set to {d} (not persisted)", .{v}); + } + + var buf: [256]u8 = undefined; + h.respondJson(conn, .ok, wp.formatJson(&buf)); +} diff --git a/src/internal/api/router.zig b/src/internal/api/router.zig index 71ce1b9..4b827f7 100644 --- a/src/internal/api/router.zig +++ b/src/internal/api/router.zig @@ -16,6 +16,7 @@ const cleaner_mod = @import("../collection_index/cleaner.zig"); const resync_mod = @import("../collection_index/resync.zig"); const host_ops_mod = @import("../host_ops.zig"); const ingest_liveness_mod = @import("../ingest_liveness.zig"); +const ws_ping_mod = @import("../ws_ping.zig"); const util = @import("../util/util.zig"); const h = @import("http.zig"); const xrpc = @import("xrpc.zig"); @@ -38,6 +39,7 @@ pub const HttpContext = struct { db_queue: *event_log_mod.DbRequestQueue, shutdown: *std.atomic.Value(bool), ingest_liveness: *ingest_liveness_mod.IngestLiveness, + ws_ping: *ws_ping_mod.WsPing, }; /// top-level HTTP request router — installed as bc.http_fallback @@ -109,6 +111,8 @@ fn handleGet(conn: *websocket.Conn, path: []const u8, query: []const u8, headers admin.handleResyncStatus(conn, headers, ctx.resyncer); } else if (std.mem.eql(u8, path, "/admin/ingest-stall")) { admin.handleIngestStallStatus(conn, headers, ctx); + } else if (std.mem.eql(u8, path, "/admin/ws-ping")) { + admin.handleWsPingStatus(conn, headers, ctx); } else if (std.mem.eql(u8, path, "/")) { h.respondText(conn, .ok, \\ _ @@ -157,6 +161,8 @@ fn handlePost(conn: *websocket.Conn, path: []const u8, query: []const u8, body: admin.handleResyncTrigger(conn, body, headers, ctx.resyncer); } else if (std.mem.eql(u8, path, "/admin/ingest-stall")) { admin.handleIngestStallUpdate(conn, body, headers, ctx); + } else if (std.mem.eql(u8, path, "/admin/ws-ping")) { + admin.handleWsPingUpdate(conn, body, headers, ctx); } else { h.respondText(conn, .not_found, "not found"); } diff --git a/src/internal/broadcaster.zig b/src/internal/broadcaster.zig index 894fd60..315bd4b 100644 --- a/src/internal/broadcaster.zig +++ b/src/internal/broadcaster.zig @@ -64,6 +64,11 @@ pub const Stats = struct { chain_breaks_since: std.atomic.Value(u64) = .{ .raw = 0 }, chain_breaks_prev_data: std.atomic.Value(u64) = .{ .raw = 0 }, pool_backpressure: std.atomic.Value(u64) = .{ .raw = 0 }, + // keepalive pinger (ws_ping.zig) — pings sent, connections torn down for + // unanswered pings, and ping writes that themselves failed + ws_pings_sent: std.atomic.Value(u64) = .{ .raw = 0 }, + ws_ping_timeout_closes: std.atomic.Value(u64) = .{ .raw = 0 }, + ws_ping_write_failures: std.atomic.Value(u64) = .{ .raw = 0 }, // host authority resolution metrics host_authority_checks: std.atomic.Value(u64) = .{ .raw = 0 }, host_authority_is_new: std.atomic.Value(u64) = .{ .raw = 0 }, @@ -1117,6 +1122,26 @@ pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, attrib stats.pool_queued_bytes.load(.acquire), }) catch return w.buffered(); + // keepalive pinger metrics (separate print to stay under 32-arg limit) + w.print( + \\# TYPE relay_ws_pings_sent_total counter + \\# HELP relay_ws_pings_sent_total keepalive pings sent to silent upstream connections + \\relay_ws_pings_sent_total {d} + \\ + \\# TYPE relay_ws_ping_timeout_closes_total counter + \\# HELP relay_ws_ping_timeout_closes_total connections closed as half-open after max unanswered keepalive pings + \\relay_ws_ping_timeout_closes_total {d} + \\ + \\# TYPE relay_ws_ping_write_failures_total counter + \\# HELP relay_ws_ping_write_failures_total keepalive ping writes that failed (connection torn down) + \\relay_ws_ping_write_failures_total {d} + \\ + , .{ + stats.ws_pings_sent.load(.acquire), + stats.ws_ping_timeout_closes.load(.acquire), + stats.ws_ping_write_failures.load(.acquire), + }) catch return w.buffered(); + w.print( \\# TYPE relay_host_authority_cache_hits_total counter \\# HELP relay_host_authority_cache_hits_total confirmed host mismatches rejected without network resolution diff --git a/src/internal/slurper.zig b/src/internal/slurper.zig index 82f2614..f574d0a 100644 --- a/src/internal/slurper.zig +++ b/src/internal/slurper.zig @@ -22,6 +22,7 @@ const resync_mod = @import("collection_index/resync.zig"); const status_checker_mod = @import("status_checker.zig"); const frame_worker_mod = @import("frame_worker.zig"); const host_ops_mod = @import("host_ops.zig"); +const ws_ping_mod = @import("ws_ping.zig"); const atproto = @import("atproto/main.zig"); const milliTimestamp = @import("util/util.zig").milliTimestamp; @@ -59,6 +60,7 @@ pub const Slurper = struct { host_ops: ?*host_ops_mod.HostOpsQueue = null, cursor_map: ?*host_ops_mod.CursorMap = null, db_queue: ?*event_log_mod.DbRequestQueue = null, + ws_ping: ?*ws_ping_mod.WsPing = null, shutdown: *std.atomic.Value(bool), options: Options, @@ -415,6 +417,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.ws_ping = self.ws_ping; if (self.cursor_map) |cm| { sub.cursor_map = cm; sub.cursor_slot = cm.register(host_id, last_seq); diff --git a/src/internal/subscriber.zig b/src/internal/subscriber.zig index 24829ae..1e43259 100644 --- a/src/internal/subscriber.zig +++ b/src/internal/subscriber.zig @@ -7,6 +7,7 @@ //! managed by the Slurper, which spawns one Subscriber per tracked host. const std = @import("std"); +const backend_config = @import("backend_config"); const websocket = @import("websocket"); const zat = @import("zat"); const broadcaster = @import("broadcaster.zig"); @@ -17,6 +18,7 @@ const resync_mod = @import("collection_index/resync.zig"); const status_checker_mod = @import("status_checker.zig"); const frame_worker_mod = @import("frame_worker.zig"); const host_ops_mod = @import("host_ops.zig"); +const ws_ping_mod = @import("ws_ping.zig"); const util = @import("util/util.zig"); const Allocator = std.mem.Allocator; @@ -223,6 +225,10 @@ pub const Subscriber = struct { /// to fall back on. A genuinely unreachable host keeps attempting, so it /// keeps stamping and stays inside the idle guard. last_progress_ms: std.atomic.Value(i64) = .{ .raw = 0 }, + /// keepalive pinger knobs, shared across all subscribers and settable at + /// runtime via /admin/ws-ping. null (tests, standalone) disables the + /// zio-side pinger; the Threaded path is unaffected either way. + ws_ping: ?*ws_ping_mod.WsPing = null, rate_limiter: RateLimiter = .{}, // per-host shutdown (e.g. FutureCursor — stops only this subscriber) @@ -412,33 +418,144 @@ pub const Subscriber = struct { } } + var conn = ConnState{ + .activity_ms = .{ .raw = milliTimestamp(self.io) }, + }; var handler = FrameHandler{ .subscriber = self, + .conn = &conn, + }; + + const heartbeat: websocket.Client.HeartbeatConfig = .{ + .interval_ms = @intCast(@max(1, if (self.ws_ping) |p| p.interval_sec.load(.acquire) else ping_interval_sec) * 1000), + .max_failures = @intCast(if (self.ws_ping) |p| p.max_failures.load(.acquire) else max_ping_failures), }; - // 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. // 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| { + if (comptime backend_config.use_zio) { + if (self.ws_ping != null) { + // under zio, readLoopWithHeartbeat's SO_RCVTIMEO mechanism is + // inert: netRead parks the fiber in epoll and the socket + // timeout never fires (the 2026-08-18/27 half-open wedges). + // run the read loop as a cancelable task and drive liveness + // from this fiber instead. cancellation (not a cross-fiber + // socket close) is the teardown: the reader observes + // cancel_requested before ever touching the fd again, so + // there is no closed-fd-number-reuse race, and the fd is + // closed exactly once by client.deinit after the task ends. + var read_future = self.io.concurrent(readLoopTask, .{ self, &client, &handler, heartbeat, &conn }) catch |e| { + _ = self.bc.stats.subscriber_disconnect_read_loop.fetchAdd(1, .monotonic); + return e; + }; + self.pingLoop(&client, &conn); + return read_future.cancel(self.io); + } + } + + // Threaded backend (and pinger-less standalone/test use): SO_RCVTIMEO + // works, so the merged keepalive + read loop covers liveness alone. + client.readLoopWithHeartbeat(&handler, heartbeat) catch |e| { + _ = self.bc.stats.subscriber_disconnect_read_loop.fetchAdd(1, .monotonic); + return e; + }; + } + + fn readLoopTask(self: *Subscriber, client: *websocket.Client, handler: *FrameHandler, heartbeat: websocket.Client.HeartbeatConfig, conn: *ConnState) anyerror!void { + defer conn.read_done.store(true, .release); + client.readLoopWithHeartbeat(handler, heartbeat) catch |e| { _ = self.bc.stats.subscriber_disconnect_read_loop.fetchAdd(1, .monotonic); return e; }; } + + /// zio keepalive: runs in the subscriber fiber while readLoopTask reads. + /// pings when the connection has been silent past the interval, tears it + /// down after max_failures unanswered pings. see ws_ping.zig for why + /// silence alone never kills and how the reader-in-handler guard works. + /// + /// a ping write that parks on a full send buffer stalls only this + /// connection's pinger, and the kernel's retransmission limit + /// (tcp_retries2, ~15-25 min) errors the socket underneath it — the + /// unanswered bytes themselves are the backstop. + fn pingLoop(self: *Subscriber, client: *websocket.Client, conn: *ConnState) void { + const cfg = self.ws_ping.?; + var pending: u64 = 0; + var empty: [0]u8 = .{}; + outer: while (true) { + // sleep the interval in 1s chunks so read-loop exit and shutdown + // are noticed promptly (bounds reconnect latency to ~1s) + var slept: u64 = 0; + while (slept < @max(1, cfg.interval_sec.load(.acquire))) : (slept += 1) { + if (conn.read_done.load(.acquire) or self.shouldStop()) break :outer; + self.io.sleep(Io.Duration.fromSeconds(1), .awake) catch {}; + } + switch (ws_ping_mod.decide( + cfg, + milliTimestamp(self.io), + conn.activity_ms.load(.acquire), + conn.in_handler.load(.acquire), + pending, + )) { + .idle => pending = 0, + .ping => { + client.writePing(&empty) catch { + _ = self.bc.stats.ws_ping_write_failures.fetchAdd(1, .monotonic); + log.info("host {s}: keepalive ping write failed, closing", .{self.options.hostname}); + break; + }; + _ = self.bc.stats.ws_pings_sent.fetchAdd(1, .monotonic); + pending += 1; + }, + .kill => { + _ = self.bc.stats.ws_ping_timeout_closes.fetchAdd(1, .monotonic); + log.warn("host {s}: {d} keepalive pings unanswered, closing half-open connection", .{ self.options.hostname, pending }); + break; + }, + } + } + } +}; + +/// per-connection liveness state shared between the read-loop task and the +/// pinger. all fields are atomics for the Threaded test path; under zio both +/// sides are fibers on the one loop thread. +const ConnState = struct { + /// ms timestamp of the last frame received on THIS connection (any type — + /// data, pong). distinct from Subscriber.last_progress_ms, which spans + /// reconnects and stamps connect attempts for the admin reconnect fiber. + activity_ms: std.atomic.Value(i64), + /// reader is inside serverMessage (possibly blocked in rate-limit or + /// pool backpressure for a long time, by design) — pongs cannot be read, + /// so the pinger must not judge the connection. + in_handler: std.atomic.Value(bool) = .{ .raw = false }, + read_done: std.atomic.Value(bool) = .{ .raw = false }, }; const FrameHandler = struct { subscriber: *Subscriber, + conn: *ConnState, + + /// any frame from the peer proves the link — pongs answer our keepalive + /// pings, so this is what resets the pinger's failure count. + pub fn serverPong(self: *FrameHandler, data: []u8) !void { + _ = data; + self.conn.activity_ms.store(milliTimestamp(self.subscriber.io), .release); + } pub fn serverMessage(self: *FrameHandler, data: []const u8) !void { const sub = self.subscriber; const io = sub.io; + self.conn.activity_ms.store(milliTimestamp(io), .release); + self.conn.in_handler.store(true, .release); + defer { + // time blocked inside the handler must not count as silence + self.conn.in_handler.store(false, .release); + self.conn.activity_ms.store(milliTimestamp(io), .release); + } + // lightweight header decode for cursor tracking + routing var arena = std.heap.ArenaAllocator.init(sub.allocator); defer arena.deinit(); diff --git a/src/internal/ws_ping.zig b/src/internal/ws_ping.zig new file mode 100644 index 0000000..0865067 --- /dev/null +++ b/src/internal/ws_ping.zig @@ -0,0 +1,162 @@ +//! 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), + ); +} diff --git a/src/main.zig b/src/main.zig index 5e73335..49520f1 100644 --- a/src/main.zig +++ b/src/main.zig @@ -18,6 +18,7 @@ //! /admin/hosts/block — block a host (POST, admin) //! /admin/hosts/unblock — unblock a host (POST, admin) //! /admin/ingest-stall — liveness stall detector knobs (GET/POST, admin) +//! /admin/ws-ping — websocket keepalive pinger knobs (GET/POST, admin) //! /_health, /_stats — health, stats //! //! port 3001 (RELAY_METRICS_PORT): internal metrics + health @@ -40,6 +41,7 @@ const resync_mod = @import("internal/collection_index/resync.zig"); const status_checker_mod = @import("internal/status_checker.zig"); const host_ops_mod = @import("internal/host_ops.zig"); const ingest_liveness_mod = @import("internal/ingest_liveness.zig"); +const ws_ping_mod = @import("internal/ws_ping.zig"); const api = @import("internal/api/main.zig"); const util = @import("internal/util/util.zig"); const build_options = @import("build_options"); @@ -375,6 +377,9 @@ fn runRelay() !void { log.info("ingest stall liveness: enabled={} threshold={d}fps window={d}s startup_grace={d}s", .{ ingest_liveness.enabled.load(.acquire), ingest_liveness.threshold_fps.load(.acquire), ingest_liveness.window_sec.load(.acquire), ingest_liveness.startup_grace_sec.load(.acquire) }); // validator uses pool_io — its cache LRU and host resolvers are called from worker threads + var ws_ping = ws_ping_mod.WsPing.fromEnv(); + log.info("ws keepalive pinger: enabled={} interval={d}s max_failures={d}", .{ ws_ping.enabled.load(.acquire), ws_ping.interval_sec.load(.acquire), ws_ping.max_failures.load(.acquire) }); + var val = validator_mod.Validator.init(allocator, &bc.stats, pool_io); defer val.deinit(); try val.start(); @@ -489,6 +494,7 @@ fn runRelay() !void { slurper.host_ops = &host_ops_queue; slurper.cursor_map = &cursor_map; slurper.db_queue = &db_queue; + slurper.ws_ping = &ws_ping; // start: loads active hosts from DB, spawns subscriber threads. // pullHosts runs on its own std.Thread (parallel, not gating spawnWorkers). @@ -540,6 +546,7 @@ fn runRelay() !void { .db_queue = &db_queue, .shutdown = &shutdown_flag, .ingest_liveness = &ingest_liveness, + .ws_ping = &ws_ping, }; bc.http_fallback = api.handleHttpRequest; bc.http_fallback_ctx = @ptrCast(&http_context); diff --git a/src/tests.zig b/src/tests.zig index bc0d2bf..cc13862 100644 --- a/src/tests.zig +++ b/src/tests.zig @@ -15,6 +15,7 @@ test { _ = @import("internal/event_log.zig"); _ = @import("internal/slurper.zig"); _ = @import("internal/ingest_liveness.zig"); + _ = @import("internal/ws_ping.zig"); _ = @import("internal/frame_worker.zig"); _ = @import("internal/status_checker.zig"); _ = @import("internal/collection_index/index.zig");