atproto relay in zig zlay.waow.tech
relay zig atproto
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163//! 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), );}