diff --git a/src/api/xrpc.zig b/src/api/xrpc.zig index c6db9ce..579232e 100644 --- a/src/api/xrpc.zig +++ b/src/api/xrpc.zig @@ -641,6 +641,10 @@ pub fn handleGetHostStatus(conn: *h.Conn, query: []const u8, ctx: *HttpContext) // --- requestCrawl --- pub fn handleRequestCrawl(conn: *h.Conn, body: []const u8, ctx: *HttpContext) void { + // log on entry so we can tell from logs whether requests are being received + // even if downstream waits wedge. body is small (~32 bytes) so this is cheap. + log.info("requestCrawl received (body_len={d})", .{body.len}); + const parsed = std.json.parseFromSlice(struct { hostname: []const u8 }, ctx.persist.allocator, body, .{ .ignore_unknown_fields = true }) catch { h.respondJson(conn, .bad_request, "{\"error\":\"InvalidRequest\",\"message\":\"invalid JSON, expected {\\\"hostname\\\":\\\"...\\\"}\"}"); return; diff --git a/src/broadcaster.zig b/src/broadcaster.zig index 70c1417..01030c5 100644 --- a/src/broadcaster.zig +++ b/src/broadcaster.zig @@ -907,6 +907,9 @@ pub const AttributionMetrics = struct { evtbuf_cap: usize = 0, outbuf_cap: usize = 0, workers_count: usize = 0, + // db_queue liveness — handled stops advancing while depth > 0 = workers wedged + db_queue_depth: u32 = 0, + db_queue_handled_total: u64 = 0, }; pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, attribution: AttributionMetrics, data_dir: []const u8, buf: []u8, io: Io) []const u8 { @@ -1156,6 +1159,14 @@ pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, attrib \\# HELP relay_workers_count active subscriber worker threads \\relay_workers_count {d} \\ + \\# TYPE relay_db_queue_depth gauge + \\# HELP relay_db_queue_depth current depth of the DbRequestQueue (tail - head). non-zero with stagnant relay_db_queue_handled_total = workers wedged + \\relay_db_queue_depth {d} + \\ + \\# TYPE relay_db_queue_handled_total counter + \\# HELP relay_db_queue_handled_total total DbRequest callbacks executed across worker threads + \\relay_db_queue_handled_total {d} + \\ , .{ attribution.validator_cache_map_cap, attribution.did_cache_map_cap, @@ -1163,6 +1174,8 @@ pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, attrib attribution.evtbuf_cap, attribution.outbuf_cap, attribution.workers_count, + attribution.db_queue_depth, + attribution.db_queue_handled_total, }) catch {}; // linux-only process metrics from /proc diff --git a/src/event_log.zig b/src/event_log.zig index 9f9b6c1..d9bdec6 100644 --- a/src/event_log.zig +++ b/src/event_log.zig @@ -123,9 +123,21 @@ pub const DbRequestQueue = struct { head: std.atomic.Value(u32) = .{ .raw = 0 }, tail: std.atomic.Value(u32) = .{ .raw = 0 }, push_lock: std.atomic.Value(u32) = .{ .raw = 0 }, // producer spinlock + /// total requests handled across all worker threads. exposed as + /// relay_db_queue_handled_total — if this stops advancing while + /// the queue has depth, the workers are dead or stuck. + handled: std.atomic.Value(u64) = .{ .raw = 0 }, shutdown: *std.atomic.Value(bool), persist: *DiskPersist, + /// current depth = tail - head. exposed as relay_db_queue_depth + /// so operators can correlate handler hangs with queue backup. + pub fn depth(self: *const DbRequestQueue) u32 { + const t = self.tail.load(.acquire); + const h = self.head.load(.acquire); + return t -% h; + } + pub fn push(self: *DbRequestQueue, req: *DbRequest) void { // acquire producer spinlock while (self.push_lock.cmpxchgWeak(0, 1, .acquire, .monotonic) != null) { @@ -161,12 +173,26 @@ pub const DbRequestQueue = struct { } pub fn run(self: *DbRequestQueue, pool_io: Io) void { + log.info("db_queue worker: starting", .{}); + var consecutive_sleep_errors: u32 = 0; while (!self.shutdown.load(.acquire)) { if (self.pop()) |req| { req.callback(req, self.persist); req.done.store(true, .release); + _ = self.handled.fetchAdd(1, .monotonic); + consecutive_sleep_errors = 0; } else { - pool_io.sleep(Io.Duration.fromMilliseconds(5), .awake) catch return; + pool_io.sleep(Io.Duration.fromMilliseconds(5), .awake) catch |e| { + // previous behavior: silently `return`, killing the worker + // forever. one of these errors during a transient io hiccup + // would wedge requestCrawl + every other DB-bound handler. + // log + continue. spin-yield to avoid a tight error loop. + consecutive_sleep_errors += 1; + if (consecutive_sleep_errors == 1 or consecutive_sleep_errors % 1000 == 0) { + log.warn("db_queue worker: pool_io.sleep failed ({s}), continuing (consecutive={d})", .{ @errorName(e), consecutive_sleep_errors }); + } + std.atomic.spinLoopHint(); + }; } } // shutdown drain — signal error on unprocessed requests @@ -174,6 +200,7 @@ pub const DbRequestQueue = struct { req.err = error.ShuttingDown; req.done.store(true, .release); } + log.info("db_queue worker: exiting (shutdown=true, handled_total={d})", .{self.handled.load(.acquire)}); } }; diff --git a/src/main.zig b/src/main.zig index 4eadc51..b971aac 100644 --- a/src/main.zig +++ b/src/main.zig @@ -129,6 +129,8 @@ const MetricsServer = struct { .evtbuf_cap = self.persist.evtbufCap(), .outbuf_cap = self.persist.outbufCap(), .workers_count = self.slurper.workerCount(), + .db_queue_depth = if (self.bc.db_queue) |q| q.depth() else 0, + .db_queue_handled_total = if (self.bc.db_queue) |q| q.handled.load(.acquire) else 0, }; var metrics_buf: [65536]u8 = undefined;