Something went wrong. Try again.
declarative relay deployment on hetzner relay-eval.waow.tech
atproto relay
Something went wrong. Try again.
76 kB · 1909 lines
Zig
at main
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910const 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", \\<svg xmlns="http://www.w3.org/2000/svg" width="1200" height="630" viewBox="0 0 1200 630"> \\<rect width="1200" height="630" fill="#0d1117"/> \\<text x="600" y="300" text-anchor="middle" fill="#c9d1d9" font-family="monospace" font-size="36">relay-eval</text> \\<text x="600" y="340" text-anchor="middle" fill="#8b949e" font-family="monospace" font-size="16">no data yet</text> \\</svg> ); 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, \\<svg xmlns="http://www.w3.org/2000/svg" width="1200" height="630" viewBox="0 0 1200 630"> \\<defs><linearGradient id="bg" x1="0" y1="0" x2="1" y2="1"> \\<stop offset="0%" stop-color="#3fb950" stop-opacity="0.06"/> \\<stop offset="100%" stop-color="#bc8cff" stop-opacity="0.06"/> \\</linearGradient></defs> \\<rect width="1200" height="630" fill="#0d1117"/> \\<rect width="1200" height="630" fill="url(#bg)"/> \\<text x="60" y="55" fill="#e6edf3" font-family="monospace" font-size="36" font-weight="bold">relay-eval</text> \\<text x="60" y="85" fill="#8b949e" font-family="monospace" font-size="16">comparing what each atproto relay sees</text> );
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, "<text x=\"60\" y=\"{d}\" fill=\"#c9d1d9\" font-family=\"monospace\" font-size=\"15\">{s}</text>", .{ y + 18, s.host }); if (bar_w > 0) { try svg.print(allocator, "<rect x=\"400\" y=\"{d}\" width=\"{d}\" height=\"24\" rx=\"3\" fill=\"{s}\" opacity=\"0.7\"/>", .{ y + 1, bar_w, color }); } try svg.print(allocator, "<text x=\"1140\" y=\"{d}\" text-anchor=\"end\" fill=\"#e6edf3\" font-family=\"monospace\" font-size=\"14\">{d}.{d}%</text>", .{ y + 18, pct_int, pct_frac }); }
// footer const window_min: i32 = @divTrunc(run_meta.window_seconds, 60); try svg.appendSlice(allocator, "<text x=\"60\" y=\"590\" fill=\"#8b949e\" font-family=\"monospace\" font-size=\"14\">"); try appendFormattedInt(&svg, allocator, @intCast(effective_union)); try svg.print(allocator, " active accounts · {d} minute window</text></svg>", .{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=<host>&limit=<n>// GET /api/relays/history?name=<host>&since=<iso>&until=<iso>//// 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 cadenceconst 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=<iso>&until=<iso>&name=<host>//// 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();}