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",
\\
);
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,
\\", .{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();
}