atproto relay in zig zlay.waow.tech
relay zig atproto
Something went wrong. Try again.
Zig
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744//! admin endpoint handlers for relay management.//!//! all handlers require Bearer token auth against RELAY_ADMIN_PASSWORD.//! includes host blocking/unblocking, account bans, and backfill control.//!//! DB-accessing handlers use DbRequest + DbRequestQueue to route queries//! through pool_io workers.
const std = @import("std");const Io = std.Io;const h = @import("http.zig");const router = @import("router.zig");const websocket = @import("websocket");const event_log_mod = @import("../event_log.zig");const backfill_mod = @import("../collection_index/backfill.zig");const cleaner_mod = @import("../collection_index/cleaner.zig");const resync_mod = @import("../collection_index/resync.zig");const util = @import("../util/util.zig");
const log = std.log.scoped(.relay);
/// default idle window for /admin/hosts/reconnect. Long enough that an/// ordinarily quiet PDS is not mistaken for a deaf one, short enough to be/// useful during an incident. Overridable per call.const default_reconnect_min_idle_s: i64 = 60;const getenv = util.getenv;const formatTimestamp = util.formatTimestamp;
const HttpContext = router.HttpContext;const DbRequest = event_log_mod.DbRequest;const DiskPersist = event_log_mod.DiskPersist;
/// check admin auth via headers, send error response if not authorized. returns true if authorized.pub fn checkAdmin(conn: *h.Conn, headers: ?*const websocket.Handshake.KeyValue) bool { const admin_pw = getenv("RELAY_ADMIN_PASSWORD") orelse { h.respondJson(conn, .forbidden, "{\"error\":\"admin endpoint not configured\"}"); return false; };
const kv = headers orelse { h.respondJson(conn, .unauthorized, "{\"error\":\"missing authorization header\"}"); return false; };
// handshake parser lowercases all header names const auth_value = kv.get("authorization") orelse { h.respondJson(conn, .unauthorized, "{\"error\":\"missing authorization header\"}"); return false; };
const bearer_prefix = "Bearer "; if (!std.mem.startsWith(u8, auth_value, bearer_prefix)) { h.respondJson(conn, .unauthorized, "{\"error\":\"invalid authorization scheme\"}"); return false; } const token = auth_value[bearer_prefix.len..]; if (!std.mem.eql(u8, token, admin_pw)) { h.respondJson(conn, .unauthorized, "{\"error\":\"invalid token\"}"); return false; } return true;}
pub fn handleBan(conn: *h.Conn, body: []const u8, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { if (!checkAdmin(conn, headers)) return;
const parsed = std.json.parseFromSlice(struct { did: []const u8 }, ctx.persist.allocator, body, .{ .ignore_unknown_fields = true }) catch { h.respondJson(conn, .bad_request, "{\"error\":\"invalid JSON, expected {\\\"did\\\":\\\"...\\\"}\"}"); return; }; defer parsed.deinit(); const did = parsed.value.did;
// resolve DID → UID via DbRequestQueue const UidReq = struct { base: DbRequest = .{ .callback = &execute }, did_buf: [256]u8 = undefined, did_len: usize = 0, uid: u64 = 0,
fn execute(b: *DbRequest, dp: *DiskPersist) void { const self: *@This() = @fieldParentPtr("base", b); const d = self.did_buf[0..self.did_len]; // check database if (dp.db.rowUnsafe("SELECT uid FROM account WHERE did = $1", .{d}) catch null) |row| { var r = row; defer r.deinit() catch {}; self.uid = @intCast(r.get(i64, 0)); return; } // create new account row _ = dp.db.exec("INSERT INTO account (did) VALUES ($1) ON CONFLICT (did) DO NOTHING", .{d}) catch { b.err = error.DatabaseError; return; }; var row = dp.db.rowUnsafe("SELECT uid FROM account WHERE did = $1", .{d}) catch { b.err = error.DatabaseError; return; } orelse { b.err = error.AccountCreationFailed; return; }; defer row.deinit() catch {}; self.uid = @intCast(row.get(i64, 0)); } }; var uid_req: UidReq = .{}; const copy_len = @min(did.len, uid_req.did_buf.len); @memcpy(uid_req.did_buf[0..copy_len], did[0..copy_len]); uid_req.did_len = copy_len; ctx.db_queue.push(&uid_req.base); uid_req.base.wait(ctx.io, ctx.shutdown);
if (uid_req.base.err != null) { h.respondJson(conn, .internal_server_error, "{\"error\":\"failed to resolve DID\"}"); return; }
// remove from collection index so banned accounts don't appear in listReposByCollection ctx.collection_index.removeAll(did) catch |err| { log.debug("collection removeAll after ban failed: {s}", .{@errorName(err)}); };
// build CBOR #account frame and route takedown + persist + broadcast // through host_ops queue (pool_io thread) — fire and forget. const host_ops_mod = @import("../host_ops.zig"); var td: host_ops_mod.HostOp.Payload.Takedown = .{ .uid = uid_req.uid };
if (buildAccountFrame(ctx.persist.allocator, did)) |frame_bytes| { defer ctx.persist.allocator.free(frame_bytes); if (frame_bytes.len <= td.frame_buf.len) { @memcpy(td.frame_buf[0..frame_bytes.len], frame_bytes); td.frame_len = @intCast(frame_bytes.len); } }
ctx.host_ops.push(.{ .host_id = 0, // not host-specific .kind = .takedown_user, .payload = .{ .takedown = td }, });
log.info("admin: banned {s} (uid={d}), takedown enqueued", .{ did, uid_req.uid }); h.respondJson(conn, .ok, "{\"success\":true}");}
pub fn handleAdminListHosts(conn: *h.Conn, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { if (!checkAdmin(conn, headers)) return;
// list all hosts via DbRequestQueue const ListAllHostsReq = struct { base: DbRequest = .{ .callback = &execute }, alloc: std.mem.Allocator, result: ?[]event_log_mod.DiskPersist.Host = null,
fn execute(b: *DbRequest, dp: *DiskPersist) void { const self: *@This() = @fieldParentPtr("base", b); self.result = dp.listAllHosts(self.alloc) catch |e| { b.err = e; return; }; } }; var list_req: ListAllHostsReq = .{ .alloc = ctx.persist.allocator }; ctx.db_queue.push(&list_req.base); list_req.base.wait(ctx.io, ctx.shutdown);
if (list_req.base.err != null or list_req.result == null) { h.respondJson(conn, .internal_server_error, "{\"error\":\"DatabaseError\",\"message\":\"query failed\"}"); return; }
const hosts = list_req.result.?; defer { for (hosts) |host| { ctx.persist.allocator.free(host.hostname); ctx.persist.allocator.free(host.status); } ctx.persist.allocator.free(hosts); }
var aw: Io.Writer.Allocating = .init(ctx.persist.allocator); defer aw.deinit(); const w = &aw.writer;
w.writeAll("{\"hosts\":[") catch return;
for (hosts, 0..) |host, i| { if (i > 0) w.writeByte(',') catch return; if (host.account_limit) |limit| { w.print("{{\"id\":{d},\"hostname\":\"{s}\",\"status\":\"{s}\",\"last_seq\":{d},\"failed_attempts\":{d},\"account_limit\":{d}}}", .{ host.id, host.hostname, host.status, host.last_seq, host.failed_attempts, limit, }) catch return; } else { w.print("{{\"id\":{d},\"hostname\":\"{s}\",\"status\":\"{s}\",\"last_seq\":{d},\"failed_attempts\":{d},\"account_limit\":null}}", .{ host.id, host.hostname, host.status, host.last_seq, host.failed_attempts, }) catch return; } }
w.print("],\"active_workers\":{d}}}", .{ctx.slurper.workerCount()}) catch return; h.respondJson(conn, .ok, aw.written());}
pub fn handleAdminBlockHost(conn: *h.Conn, body: []const u8, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { if (!checkAdmin(conn, headers)) return;
const parsed = std.json.parseFromSlice(struct { hostname: []const u8 }, ctx.persist.allocator, body, .{ .ignore_unknown_fields = true }) catch { h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"invalid JSON\"}"); return; }; defer parsed.deinit();
const BlockHostReq = struct { base: DbRequest = .{ .callback = &execute }, hostname_buf: [256]u8 = undefined, hostname_len: usize = 0, host_id: u64 = 0,
fn execute(b: *DbRequest, dp: *DiskPersist) void { const self: *@This() = @fieldParentPtr("base", b); const hn = self.hostname_buf[0..self.hostname_len]; const info = dp.getOrCreateHost(hn) catch |e| { b.err = e; return; }; self.host_id = info.id; dp.updateHostStatus(info.id, "blocked") catch |e| { b.err = e; return; }; } }; var req: BlockHostReq = .{}; const copy_len = @min(parsed.value.hostname.len, req.hostname_buf.len); @memcpy(req.hostname_buf[0..copy_len], parsed.value.hostname[0..copy_len]); req.hostname_len = copy_len; ctx.db_queue.push(&req.base); req.base.wait(ctx.io, ctx.shutdown);
if (req.base.err != null) { h.respondJson(conn, .internal_server_error, "{\"error\":\"DatabaseError\",\"message\":\"operation failed\"}"); return; }
log.info("admin: blocked host {s} (id={d})", .{ parsed.value.hostname, req.host_id }); h.respondJson(conn, .ok, "{\"success\":true}");}
pub fn handleAdminUnblockHost(conn: *h.Conn, body: []const u8, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { if (!checkAdmin(conn, headers)) return;
const parsed = std.json.parseFromSlice(struct { hostname: []const u8 }, ctx.persist.allocator, body, .{ .ignore_unknown_fields = true }) catch { h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"invalid JSON\"}"); return; }; defer parsed.deinit();
const UnblockHostReq = struct { base: DbRequest = .{ .callback = &execute }, hostname_buf: [256]u8 = undefined, hostname_len: usize = 0, host_id: u64 = 0,
fn execute(b: *DbRequest, dp: *DiskPersist) void { const self: *@This() = @fieldParentPtr("base", b); const hn = self.hostname_buf[0..self.hostname_len]; const info = dp.getOrCreateHost(hn) catch |e| { b.err = e; return; }; self.host_id = info.id; dp.updateHostStatus(info.id, "active") catch |e| { b.err = e; return; }; dp.resetHostFailures(info.id) catch {}; } }; var req: UnblockHostReq = .{}; const copy_len = @min(parsed.value.hostname.len, req.hostname_buf.len); @memcpy(req.hostname_buf[0..copy_len], parsed.value.hostname[0..copy_len]); req.hostname_len = copy_len; ctx.db_queue.push(&req.base); req.base.wait(ctx.io, ctx.shutdown);
if (req.base.err != null) { h.respondJson(conn, .internal_server_error, "{\"error\":\"DatabaseError\",\"message\":\"operation failed\"}"); return; }
log.info("admin: unblocked host {s} (id={d})", .{ parsed.value.hostname, req.host_id }); h.respondJson(conn, .ok, "{\"success\":true}");}
/// force a host's worker to drop and re-establish its connection.////// In-band recovery for a worker that is alive but deaf. Unlike/// block+unblock this never changes the host's status, so a host cannot be/// left stranded as blocked if the operator's second call does not land.////// POST /admin/hosts/reconnect?min_idle_seconds=60&dry_run=true/// body: {"hostname": "..."}////// `min_idle_seconds` guards against tearing down a healthy worker: a deaf/// host and a merely quiet host are indistinguishable from seq sampling, so/// reconnecting on silence alone manufactures the connection churn it is/// trying to recover from. It is a query parameter, not a constant, so the/// window can be widened mid-incident without a redeploy.pub fn handleAdminReconnectHost(conn: *h.Conn, query: []const u8, body: []const u8, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { if (!checkAdmin(conn, headers)) return;
const min_idle_s = blk: { const raw = h.queryParam(query, "min_idle_seconds") orelse break :blk default_reconnect_min_idle_s; break :blk std.fmt.parseInt(i64, raw, 10) catch { h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"min_idle_seconds must be an integer\"}"); return; }; }; if (min_idle_s < 0) { h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"min_idle_seconds must not be negative\"}"); return; } const dry_run = if (h.queryParam(query, "dry_run")) |v| std.mem.eql(u8, v, "true") else false;
const parsed = std.json.parseFromSlice(struct { hostname: []const u8 }, ctx.persist.allocator, body, .{ .ignore_unknown_fields = true }) catch { h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"invalid JSON\"}"); return; }; defer parsed.deinit();
// resolve the host WITHOUT creating it: reconnecting a host we have never // seen is a caller error, not a reason to add one. const LookupReq = struct { base: DbRequest = .{ .callback = &execute }, hostname_buf: [256]u8 = undefined, hostname_len: usize = 0, host_id: u64 = 0, found: bool = false,
fn execute(b: *DbRequest, dp: *DiskPersist) void { const self: *@This() = @fieldParentPtr("base", b); const hn = self.hostname_buf[0..self.hostname_len]; const maybe_id = dp.getHostIdForHostname(hn) catch |e| { b.err = e; return; }; if (maybe_id) |id| { self.host_id = id; self.found = true; } } }; var req: LookupReq = .{}; const copy_len = @min(parsed.value.hostname.len, req.hostname_buf.len); @memcpy(req.hostname_buf[0..copy_len], parsed.value.hostname[0..copy_len]); req.hostname_len = copy_len; ctx.db_queue.push(&req.base); req.base.wait(ctx.io, ctx.shutdown);
if (req.base.err != null) { h.respondJson(conn, .internal_server_error, "{\"error\":\"DatabaseError\",\"message\":\"operation failed\"}"); return; } if (!req.found) { h.respondJson(conn, .not_found, "{\"error\":\"NotFound\",\"message\":\"unknown host\"}"); return; }
const min_idle_ms = min_idle_s * std.time.ms_per_s; var buf: [256]u8 = undefined;
if (dry_run) { // report the decision, and for a skip the idle time behind it -- the // skips are what tell an operator whether their silent-set is real. const outcome = ctx.slurper.inspectReconnect(req.host_id, min_idle_ms); const payload = switch (outcome) { .reconnected => std.fmt.bufPrint(&buf, "{{\"dry_run\":true,\"would\":\"reconnect\"}}", .{}), .skipped_recent_frame => |idle_ms| std.fmt.bufPrint(&buf, "{{\"dry_run\":true,\"would\":\"skip\",\"reason\":\"recent_frame\",\"idle_seconds\":{d},\"min_idle_seconds\":{d}}}", .{ @divFloor(idle_ms, std.time.ms_per_s), min_idle_s }), .no_worker => std.fmt.bufPrint(&buf, "{{\"dry_run\":true,\"would\":\"skip\",\"reason\":\"no_worker\"}}", .{}), .respawn_failed => std.fmt.bufPrint(&buf, "{{\"dry_run\":true,\"would\":\"reconnect\"}}", .{}), } catch { h.respondJson(conn, .internal_server_error, "{\"error\":\"Internal\"}"); return; }; h.respondJson(conn, .ok, payload); return; }
switch (ctx.slurper.forceReconnect(req.host_id, min_idle_ms)) { .reconnected => { log.info("admin: forced reconnect for host {s} (id={d})", .{ parsed.value.hostname, req.host_id }); h.respondJson(conn, .ok, "{\"success\":true,\"reconnected\":true}"); }, .skipped_recent_frame => |idle_ms| { const payload = std.fmt.bufPrint(&buf, "{{\"success\":true,\"reconnected\":false,\"reason\":\"recent_frame\",\"idle_seconds\":{d},\"min_idle_seconds\":{d}}}", .{ @divFloor(idle_ms, std.time.ms_per_s), min_idle_s }) catch "{\"success\":true,\"reconnected\":false,\"reason\":\"recent_frame\"}"; h.respondJson(conn, .ok, payload); }, .no_worker => h.respondJson(conn, .not_found, "{\"error\":\"NoWorker\",\"message\":\"host has no active worker\"}"), .respawn_failed => h.respondJson(conn, .internal_server_error, "{\"error\":\"RespawnFailed\",\"message\":\"worker torn down but respawn could not be queued\"}"), }}
/// set or clear the account_limit override for a host.pub fn handleAdminChangeLimits(conn: *h.Conn, body: []const u8, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { if (!checkAdmin(conn, headers)) return;
const parsed = std.json.parseFromSlice( struct { host: []const u8, account_limit: ?u64 }, ctx.persist.allocator, body, .{ .ignore_unknown_fields = true }, ) catch { h.respondJson(conn, .bad_request, "{\"error\":\"invalid JSON, expected {\\\"host\\\":\\\"...\\\",\\\"account_limit\\\":...}\"}"); return; }; defer parsed.deinit();
const ChangeLimitsReq = struct { base: DbRequest = .{ .callback = &execute }, hostname_buf: [256]u8 = undefined, hostname_len: usize = 0, new_limit: ?u64, host_id: ?u64 = null, effective: u64 = 0,
fn execute(b: *DbRequest, dp: *DiskPersist) void { const self: *@This() = @fieldParentPtr("base", b); const hn = self.hostname_buf[0..self.hostname_len]; self.host_id = dp.getHostIdForHostname(hn) catch |e| { b.err = e; return; }; const hid = self.host_id orelse return; dp.setHostAccountLimit(hid, self.new_limit) catch |e| { b.err = e; return; }; self.effective = if (self.new_limit) |l| l else dp.getHostAccountCount(hid); } }; var req: ChangeLimitsReq = .{ .new_limit = parsed.value.account_limit }; const copy_len = @min(parsed.value.host.len, req.hostname_buf.len); @memcpy(req.hostname_buf[0..copy_len], parsed.value.host[0..copy_len]); req.hostname_len = copy_len; ctx.db_queue.push(&req.base); req.base.wait(ctx.io, ctx.shutdown);
if (req.base.err != null) { h.respondJson(conn, .internal_server_error, "{\"error\":\"database error\"}"); return; } const host_id = req.host_id orelse { h.respondJson(conn, .not_found, "{\"error\":\"host not found\"}"); return; };
// update running subscriber's rate limits immediately ctx.slurper.updateHostLimits(host_id, req.effective);
if (parsed.value.account_limit) |limit| { log.info("admin: set account_limit for {s} (id={d}): {d}", .{ parsed.value.host, host_id, limit }); } else { log.info("admin: cleared account_limit for {s} (id={d}), reverted to COUNT(*)", .{ parsed.value.host, host_id }); } h.respondJson(conn, .ok, "{\"success\":true}");}
pub fn handleAdminBackfillTrigger(conn: *h.Conn, query: []const u8, headers: *const websocket.Handshake.KeyValue, backfiller: *backfill_mod.Backfiller) void { if (!checkAdmin(conn, headers)) return;
const source = h.queryParam(query, "source") orelse "bsky.network";
backfiller.start(source) catch |err| { switch (err) { error.AlreadyRunning => { h.respondJson(conn, .conflict, "{\"error\":\"backfill already in progress\"}"); }, else => { h.respondJson(conn, .internal_server_error, "{\"error\":\"failed to start backfill\"}"); }, } return; };
var buf: [256]u8 = undefined; const resp_body = std.fmt.bufPrint(&buf, "{{\"status\":\"started\",\"source\":\"{s}\"}}", .{source}) catch { h.respondJson(conn, .ok, "{\"status\":\"started\"}"); return; }; h.respondJson(conn, .ok, resp_body);}
pub fn handleAdminBackfillStatus(conn: *h.Conn, headers: *const websocket.Handshake.KeyValue, backfiller: *backfill_mod.Backfiller) void { if (!checkAdmin(conn, headers)) return;
const body = backfiller.getStatus(backfiller.allocator) catch { h.respondJson(conn, .internal_server_error, "{\"error\":\"failed to query backfill status\"}"); return; }; defer backfiller.allocator.free(body);
h.respondJson(conn, .ok, body);}
pub fn handleCleanupTrigger(conn: *h.Conn, headers: *const websocket.Handshake.KeyValue, cleaner: *cleaner_mod.Cleaner) void { if (!checkAdmin(conn, headers)) return;
cleaner.start() catch |err| { switch (err) { error.AlreadyRunning => { h.respondJson(conn, .conflict, "{\"error\":\"cleanup already in progress\"}"); }, else => { h.respondJson(conn, .internal_server_error, "{\"error\":\"failed to start cleanup\"}"); }, } return; };
h.respondJson(conn, .ok, "{\"status\":\"started\"}");}
pub fn handleCleanupStatus(conn: *h.Conn, headers: *const websocket.Handshake.KeyValue, cleaner: *cleaner_mod.Cleaner) void { if (!checkAdmin(conn, headers)) return;
const status = cleaner.getStatus(); var buf: [256]u8 = undefined; const body = std.fmt.bufPrint(&buf, "{{\"running\":{},\"scanned\":{d},\"removed\":{d}}}", .{ status.running, status.scanned, status.removed, }) catch { h.respondJson(conn, .internal_server_error, "{\"error\":\"format failed\"}"); return; }; h.respondJson(conn, .ok, body);}
pub fn handleResyncStatus(conn: *h.Conn, headers: *const websocket.Handshake.KeyValue, resyncer: *resync_mod.Resyncer) void { if (!checkAdmin(conn, headers)) return;
var buf: [256]u8 = undefined; const body = std.fmt.bufPrint(&buf, "{{\"processed\":{d},\"failed\":{d},\"dropped\":{d},\"queue_depth\":{d}}}", .{ resyncer.processed.load(.monotonic), resyncer.failed.load(.monotonic), resyncer.dropped.load(.monotonic), resyncer.queueDepth(), }) catch { h.respondJson(conn, .internal_server_error, "{\"error\":\"format failed\"}"); return; }; h.respondJson(conn, .ok, body);}
pub fn handleResyncTrigger(conn: *h.Conn, body: []const u8, headers: *const websocket.Handshake.KeyValue, resyncer: *resync_mod.Resyncer) void { if (!checkAdmin(conn, headers)) return;
const parsed = std.json.parseFromSlice( struct { did: []const u8, hostname: []const u8 }, std.heap.c_allocator, body, .{ .ignore_unknown_fields = true }, ) catch { h.respondJson(conn, .bad_request, "{\"error\":\"invalid JSON, expected {\\\"did\\\":\\\"...\\\",\\\"hostname\\\":\\\"...\\\"}\"}"); return; }; defer parsed.deinit();
resyncer.enqueue(parsed.value.did, parsed.value.hostname); h.respondJson(conn, .ok, "{\"status\":\"enqueued\"}");}
// --- protocol helpers (used only by handleBan) ---
/// build a CBOR #account frame for a takedown event.fn buildAccountFrame(allocator: std.mem.Allocator, did: []const u8) ?[]const u8 { const zat = @import("zat"); const cbor = zat.cbor;
const header: cbor.Value = .{ .map = &.{ .{ .key = "op", .value = .{ .unsigned = 1 } }, .{ .key = "t", .value = .{ .text = "#account" } }, } };
var time_buf: [24]u8 = undefined; const time_str = formatTimestamp(&time_buf);
const payload: cbor.Value = .{ .map = &.{ .{ .key = "seq", .value = .{ .unsigned = 0 } }, .{ .key = "did", .value = .{ .text = did } }, .{ .key = "time", .value = .{ .text = time_str } }, .{ .key = "active", .value = .{ .boolean = false } }, .{ .key = "status", .value = .{ .text = "takendown" } }, } };
const header_bytes = cbor.encodeAlloc(allocator, header) catch return null; const payload_bytes = cbor.encodeAlloc(allocator, payload) catch { allocator.free(header_bytes); return null; };
var frame = allocator.alloc(u8, header_bytes.len + payload_bytes.len) catch { allocator.free(header_bytes); allocator.free(payload_bytes); return null; }; @memcpy(frame[0..header_bytes.len], header_bytes); @memcpy(frame[header_bytes.len..], payload_bytes);
allocator.free(header_bytes); allocator.free(payload_bytes);
return frame;}
/// GET /admin/ingest-stall — the liveness stall detector's knobs and verdict.pub fn handleIngestStallStatus(conn: *h.Conn, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { if (!checkAdmin(conn, headers)) return; var buf: [512]u8 = undefined; h.respondJson(conn, .ok, ctx.ingest_liveness.formatJson(&buf, util.milliTimestamp(ctx.io)));}
/// POST /admin/ingest-stall/// body: any subset of {"enabled":bool,"threshold_fps":N,"window_sec":N,"startup_grace_sec":N}////// takes effect immediately and is NOT persisted: a restart returns to the/// RELAY_INGEST_STALL_* env defaults. `{"enabled":false}` is the kill switch/// for capture-before-restart forensics on a wedged pod.pub fn handleIngestStallUpdate(conn: *h.Conn, body: []const u8, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { if (!checkAdmin(conn, headers)) return;
const Update = struct { enabled: ?bool = null, threshold_fps: ?u64 = null, window_sec: ?u64 = null, startup_grace_sec: ?u64 = null, }; const parsed = std.json.parseFromSlice(Update, ctx.persist.allocator, body, .{ .ignore_unknown_fields = true }) catch { h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"expected JSON with any of enabled, threshold_fps, window_sec, startup_grace_sec\"}"); return; }; defer parsed.deinit(); const u = parsed.value;
if (u.enabled == null and u.threshold_fps == null and u.window_sec == null and u.startup_grace_sec == null) { h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"no fields to update\"}"); return; } if (u.window_sec) |w| if (w == 0) { h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"window_sec must be positive\"}"); return; };
const il = ctx.ingest_liveness; if (u.enabled) |v| { il.enabled.store(v, .release); log.warn("admin: ingest stall liveness {s} (not persisted; env default returns on restart)", .{if (v) "enabled" else "DISABLED"}); } if (u.threshold_fps) |v| { il.threshold_fps.store(v, .release); log.warn("admin: ingest stall threshold set to {d} fps (not persisted)", .{v}); } if (u.window_sec) |v| { il.window_sec.store(v, .release); log.warn("admin: ingest stall window set to {d}s (not persisted)", .{v}); } if (u.startup_grace_sec) |v| { il.startup_grace_sec.store(v, .release); log.warn("admin: ingest stall startup grace set to {d}s (not persisted)", .{v}); }
var buf: [512]u8 = undefined; h.respondJson(conn, .ok, il.formatJson(&buf, util.milliTimestamp(ctx.io)));}
/// GET /admin/ws-ping — the websocket keepalive pinger's knobs.pub fn handleWsPingStatus(conn: *h.Conn, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { if (!checkAdmin(conn, headers)) return; var buf: [256]u8 = undefined; h.respondJson(conn, .ok, ctx.ws_ping.formatJson(&buf));}
/// POST /admin/ws-ping/// body: any subset of {"enabled":bool,"interval_sec":N,"max_failures":N}////// takes effect on each pinger's next tick and is NOT persisted: a restart/// returns to the RELAY_WS_PING_* env defaults. `{"enabled":false}` is the/// kill switch — connections revert to no-keepalive behavior.pub fn handleWsPingUpdate(conn: *h.Conn, body: []const u8, headers: *const websocket.Handshake.KeyValue, ctx: *HttpContext) void { if (!checkAdmin(conn, headers)) return;
const Update = struct { enabled: ?bool = null, interval_sec: ?u64 = null, max_failures: ?u64 = null, }; const parsed = std.json.parseFromSlice(Update, ctx.persist.allocator, body, .{ .ignore_unknown_fields = true }) catch { h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"expected JSON with any of enabled, interval_sec, max_failures\"}"); return; }; defer parsed.deinit(); const u = parsed.value;
if (u.enabled == null and u.interval_sec == null and u.max_failures == null) { h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"no fields to update\"}"); return; } if (u.interval_sec) |v| if (v == 0) { h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"interval_sec must be positive\"}"); return; }; if (u.max_failures) |v| if (v == 0) { h.respondJson(conn, .bad_request, "{\"error\":\"BadRequest\",\"message\":\"max_failures must be positive\"}"); return; };
const wp = ctx.ws_ping; if (u.enabled) |v| { wp.enabled.store(v, .release); log.warn("admin: ws keepalive pinger {s} (not persisted; env default returns on restart)", .{if (v) "enabled" else "DISABLED"}); } if (u.interval_sec) |v| { wp.interval_sec.store(v, .release); log.warn("admin: ws keepalive interval set to {d}s (not persisted)", .{v}); } if (u.max_failures) |v| { wp.max_failures.store(v, .release); log.warn("admin: ws keepalive max_failures set to {d} (not persisted)", .{v}); }
var buf: [256]u8 = undefined; h.respondJson(conn, .ok, wp.formatJson(&buf));}