diff --git a/src/broadcaster.zig b/src/broadcaster.zig index ca02d96..4cac682 100644 --- a/src/broadcaster.zig +++ b/src/broadcaster.zig @@ -683,8 +683,13 @@ fn appendProcMetrics(w: anytype) void { } else |_| {} // glibc malloc arena stats — distinguishes in-use heap from fragmentation - // mallinfo uses int fields (cap at 2 GiB) but sufficient for trend analysis + // mallinfo returns c_int (i32) fields which overflow at 2 GiB; bitcast to u32 + // extends useful range to 4 GiB per field const mi = malloc_h.mallinfo(); + const arena: u64 = @as(u32, @bitCast(mi.arena)); + const in_use: u64 = @as(u32, @bitCast(mi.uordblks)); + const free_bytes: u64 = @as(u32, @bitCast(mi.fordblks)); + const mmap_bytes: u64 = @as(u32, @bitCast(mi.hblkhd)); std.fmt.format(w, \\# TYPE relay_malloc_arena_bytes gauge \\# HELP relay_malloc_arena_bytes total bytes claimed from OS by malloc @@ -702,7 +707,7 @@ fn appendProcMetrics(w: anytype) void { \\# HELP relay_malloc_mmap_bytes bytes allocated via mmap (large blocks) \\relay_malloc_mmap_bytes {d} \\ - , .{ mi.arena, mi.uordblks, mi.fordblks, mi.hblkhd }) catch {}; + , .{ arena, in_use, free_bytes, mmap_bytes }) catch {}; } const posix_vfs = @cImport(@cInclude("sys/statvfs.h")); diff --git a/src/validator.zig b/src/validator.zig index 3402c0c..9d40a1e 100644 --- a/src/validator.zig +++ b/src/validator.zig @@ -55,6 +55,8 @@ pub const Validator = struct { cache_mutex: std.Thread.Mutex = .{}, // background resolve queue queue: std.ArrayListUnmanaged([]const u8) = .{}, + // in-flight set — prevents duplicate DID entries in the queue + queued_set: std.StringHashMapUnmanaged(void) = .{}, // migration validation queue migration_queue: std.ArrayListUnmanaged(MigrationCheck) = .{}, queue_mutex: std.Thread.Mutex = .{}, @@ -65,6 +67,10 @@ pub const Validator = struct { const max_resolver_threads = 8; const default_resolver_threads = 4; + // process 1 migration check per N DID resolutions to prevent migration queue starvation + const migration_interleave_interval = 10; + // recreate http resolver after N resolutions to bound connection pool growth + const resolver_recycle_interval = 1000; pub fn init(allocator: Allocator, stats: *broadcaster.Stats) Validator { return initWithConfig(allocator, stats, .{}); @@ -100,6 +106,7 @@ pub const Validator = struct { self.allocator.free(did); } self.queue.deinit(self.allocator); + self.queued_set.deinit(self.allocator); // free migration queue for (self.migration_queue.items) |mc| { @@ -387,10 +394,18 @@ pub const Validator = struct { self.queue_mutex.lock(); defer self.queue_mutex.unlock(); + + // skip if already queued (prevents duplicate in-flight resolutions) + if (self.queued_set.contains(duped)) { + self.allocator.free(duped); + return; + } + self.queue.append(self.allocator, duped) catch { self.allocator.free(duped); return; }; + self.queued_set.put(self.allocator, duped, {}) catch {}; self.queue_cond.signal(); } @@ -408,10 +423,11 @@ pub const Validator = struct { while (self.queue.items.len == 0 and self.migration_queue.items.len == 0 and self.alive.load(.acquire)) { self.queue_cond.timedWait(&self.queue_mutex, 1 * std.time.ns_per_s) catch {}; } - if (self.queue.items.len > 0) { - did = self.queue.orderedRemove(0); - } else if (self.migration_queue.items.len > 0) { + if (self.migration_queue.items.len > 0 and (self.queue.items.len == 0 or resolve_count % migration_interleave_interval == 0)) { migration = self.migration_queue.orderedRemove(0); + } else if (self.queue.items.len > 0) { + did = self.queue.orderedRemove(0); + _ = self.queued_set.remove(did.?); } } @@ -427,7 +443,7 @@ pub const Validator = struct { // periodically recreate resolver to free accumulated http.Client state resolve_count += 1; - if (resolve_count % 1000 == 0) { + if (resolve_count % resolver_recycle_interval == 0) { resolver.deinit(); resolver = zat.DidResolver.init(self.allocator); }