diff --git a/src/broadcaster.zig b/src/broadcaster.zig index 094fd1c..757254b 100644 --- a/src/broadcaster.zig +++ b/src/broadcaster.zig @@ -480,6 +480,19 @@ pub const Broadcaster = struct { defer self.consumers_mutex.unlock(); return self.consumers.items.len; } + + /// sum of all consumer send buffer depths (for metrics) + pub fn consumerQueueDepth(self: *Broadcaster) usize { + self.consumers_mutex.lock(); + defer self.consumers_mutex.unlock(); + var total: usize = 0; + for (self.consumers.items) |c| { + c.mutex.lock(); + total += c.buf_len; + c.mutex.unlock(); + } + return total; + } }; // --- websocket handler --- @@ -546,7 +559,14 @@ pub const Handler = struct { } }; -pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, migration_queue_len: usize, data_dir: []const u8, buf: []u8) []const u8 { +pub const AttributionMetrics = struct { + history_entries: usize = 0, + evtbuf_entries: usize = 0, + did_cache_entries: usize = 0, + consumer_queue_depth: usize = 0, +}; + +pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, migration_queue_len: usize, attribution: AttributionMetrics, data_dir: []const u8, buf: []u8) []const u8 { const uptime: i64 = std.time.timestamp() - stats.start_time; var fbs = std.io.fixedBufferStream(buf); const w = fbs.writer(); @@ -600,6 +620,22 @@ pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, migrat \\# TYPE relay_validator_cache_evictions_total counter \\relay_validator_cache_evictions_total {d} \\ + \\# TYPE relay_history_entries gauge + \\# HELP relay_history_entries in-memory frame history ring buffer entries + \\relay_history_entries {d} + \\ + \\# TYPE relay_evtbuf_entries gauge + \\# HELP relay_evtbuf_entries pending flush buffer entries + \\relay_evtbuf_entries {d} + \\ + \\# TYPE relay_did_cache_entries gauge + \\# HELP relay_did_cache_entries DID-to-UID cache entries + \\relay_did_cache_entries {d} + \\ + \\# TYPE relay_consumer_queue_depth gauge + \\# HELP relay_consumer_queue_depth total frames queued across all consumer send buffers + \\relay_consumer_queue_depth {d} + \\ , .{ stats.frames_in.load(.acquire), stats.frames_out.load(.acquire), @@ -619,6 +655,10 @@ pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, migrat cache_entries, migration_queue_len, stats.cache_evictions.load(.acquire), + attribution.history_entries, + attribution.evtbuf_entries, + attribution.did_cache_entries, + attribution.consumer_queue_depth, }) catch return fbs.getWritten(); // linux-only process metrics from /proc @@ -684,9 +724,8 @@ fn appendProcMetrics(w: anytype) void { } } else |_| {} - // glibc malloc arena stats — distinguishes in-use heap from fragmentation - // mallinfo returns c_int (i32) fields which overflow at 2 GiB; bitcast to u32 - // extends useful range to 4 GiB per field + // glibc malloc arena stats — mallinfo returns c_int (i32) fields which + // overflow at 2 GiB; bitcast to u32 extends useful range to 4 GiB const mi = malloc_h.mallinfo(); const arena: u64 = @as(u32, @bitCast(mi.arena)); const in_use: u64 = @as(u32, @bitCast(mi.uordblks)); @@ -847,8 +886,8 @@ test "formatPrometheusMetrics produces valid output" { stats.cache_hits.store(400, .release); stats.cache_misses.store(100, .release); - var buf: [8192]u8 = undefined; - const output = formatPrometheusMetrics(&stats, 42, 3, "/tmp", &buf); + var buf: [12288]u8 = undefined; + const output = formatPrometheusMetrics(&stats, 42, 3, .{}, "/tmp", &buf); try std.testing.expect(std.mem.indexOf(u8, output, "relay_frames_received_total 10000") != null); try std.testing.expect(std.mem.indexOf(u8, output, "relay_frames_broadcast_total 9000") != null); diff --git a/src/event_log.zig b/src/event_log.zig index 9b49bc3..46cd0c8 100644 --- a/src/event_log.zig +++ b/src/event_log.zig @@ -105,6 +105,20 @@ pub const DiskPersist = struct { alive: std.atomic.Value(bool) = .{ .raw = true }, flush_cond: std.Thread.Condition = .{}, + /// current evtbuf entry count (for metrics) + pub fn evtbufLen(self: *DiskPersist) usize { + self.mutex.lock(); + defer self.mutex.unlock(); + return self.evtbuf.items.len; + } + + /// current DID cache entry count (for metrics) + pub fn didCacheLen(self: *DiskPersist) usize { + self.did_cache_mutex.lock(); + defer self.did_cache_mutex.unlock(); + return self.did_cache.count(); + } + pub fn init(allocator: Allocator, dir_path: []const u8, database_url: []const u8) !DiskPersist { // ensure directory exists std.fs.cwd().makePath(dir_path) catch |err| switch (err) { diff --git a/src/main.zig b/src/main.zig index 45fe4a9..b02890f 100644 --- a/src/main.zig +++ b/src/main.zig @@ -50,6 +50,7 @@ const MetricsServer = struct { validator: *validator_mod.Validator, data_dir: []const u8, persist: *event_log_mod.DiskPersist, + bc: *broadcaster.Broadcaster, fn run(self: *MetricsServer) void { while (!shutdown_flag.load(.acquire)) { @@ -58,12 +59,12 @@ const MetricsServer = struct { log.debug("metrics accept error: {s}", .{@errorName(err)}); continue; }; - handleMetricsConn(conn.stream, self.stats, self.validator, self.data_dir, self.persist); + handleMetricsConn(conn.stream, self.stats, self.validator, self.data_dir, self.persist, self.bc); } } }; -fn handleMetricsConn(stream: std.net.Stream, stats: *broadcaster.Stats, validator: *validator_mod.Validator, data_dir: []const u8, persist: *event_log_mod.DiskPersist) void { +fn handleMetricsConn(stream: std.net.Stream, stats: *broadcaster.Stats, validator: *validator_mod.Validator, data_dir: []const u8, persist: *event_log_mod.DiskPersist, bc: *broadcaster.Broadcaster) void { defer stream.close(); var recv_buf: [4096]u8 = undefined; @@ -86,9 +87,15 @@ fn handleMetricsConn(stream: std.net.Stream, stats: *broadcaster.Stats, validato } else if (std.mem.eql(u8, path, "/metrics")) { const cache_entries = validator.cacheSize(); const migration_queue_len = validator.migrationQueueLen(); + const attribution = broadcaster.AttributionMetrics{ + .history_entries = bc.history.count(), + .evtbuf_entries = persist.evtbufLen(), + .did_cache_entries = persist.didCacheLen(), + .consumer_queue_depth = bc.consumerQueueDepth(), + }; - var metrics_buf: [8192]u8 = undefined; - const body = broadcaster.formatPrometheusMetrics(stats, cache_entries, migration_queue_len, data_dir, &metrics_buf); + var metrics_buf: [12288]u8 = undefined; + const body = broadcaster.formatPrometheusMetrics(stats, cache_entries, migration_queue_len, attribution, data_dir, &metrics_buf); request.respond(body, .{ .status = .ok, .keep_alive = false, .extra_headers = &.{ .{ .name = "content-type", .value = "text/plain; version=0.0.4; charset=utf-8" }, .{ .name = "server", .value = "zlay (atproto-relay)" }, @@ -201,6 +208,7 @@ pub fn main() !void { .validator = &val, .data_dir = data_dir, .persist = &dp, + .bc = &bc, }; const metrics_thread = try std.Thread.spawn(.{ .stack_size = default_stack_size }, MetricsServer.run, .{&metrics_srv}); @@ -240,6 +248,9 @@ pub fn main() !void { log.info("relay stopped cleanly", .{}); } +const builtin = @import("builtin"); +const malloc_h = if (builtin.os.tag == .linux) @cImport(@cInclude("malloc.h")) else struct {}; + fn gcLoop(dp: *event_log_mod.DiskPersist) void { const gc_interval: u64 = 10 * 60; // 10 minutes in seconds while (!shutdown_flag.load(.acquire)) { @@ -255,6 +266,13 @@ fn gcLoop(dp: *event_log_mod.DiskPersist) void { dp.gc() catch |err| { log.warn("event log GC failed: {s}", .{@errorName(err)}); }; + + // return free heap pages to the OS — glibc retains freed pages in + // per-thread arenas indefinitely; with ~2,700 threads this accumulates + // significant RSS that the application no longer needs. + if (comptime builtin.os.tag == .linux) { + _ = malloc_h.malloc_trim(0); + } } }