const std = @import("std"); const Store = @import("store.zig").Store; const live_rate = @import("live_rate.zig"); const coverage = @import("coverage.zig"); const log = std.log.scoped(.server); const net = std.Io.net; const index_html = @embedFile("static/index.html"); const sonify_html = @embedFile("static/sonify.html"); const sonify_og_svg = @embedFile("static/sonify-og.svg"); const diffs_api_html = @embedFile("static/diffs-api.html"); const llms_txt = @embedFile("static/llms.txt"); const sonify_stream_dir = "/var/lib/relay-eval-sonify-stream"; const sonify_playlist_url = "https://relay-eval.waow.tech/sonify/live/index.m3u8"; const sonify_artwork_url = "https://relay-eval.waow.tech/sonify/live/cover.png"; const static_cache_control = "public, max-age=0, s-maxage=300, stale-while-revalidate=60"; const snapshot_cache_control = "public, max-age=0, s-maxage=15, stale-while-revalidate=30"; const sonify_segment_max_age_seconds = 15; pub fn run( allocator: std.mem.Allocator, io: std.Io, db_path: [:0]const u8, port: u16, trend_limit: u32, meters: []live_rate.LiveRateMeter, ) !void { var store = try Store.open(db_path); defer store.close(); // derive og.png path from db directory const db_dir = std.fs.path.dirname(db_path) orelse "."; const og_png_path = try std.fmt.allocPrint(allocator, "{s}/og.png", .{db_dir}); defer allocator.free(og_png_path); const addr = net.IpAddress.parseIp4("0.0.0.0", port) catch unreachable; var listener = try addr.listen(io, .{ .reuse_address = true }); defer listener.deinit(io); log.info("listening on :{d}", .{port}); while (true) { const stream = listener.accept(io) catch |err| { log.err("accept: {s}", .{@errorName(err)}); continue; }; handleConnection(allocator, io, stream, &store, trend_limit, og_png_path, meters) catch |err| { log.debug("request error: {s}", .{@errorName(err)}); }; stream.close(io); } } fn handleConnection( allocator: std.mem.Allocator, io: std.Io, stream: net.Stream, store: *Store, trend_limit: u32, og_png_path: []const u8, meters: []live_rate.LiveRateMeter, ) !void { var buf: [4096]u8 = undefined; var bufs: [1][]u8 = .{&buf}; const n = try io.vtable.netRead(io.userdata, stream.socket.handle, &bufs); if (n == 0) return; const request = buf[0..n]; // parse method + path from first line const first_line_end = std.mem.indexOf(u8, request, "\r\n") orelse return; const first_line = request[0..first_line_end]; var parts = std.mem.splitScalar(u8, first_line, ' '); const method = parts.next() orelse return; const full_path = parts.next() orelse return; if (!std.mem.eql(u8, method, "GET")) { try sendResponse(stream, "405 Method Not Allowed", "text/plain", "method not allowed"); return; } // split path from query string. well-behaved clients percent-encode // reserved characters in query values (httpx sends since=2026-07-31T00%3A00%3A00Z), // so decode before any handler validates or compares raw bytes. const qs_sep = std.mem.indexOf(u8, full_path, "?"); const path = if (qs_sep) |i| full_path[0..i] else full_path; var query_buf: [2048]u8 = undefined; const query = if (qs_sep) |i| percentDecodeQuery(&query_buf, full_path[i + 1 ..]) else ""; if (std.mem.eql(u8, path, "/") or std.mem.eql(u8, path, "/chart") or std.mem.startsWith(u8, path, "/chart/")) { try sendResponseCached(stream, "200 OK", "text/html", index_html, static_cache_control); } else if (std.mem.eql(u8, path, "/diffs-api")) { try sendResponseCached(stream, "200 OK", "text/html", diffs_api_html, static_cache_control); } else if (std.mem.eql(u8, path, "/llms.txt")) { try sendResponseCached(stream, "200 OK", "text/plain", llms_txt, static_cache_control); } else if (std.mem.eql(u8, path, "/sonify")) { try sendResponseCached(stream, "200 OK", "text/html", sonify_html, static_cache_control); } else if (std.mem.eql(u8, path, "/sonify/og.svg")) { try sendResponseCached(stream, "200 OK", "image/svg+xml", sonify_og_svg, static_cache_control); } else if (std.mem.eql(u8, path, "/sonify/live/index.m3u8")) { try serveSonifyPlaylist(allocator, stream); } else if (std.mem.eql(u8, path, "/sonify/live/cover.png")) { try serveSonifyArtwork(allocator, stream); } else if (std.mem.startsWith(u8, path, "/sonify/live/")) { try serveSonifySegment(allocator, stream, path["/sonify/live/".len..]); } else if (std.mem.eql(u8, path, "/api/sonify/live")) { try serveSonifyStatus(allocator, io, stream); } else if (std.mem.eql(u8, path, "/api/runs")) { try serveRuns(allocator, stream, store); } else if (std.mem.eql(u8, path, "/api/runs/count")) { try serveRunCount(stream, store); } else if (std.mem.eql(u8, path, "/api/latest")) { try serveLatest(allocator, stream, store); } else if (std.mem.eql(u8, path, "/api/latest/diffs/summary")) { const run_id = try store.getLatestRunId() orelse { try sendResponse(stream, "200 OK", "application/json", "{\"run_id\":null,\"total\":0,\"checked\":0,\"facets\":{\"relays\":[],\"classifications\":[],\"pds\":[]}}"); return; }; try serveRunDiffSummary(allocator, stream, store, run_id, query); } else if (std.mem.eql(u8, path, "/api/latest/diffs")) { const run_id = try store.getLatestRunId() orelse { try sendResponse(stream, "200 OK", "application/json", "{\"run_id\":null,\"diffs\":[],\"next_cursor\":null}"); return; }; try serveRunDiffs(allocator, stream, store, run_id, query); } else if (std.mem.startsWith(u8, path, "/api/runs/")) { const id_str = path["/api/runs/".len..]; if (std.mem.endsWith(u8, id_str, "/diffs/summary")) { const run_id_str = id_str[0 .. id_str.len - "/diffs/summary".len]; const id = std.fmt.parseInt(i64, run_id_str, 10) catch { try sendResponse(stream, "400 Bad Request", "text/plain", "invalid run id"); return; }; try serveRunDiffSummary(allocator, stream, store, id, query); return; } else if (std.mem.endsWith(u8, id_str, "/diffs")) { const run_id_str = id_str[0 .. id_str.len - "/diffs".len]; const id = std.fmt.parseInt(i64, run_id_str, 10) catch { try sendResponse(stream, "400 Bad Request", "text/plain", "invalid run id"); return; }; try serveRunDiffs(allocator, stream, store, id, query); return; } const id = std.fmt.parseInt(i64, id_str, 10) catch { try sendResponse(stream, "400 Bad Request", "text/plain", "invalid run id"); return; }; try serveRunDetail(allocator, stream, store, id); } else if (std.mem.eql(u8, path, "/api/trend")) { const limit = parseLimitParam(query, trend_limit, store); const quorum = parseQuorumParam(query, 2); try serveTrend(allocator, stream, store, query, limit, quorum); } else if (std.mem.eql(u8, path, "/api/status")) { try serveStatus(allocator, stream, store); } else if (std.mem.eql(u8, path, "/api/relays")) { try serveRelays(allocator, stream, store); } else if (std.mem.eql(u8, path, "/api/relays/history")) { try serveRelaysHistory(allocator, stream, store, query); } else if (std.mem.eql(u8, path, "/api/relays/events")) { try serveRelaysEvents(allocator, stream, store, query); } else if (std.mem.eql(u8, path, "/api/relays/live")) { try serveRelaysLive(allocator, stream, meters); } else if (std.mem.eql(u8, path, "/api/rate-profile")) { try serveRateProfile(allocator, io, stream, store, query); } else if (std.mem.eql(u8, path, "/og")) { try serveOgPng(allocator, stream, og_png_path); } else if (std.mem.eql(u8, path, "/og.svg")) { try serveOgSvg(allocator, stream, store); } else { try sendResponse(stream, "404 Not Found", "text/plain", "not found"); } } fn parseLimitParam(query: []const u8, default: u32, store: *Store) u32 { const max = store.getRunCount() catch default; if (query.len == 0) return @min(default, max); var it = std.mem.splitScalar(u8, query, '&'); while (it.next()) |param| { if (std.mem.startsWith(u8, param, "limit=")) { const val = std.fmt.parseInt(u32, param["limit=".len..], 10) catch return @min(default, max); const upper = if (max > 0) max else default; return std.math.clamp(val, @min(2, upper), upper); } } return @min(default, max); } /// quorum for the trend denominator: coverage is computed against DIDs seen by /// at least this many relays. default 2 strips any single relay's idiosyncratic /// DIDs from the shared denominator. quorum=1 reproduces the raw-union history. fn parseQuorumParam(query: []const u8, default: u16) u16 { if (query.len == 0) return default; var it = std.mem.splitScalar(u8, query, '&'); while (it.next()) |param| { if (std.mem.startsWith(u8, param, "quorum=")) { return std.fmt.parseInt(u16, param["quorum=".len..], 10) catch default; } } return default; } fn serveRunCount(stream: net.Stream, store: *Store) !void { const count = try store.getRunCount(); var buf: [64]u8 = undefined; const json = std.fmt.bufPrint(&buf, "{{\"count\":{d}}}", .{count}) catch return; try sendResponseCached(stream, "200 OK", "application/json", json, snapshot_cache_control); } fn serveRuns(allocator: std.mem.Allocator, stream: net.Stream, store: *Store) !void { const runs = try store.getRecentRuns(allocator, 50); defer { for (runs) |r| allocator.free(r.timestamp); allocator.free(runs); } var json: std.ArrayList(u8) = .empty; defer json.deinit(allocator); try json.appendSlice(allocator, "["); for (runs, 0..) |r, i| { if (i > 0) try json.appendSlice(allocator, ","); try json.print(allocator, "{{\"id\":{d},\"timestamp\":\"{s}\",\"window_seconds\":{d},\"union_dids\":{d}}}", .{ r.id, r.timestamp, r.window_seconds, r.union_dids, }); } try json.appendSlice(allocator, "]"); try sendResponseCached(stream, "200 OK", "application/json", json.items, snapshot_cache_control); } fn serveLatest(allocator: std.mem.Allocator, stream: net.Stream, store: *Store) !void { const run_id = try store.getLatestRunId() orelse { try sendResponse(stream, "200 OK", "application/json", "null"); return; }; try serveRunDetail(allocator, stream, store, run_id); } fn serveRunDetail(allocator: std.mem.Allocator, stream: net.Stream, store: *Store, run_id: i64) !void { const run_meta = try store.getRun(allocator, run_id) orelse { try sendResponse(stream, "404 Not Found", "text/plain", "run not found"); return; }; defer allocator.free(run_meta.timestamp); const stats = try store.getRunStats(allocator, run_id); defer { for (stats) |s| { allocator.free(s.host); if (s.overlap) |o| allocator.free(o); } allocator.free(stats); } const diff_counts = try store.getDiffCounts(allocator, run_id); defer { for (diff_counts) |dc| { allocator.free(dc.relay); allocator.free(dc.classification); } allocator.free(diff_counts); } const diff_samples = try store.getDiffSamples(allocator, run_id, 30); defer { for (diff_samples) |d| { allocator.free(d.relay); allocator.free(d.did); allocator.free(d.classification); if (d.pds) |p| allocator.free(p); } allocator.free(diff_samples); } var json: std.ArrayList(u8) = .empty; defer json.deinit(allocator); var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const a = arena.allocator(); const parsed = try a.alloc(?[]u32, stats.len); var all_present = stats.len > 0; var relay_count: usize = 0; for (stats, 0..) |s, i| { parsed[i] = if (s.overlap) |o| coverage.parseHistogram(a, o) else null; if (parsed[i] == null) all_present = false; if (s.source_type == .relay) relay_count += 1; } const coverage_quorum: u16 = if (all_present and relay_count > 1) 2 else 1; var coverage_union_dids: i32 = run_meta.union_dids; if (coverage_quorum > 1) { const hslice = try a.alloc([]const u32, relay_count); var relay_idx: usize = 0; for (stats, 0..) |s, i| { if (s.source_type != .relay) continue; hslice[relay_idx] = parsed[i].?; relay_idx += 1; } coverage_union_dids = @intCast(coverage.unionAtQuorum(hslice, coverage_quorum)); } try json.print(allocator, "{{\"id\":{d},\"timestamp\":\"{s}\",\"window_seconds\":{d},\"union_dids\":{d},\"coverage_quorum\":{d},\"coverage_union_dids\":{d},\"stats\":[", .{ run_id, run_meta.timestamp, run_meta.window_seconds, run_meta.union_dids, coverage_quorum, coverage_union_dids, }); for (stats, 0..) |s, i| { if (i > 0) try json.appendSlice(allocator, ","); const coverage_dids: i32 = if (coverage_quorum > 1) @intCast(coverage.numeratorAtQuorum(parsed[i].?, coverage_quorum)) else @intCast(s.unique_dids); try json.print(allocator, "{{\"host\":\"{s}\",\"source_type\":\"{s}\",\"events\":{d},\"unique_dids\":{d},\"coverage_dids\":{d},\"connected\":{}}}", .{ s.host, s.source_type.text(), s.events, s.unique_dids, coverage_dids, s.connected, }); } try json.appendSlice(allocator, "],\"diff_counts\":["); for (diff_counts, 0..) |dc, i| { if (i > 0) try json.appendSlice(allocator, ","); try json.print(allocator, "{{\"relay\":\"{s}\",\"classification\":\"{s}\",\"classification_version\":{d},\"count\":{d}}}", .{ dc.relay, dc.classification, dc.classification_version, dc.count, }); } try json.appendSlice(allocator, "],\"diff_samples\":["); for (diff_samples, 0..) |d, i| { if (i > 0) try json.appendSlice(allocator, ","); try json.print(allocator, "{{\"relay\":\"{s}\",\"did\":\"{s}\",\"classification\":\"{s}\",\"classification_version\":{d}", .{ d.relay, d.did, d.classification, d.classification_version, }); if (d.pds) |pds| try json.print(allocator, ",\"pds\":\"{s}\"", .{pds}); try json.append(allocator, '}'); } try json.appendSlice(allocator, "]}"); try sendResponseCached(stream, "200 OK", "application/json", json.items, snapshot_cache_control); } const default_diff_limit: u32 = 1000; const max_diff_limit: u32 = 5000; fn isValidClassification(s: []const u8) bool { return s.len == 0 or std.mem.eql(u8, s, "active_missing") or std.mem.eql(u8, s, "no_pds_endpoint") or std.mem.eql(u8, s, "malformed_did") or std.mem.eql(u8, s, "unsupported_did_method") or std.mem.eql(u8, s, "did_resolution_failed") or std.mem.eql(u8, s, "invalid_did_document") or std.mem.eql(u8, s, "classification_not_attempted") or std.mem.eql(u8, s, "coverage_gap") or std.mem.eql(u8, s, "unresolvable") or std.mem.eql(u8, s, "deactivated"); } fn isValidLane(s: []const u8) bool { return s.len == 0 or std.mem.eql(u8, s, "live") or std.mem.eql(u8, s, "review") or std.mem.eql(u8, s, "inactive") or std.mem.eql(u8, s, "unchecked") or std.mem.eql(u8, s, "checked"); } fn parseDiffLimit(query: []const u8) u32 { if (parseQueryString(query, "limit")) |raw| { return std.math.clamp(std.fmt.parseInt(u32, raw, 10) catch default_diff_limit, 1, max_diff_limit); } return default_diff_limit; } fn parseDiffCursor(query: []const u8) i64 { if (parseQueryString(query, "cursor")) |raw| { return std.fmt.parseInt(i64, raw, 10) catch 0; } return 0; } fn isValidDidSearch(s: []const u8) bool { if (s.len > 128) return false; for (s) |ch| { if ((ch >= 'a' and ch <= 'z') or (ch >= 'A' and ch <= 'Z') or (ch >= '0' and ch <= '9') or ch == ':' or ch == '.' or ch == '-' or ch == '_' or ch == '*') { continue; } return false; } return true; } fn isValidPdsHost(s: []const u8) bool { if (s.len > 128) return false; for (s) |ch| { if ((ch >= 'a' and ch <= 'z') or (ch >= 'A' and ch <= 'Z') or (ch >= '0' and ch <= '9') or ch == '.' or ch == '-' or ch == '_') { continue; } return false; } return true; } fn parseDiffFilters(stream: net.Stream, query: []const u8) !?struct { relay: []const u8, classification: []const u8, lane: []const u8, did_search: []const u8, pds_host: []const u8, } { const relay = parseQueryString(query, "relay") orelse ""; if (relay.len > 0 and !isValidHostname(relay)) { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"invalid relay\"}"); return null; } const classification = parseQueryString(query, "classification") orelse ""; if (!isValidClassification(classification)) { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"invalid classification\"}"); return null; } const lane = parseQueryString(query, "lane") orelse ""; if (!isValidLane(lane)) { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"invalid lane\"}"); return null; } const did_search = parseQueryString(query, "q") orelse ""; if (!isValidDidSearch(did_search)) { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"invalid DID search\"}"); return null; } const pds_host = parseQueryString(query, "pds_host") orelse ""; if (!isValidPdsHost(pds_host)) { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"invalid pds_host\"}"); return null; } return .{ .relay = relay, .classification = classification, .lane = lane, .did_search = did_search, .pds_host = pds_host, }; } fn serveRunDiffs(allocator: std.mem.Allocator, stream: net.Stream, store: *Store, run_id: i64, query: []const u8) !void { const run_meta = try store.getRun(allocator, run_id) orelse { try sendResponse(stream, "404 Not Found", "application/json", "{\"error\":\"run not found\"}"); return; }; allocator.free(run_meta.timestamp); const filters = (try parseDiffFilters(stream, query)) orelse return; const limit = parseDiffLimit(query); const cursor = parseDiffCursor(query); const diffs = try store.getDiffPage(allocator, run_id, filters.relay, filters.classification, filters.lane, filters.did_search, filters.pds_host, cursor, limit + 1); defer { for (diffs) |d| { allocator.free(d.relay); allocator.free(d.did); allocator.free(d.classification); if (d.pds) |p| allocator.free(p); } allocator.free(diffs); } const limit_usize: usize = @intCast(limit); const emit_len = @min(diffs.len, limit_usize); const next_cursor: ?i64 = if (diffs.len > limit_usize) diffs[emit_len - 1].id else null; var json: std.ArrayList(u8) = .empty; defer json.deinit(allocator); try json.print(allocator, "{{\"run_id\":{d},\"limit\":{d}", .{ run_id, limit }); if (filters.relay.len > 0) try json.print(allocator, ",\"relay\":\"{s}\"", .{filters.relay}); if (filters.classification.len > 0) try json.print(allocator, ",\"classification\":\"{s}\"", .{filters.classification}); if (filters.lane.len > 0) try json.print(allocator, ",\"lane\":\"{s}\"", .{filters.lane}); if (filters.did_search.len > 0) try json.print(allocator, ",\"q\":\"{s}\"", .{filters.did_search}); if (filters.pds_host.len > 0) try json.print(allocator, ",\"pds_host\":\"{s}\"", .{filters.pds_host}); try json.appendSlice(allocator, ",\"diffs\":["); for (diffs[0..emit_len], 0..) |d, i| { if (i > 0) try json.append(allocator, ','); try json.print(allocator, "{{\"cursor\":{d},\"relay\":\"{s}\",\"did\":\"{s}\",\"classification\":\"{s}\",\"classification_version\":{d}", .{ d.id, d.relay, d.did, d.classification, d.classification_version, }); if (d.pds) |pds| try json.print(allocator, ",\"pds\":\"{s}\"", .{pds}); try json.append(allocator, '}'); } try json.appendSlice(allocator, "],\"next_cursor\":"); if (next_cursor) |c| { try json.print(allocator, "{d}", .{c}); } else { try json.appendSlice(allocator, "null"); } try json.append(allocator, '}'); try sendResponse(stream, "200 OK", "application/json", json.items); } fn appendPdsHost(json: *std.ArrayList(u8), allocator: std.mem.Allocator, pds: []const u8) !void { var host = pds; if (std.mem.startsWith(u8, host, "https://")) host = host["https://".len..]; if (std.mem.startsWith(u8, host, "http://")) host = host["http://".len..]; if (std.mem.indexOfScalar(u8, host, '/')) |slash| host = host[0..slash]; try json.print(allocator, "\"{s}\"", .{host}); } fn serveRunDiffSummary(allocator: std.mem.Allocator, stream: net.Stream, store: *Store, run_id: i64, query: []const u8) !void { const run_meta = try store.getRun(allocator, run_id) orelse { try sendResponse(stream, "404 Not Found", "application/json", "{\"error\":\"run not found\"}"); return; }; allocator.free(run_meta.timestamp); const filters = (try parseDiffFilters(stream, query)) orelse return; const pds_limit = std.math.clamp(parseDiffLimit(query), 1, 100); const summary = try store.getDiffSummary(run_id, filters.relay, filters.classification, filters.lane, filters.did_search, filters.pds_host); const relays = try store.getDiffRelayFacets(allocator, run_id, filters.relay, filters.classification, filters.lane, filters.did_search, filters.pds_host); defer { for (relays) |r| allocator.free(r.relay); allocator.free(relays); } const classes = try store.getDiffClassificationFacets(allocator, run_id, filters.relay, filters.classification, filters.lane, filters.did_search, filters.pds_host); defer { for (classes) |c| allocator.free(c.classification); allocator.free(classes); } const pds = try store.getDiffPdsFacets(allocator, run_id, filters.relay, filters.classification, filters.lane, filters.did_search, filters.pds_host, pds_limit); defer { for (pds) |p| allocator.free(p.pds); allocator.free(pds); } var live_lower_bound: u32 = 0; for (classes) |c| { if (std.mem.eql(u8, c.classification, "active_missing") or std.mem.eql(u8, c.classification, "coverage_gap")) { live_lower_bound += c.count; } } const unchecked = summary.total - summary.checked; const checked_limited = unchecked > 0; var json: std.ArrayList(u8) = .empty; defer json.deinit(allocator); try json.print(allocator, "{{\"run_id\":{d},\"filters\":{{", .{run_id}); var filter_first = true; if (filters.relay.len > 0) { try json.print(allocator, "\"relay\":\"{s}\"", .{filters.relay}); filter_first = false; } if (filters.classification.len > 0) { if (!filter_first) try json.append(allocator, ','); try json.print(allocator, "\"classification\":\"{s}\"", .{filters.classification}); filter_first = false; } if (filters.lane.len > 0) { if (!filter_first) try json.append(allocator, ','); try json.print(allocator, "\"lane\":\"{s}\"", .{filters.lane}); filter_first = false; } if (filters.did_search.len > 0) { if (!filter_first) try json.append(allocator, ','); try json.print(allocator, "\"q\":\"{s}\"", .{filters.did_search}); filter_first = false; } if (filters.pds_host.len > 0) { if (!filter_first) try json.append(allocator, ','); try json.print(allocator, "\"pds_host\":\"{s}\"", .{filters.pds_host}); } try json.print(allocator, "}},\"total\":{d},\"checked\":{d},\"unchecked\":{d}", .{ summary.total, summary.checked, unchecked }); try json.print(allocator, ",\"estimate\":{{\"live_gap_lower_bound\":{d},\"live_gap_display\":\"{d}{s}\",\"live_gap_more_possible\":{},\"checked_limited\":{}}}", .{ live_lower_bound, live_lower_bound, if (checked_limited) "+" else "", checked_limited, checked_limited, }); try json.appendSlice(allocator, ",\"facets\":{\"relays\":["); for (relays, 0..) |r, i| { if (i > 0) try json.append(allocator, ','); const relay_limited = r.total > r.checked; try json.print(allocator, "{{\"relay\":\"{s}\",\"total\":{d},\"checked\":{d},\"unchecked\":{d},\"live_lower_bound\":{d},\"live_display\":\"{d}{s}\",\"review\":{d},\"inactive\":{d}}}", .{ r.relay, r.total, r.checked, r.total - r.checked, r.live, r.live, if (relay_limited) "+" else "", r.review, r.inactive, }); } try json.appendSlice(allocator, "],\"classifications\":["); for (classes, 0..) |c, i| { if (i > 0) try json.append(allocator, ','); try json.print(allocator, "{{\"classification\":\"{s}\",\"classification_version\":{d},\"count\":{d}}}", .{ c.classification, c.classification_version, c.count, }); } try json.appendSlice(allocator, "],\"pds\":["); for (pds, 0..) |p, i| { if (i > 0) try json.append(allocator, ','); try json.print(allocator, "{{\"pds\":\"{s}\",\"pds_host\":", .{p.pds}); try appendPdsHost(&json, allocator, p.pds); try json.print(allocator, ",\"total\":{d},\"live\":{d},\"review\":{d},\"inactive\":{d}}}", .{ p.total, p.live, p.review, p.inactive, }); } try json.appendSlice(allocator, "]}}"); try sendResponse(stream, "200 OK", "application/json", json.items); } fn parseTrendTimestamp(raw: []const u8, buf: []u8) ?[]const u8 { if (raw.len == 0) return null; var all_digits = true; for (raw) |ch| { if (ch < '0' or ch > '9') { all_digits = false; break; } } if (all_digits) { var epoch = std.fmt.parseInt(i64, raw, 10) catch return null; // Accept Unix milliseconds as a convenience, but emit stored UTC seconds. if (epoch > 100_000_000_000) epoch = @divTrunc(epoch, 1000); return formatIsoEpoch(epoch, buf); } if (!isValidIsoTimestamp(raw)) return null; return raw; } fn relaySelected(relays_csv: []const u8, host: []const u8) bool { if (relays_csv.len == 0) return true; var it = std.mem.splitScalar(u8, relays_csv, ','); while (it.next()) |relay| { if (std.mem.eql(u8, relay, host)) return true; } return false; } fn serveTrend(allocator: std.mem.Allocator, stream: net.Stream, store: *Store, query: []const u8, limit: u32, quorum: u16) !void { var since_buf: [32]u8 = undefined; var until_buf: [32]u8 = undefined; const since_raw = parseQueryString(query, "since"); const until_raw = parseQueryString(query, "until"); const use_range = since_raw != null or until_raw != null; const since = if (since_raw) |raw| parseTrendTimestamp(raw, &since_buf) orelse { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"invalid since\"}"); return; } else "0000"; const until = if (until_raw) |raw| parseTrendTimestamp(raw, &until_buf) orelse { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"invalid until\"}"); return; } else "9999"; const max_runs: u32 = if (parseQueryString(query, "points")) |raw| std.fmt.parseInt(u32, raw, 10) catch 0 else 0; const points = if (use_range) try store.getTrendDataRange(allocator, since, until, max_runs) else try store.getTrendData(allocator, limit); defer { for (points) |p| { allocator.free(p.timestamp); allocator.free(p.host); allocator.free(p.overlap); } allocator.free(points); } const relays = parseQueryString(query, "relays") orelse ""; // parse each row's overlap histogram once (arena-scoped to this request). var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const a = arena.allocator(); const parsed = try a.alloc(?[]u32, points.len); for (points, 0..) |p, i| parsed[i] = coverage.parseHistogram(a, p.overlap); var json: std.ArrayList(u8) = .empty; defer json.deinit(allocator); try json.appendSlice(allocator, "["); // points arrive grouped by run (run_id ASC, then host). per run, recompute // union(q) and each relay's numerator(q) from the histograms. a run missing // any histogram (old data) — or quorum=1 — falls back to the stored raw // union, so the response shape and q=1 history are unchanged. var first = true; var i: usize = 0; while (i < points.len) { const run_id = points[i].run_id; var j = i; while (j < points.len and points[j].run_id == run_id) j += 1; const group = points[i..j]; const parsed_group = parsed[i..j]; var all_present = true; var relay_count: usize = 0; for (parsed_group) |h| { if (h == null) { all_present = false; break; } } for (group) |p| if (p.source_type == .relay) { relay_count += 1; }; // a DID can be seen by at most `group.len` relays; clamp so an absurd // quorum can't zero out the denominator. const eff_q: u16 = if (all_present) @min(@max(quorum, 1), @as(u16, @intCast(relay_count))) else 1; const use_quorum = all_present and eff_q > 1; var union_val: i64 = points[i].union_dids; if (use_quorum) { const hslice = try a.alloc([]const u32, relay_count); var relay_idx: usize = 0; for (group, 0..) |p, k| { if (p.source_type != .relay) continue; hslice[relay_idx] = parsed_group[k].?; relay_idx += 1; } union_val = coverage.unionAtQuorum(hslice, eff_q); } for (group, 0..) |p, k| { if (!relaySelected(relays, p.host)) continue; if (!first) try json.appendSlice(allocator, ","); first = false; const dids: i64 = if (use_quorum) coverage.numeratorAtQuorum(parsed_group[k].?, eff_q) else p.unique_dids; try json.print(allocator, "{{\"ts\":\"{s}\",\"union\":{d},\"host\":\"{s}\",\"source_type\":\"{s}\",\"dids\":{d},\"connected\":{},\"connected_runs\":{d},\"range_runs\":{d},\"quorum\":{d}}}", .{ p.timestamp, union_val, p.host, p.source_type.text(), dids, p.connected, p.connected_runs, p.range_runs, eff_q, }); } i = j; } try json.appendSlice(allocator, "]"); try sendResponseCached(stream, "200 OK", "application/json", json.items, snapshot_cache_control); } // --- "behind lately" status endpoint --- // // GET /api/status // // the lay-visitor question: is any relay behind the network right now-ish? // scores each relay's last `status_window_runs` valid runs against the run's // quorum-2 consensus union (verdict math lives in coverage.zig: a run is // "behind" below 85% of consensus; "behind lately" = behind in ≥ 1/3 of scored // runs). network-absolute by design — the self-relative signal is /api/relays. // // runs predating overlap histograms fall back to the raw union (same rule as // /api/trend); scoped-relay carve-outs are a UI concern, every host is scored. const status_window_runs: u32 = 24; fn serveStatus(allocator: std.mem.Allocator, stream: net.Stream, store: *Store) !void { const points = try store.getTrendData(allocator, status_window_runs); defer { for (points) |p| { allocator.free(p.timestamp); allocator.free(p.host); allocator.free(p.overlap); } allocator.free(points); } var arena = std.heap.ArenaAllocator.init(allocator); defer arena.deinit(); const a = arena.allocator(); const parsed = try a.alloc(?[]u32, points.len); for (points, 0..) |p, i| parsed[i] = coverage.parseHistogram(a, p.overlap); const HostAcc = struct { runs: u32 = 0, behind_runs: u32 = 0, ratio_sum: f64 = 0, latest_ts: []const u8 = "", latest_ratio: f64 = 0, latest_behind: bool = false, latest_connected: bool = false, }; var acc: std.StringArrayHashMapUnmanaged(HostAcc) = .empty; var since: []const u8 = ""; var until: []const u8 = ""; var total_runs: u32 = 0; // points arrive grouped by run; per run recompute the quorum-2 union and // each host's numerator (identical grouping to serveTrend). var i: usize = 0; while (i < points.len) { const run_id = points[i].run_id; var j = i; while (j < points.len and points[j].run_id == run_id) j += 1; const group = points[i..j]; const parsed_group = parsed[i..j]; var all_present = true; for (parsed_group) |h| { if (h == null) { all_present = false; break; } } var relay_count: usize = 0; for (group) |p| if (p.source_type == .relay) { relay_count += 1; }; const use_quorum = all_present and relay_count > 1; var union_val: i64 = points[i].union_dids; if (use_quorum) { const hslice = try a.alloc([]const u32, relay_count); var relay_idx: usize = 0; for (group, 0..) |p, k| { if (p.source_type != .relay) continue; hslice[relay_idx] = parsed_group[k].?; relay_idx += 1; } union_val = coverage.unionAtQuorum(hslice, 2); } if (since.len == 0) since = points[i].timestamp; until = points[i].timestamp; total_runs += 1; for (group, 0..) |p, k| { const dids: i64 = if (use_quorum) coverage.numeratorAtQuorum(parsed_group[k].?, 2) else p.unique_dids; const gop = try acc.getOrPut(a, p.host); if (!gop.found_existing) gop.value_ptr.* = .{}; const e = gop.value_ptr; const ratio: f64 = if (union_val > 0) @as(f64, @floatFromInt(dids)) / @as(f64, @floatFromInt(union_val)) else 0; e.runs += 1; e.ratio_sum += @min(ratio, 1.0); if (coverage.runIsBehind(dids, union_val)) e.behind_runs += 1; e.latest_ts = p.timestamp; e.latest_ratio = ratio; e.latest_behind = coverage.runIsBehind(dids, union_val); e.latest_connected = p.connected; } i = j; } var json: std.ArrayList(u8) = .empty; defer json.deinit(allocator); try json.print(allocator, "{{\"window\":{{\"runs\":{d},\"since\":\"{s}\",\"until\":\"{s}\"}},\"relays\":[", .{ total_runs, since, until, }); var first = true; var it = acc.iterator(); while (it.next()) |entry| { const host = entry.key_ptr.*; const e = entry.value_ptr.*; if (!first) try json.appendSlice(allocator, ","); first = false; const avg_pct_h: u32 = if (e.runs > 0) @intFromFloat(std.math.clamp(e.ratio_sum / @as(f64, @floatFromInt(e.runs)), 0, 1.0) * 10000.0 + 0.5) else 0; const latest_pct_h: u32 = @intFromFloat(std.math.clamp(e.latest_ratio, 0, 10.0) * 10000.0 + 0.5); try json.print( allocator, "{{\"host\":\"{s}\",\"behind_lately\":{},\"behind_runs\":{d},\"runs\":{d},\"avg_coverage_pct\":{d}.{d:0>2}," ++ "\"latest\":{{\"ts\":\"{s}\",\"coverage_pct\":{d}.{d:0>2},\"behind\":{},\"connected\":{}}}}}", .{ host, coverage.isBehindLately(e.behind_runs, e.runs), e.behind_runs, e.runs, avg_pct_h / 100, avg_pct_h % 100, e.latest_ts, latest_pct_h / 100, latest_pct_h % 100, e.latest_behind, e.latest_connected, }, ); } try json.appendSlice(allocator, "]}"); try sendResponseCached(stream, "200 OK", "application/json", json.items, snapshot_cache_control); } fn serveOgPng(allocator: std.mem.Allocator, stream: net.Stream, og_png_path: []const u8) !void { const io = std.Options.debug_io; const file = std.Io.Dir.cwd().openFile(io, og_png_path, .{}) catch { try sendResponse(stream, "404 Not Found", "text/plain", "og image not generated yet"); return; }; defer file.close(io); var file_buf: [4096]u8 = undefined; var reader = file.reader(io, &file_buf); const body = reader.interface.allocRemaining(allocator, .limited(2 * 1024 * 1024)) catch { try sendResponse(stream, "500 Internal Server Error", "text/plain", "failed to read og image"); return; }; defer allocator.free(body); try sendResponse(stream, "200 OK", "image/png", body); } fn readFile( allocator: std.mem.Allocator, path: []const u8, max_bytes: usize, ) !?[]u8 { const io = std.Options.debug_io; const file = std.Io.Dir.cwd().openFile(io, path, .{}) catch return null; defer file.close(io); var file_buf: [4096]u8 = undefined; var reader = file.reader(io, &file_buf); return reader.interface.allocRemaining(allocator, .limited(max_bytes)) catch return error.FileTooLarge; } fn serveSonifyPlaylist(allocator: std.mem.Allocator, stream: net.Stream) !void { const path = sonify_stream_dir ++ "/index.m3u8"; const body = readFile(allocator, path, 256 * 1024) catch { try sendResponse(stream, "500 Internal Server Error", "text/plain", "failed to read live playlist"); return; } orelse { try sendResponseCached( stream, "503 Service Unavailable", "application/json", "{\"live\":false,\"error\":\"live playlist is not ready\"}", "no-store", ); return; }; defer allocator.free(body); try sendResponseCached( stream, "200 OK", "application/vnd.apple.mpegurl", body, "public, max-age=0, s-maxage=2", ); } fn serveSonifyArtwork(allocator: std.mem.Allocator, stream: net.Stream) !void { const path = sonify_stream_dir ++ "/cover.png"; const body = readFile(allocator, path, 4 * 1024 * 1024) catch { try sendResponse(stream, "500 Internal Server Error", "text/plain", "failed to read live artwork"); return; } orelse { try sendResponseCached( stream, "503 Service Unavailable", "application/json", "{\"live\":false,\"error\":\"live artwork is not ready\"}", "no-store", ); return; }; defer allocator.free(body); try sendResponseCached( stream, "200 OK", "image/png", body, "public, max-age=0, s-maxage=15", ); } fn validSonifySegmentName(name: []const u8) bool { if (!std.mem.startsWith(u8, name, "segment-") or !std.mem.endsWith(u8, name, ".ts") or name.len > 96) { return false; } for (name) |char| { if (!std.ascii.isAlphanumeric(char) and char != '-' and char != '.') { return false; } } return true; } fn lastPlaylistSegment(playlist: []const u8) ?[]const u8 { var lines = std.mem.splitBackwardsScalar(u8, playlist, '\n'); while (lines.next()) |raw_line| { const line = std.mem.trim(u8, raw_line, " \t\r"); if (line.len == 0 or line[0] == '#') continue; return if (validSonifySegmentName(line)) line else null; } return null; } test "live sonification segment paths are constrained to generated names" { try std.testing.expect(validSonifySegmentName( "segment-20260730T190652Z-000000004.ts", )); try std.testing.expect(!validSonifySegmentName("../status.json")); try std.testing.expect(!validSonifySegmentName("segment-nested/file.ts")); try std.testing.expect(!validSonifySegmentName("index.m3u8")); } test "live sonification health follows the playlist's active segment" { const playlist = \\#EXTM3U \\#EXT-X-MEDIA-SEQUENCE:41 \\#EXTINF:4.0, \\segment-run-000000041.ts \\#EXTINF:4.0, \\segment-run-000000042.ts ; try std.testing.expectEqualStrings( "segment-run-000000042.ts", lastPlaylistSegment(playlist).?, ); try std.testing.expect(lastPlaylistSegment( "#EXTM3U\n../outside.ts\n", ) == null); } fn serveSonifySegment( allocator: std.mem.Allocator, stream: net.Stream, name: []const u8, ) !void { if (!validSonifySegmentName(name)) { try sendResponse(stream, "404 Not Found", "text/plain", "not found"); return; } const path = try std.fmt.allocPrint( allocator, "{s}/{s}", .{ sonify_stream_dir, name }, ); defer allocator.free(path); const body = readFile(allocator, path, 4 * 1024 * 1024) catch { try sendResponse(stream, "500 Internal Server Error", "text/plain", "failed to read live segment"); return; } orelse { try sendResponse(stream, "404 Not Found", "text/plain", "segment not found"); return; }; defer allocator.free(body); try sendResponseCached( stream, "200 OK", "video/mp2t", body, "public, max-age=31536000, immutable", ); } fn serveSonifyStatus( allocator: std.mem.Allocator, io: std.Io, stream: net.Stream, ) !void { if (!try sonifyStreamIsLive(allocator, io)) { try sendSonifyOfflineStatus(stream); return; } const path = sonify_stream_dir ++ "/status.json"; const body = readFile(allocator, path, 64 * 1024) catch { try sendSonifyLiveStatus(stream); return; } orelse { try sendSonifyLiveStatus(stream); return; }; defer allocator.free(body); try sendResponseCached(stream, "200 OK", "application/json", body, "no-store"); } fn sonifyStreamIsLive(allocator: std.mem.Allocator, io: std.Io) !bool { const playlist_path = sonify_stream_dir ++ "/index.m3u8"; const playlist = readFile(allocator, playlist_path, 256 * 1024) catch return false; defer if (playlist) |body| allocator.free(body); const body = playlist orelse return false; const segment_name = lastPlaylistSegment(body) orelse return false; const segment_path = try std.fmt.allocPrint( allocator, "{s}/{s}", .{ sonify_stream_dir, segment_name }, ); defer allocator.free(segment_path); const stat = std.Io.Dir.cwd().statFile(io, segment_path, .{}) catch return false; const now = std.Io.Timestamp.now(io, .real); const age_seconds = @max( @as(i64, 0), stat.mtime.durationTo(now).toSeconds(), ); return age_seconds <= sonify_segment_max_age_seconds; } fn sendSonifyLiveStatus(stream: net.Stream) !void { try sendResponseCached( stream, "200 OK", "application/json", "{\"live\":true,\"waveform\":\"sine\",\"playlist_url\":\"" ++ sonify_playlist_url ++ "\",\"artwork_url\":\"" ++ sonify_artwork_url ++ "\"}", "no-store", ); } fn sendSonifyOfflineStatus(stream: net.Stream) !void { try sendResponseCached( stream, "200 OK", "application/json", "{\"live\":false,\"waveform\":\"sine\",\"playlist_url\":\"" ++ sonify_playlist_url ++ "\",\"artwork_url\":\"" ++ sonify_artwork_url ++ "\"}", "no-store", ); } fn serveOgSvg(allocator: std.mem.Allocator, stream: net.Stream, store: *Store) !void { const run_id = try store.getLatestRunId() orelse { try sendResponse(stream, "200 OK", "image/svg+xml", \\ \\ \\relay-eval \\no data yet \\ ); return; }; const run_meta = try store.getRun(allocator, run_id) orelse return; defer allocator.free(run_meta.timestamp); const stats = try store.getRunStats(allocator, run_id); defer { for (stats) |s| allocator.free(s.host); allocator.free(stats); } var svg: std.ArrayList(u8) = .empty; defer svg.deinit(allocator); try svg.appendSlice(allocator, \\ \\ \\ \\ \\ \\ \\ \\relay-eval \\comparing what each atproto relay sees ); const max_rows: usize = @min(stats.len, 10); // outlier detection: if a relay has >1.5x median events, it's likely replaying; // use the max unique_dids among non-outliers as the effective union var events_sorted: [32]u64 = undefined; const n_events = @min(stats.len, 32); for (stats[0..n_events], 0..) |s, idx| events_sorted[idx] = s.events; std.mem.sort(u64, events_sorted[0..n_events], {}, std.sort.asc(u64)); const median_events = if (n_events > 0) events_sorted[n_events / 2] else 0; const threshold = median_events + median_events * 3 / 10; // 1.3x var effective_union: u64 = @intCast(@max(run_meta.union_dids, 1)); if (median_events > 0) { var max_non_outlier: u64 = 1; var has_outlier = false; for (stats) |s| { if (s.events > threshold) { has_outlier = true; } else if (s.events > 0 and s.unique_dids > max_non_outlier) { max_non_outlier = s.unique_dids; } } if (has_outlier) effective_union = max_non_outlier; } for (stats[0..max_rows], 0..) |s, i| { const y: u32 = 130 + @as(u32, @intCast(i)) * 42; const pct_x10: u64 = @as(u64, s.unique_dids) * 1000 / effective_union; const pct_int: u32 = @intCast(pct_x10 / 10); const pct_frac: u32 = @intCast(pct_x10 % 10); const bar_w: u32 = @intCast(@min(@as(u64, s.unique_dids) * 600 / effective_union, 600)); const color: []const u8 = if (pct_x10 >= 990) "#3fb950" else if (pct_x10 >= 950) "#58a6ff" else if (pct_x10 >= 800) "#bc8cff" else "#8b949e"; try svg.print(allocator, "{s}", .{ y + 18, s.host }); if (bar_w > 0) { try svg.print(allocator, "", .{ y + 1, bar_w, color }); } try svg.print(allocator, "{d}.{d}%", .{ y + 18, pct_int, pct_frac }); } // footer const window_min: i32 = @divTrunc(run_meta.window_seconds, 60); try svg.appendSlice(allocator, ""); try appendFormattedInt(&svg, allocator, @intCast(effective_union)); try svg.print(allocator, " active accounts · {d} minute window", .{window_min}); try sendResponse(stream, "200 OK", "image/svg+xml", svg.items); } fn appendFormattedInt(svg: *std.ArrayList(u8), allocator: std.mem.Allocator, value: i32) !void { if (value < 0) { try svg.append(allocator, '-'); return appendFormattedInt(svg, allocator, -value); } if (value >= 1000) { try appendFormattedInt(svg, allocator, @divTrunc(value, 1000)); const r: u32 = @intCast(@rem(value, 1000)); try svg.append(allocator, ','); try svg.append(allocator, '0' + @as(u8, @intCast(r / 100))); try svg.append(allocator, '0' + @as(u8, @intCast(r / 10 % 10))); try svg.append(allocator, '0' + @as(u8, @intCast(r % 10))); } else { var buf: [16]u8 = undefined; const s = std.fmt.bufPrint(&buf, "{d}", .{value}) catch unreachable; try svg.appendSlice(allocator, s); } } // --- phi monitors endpoint --- // // shape is stable and documented — phi (bluesky bot) polls this and reports // meaningful status transitions. it does not interpret; it reports what's // written in `headline`. see docs for details. // // thresholds (a host's short-window coverage vs baseline coverage): // ratio < 0.70 or disconnected in majority of recent runs → "critical" // ratio < 0.90 → "degraded" // otherwise → "nominal" // // short window = last 3 valid runs (~15 min at 5-min cadence). // baseline window = last 24 valid runs (~2h). // coverages are clamped to [0, 1] for comparison to avoid replay-induced // inflation flipping healthy relays into "critical". const short_window_runs: u32 = 3; const baseline_window_runs: u32 = 24; const Status = enum { nominal, degraded, critical, fn text(self: Status) []const u8 { return switch (self) { .nominal => "nominal", .degraded => "degraded", .critical => "critical", }; } }; fn classifyMonitor(short: Store.MonitorSnapshot, baseline_opt: ?Store.MonitorSnapshot) Status { // connectivity check dominates: if the relay isn't actually delivering, // nothing else matters. if (short.n_runs > 0 and short.connected_runs == 0) return .critical; if (short.n_runs >= 3 and short.connected_runs <= short.n_runs / 2) return .critical; // coverage ratio vs baseline. clamp both sides to 1.0 so a replay-heavy // run doesn't push a healthy relay into "degraded" via inflated baseline. const s_cov = @min(short.coverage, 1.0); const base_cov = if (baseline_opt) |b| @min(b.coverage, 1.0) else s_cov; if (base_cov <= 0) return .nominal; const ratio = s_cov / base_cov; if (ratio < 0.70) return .critical; if (ratio < 0.90) return .degraded; return .nominal; } fn findSnapshot(list: []const Store.MonitorSnapshot, host: []const u8) ?Store.MonitorSnapshot { for (list) |s| { if (std.mem.eql(u8, s.host, host)) return s; } return null; } fn formatIsoNow(buf: []u8) []const u8 { const now = currentEpochSeconds(); const es = std.time.epoch.EpochSeconds{ .secs = @intCast(now) }; const day = es.getEpochDay(); const yd = day.calculateYearDay(); const md = yd.calculateMonthDay(); const ds = es.getDaySeconds(); return std.fmt.bufPrint(buf, "{d}-{d:0>2}-{d:0>2}T{d:0>2}:{d:0>2}:{d:0>2}Z", .{ yd.year, @as(u32, @intFromEnum(md.month)), @as(u32, md.day_index) + 1, ds.getHoursIntoDay(), ds.getMinutesIntoHour(), ds.getSecondsIntoMinute(), }) catch "1970-01-01T00:00:00Z"; } fn currentEpochSeconds() i64 { return @intCast(@divFloor(std.Io.Timestamp.now(std.Options.debug_io, .real).nanoseconds, std.time.ns_per_s)); } fn pctInt(v: f64) u32 { const clamped = std.math.clamp(v, 0.0, 10.0); // allow >1.0 to be visible in headlines return @intFromFloat(clamped * 100.0 + 0.5); } fn writeHeadline( out: *std.ArrayList(u8), allocator: std.mem.Allocator, host: []const u8, status: Status, short: Store.MonitorSnapshot, baseline_opt: ?Store.MonitorSnapshot, ) !void { const short_pct = pctInt(short.coverage); const base_pct = if (baseline_opt) |b| pctInt(b.coverage) else short_pct; switch (status) { .nominal => { try out.print(allocator, "{s}: {d}% coverage over last {d} eval runs", .{ host, short_pct, short.n_runs }); }, .degraded => { try out.print(allocator, "{s}: coverage at {d}% over last {d} runs (baseline {d}%)", .{ host, short_pct, short.n_runs, base_pct }); }, .critical => { const disconnected = short.n_runs - short.connected_runs; if (short.connected_runs == 0 and short.n_runs > 0) { try out.print(allocator, "{s}: no firehose connection in last {d} eval runs", .{ host, short.n_runs }); } else if (disconnected > short.n_runs / 2) { try out.print(allocator, "{s}: firehose disconnected in {d} of last {d} eval runs", .{ host, disconnected, short.n_runs }); } else { try out.print(allocator, "{s}: coverage at {d}% over last {d} runs (baseline {d}%)", .{ host, short_pct, short.n_runs, base_pct }); } }, } } fn serveRelays(allocator: std.mem.Allocator, stream: net.Stream, store: *Store) !void { const short = try store.getPerHostStats(allocator, short_window_runs); defer { for (short) |s| allocator.free(s.host); allocator.free(short); } const baseline = try store.getPerHostStats(allocator, baseline_window_runs); defer { for (baseline) |s| allocator.free(s.host); allocator.free(baseline); } var ts_buf: [32]u8 = undefined; const checked_at = formatIsoNow(&ts_buf); var json: std.ArrayList(u8) = .empty; defer json.deinit(allocator); try json.append(allocator, '['); // reusable scratch for building headlines (rebuilt per host) var headline_buf: std.ArrayList(u8) = .empty; defer headline_buf.deinit(allocator); for (short, 0..) |host_short, i| { if (i > 0) try json.append(allocator, ','); const host_baseline = findSnapshot(baseline, host_short.host); const status = classifyMonitor(host_short, host_baseline); // build the headline into a scratch buffer so it can be both appended // to the response AND used as the headline on any transition row. headline_buf.clearRetainingCapacity(); try writeHeadline(&headline_buf, allocator, host_short.host, status, host_short, host_baseline); const headline = headline_buf.items; // look up prior monitor_state. three cases: // 1. no prior state → seed, no transition logged // 2. prior state == current → keep existing last_changed // 3. prior state != current → update last_changed, log transition var last_changed_buf: [32]u8 = undefined; const last_changed: []const u8 = blk: { const prior = try store.getMonitorState(allocator, host_short.host); if (prior) |p| { defer { allocator.free(p.status); allocator.free(p.last_changed); } if (std.mem.eql(u8, p.status, status.text())) { // status unchanged — keep existing last_changed const copy = std.fmt.bufPrint(&last_changed_buf, "{s}", .{p.last_changed}) catch p.last_changed; break :blk copy; } // status transitioned — append to event log AND update state try store.insertTransition(checked_at, host_short.host, p.status, status.text(), headline); try store.upsertMonitorState(host_short.host, status.text(), checked_at); const copy = std.fmt.bufPrint(&last_changed_buf, "{s}", .{checked_at}) catch checked_at; break :blk copy; } else { // first time seeing this host — seed without logging a transition // (there's no "from" state to report) try store.upsertMonitorState(host_short.host, status.text(), checked_at); const copy = std.fmt.bufPrint(&last_changed_buf, "{s}", .{checked_at}) catch checked_at; break :blk copy; } }; // emit monitor object try json.print(allocator, "{{\"name\":\"{s}\",\"status\":\"{s}\",\"headline\":\"{s}\",\"metrics\":{{", .{ host_short.host, status.text(), headline, }); // coverage values are fractions (0..~1 for normal runs, can exceed 1 // transiently during replay). emit as percent with 2 decimals: a // 0.9784 fraction → "97.84". const short_pct_hundredths: u32 = @intFromFloat(std.math.clamp(host_short.coverage, 0, 10.0) * 10000.0 + 0.5); try json.print(allocator, "\"coverage_short_pct\":{d}.{d:0>2}", .{ short_pct_hundredths / 100, short_pct_hundredths % 100 }); if (host_baseline) |b| { const base_pct_hundredths: u32 = @intFromFloat(std.math.clamp(b.coverage, 0, 10.0) * 10000.0 + 0.5); try json.print(allocator, ",\"coverage_baseline_pct\":{d}.{d:0>2},\"baseline_runs\":{d}", .{ base_pct_hundredths / 100, base_pct_hundredths % 100, b.n_runs }); } try json.print( allocator, ",\"events_short\":{d},\"dids_short\":{d},\"connected_runs\":{d},\"short_runs\":{d}}},\"last_changed\":\"{s}\",\"checked_at\":\"{s}\"}}", .{ host_short.total_events, host_short.total_dids, host_short.connected_runs, host_short.n_runs, last_changed, checked_at }, ); } try json.append(allocator, ']'); try sendResponse(stream, "200 OK", "application/json", json.items); } // --- relays history endpoint --- // // GET /api/relays/history?name=&limit= // GET /api/relays/history?name=&since=&until= // // returns per-run coverage history for a single host, oldest → newest, with // precomputed coverage_pct. includes a summary block with mean/min/max // coverage + connected-run count for the window. // // when since/until are set, returns every point in the inclusive range // (no cap). when they're absent, returns the last N runs (default 288 // ≈ 24h, max 2016 ≈ 7d). limit is ignored when since/until are set. // // coverage semantics match /api/relays: // coverage_pct = unique_dids / max(unique_dids in that run) × 100 // self-normalizing — replay inflates only the replaying host's own numbers. const default_history_limit: u32 = 288; // ~24h at 5-min eval cadence const max_history_limit: u32 = 2016; // 7d // decode %XX escapes (and '+' as space) in a query string, preserving the // param structure: '&' and '=' separators are literal in the input, and any // *encoded* '&'/'=' inside a value is left encoded so it cannot change how // the query splits. invalid escapes pass through untouched. truncates at // buf.len (far above any real query here). fn percentDecodeQuery(buf: []u8, query: []const u8) []const u8 { var out: usize = 0; var i: usize = 0; while (i < query.len and out < buf.len) { const ch = query[i]; if (ch == '%' and i + 2 < query.len) { const hi = std.fmt.charToDigit(query[i + 1], 16) catch null; const lo = std.fmt.charToDigit(query[i + 2], 16) catch null; if (hi != null and lo != null) { const decoded: u8 = @intCast(hi.? * 16 + lo.?); if (decoded != '&' and decoded != '=') { buf[out] = decoded; out += 1; i += 3; continue; } } } buf[out] = if (ch == '+') ' ' else ch; out += 1; i += 1; } return buf[0..out]; } test "percentDecodeQuery decodes iso timestamps the way httpx sends them" { var buf: [2048]u8 = undefined; try std.testing.expectEqualStrings( "since=2026-07-31T00:00:00Z&name=atproto.africa", percentDecodeQuery(&buf, "since=2026-07-31T00%3A00%3A00Z&name=atproto.africa"), ); } test "percentDecodeQuery leaves separators and invalid escapes alone" { var buf: [2048]u8 = undefined; // an encoded '&' must not split the value; a bare '%' passes through try std.testing.expectEqualStrings( "q=a%26b&x=100%", percentDecodeQuery(&buf, "q=a%26b&x=100%"), ); try std.testing.expectEqualStrings("q=did:plc:abc", percentDecodeQuery(&buf, "q=did%3Aplc%3Aabc")); } fn parseQueryString(query: []const u8, key: []const u8) ?[]const u8 { var it = std.mem.splitScalar(u8, query, '&'); while (it.next()) |param| { if (std.mem.startsWith(u8, param, key) and param.len > key.len and param[key.len] == '=') { return param[key.len + 1 ..]; } } return null; } // permissive check that a query value looks like an ISO 8601 UTC timestamp // (e.g. "2026-04-17T06:36:30Z"). string-comparable against stored // `runs.timestamp` values which use the same format. fn isValidIsoTimestamp(s: []const u8) bool { if (s.len == 0 or s.len > 32) return false; for (s) |ch| { const ok = (ch >= '0' and ch <= '9') or ch == '-' or ch == ':' or ch == 'T' or ch == 'Z' or ch == '.' or ch == '+'; if (!ok) return false; } return true; } // hostnames we might plausibly track: alphanumerics, dots, hyphens. // cheap validation before binding into SQL. fn isValidHostname(s: []const u8) bool { if (s.len == 0 or s.len > 253) return false; for (s) |ch| { const ok = (ch >= 'a' and ch <= 'z') or (ch >= 'A' and ch <= 'Z') or (ch >= '0' and ch <= '9') or ch == '.' or ch == '-'; if (!ok) return false; } return true; } fn serveRelaysHistory(allocator: std.mem.Allocator, stream: net.Stream, store: *Store, query: []const u8) !void { const name = parseQueryString(query, "name") orelse { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"missing required query parameter: name\"}"); return; }; if (!isValidHostname(name)) { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"invalid name\"}"); return; } const since_raw = parseQueryString(query, "since") orelse ""; const until_raw = parseQueryString(query, "until") orelse ""; if (since_raw.len > 0 and !isValidIsoTimestamp(since_raw)) { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"invalid since (expected ISO 8601)\"}"); return; } if (until_raw.len > 0 and !isValidIsoTimestamp(until_raw)) { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"invalid until (expected ISO 8601)\"}"); return; } const limit: u32 = if (parseQueryString(query, "limit")) |l| std.math.clamp(std.fmt.parseInt(u32, l, 10) catch default_history_limit, 2, max_history_limit) else default_history_limit; const points = try store.getHostHistory(allocator, name, limit, since_raw, until_raw); defer { for (points) |p| allocator.free(p.timestamp); allocator.free(points); } // summary: mean/min/max of clamped coverage over connected points; // connected-run count stands on its own so phi can say "alive 280/288" var sum_cov: f64 = 0; var connected_count: u32 = 0; var min_cov: f64 = std.math.inf(f64); var max_cov: f64 = 0; for (points) |p| { if (!p.connected) continue; connected_count += 1; const c_clamped = @min(p.coverage, 1.0); sum_cov += c_clamped; if (c_clamped < min_cov) min_cov = c_clamped; if (c_clamped > max_cov) max_cov = c_clamped; } if (connected_count == 0) { min_cov = 0; max_cov = 0; } const mean_cov: f64 = if (connected_count > 0) sum_cov / @as(f64, @floatFromInt(connected_count)) else 0; var json: std.ArrayList(u8) = .empty; defer json.deinit(allocator); try json.print(allocator, "{{\"name\":\"{s}\"", .{name}); if (since_raw.len > 0 or until_raw.len > 0) { try json.print(allocator, ",\"since\":\"{s}\",\"until\":\"{s}\"", .{ since_raw, until_raw }); } else { try json.print(allocator, ",\"limit\":{d}", .{limit}); } try json.appendSlice(allocator, ",\"points\":["); for (points, 0..) |p, i| { if (i > 0) try json.append(allocator, ','); const pct_hundredths: u32 = @intFromFloat(std.math.clamp(p.coverage, 0, 10.0) * 10000.0 + 0.5); try json.print( allocator, "{{\"ts\":\"{s}\",\"coverage_pct\":{d}.{d:0>2},\"events\":{d},\"dids\":{d},\"connected\":{}}}", .{ p.timestamp, pct_hundredths / 100, pct_hundredths % 100, p.events, p.unique_dids, p.connected }, ); } try json.appendSlice(allocator, "],\"summary\":{"); const mean_h: u32 = @intFromFloat(mean_cov * 10000.0 + 0.5); const min_h: u32 = @intFromFloat(min_cov * 10000.0 + 0.5); const max_h: u32 = @intFromFloat(max_cov * 10000.0 + 0.5); try json.print( allocator, "\"mean_coverage_pct\":{d}.{d:0>2},\"min_coverage_pct\":{d}.{d:0>2},\"max_coverage_pct\":{d}.{d:0>2},\"connected_runs\":{d},\"total_runs\":{d}}}}}", .{ mean_h / 100, mean_h % 100, min_h / 100, min_h % 100, max_h / 100, max_h % 100, connected_count, points.len, }, ); try sendResponse(stream, "200 OK", "application/json", json.items); } // --- relays events endpoint --- // // GET /api/relays/events // GET /api/relays/events?since=&until=&name= // // returns status-transition log rows, oldest → newest. complements the // snapshot (/api/relays) and per-host timeseries (/api/relays/history) by // exposing the "what changed, when" axis directly. headlines are stored // at transition time and returned verbatim so callers can quote them. // // default window is last 24h. name filter is optional; omitted = all relays. fn formatIsoEpoch(epoch: i64, buf: []u8) []const u8 { const es = std.time.epoch.EpochSeconds{ .secs = @intCast(epoch) }; const day = es.getEpochDay(); const yd = day.calculateYearDay(); const md = yd.calculateMonthDay(); const ds = es.getDaySeconds(); return std.fmt.bufPrint(buf, "{d}-{d:0>2}-{d:0>2}T{d:0>2}:{d:0>2}:{d:0>2}Z", .{ yd.year, @as(u32, @intFromEnum(md.month)), @as(u32, md.day_index) + 1, ds.getHoursIntoDay(), ds.getMinutesIntoHour(), ds.getSecondsIntoMinute(), }) catch "1970-01-01T00:00:00Z"; } fn serveRelaysEvents(allocator: std.mem.Allocator, stream: net.Stream, store: *Store, query: []const u8) !void { var since_buf: [32]u8 = undefined; var until_buf: [32]u8 = undefined; const since = blk: { if (parseQueryString(query, "since")) |raw| { if (!isValidIsoTimestamp(raw)) { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"invalid since (expected ISO 8601)\"}"); return; } break :blk raw; } // default: 24h ago break :blk formatIsoEpoch(currentEpochSeconds() - 86400, &since_buf); }; const until = blk: { if (parseQueryString(query, "until")) |raw| { if (!isValidIsoTimestamp(raw)) { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"invalid until (expected ISO 8601)\"}"); return; } break :blk raw; } // default: now break :blk formatIsoEpoch(currentEpochSeconds(), &until_buf); }; const name = parseQueryString(query, "name") orelse ""; if (name.len > 0 and !isValidHostname(name)) { try sendResponse(stream, "400 Bad Request", "application/json", "{\"error\":\"invalid name\"}"); return; } const transitions = try store.getTransitions(allocator, since, until, name); defer { for (transitions) |t| { allocator.free(t.ts); allocator.free(t.name); if (t.from_status) |fs| allocator.free(fs); allocator.free(t.to_status); allocator.free(t.headline); } allocator.free(transitions); } var json: std.ArrayList(u8) = .empty; defer json.deinit(allocator); try json.append(allocator, '['); for (transitions, 0..) |t, i| { if (i > 0) try json.append(allocator, ','); try json.print(allocator, "{{\"ts\":\"{s}\",\"name\":\"{s}\",", .{ t.ts, t.name }); if (t.from_status) |fs| { try json.print(allocator, "\"from_status\":\"{s}\",", .{fs}); } else { try json.appendSlice(allocator, "\"from_status\":null,"); } try json.print(allocator, "\"to_status\":\"{s}\",\"headline\":\"{s}\"}}", .{ t.to_status, t.headline }); } try json.append(allocator, ']'); try sendResponse(stream, "200 OK", "application/json", json.items); } // --- live rate endpoint --- // // GET /api/relays/live // // returns the current rolling 10s event rate per relay, sourced from // persistent rate-meter threads (one per --relays host passed to `serve`). // shape: [{name, ev_per_sec, connected, window_seconds}, ...] // // frontend polls this at ~1Hz to drive the audio sonification — each relay's // ev/s maps to a tone frequency in Hz (audible-frequency range, fig's idea). // returns [] if no --relays were configured (the binary then falls back to // snapshot-only behavior, but the audio UI hides itself). fn serveRelaysLive(allocator: std.mem.Allocator, stream: net.Stream, meters: []live_rate.LiveRateMeter) !void { var json: std.ArrayList(u8) = .empty; defer json.deinit(allocator); try json.append(allocator, '['); for (meters, 0..) |*m, i| { if (i > 0) try json.append(allocator, ','); const rate = m.ratePerSec(); const rate_int: u64 = @intFromFloat(rate * 100.0 + 0.5); try json.print( allocator, "{{\"name\":\"{s}\",\"source_type\":\"{s}\",\"ev_per_sec\":{d}.{d:0>2},\"connected\":{},\"window_seconds\":{d}}}", .{ m.host, m.source_type.text(), rate_int / 100, rate_int % 100, m.connected.load(.acquire), live_rate.window_seconds, }, ); } try json.append(allocator, ']'); try sendResponse(stream, "200 OK", "application/json", json.items); } /// hour-of-day event-rate profile for one host over a trailing window: /// per-utc-hour median and p10/p25/p75/p90 of events-per-second across /// eval windows. computed from the store on each request, so it stays /// current as history accumulates. ?host= (default bsky.network), /// ?days= (default 90, clamped 7..365). fn serveRateProfile(allocator: std.mem.Allocator, io: std.Io, stream: net.Stream, store: *Store, query: []const u8) !void { const host = parseQueryString(query, "host") orelse "bsky.network"; const days: u32 = if (parseQueryString(query, "days")) |raw| std.math.clamp(std.fmt.parseInt(u32, raw, 10) catch 90, 7, 365) else 90; // ISO date prefix for the lexical timestamp comparison const now_s: i64 = @intCast(@divFloor(std.Io.Timestamp.now(io, .real).nanoseconds, std.time.ns_per_s)); const cutoff_epoch = now_s - @as(i64, days) * 86_400; const epoch_day = std.math.divFloor(i64, cutoff_epoch, 86_400) catch unreachable; const z = epoch_day + 719_468; const era = @divFloor(if (z >= 0) z else z - 146_096, 146_097); const doe = z - era * 146_097; const yoe = @divFloor(doe - @divFloor(doe, 1460) + @divFloor(doe, 36_524) - @divFloor(doe, 146_096), 365); const y = yoe + era * 400; const doy = doe - (365 * yoe + @divFloor(yoe, 4) - @divFloor(yoe, 100)); const mp = @divFloor(5 * doy + 2, 153); const day = doy - @divFloor(153 * mp + 2, 5) + 1; const month: i64 = mp + (if (mp < 10) @as(i64, 3) else @as(i64, -9)); const year: i64 = y + (if (month <= 2) @as(i64, 1) else @as(i64, 0)); var since_buf: [16]u8 = undefined; const since = try std.fmt.bufPrint(&since_buf, "{d:0>4}-{d:0>2}-{d:0>2}", .{ @as(u32, @intCast(year)), @as(u32, @intCast(month)), @as(u32, @intCast(day)), }); const rows = store.getHourlyRates(allocator, host, since) catch { try sendResponse(stream, "500 Internal Server Error", "application/json", "{\"error\":\"query failed\"}"); return; }; defer allocator.free(rows); var buckets: [24]std.ArrayList(f64) = @splat(.empty); defer for (&buckets) |*b| b.deinit(allocator); for (rows) |r| try buckets[r.hour].append(allocator, r.rate); var json: std.ArrayList(u8) = .empty; defer json.deinit(allocator); try json.print(allocator, "{{\"host\":\"{s}\",\"days\":{d},\"windows\":{d},\"hours\":[", .{ host, days, rows.len }); for (&buckets, 0..) |*b, h| { std.mem.sort(f64, b.items, {}, std.sort.asc(f64)); const n = b.items.len; const q = struct { fn at(items: []const f64, num: usize, den: usize) u32 { if (items.len == 0) return 0; return @intFromFloat(@round(items[@min(items.len - 1, items.len * num / den)])); } }; try json.print(allocator, "{s}{{\"h\":{d},\"n\":{d},\"med\":{d},\"p10\":{d},\"p25\":{d},\"p75\":{d},\"p90\":{d}}}", .{ if (h == 0) "" else ",", h, n, q.at(b.items, 1, 2), q.at(b.items, 1, 10), q.at(b.items, 1, 4), q.at(b.items, 3, 4), q.at(b.items, 9, 10), }); } try json.appendSlice(allocator, "]}"); try sendResponseCached(stream, "200 OK", "application/json", json.items, "public, max-age=3600"); } fn sendResponse(stream: net.Stream, status: []const u8, content_type: []const u8, body: []const u8) !void { return sendResponseCached(stream, status, content_type, body, "no-cache"); } fn sendResponseCached( stream: net.Stream, status: []const u8, content_type: []const u8, body: []const u8, cache_control: []const u8, ) !void { var header_buf: [512]u8 = undefined; const header = std.fmt.bufPrint( &header_buf, "HTTP/1.1 {s}\r\n" ++ "Content-Type: {s}\r\n" ++ "Content-Length: {d}\r\n" ++ "Access-Control-Allow-Origin: *\r\n" ++ "Cache-Control: {s}\r\n" ++ "Connection: close\r\n" ++ "\r\n", .{ status, content_type, body.len, cache_control }, ) catch return error.HeaderTooLong; var writer_buf: [4096]u8 = undefined; var writer = net.Stream.writer(stream, std.Options.debug_io, &writer_buf); try writer.interface.writeAll(header); try writer.interface.writeAll(body); try writer.interface.flush(); }