diff --git a/src/api/router.zig b/src/api/router.zig index 37db586..b15015c 100644 --- a/src/api/router.zig +++ b/src/api/router.zig @@ -61,11 +61,11 @@ fn handleGet(conn: *websocket.Conn, path: []const u8, query: []const u8, headers // trivial liveness — process is alive, constant-time, no dependencies h.respondJson(conn, .ok, "{\"status\":\"ok\"}"); } else if (std.mem.eql(u8, path, "/_readyz") or std.mem.eql(u8, path, "/_health") or std.mem.eql(u8, path, "/xrpc/_health")) { - // readiness — checks DB dependency - _ = ctx.persist.db.exec("SELECT 1", .{}) catch { + // readiness — use atomic health flag (pg.Pool runs on Threaded, HTTP handlers are Evented) + if (!ctx.persist.isDbHealthy()) { h.respondJson(conn, .internal_server_error, "{\"status\":\"error\",\"msg\":\"database unavailable\"}"); return; - }; + } h.respondJson(conn, .ok, "{\"status\":\"ok\"}"); } else if (std.mem.eql(u8, path, "/_stats")) { var stats_buf: [4096]u8 = undefined; diff --git a/src/event_log.zig b/src/event_log.zig index f08f78b..eac0bdc 100644 --- a/src/event_log.zig +++ b/src/event_log.zig @@ -109,6 +109,10 @@ pub const DiskPersist = struct { io: Io, + /// last successful DB interaction (epoch seconds, set by Threaded workers). + /// read by metrics server to report health without cross-Io pg.Pool access. + last_db_success: std.atomic.Value(i64) = .{ .raw = 0 }, + /// current evtbuf entry count (for metrics — non-blocking, returns 0 if lock is contended) pub fn evtbufLen(self: *DiskPersist) usize { if (!self.mutex.tryLock()) return 0; @@ -306,6 +310,22 @@ pub const DiskPersist = struct { return .{ .uid = uid }; } + /// check DB health without touching pg.Pool — safe from any Io context. + /// returns true if a Threaded worker successfully queried the DB within the last 30s. + pub fn isDbHealthy(self: *DiskPersist) bool { + const last = self.last_db_success.load(.acquire); + if (last == 0) return false; + var ts: std.c.timespec = undefined; + _ = std.c.clock_gettime(.REALTIME, &ts); + return (@as(i64, ts.sec) - last) < 30; + } + + fn markDbSuccess(self: *DiskPersist) void { + var ts: std.c.timespec = undefined; + _ = std.c.clock_gettime(.REALTIME, &ts); + self.last_db_success.store(@as(i64, ts.sec), .release); + } + /// resolve a DID to a numeric UID. creates a new account row on first encounter. /// matches indigo's Relay.DidToUid → Account.UID mapping. pub fn uidForDid(self: *DiskPersist, did: []const u8) !u64 { @@ -321,6 +341,7 @@ pub const DiskPersist = struct { defer r.deinit() catch {}; const uid: u64 = @intCast(r.get(i64, 0)); self.didCachePut(did, uid); + self.markDbSuccess(); return uid; } diff --git a/src/main.zig b/src/main.zig index 2175b50..d963868 100644 --- a/src/main.zig +++ b/src/main.zig @@ -106,7 +106,8 @@ const MetricsServer = struct { .{ .name = "server", .value = "zlay (atproto-relay)" }, } }) catch {}; } else if (std.mem.eql(u8, path, "/_health") or std.mem.eql(u8, path, "/_readyz")) { - const db_ok = if (self.persist.db.exec("SELECT 1", .{})) |_| true else |_| false; + // use atomic health flag — pg.Pool runs on Threaded, metrics server is Evented + const db_ok = self.persist.isDbHealthy(); const status: http.Status = if (db_ok) .ok else .internal_server_error; const body = if (db_ok) "{\"status\":\"ok\"}" else "{\"status\":\"error\",\"msg\":\"database unavailable\"}"; request.respond(body, .{ .status = status, .keep_alive = false, .extra_headers = &.{ @@ -289,9 +290,14 @@ pub fn main() !void { var broadcast_future = try io.concurrent(broadcaster.Broadcaster.runBroadcastLoop, .{&bc}); defer _ = broadcast_future.cancel(io); - // start GC loop (runs as background task — does disk I/O + malloc_trim) - var gc_future = try io.concurrent(gcLoop, .{ &dp, io }); - defer _ = gc_future.cancel(io); + // start GC loop on a plain thread — dp.gc() uses pool_io (Threaded) mutex + // and pg.Pool. MUST NOT run as Evented fiber: Threaded futex on Evented + // fiber dereferences NULL Thread.current() threadlocal → heap corruption. + const gc_thread = std.Thread.spawn(.{}, gcLoop, .{ &dp, pool_io }) catch |err| { + log.err("failed to start GC thread: {s}", .{@errorName(err)}); + return err; + }; + gc_thread.detach(); // wire HTTP fallback into broadcaster (all API endpoints served on WS port) var http_context = api.HttpContext{ @@ -360,8 +366,7 @@ pub fn main() !void { ws_listener.deinit(io); server_future.cancel(io); - // cancel GC task - gc_future.cancel(io); + // GC thread is detached and checks shutdown_flag — no cancel needed // cancel broadcaster fiber (shutdown flag already set, it will drain remaining) broadcast_future.cancel(io);