atproto relay in zig zlay.waow.tech
relay zig atproto
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387//! ingest liveness — the process-wide stall detector behind /_healthz.//!//! samples relay_frames_received_total on a fixed interval and calls the//! process stalled once the whole-process ingest rate has stayed below a//! threshold for a full window. the judgement is process-wide on purpose://! of ~3,000 connected hosts only ~120 deliver anything in a given 30s, so//! per-worker silence is the normal state of a small PDS, and a sweep that//! judged workers individually tore production down on 2026-08-19 (9730b66).//!//! every knob is an atomic so /admin/ingest-stall can change it at runtime//! without a restart. nothing is persisted: a restart returns to the env//! defaults. that is deliberate — the kill switch exists so an operator can//! hold a wedged pod up for forensics, and a forgotten "disabled" must not//! outlive the incident.
const std = @import("std");const Io = std.Io;const broadcaster = @import("broadcaster.zig");const util = @import("util/util.zig");
const log = std.log.scoped(.relay);
pub const sample_interval_sec: u64 = 10;
/// a sample older than this is treated as no sample at all, so a dead/// sampler thread can never leave a stale "stalled" verdict in place.const sample_stale_ms: i64 = 3 * sample_interval_sec * std.time.ms_per_s;
pub const Verdict = struct { stalled: bool, enabled: bool, rate_milli_fps: u64, below_sec: u64, threshold_fps: u64, window_sec: u64,};
pub const IngestLiveness = struct { enabled: std.atomic.Value(bool), threshold_fps: std.atomic.Value(u64), window_sec: std.atomic.Value(u64), startup_grace_sec: std.atomic.Value(u64), started_ms: i64,
last_sample_ms: std.atomic.Value(i64) = .{ .raw = 0 }, last_frames: std.atomic.Value(u64) = .{ .raw = 0 }, below_since_ms: std.atomic.Value(i64) = .{ .raw = 0 }, rate_milli_fps: std.atomic.Value(u64) = .{ .raw = 0 },
pub fn init(started_ms: i64, enabled: bool, threshold_fps: u64, window_sec: u64, startup_grace_sec: u64) IngestLiveness { return .{ .enabled = .{ .raw = enabled }, .threshold_fps = .{ .raw = threshold_fps }, .window_sec = .{ .raw = window_sec }, .startup_grace_sec = .{ .raw = startup_grace_sec }, .started_ms = started_ms, }; }
/// env defaults. the threshold matches the ZlayIngestStalled prometheus /// rule (<50 fps over 10m, for 10m) and the window is longer than that /// rule's total lead time, so a human is paged before the pod is killed. pub fn fromEnv(started_ms: i64) IngestLiveness { const enabled = if (util.getenv("RELAY_INGEST_STALL_ENABLED")) |v| parseBool(v) orelse true else true; return init( started_ms, enabled, util.parseEnvInt(u64, "RELAY_INGEST_STALL_THRESHOLD_FPS", 50), util.parseEnvInt(u64, "RELAY_INGEST_STALL_WINDOW_SEC", 900), util.parseEnvInt(u64, "RELAY_INGEST_STALL_STARTUP_GRACE_SEC", 900), ); }
/// feed one sample of the process-wide frame counter. the sampler thread /// is the only caller in production; tests drive it directly. pub fn observe(self: *IngestLiveness, now_ms: i64, frames_total: u64) void { const prev_ms = self.last_sample_ms.load(.acquire); const prev_frames = self.last_frames.load(.acquire); self.last_frames.store(frames_total, .release); self.last_sample_ms.store(now_ms, .release);
if (!self.enabled.load(.acquire) or self.inStartupGrace(now_ms)) { self.below_since_ms.store(0, .release); return; } if (prev_ms == 0 or now_ms <= prev_ms) return; if (now_ms - prev_ms > sample_stale_ms) { self.below_since_ms.store(0, .release); return; }
const elapsed_ms: u64 = @intCast(now_ms - prev_ms); const delta = frames_total -| prev_frames; const rate_milli = delta * std.time.ms_per_s * std.time.ms_per_s / elapsed_ms; self.rate_milli_fps.store(rate_milli, .release);
if (rate_milli < self.threshold_fps.load(.acquire) * std.time.ms_per_s) { if (self.below_since_ms.load(.acquire) == 0) self.below_since_ms.store(prev_ms, .release); } else { self.below_since_ms.store(0, .release); } }
/// what /_healthz acts on. the window is applied here rather than at /// sample time so widening it mid-incident takes effect immediately. pub fn verdict(self: *const IngestLiveness, now_ms: i64) Verdict { const enabled = self.enabled.load(.acquire); const threshold = self.threshold_fps.load(.acquire); const window = self.window_sec.load(.acquire); const below_since = self.below_since_ms.load(.acquire); const last_sample = self.last_sample_ms.load(.acquire);
var below_sec: u64 = 0; if (below_since != 0 and now_ms > below_since) below_sec = @intCast(@divFloor(now_ms - below_since, std.time.ms_per_s));
const fresh = last_sample != 0 and now_ms - last_sample <= sample_stale_ms; const stalled = enabled and fresh and !self.inStartupGrace(now_ms) and below_since != 0 and below_sec >= window;
return .{ .stalled = stalled, .enabled = enabled, .rate_milli_fps = self.rate_milli_fps.load(.acquire), .below_sec = below_sec, .threshold_fps = threshold, .window_sec = window, }; }
fn inStartupGrace(self: *const IngestLiveness, now_ms: i64) bool { const grace_ms: i64 = @intCast(self.startup_grace_sec.load(.acquire) * std.time.ms_per_s); return now_ms - self.started_ms < grace_ms; }
/// sampler loop — a plain thread on pool_io, like gcLoop. pub fn run(self: *IngestLiveness, io: Io, stats: *const broadcaster.Stats, shutdown: *std.atomic.Value(bool)) void { var was_stalled = false; while (!shutdown.load(.acquire)) { var remaining: u64 = sample_interval_sec; while (remaining > 0 and !shutdown.load(.acquire)) { io.sleep(Io.Duration.fromSeconds(1), .awake) catch return; remaining -= 1; } if (shutdown.load(.acquire)) return;
const now_ms = util.milliTimestamp(io); self.observe(now_ms, stats.frames_in.load(.monotonic)); const v = self.verdict(now_ms); if (v.stalled != was_stalled) { if (v.stalled) { log.err("ingest stalled: {d}.{d:0>3} fps below {d} fps for {d}s (window {d}s); /_healthz now 503", .{ v.rate_milli_fps / 1000, v.rate_milli_fps % 1000, v.threshold_fps, v.below_sec, v.window_sec }); } else { log.info("ingest stall cleared: {d}.{d:0>3} fps", .{ v.rate_milli_fps / 1000, v.rate_milli_fps % 1000 }); } was_stalled = v.stalled; } } }
/// JSON for /admin/ingest-stall and the 503 body. pub fn formatJson(self: *const IngestLiveness, buf: []u8, now_ms: i64) []const u8 { const v = self.verdict(now_ms); const uptime_sec = @divFloor(now_ms - self.started_ms, std.time.ms_per_s); return std.fmt.bufPrint(buf, "{{\"enabled\":{},\"stalled\":{},\"rate_fps\":{d}.{d:0>3},\"below_threshold_sec\":{d},\"threshold_fps\":{d},\"window_sec\":{d},\"startup_grace_sec\":{d},\"uptime_sec\":{d},\"persisted\":false}}", .{ v.enabled, v.stalled, v.rate_milli_fps / 1000, v.rate_milli_fps % 1000, v.below_sec, v.threshold_fps, v.window_sec, self.startup_grace_sec.load(.acquire), uptime_sec, }) catch "{\"error\":\"Internal\"}"; }
pub fn formatPrometheus(self: *const IngestLiveness, buf: []u8, now_ms: i64) []const u8 { const v = self.verdict(now_ms); return std.fmt.bufPrint(buf, \\# HELP relay_ingest_stalled 1 when process-wide ingest has been below relay_ingest_stall_threshold_fps for relay_ingest_stall_window_seconds; /_healthz returns 503 while set \\# TYPE relay_ingest_stalled gauge \\relay_ingest_stalled {d} \\# HELP relay_ingest_stall_enabled 0 = kill switch thrown via /admin/ingest-stall; /_healthz never reports a stall \\# TYPE relay_ingest_stall_enabled gauge \\relay_ingest_stall_enabled {d} \\# HELP relay_ingest_rate_fps process-wide ingest rate as sampled by the stall detector \\# TYPE relay_ingest_rate_fps gauge \\relay_ingest_rate_fps {d}.{d:0>3} \\# HELP relay_ingest_below_threshold_seconds how long ingest has been continuously below the threshold (0 = not below) \\# TYPE relay_ingest_below_threshold_seconds gauge \\relay_ingest_below_threshold_seconds {d} \\# TYPE relay_ingest_stall_threshold_fps gauge \\relay_ingest_stall_threshold_fps {d} \\# TYPE relay_ingest_stall_window_seconds gauge \\relay_ingest_stall_window_seconds {d} \\ , .{ @as(u8, @intFromBool(v.stalled)), @as(u8, @intFromBool(v.enabled)), v.rate_milli_fps / 1000, v.rate_milli_fps % 1000, v.below_sec, v.threshold_fps, v.window_sec, }) catch ""; }};
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;
fn testDetector() IngestLiveness { return IngestLiveness.init(0, true, 50, 600, 300);}
fn feed(d: *IngestLiveness, from_ms: i64, to_ms: i64, fps: u64, frames: *u64) void { var t = from_ms; while (t < to_ms) : (t += sample_interval_sec * ms) { d.observe(t, frames.*); frames.* += fps * sample_interval_sec; } d.observe(to_ms, frames.*);}
test "healthy ingest never stalls" { var d = testDetector(); var frames: u64 = 0; feed(&d, 0, 3600 * ms, 400, &frames); const v = d.verdict(3600 * ms); try std.testing.expect(!v.stalled); try std.testing.expectEqual(@as(u64, 0), v.below_sec); try std.testing.expectEqual(@as(u64, 400_000), v.rate_milli_fps);}
test "below threshold for less than the window stays healthy" { var d = testDetector(); var frames: u64 = 0; feed(&d, 0, 1000 * ms, 400, &frames); feed(&d, 1000 * ms, 1500 * ms, 10, &frames); const v = d.verdict(1500 * ms); try std.testing.expect(!v.stalled); try std.testing.expect(v.below_sec > 0 and v.below_sec < 600);}
test "below threshold for the whole window is a stall" { var d = testDetector(); var frames: u64 = 0; feed(&d, 0, 1000 * ms, 400, &frames); feed(&d, 1000 * ms, 1700 * ms, 10, &frames); const v = d.verdict(1700 * ms); try std.testing.expect(v.stalled); try std.testing.expect(v.below_sec >= 600); try std.testing.expectEqual(@as(u64, 10_000), v.rate_milli_fps);}
test "total silence is a stall" { var d = testDetector(); var frames: u64 = 0; feed(&d, 0, 1000 * ms, 400, &frames); feed(&d, 1000 * ms, 1700 * ms, 0, &frames); try std.testing.expect(d.verdict(1700 * ms).stalled);}
test "recovery clears the stall" { var d = testDetector(); var frames: u64 = 0; feed(&d, 0, 1000 * ms, 400, &frames); feed(&d, 1000 * ms, 1700 * ms, 0, &frames); try std.testing.expect(d.verdict(1700 * ms).stalled); feed(&d, 1700 * ms, 1720 * ms, 400, &frames); try std.testing.expect(!d.verdict(1720 * ms).stalled);}
test "disabled never stalls and re-enabling restarts the window" { var d = testDetector(); var frames: u64 = 0; d.enabled.store(false, .release); feed(&d, 0, 1000 * ms, 400, &frames); feed(&d, 1000 * ms, 2000 * ms, 0, &frames); try std.testing.expect(!d.verdict(2000 * ms).stalled);
d.enabled.store(true, .release); feed(&d, 2000 * ms, 2300 * ms, 0, &frames); try std.testing.expect(!d.verdict(2300 * ms).stalled); feed(&d, 2300 * ms, 2700 * ms, 0, &frames); try std.testing.expect(d.verdict(2700 * ms).stalled);}
test "kill switch takes effect without waiting for a sample" { var d = testDetector(); var frames: u64 = 0; feed(&d, 0, 1000 * ms, 400, &frames); feed(&d, 1000 * ms, 1700 * ms, 0, &frames); try std.testing.expect(d.verdict(1700 * ms).stalled); d.enabled.store(false, .release); try std.testing.expect(!d.verdict(1700 * ms).stalled);}
test "startup grace covers cold start and the window starts after it" { var d = testDetector(); var frames: u64 = 0; feed(&d, 0, 290 * ms, 0, &frames); try std.testing.expect(!d.verdict(290 * ms).stalled); try std.testing.expectEqual(@as(u64, 0), d.verdict(290 * ms).below_sec); feed(&d, 290 * ms, 800 * ms, 0, &frames); try std.testing.expect(!d.verdict(800 * ms).stalled); feed(&d, 800 * ms, 1000 * ms, 0, &frames); try std.testing.expect(d.verdict(1000 * ms).stalled);}
test "widening the window at runtime clears a stall immediately" { var d = testDetector(); var frames: u64 = 0; feed(&d, 0, 1000 * ms, 400, &frames); feed(&d, 1000 * ms, 1700 * ms, 0, &frames); try std.testing.expect(d.verdict(1700 * ms).stalled); d.window_sec.store(3600, .release); try std.testing.expect(!d.verdict(1700 * ms).stalled);}
test "a stale sample never reads as stalled" { var d = testDetector(); var frames: u64 = 0; feed(&d, 0, 1000 * ms, 400, &frames); feed(&d, 1000 * ms, 1700 * ms, 0, &frames); try std.testing.expect(d.verdict(1700 * ms).stalled); try std.testing.expect(!d.verdict(1700 * ms + sample_stale_ms + 1).stalled);}
test "sampling resumes cleanly after a gap" { var d = testDetector(); var frames: u64 = 0; feed(&d, 0, 1000 * ms, 400, &frames); d.observe(1000 * ms + sample_stale_ms + ms, frames); try std.testing.expect(!d.verdict(1000 * ms + sample_stale_ms + ms).stalled); try std.testing.expectEqual(@as(i64, 0), d.below_since_ms.load(.acquire));}
test "formatJson and formatPrometheus carry the verdict" { var d = testDetector(); var frames: u64 = 0; feed(&d, 0, 1000 * ms, 400, &frames); var buf: [1024]u8 = undefined; const json = d.formatJson(&buf, 1000 * ms); try std.testing.expect(std.mem.indexOf(u8, json, "\"stalled\":false") != null); try std.testing.expect(std.mem.indexOf(u8, json, "\"rate_fps\":400.000") != null); try std.testing.expect(std.mem.indexOf(u8, json, "\"persisted\":false") != null); const prom = d.formatPrometheus(&buf, 1000 * ms); try std.testing.expect(std.mem.indexOf(u8, prom, "relay_ingest_stalled 0\n") != null); try std.testing.expect(std.mem.indexOf(u8, prom, "relay_ingest_rate_fps 400.000\n") != null); try std.testing.expect(std.mem.indexOf(u8, prom, "relay_ingest_stall_window_seconds 600\n") != null);}
test "parseBool" { try std.testing.expectEqual(@as(?bool, true), parseBool("1")); try std.testing.expectEqual(@as(?bool, true), parseBool("TRUE")); try std.testing.expectEqual(@as(?bool, false), parseBool("0")); try std.testing.expectEqual(@as(?bool, false), parseBool("false")); try std.testing.expectEqual(@as(?bool, null), parseBool("maybe"));}
/// the /_healthz 503 body: a short reason, the sampled rate, and the window.pub fn formatHealthzStalled(buf: []u8, v: Verdict) []const u8 { return std.fmt.bufPrint(buf, "{{\"status\":\"error\",\"msg\":\"ingest stalled: {d}.{d:0>3} fps below {d} fps for {d}s (window {d}s)\",\"rate_fps\":{d}.{d:0>3},\"threshold_fps\":{d},\"below_threshold_sec\":{d},\"window_sec\":{d}}}", .{ v.rate_milli_fps / 1000, v.rate_milli_fps % 1000, v.threshold_fps, v.below_sec, v.window_sec, v.rate_milli_fps / 1000, v.rate_milli_fps % 1000, v.threshold_fps, v.below_sec, v.window_sec, }) catch "{\"status\":\"error\",\"msg\":\"ingest stalled\"}";}
test "formatHealthzStalled" { var buf: [512]u8 = undefined; const body = formatHealthzStalled(&buf, .{ .stalled = true, .enabled = true, .rate_milli_fps = 1_500, .below_sec = 930, .threshold_fps = 50, .window_sec = 900 }); try std.testing.expectEqualStrings("{\"status\":\"error\",\"msg\":\"ingest stalled: 1.500 fps below 50 fps for 930s (window 900s)\",\"rate_fps\":1.500,\"threshold_fps\":50,\"below_threshold_sec\":930,\"window_sec\":900}", body);}