From ea663ea1686b34c1e39ba2005cedfcbd08e30276 Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Mon, 8 Jun 2026 20:06:35 -0500 Subject: [PATCH] relay-eval: expose missing repo diff API --- AGENTS.md | 5 + docs/relay-eval-recipes.md | 77 ++- relay-eval/build.zig | 8 +- relay-eval/build.zig.zon | 10 +- relay-eval/justfile | 2 +- relay-eval/src/classifier.zig | 67 ++- relay-eval/src/collector.zig | 18 +- relay-eval/src/live_rate.zig | 26 +- relay-eval/src/main.zig | 103 ++-- relay-eval/src/server.zig | 462 +++++++++++++-- relay-eval/src/static/diffs-api.html | 153 +++++ relay-eval/src/static/index.html | 827 +++++++++++++++++++++++++-- relay-eval/src/store.zig | 400 ++++++++++++- 13 files changed, 1980 insertions(+), 178 deletions(-) create mode 100644 AGENTS.md create mode 100644 relay-eval/src/static/diffs-api.html diff --git a/AGENTS.md b/AGENTS.md new file mode 100644 index 0000000..0bffb89 --- /dev/null +++ b/AGENTS.md @@ -0,0 +1,5 @@ +## Zig + +Use Zig 0.16 for this repository. Do not suggest or use Zig 0.15 as a +fallback. If a build, dependency, or source file is incompatible with Zig +0.16, migrate it forward to Zig 0.16. diff --git a/docs/relay-eval-recipes.md b/docs/relay-eval-recipes.md index a38f798..35e1658 100644 --- a/docs/relay-eval-recipes.md +++ b/docs/relay-eval-recipes.md @@ -104,6 +104,71 @@ curl -sS 'https://relay-eval.waow.tech/api/trend?limit=50' \ curl -sS 'https://relay-eval.waow.tech/api/latest' | jq . ``` +**summarize missed accounts for the latest run**: +```bash +curl -sS 'https://relay-eval.waow.tech/api/latest/diffs/summary?limit=20' \ + | jq '{total, checked, unchecked, estimate, relays: .facets.relays[0:5], pds: .facets.pds[0:5]}' +``` + +`estimate.live_gap_lower_bound` is the number of checked rows that look like +live relay gaps. When `estimate.checked_limited` is true, treat it as a lower +bound; `estimate.live_gap_display` includes the `+` marker for UI/scripts. + +**page likely live gaps for one relay**: +```bash +cursor=0 +while :; do + page=$(curl -sS "https://relay-eval.waow.tech/api/latest/diffs?relay=zlay.waow.tech&lane=live&limit=1000&cursor=$cursor") + echo "$page" | jq -r '.diffs[].did' + cursor=$(echo "$page" | jq -r '.next_cursor // empty') + [ -z "$cursor" ] && break +done +``` + +Use `/api/runs//diffs` for a specific run. Omit `relay` and/or +`classification` to fetch a broader diff list. Both `/diffs` and +`/diffs/summary` accept: + +- `relay=` — one missing relay. +- `classification=` — exact stored classifier label. +- `lane=live|review|inactive|unchecked|checked` — semantic bucket. +- `q=` — DID substring search. +- `pds_host=` — rows whose resolved PDS host matches. +- `limit=` and `cursor=` — page raw `/diffs` rows. + +new rows are written with `classification_version=2` and these labels: + +- `active_missing` — DID resolves, has a PDS endpoint, and appears active; + this is the lane to treat as likely relay coverage debt. +- `no_pds_endpoint` — DID resolved, but no usable PDS service endpoint was + present in the DID document. +- `malformed_did` — syntactically invalid DID. +- `unsupported_did_method` — not a DID method relay-eval can resolve. +- `did_resolution_failed` — resolver/network failure while looking up the DID. +- `invalid_did_document` — resolver returned a DID document relay-eval could + not interpret safely. +- `classification_not_attempted` — the row was outside the classification + budget for the run. + +old rows may still use `classification_version=1` labels: +`coverage_gap`, `unresolvable`, and `deactivated`. Keep them available for +historical queries, but avoid over-interpreting `unresolvable`: older runs +mixed true resolver failures with DIDs that were never attempted after the +classifier cap was exhausted. + +**time-window chart route**: +```text +https://relay-eval.waow.tech/chart/// +``` + +`since` and `until` may be Unix seconds, Unix milliseconds, or ISO-ish +timestamps accepted by `/api/trend`. `relay_csv` is optional; when present it +selects the relays rendered in the chart: + +```text +https://relay-eval.waow.tech/chart/1780700000/1780786400/zlay.waow.tech,bsky.network,relay.waow.tech +``` + **raw trend JSON for loading into pandas / duckdb / sqlite**: ```bash curl -sS 'https://relay-eval.waow.tech/api/trend?limit=500' > trend.json @@ -118,12 +183,12 @@ curl -sS 'https://relay-eval.waow.tech/api/trend?limit=500' > trend.json trend API just gives you the raw numbers. if a specific relay is >1.3× the median, treat it as "replaying" not "live". - `diffs.classification` (via `/api/latest`) separates missed DIDs - into `coverage_gap` (real bug — PDS was reachable and resolvable) - vs `unresolvable` (DID lookup failed — ambiguous) vs `deactivated` - (account is dead). `coverage_gap` is the number to fixate on. -- `/api/latest` truncates to 30 `diff_samples` per relay. for the - full missed-DID set you'd need to query sqlite directly, but you - probably don't need that for the canary-2 diagnostic work. + into v2 lanes. `active_missing` is the likely relay issue; resolver, + DID-document, and budget lanes are useful audit context, not the same kind + of action item. +- `/api/latest` includes only 30 `diff_samples` per relay to keep the + dashboard response small. use `/api/latest/diffs` or + `/api/runs//diffs` to page through the complete missed-DID set. ## using this as a canary acceptance signal diff --git a/relay-eval/build.zig b/relay-eval/build.zig index 779df07..f4b0e4e 100644 --- a/relay-eval/build.zig +++ b/relay-eval/build.zig @@ -23,14 +23,14 @@ pub fn build(b: *std.Build) void { .target = target, .optimize = optimize, .imports = imports, + .link_libc = true, }); + exe_mod.linkSystemLibrary("sqlite3", .{}); const exe = b.addExecutable(.{ .name = "relay-eval", .root_module = exe_mod, }); - exe.linkLibC(); - exe.linkSystemLibrary("sqlite3"); b.installArtifact(exe); const run = b.addRunArtifact(exe); @@ -52,12 +52,12 @@ pub fn build(b: *std.Build) void { .target = target, .optimize = optimize, .imports = imports, + .link_libc = true, }); + test_mod.linkSystemLibrary("sqlite3", .{}); const t = b.addTest(.{ .root_module = test_mod, }); - t.linkLibC(); - t.linkSystemLibrary("sqlite3"); test_step.dependOn(&b.addRunArtifact(t).step); } } diff --git a/relay-eval/build.zig.zon b/relay-eval/build.zig.zon index 925ba1c..d3544ac 100644 --- a/relay-eval/build.zig.zon +++ b/relay-eval/build.zig.zon @@ -2,15 +2,15 @@ .name = .relay_eval, .version = "0.0.1", .fingerprint = 0x5f4f6bc8058cd6d1, - .minimum_zig_version = "0.15.0", + .minimum_zig_version = "0.16.0", .dependencies = .{ .zat = .{ - .url = "https://tangled.org/zat.dev/zat/archive/v0.2.16.tar.gz", - .hash = "zat-0.2.16-5PuC7tjwBADbnwV5y8ztKUHhGHMJHh2HouvoYImnZ7y5", + .url = "https://tangled.org/zat.dev/zat/archive/v0.3.4.tar.gz", + .hash = "zat-0.3.4-5PuC7l8GCQA9xdkWh3ha_cf9fZvELT_370O14iNqjSy6", }, .websocket = .{ - .url = "https://github.com/zzstoatzz/websocket.zig/archive/395d0f4.tar.gz", - .hash = "websocket-0.1.0-ZPISdVJ8AwD7U03ARGgHclzlYSd9GeU91_WDXjRyjYdh", + .url = "https://github.com/zzstoatzz/websocket.zig/archive/refs/tags/v0.1.4.tar.gz", + .hash = "websocket-0.1.2-ZPISdXFPBAB4F0i5uz0NkdgHU4J0PWrGFnjgYKzaBwmJ", }, }, .paths = .{ diff --git a/relay-eval/justfile b/relay-eval/justfile index 67c3958..8b02013 100644 --- a/relay-eval/justfile +++ b/relay-eval/justfile @@ -1,7 +1,7 @@ # relay-eval — firehose comparison tool # required env vars: HCLOUD_TOKEN -server := "root@" + `terraform -chdir=infra output -raw server_ip 2>/dev/null || echo "NO_SERVER"` +server := "root@" + `IP=$(terraform -chdir=infra output -raw server_ip 2>/dev/null || true); if [ -n "$IP" ]; then echo "$IP"; else dig +short relay-eval.waow.tech | grep -E '^[0-9.]+' | head -1 || echo "NO_SERVER"; fi` # --- infrastructure --- diff --git a/relay-eval/src/classifier.zig b/relay-eval/src/classifier.zig index e82bcc8..280a4bf 100644 --- a/relay-eval/src/classifier.zig +++ b/relay-eval/src/classifier.zig @@ -4,33 +4,44 @@ const zat = @import("zat"); const log = std.log.scoped(.classifier); pub const Classification = enum { - coverage_gap, - unresolvable, - deactivated, + active_missing, + no_pds_endpoint, + malformed_did, + unsupported_did_method, + did_resolution_failed, + invalid_did_document, + classification_not_attempted, pub fn toString(self: Classification) []const u8 { return switch (self) { - .coverage_gap => "coverage_gap", - .unresolvable => "unresolvable", - .deactivated => "deactivated", + .active_missing => "active_missing", + .no_pds_endpoint => "no_pds_endpoint", + .malformed_did => "malformed_did", + .unsupported_did_method => "unsupported_did_method", + .did_resolution_failed => "did_resolution_failed", + .invalid_did_document => "invalid_did_document", + .classification_not_attempted => "classification_not_attempted", }; } }; +pub const classification_version: u16 = 2; + pub const ClassifiedDid = struct { did: []const u8, classification: Classification, + pds: ?[]const u8 = null, }; /// classify DIDs that appear in one relay but not another. -/// resolves each DID to determine if it's a real coverage gap, -/// an unresolvable DID, or a deactivated account. +/// `max_count` is a real classification budget: callers should store +/// classification_not_attempted for misses outside this returned set. pub fn classifyDids( allocator: std.mem.Allocator, dids: []const []const u8, max_count: usize, ) ![]ClassifiedDid { - var resolver = zat.DidResolver.init(allocator); + var resolver = zat.DidResolver.init(std.Options.debug_io, allocator); defer resolver.deinit(); var results: std.ArrayList(ClassifiedDid) = .empty; @@ -41,10 +52,11 @@ pub fn classifyDids( if (i > 0 and i % 50 == 0) { log.info("classified {d}/{d} DIDs", .{ i, limit }); } - const classification = classifyOne(&resolver, did); + const classified = classifyOne(allocator, &resolver, did); try results.append(allocator, .{ .did = did, - .classification = classification, + .classification = classified.classification, + .pds = classified.pds, }); } log.info("classified {d}/{d} DIDs", .{ limit, limit }); @@ -52,19 +64,36 @@ pub fn classifyDids( return results.toOwnedSlice(allocator); } -fn classifyOne(resolver: *zat.DidResolver, did_str: []const u8) Classification { - const parsed = zat.Did.parse(did_str) orelse return .unresolvable; +pub fn freeClassified(allocator: std.mem.Allocator, classified: []ClassifiedDid) void { + for (classified) |cd| { + if (cd.pds) |p| allocator.free(p); + } + allocator.free(classified); +} + +const OneResult = struct { + classification: Classification, + pds: ?[]const u8 = null, +}; + +fn classifyOne(allocator: std.mem.Allocator, resolver: *zat.DidResolver, did_str: []const u8) OneResult { + const parsed = zat.Did.parse(did_str) orelse return .{ .classification = .malformed_did }; var doc = resolver.resolve(parsed) catch |err| { log.debug("resolve failed for {s}: {s}", .{ did_str, @errorName(err) }); - return .unresolvable; + return .{ .classification = switch (err) { + error.UnsupportedDidMethod => .unsupported_did_method, + error.InvalidDidDocument => .invalid_did_document, + else => .did_resolution_failed, + } }; }; defer doc.deinit(); - // has PDS endpoint = real account, this is a coverage gap - if (doc.pdsEndpoint()) |_| { - return .coverage_gap; + if (doc.pdsEndpoint()) |pds| { + return .{ + .classification = .active_missing, + .pds = allocator.dupe(u8, pds) catch null, + }; } - // resolves but no PDS = deactivated - return .deactivated; + return .{ .classification = .no_pds_endpoint }; } diff --git a/relay-eval/src/collector.zig b/relay-eval/src/collector.zig index 1823767..fe23767 100644 --- a/relay-eval/src/collector.zig +++ b/relay-eval/src/collector.zig @@ -9,25 +9,27 @@ pub const Collector = struct { shutdown: std.atomic.Value(bool), event_count: std.atomic.Value(u64), did_set: DidSet, - mutex: std.Thread.Mutex, + mutex: std.Io.Mutex, connected: std.atomic.Value(bool), was_connected: std.atomic.Value(bool), error_msg: ?[]const u8, allocator: std.mem.Allocator, + io: std.Io, const DidSet = std.StringHashMap(void); - pub fn init(allocator: std.mem.Allocator, host: []const u8) Collector { + pub fn init(allocator: std.mem.Allocator, io: std.Io, host: []const u8) Collector { return .{ .host = host, .shutdown = std.atomic.Value(bool).init(false), .event_count = std.atomic.Value(u64).init(0), .did_set = DidSet.init(allocator), - .mutex = .{}, + .mutex = std.Io.Mutex.init, .connected = std.atomic.Value(bool).init(false), .was_connected = std.atomic.Value(bool).init(false), .error_msg = null, .allocator = allocator, + .io = io, }; } @@ -40,8 +42,8 @@ pub const Collector = struct { } pub fn uniqueDids(self: *Collector) u32 { - self.mutex.lock(); - defer self.mutex.unlock(); + self.mutex.lockUncancelable(self.io); + defer self.mutex.unlock(self.io); return @intCast(self.did_set.count()); } @@ -57,7 +59,7 @@ pub const Collector = struct { fn connectAndRead(self: *Collector) !void { const path = "/xrpc/com.atproto.sync.subscribeRepos"; - var client = try websocket.Client.init(self.allocator, .{ + var client = try websocket.Client.init(self.io, self.allocator, .{ .host = self.host, .port = 443, .tls = true, @@ -130,8 +132,8 @@ pub const Collector = struct { _ = self.event_count.fetchAdd(1, .monotonic); // insert into DID set (dupe the string to outlive the arena) - self.mutex.lock(); - defer self.mutex.unlock(); + self.mutex.lockUncancelable(self.io); + defer self.mutex.unlock(self.io); const gop = self.did_set.getOrPut(did) catch return; if (!gop.found_existing) { diff --git a/relay-eval/src/live_rate.zig b/relay-eval/src/live_rate.zig index 2eee68b..a4ad14a 100644 --- a/relay-eval/src/live_rate.zig +++ b/relay-eval/src/live_rate.zig @@ -23,6 +23,7 @@ const Snapshot = struct { pub const LiveRateMeter = struct { host: []const u8, allocator: std.mem.Allocator, + io: std.Io, // hot path: writer thread bumps; readers (sampler, /api) read atomically counter: std.atomic.Value(u64), @@ -32,21 +33,22 @@ pub const LiveRateMeter = struct { // sampler ring (1Hz, last 10s). `snapshot_idx` points at the newest entry. snapshots: [sample_buckets]Snapshot, snapshot_idx: usize, - snapshots_lock: std.Thread.Mutex, + snapshots_lock: std.Io.Mutex, // reconnect backoff state — only mutated by the run() thread, so no atomic backoff_ms: u64, - pub fn init(allocator: std.mem.Allocator, host: []const u8) LiveRateMeter { + pub fn init(allocator: std.mem.Allocator, io: std.Io, host: []const u8) LiveRateMeter { return .{ .host = host, .allocator = allocator, + .io = io, .counter = .init(0), .connected = .init(false), .shutdown = .init(false), .snapshots = [_]Snapshot{.{ .counter = 0, .ts_ms = 0 }} ** sample_buckets, .snapshot_idx = 0, - .snapshots_lock = .{}, + .snapshots_lock = std.Io.Mutex.init, .backoff_ms = 1000, }; } @@ -55,9 +57,9 @@ pub const LiveRateMeter = struct { /// shared sampler thread once per second. pub fn snapshot(self: *LiveRateMeter) void { const c = self.counter.load(.monotonic); - const t = std.time.milliTimestamp(); - self.snapshots_lock.lock(); - defer self.snapshots_lock.unlock(); + const t: i64 = @intCast(@divFloor(std.Io.Timestamp.now(self.io, .real).nanoseconds, std.time.ns_per_ms)); + self.snapshots_lock.lockUncancelable(self.io); + defer self.snapshots_lock.unlock(self.io); self.snapshot_idx = (self.snapshot_idx + 1) % sample_buckets; self.snapshots[self.snapshot_idx] = .{ .counter = c, .ts_ms = t }; } @@ -65,8 +67,8 @@ pub const LiveRateMeter = struct { /// events/sec averaged over the last ~10s. returns 0.0 if the ring hasn't /// filled enough to span a real window. pub fn ratePerSec(self: *LiveRateMeter) f64 { - self.snapshots_lock.lock(); - defer self.snapshots_lock.unlock(); + self.snapshots_lock.lockUncancelable(self.io); + defer self.snapshots_lock.unlock(self.io); const newest = self.snapshots[self.snapshot_idx]; const oldest_idx = (self.snapshot_idx + 1) % sample_buckets; const oldest = self.snapshots[oldest_idx]; @@ -87,7 +89,7 @@ pub const LiveRateMeter = struct { }; self.connected.store(false, .release); if (self.shutdown.load(.acquire)) break; - std.Thread.sleep(self.backoff_ms * std.time.ns_per_ms); + self.io.sleep(.{ .nanoseconds = self.backoff_ms * std.time.ns_per_ms }, .awake) catch {}; self.backoff_ms = @min(self.backoff_ms * 2, 60_000); } } @@ -95,7 +97,7 @@ pub const LiveRateMeter = struct { fn connectAndCount(self: *LiveRateMeter) !void { const path = "/xrpc/com.atproto.sync.subscribeRepos"; - var client = try websocket.Client.init(self.allocator, .{ + var client = try websocket.Client.init(self.io, self.allocator, .{ .host = self.host, .port = 443, .tls = true, @@ -134,9 +136,9 @@ pub const LiveRateMeter = struct { /// snapshots all meters at ~1 Hz. runs in its own thread for the lifetime of /// the server. this is the only writer to `LiveRateMeter.snapshots`. -pub fn sampleLoop(meters: []LiveRateMeter, shutdown: *std.atomic.Value(bool)) void { +pub fn sampleLoop(io: std.Io, meters: []LiveRateMeter, shutdown: *std.atomic.Value(bool)) void { while (!shutdown.load(.acquire)) { for (meters) |*m| m.snapshot(); - std.Thread.sleep(std.time.ns_per_s); + io.sleep(.{ .nanoseconds = std.time.ns_per_s }, .awake) catch {}; } } diff --git a/relay-eval/src/main.zig b/relay-eval/src/main.zig index d8d1adb..4b8d41c 100644 --- a/relay-eval/src/main.zig +++ b/relay-eval/src/main.zig @@ -8,13 +8,17 @@ const coverage = @import("coverage.zig"); const log = std.log.scoped(.main); -pub fn main() !void { - var gpa: std.heap.GeneralPurposeAllocator(.{}) = .init; - defer _ = gpa.deinit(); - const allocator = gpa.allocator(); - - const args = try std.process.argsAlloc(allocator); - defer std.process.argsFree(allocator, args); +pub fn main(init: std.process.Init) !void { + const allocator = init.gpa; + const io = init.io; + + var args_it = std.process.Args.Iterator.init(init.minimal.args); + var args_list: std.ArrayList([]const u8) = .empty; + defer args_list.deinit(allocator); + while (args_it.next()) |arg| { + try args_list.append(allocator, arg); + } + const args = args_list.items; if (args.len < 2) { printUsage(); @@ -23,9 +27,9 @@ pub fn main() !void { const command = args[1]; if (std.mem.eql(u8, command, "eval")) { - try runEval(allocator, args[2..]); + try runEval(allocator, io, args[2..]); } else if (std.mem.eql(u8, command, "serve")) { - try runServe(allocator, args[2..]); + try runServe(allocator, io, args[2..]); } else { printUsage(); } @@ -54,7 +58,7 @@ fn printUsage() void { , .{}); } -fn runEval(allocator: std.mem.Allocator, args: []const []const u8) !void { +fn runEval(allocator: std.mem.Allocator, io: std.Io, args: []const []const u8) !void { var relays_str: ?[]const u8 = null; var window: u32 = 300; var db_path: [:0]const u8 = "relay-eval.db"; @@ -117,14 +121,14 @@ fn runEval(allocator: std.mem.Allocator, args: []const []const u8) !void { defer allocator.free(threads); for (hosts.items, 0..) |host, idx| { - collectors[idx] = Collector.init(allocator, host); + collectors[idx] = Collector.init(allocator, io, host); threads[idx] = try std.Thread.spawn(.{}, Collector.run, .{&collectors[idx]}); log.info("spawned collector for {s}", .{host}); } // wait for collection window log.info("collecting for {d}s...", .{window}); - std.Thread.sleep(@as(u64, window) * std.time.ns_per_s); + io.sleep(.{ .nanoseconds = @as(u64, window) * std.time.ns_per_s }, .awake) catch {}; // signal shutdown for (collectors) |*col| { @@ -154,13 +158,13 @@ fn runEval(allocator: std.mem.Allocator, args: []const []const u8) !void { defer did_counts.deinit(); for (collectors) |*col| { - col.mutex.lock(); + col.mutex.lockUncancelable(io); var key_it = col.did_set.keyIterator(); while (key_it.next()) |key| { const gop = try did_counts.getOrPut(key.*); gop.value_ptr.* = if (gop.found_existing) gop.value_ptr.* + 1 else 1; } - col.mutex.unlock(); + col.mutex.unlock(io); } const union_size: u32 = @intCast(did_counts.count()); @@ -182,7 +186,7 @@ fn runEval(allocator: std.mem.Allocator, args: []const []const u8) !void { } for (collectors, 0..) |*col, idx| { - col.mutex.lock(); + col.mutex.lockUncancelable(io); var count_it = did_counts.iterator(); while (count_it.next()) |entry| { if (entry.value_ptr.* < union_quorum) continue; // not in the consensus union @@ -191,7 +195,7 @@ fn runEval(allocator: std.mem.Allocator, args: []const []const u8) !void { try relay_misses[idx].items.append(allocator, entry.key_ptr.*); } } - col.mutex.unlock(); + col.mutex.unlock(io); log.info("{s}: missing {d} DIDs from consensus union", .{ col.host, relay_misses[idx].items.items.len, @@ -211,22 +215,25 @@ fn runEval(allocator: std.mem.Allocator, args: []const []const u8) !void { } const classified = try classifier.classifyDids(allocator, missing_list.items, 500); - defer allocator.free(classified); + defer classifier.freeClassified(allocator, classified); // build DID → classification lookup var class_map: std.StringHashMap(classifier.Classification) = .init(allocator); defer class_map.deinit(); + var pds_map: std.StringHashMap([]const u8) = .init(allocator); + defer pds_map.deinit(); for (classified) |cd| { try class_map.put(cd.did, cd.classification); + if (cd.pds) |p| try pds_map.put(cd.did, p); } - // for DIDs beyond the 500 cap, mark as unresolvable (not classified) - // they'll just not be in class_map + // DIDs beyond the classification budget remain absent from class_map and + // are stored as classification_not_attempted, not as resolver failures. var store = try Store.open(db_path); defer store.close(); - const now = std.time.timestamp(); + const now: i64 = @intCast(@divFloor(std.Io.Timestamp.now(io, .real).nanoseconds, std.time.ns_per_s)); var ts_buf: [32]u8 = undefined; const timestamp = formatTimestamp(now, &ts_buf); @@ -242,12 +249,12 @@ fn runEval(allocator: std.mem.Allocator, args: []const []const u8) !void { for (collectors) |*col| { var seen_counts: std.ArrayList(u16) = .empty; defer seen_counts.deinit(allocator); - col.mutex.lock(); + col.mutex.lockUncancelable(io); var key_it = col.did_set.keyIterator(); while (key_it.next()) |key| { try seen_counts.append(allocator, did_counts.get(key.*) orelse 1); } - col.mutex.unlock(); + col.mutex.unlock(io); const hist = try coverage.buildHistogram(allocator, num_relays, seen_counts.items); defer allocator.free(hist); @@ -265,35 +272,55 @@ fn runEval(allocator: std.mem.Allocator, args: []const []const u8) !void { // insert per-relay diffs for (collectors, 0..) |*col, idx| { - var gap: u32 = 0; - var unresolvable: u32 = 0; - var deactivated: u32 = 0; + var active_missing: u32 = 0; + var no_pds_endpoint: u32 = 0; + var malformed_did: u32 = 0; + var unsupported_did_method: u32 = 0; + var did_resolution_failed: u32 = 0; + var invalid_did_document: u32 = 0; + var classification_not_attempted: u32 = 0; for (relay_misses[idx].items.items) |did| { - const classification = class_map.get(did) orelse .unresolvable; + const classification = class_map.get(did) orelse .classification_not_attempted; switch (classification) { - .coverage_gap => gap += 1, - .unresolvable => unresolvable += 1, - .deactivated => deactivated += 1, + .active_missing => active_missing += 1, + .no_pds_endpoint => no_pds_endpoint += 1, + .malformed_did => malformed_did += 1, + .unsupported_did_method => unsupported_did_method += 1, + .did_resolution_failed => did_resolution_failed += 1, + .invalid_did_document => invalid_did_document += 1, + .classification_not_attempted => classification_not_attempted += 1, } try store.insertDiff(run_id, .{ .relay = col.host, .did = did, .classification = classification.toString(), + .classification_version = classifier.classification_version, + .pds = pds_map.get(did), }); } - log.info("{s}: coverage_gap={d}, unresolvable={d}, deactivated={d}", .{ - col.host, gap, unresolvable, deactivated, - }); + log.info( + "{s}: active_missing={d}, no_pds_endpoint={d}, malformed_did={d}, unsupported_did_method={d}, did_resolution_failed={d}, invalid_did_document={d}, classification_not_attempted={d}", + .{ + col.host, + active_missing, + no_pds_endpoint, + malformed_did, + unsupported_did_method, + did_resolution_failed, + invalid_did_document, + classification_not_attempted, + }, + ); } try store.commit(); log.info("run #{d} saved to {s}", .{ run_id, db_path }); } -fn runServe(allocator: std.mem.Allocator, args: []const []const u8) !void { +fn runServe(allocator: std.mem.Allocator, io: std.Io, args: []const []const u8) !void { var db_path: [:0]const u8 = "relay-eval.db"; var db_path_alloc: ?[:0]const u8 = null; defer if (db_path_alloc) |p| allocator.free(p); @@ -301,8 +328,8 @@ fn runServe(allocator: std.mem.Allocator, args: []const []const u8) !void { var relays_str: ?[]const u8 = null; // default trend limit: env var, then 48 - var trend_limit: u32 = if (std.posix.getenv("RELAY_EVAL_TREND_LIMIT")) |v| - std.fmt.parseInt(u32, v, 10) catch 48 + var trend_limit: u32 = if (std.c.getenv("RELAY_EVAL_TREND_LIMIT")) |v| + std.fmt.parseInt(u32, std.mem.span(v), 10) catch 48 else 48; @@ -343,15 +370,15 @@ fn runServe(allocator: std.mem.Allocator, args: []const []const u8) !void { meters = try allocator.alloc(live_rate.LiveRateMeter, hosts.items.len); meter_threads = try allocator.alloc(std.Thread, hosts.items.len); for (hosts.items, 0..) |h, idx| { - meters[idx] = live_rate.LiveRateMeter.init(allocator, h); + meters[idx] = live_rate.LiveRateMeter.init(allocator, io, h); meter_threads[idx] = try std.Thread.spawn(.{}, live_rate.LiveRateMeter.run, .{&meters[idx]}); } - _ = try std.Thread.spawn(.{}, live_rate.sampleLoop, .{ meters, &sampler_shutdown }); + _ = try std.Thread.spawn(.{}, live_rate.sampleLoop, .{ io, meters, &sampler_shutdown }); log.info("spawned {d} live rate meters + sampler thread", .{meters.len}); } } - try server.run(allocator, db_path, port, trend_limit, meters); + try server.run(allocator, io, db_path, port, trend_limit, meters); // server.run is intended to run forever; on the (unexpected) clean return // path we stop the sampler. meter threads detach naturally with the diff --git a/relay-eval/src/server.zig b/relay-eval/src/server.zig index 41fc3d7..faebbf2 100644 --- a/relay-eval/src/server.zig +++ b/relay-eval/src/server.zig @@ -4,13 +4,16 @@ 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"); pub fn run( allocator: std.mem.Allocator, + io: std.Io, db_path: [:0]const u8, port: u16, trend_limit: u32, @@ -24,34 +27,36 @@ pub fn run( const og_png_path = try std.fmt.allocPrint(allocator, "{s}/og.png", .{db_dir}); defer allocator.free(og_png_path); - const addr = std.net.Address.parseIp4("0.0.0.0", port) catch unreachable; - var listener = try addr.listen(.{ .reuse_address = true }); - defer listener.deinit(); + 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 conn = listener.accept() catch |err| { + const stream = listener.accept(io) catch |err| { log.err("accept: {s}", .{@errorName(err)}); continue; }; - handleConnection(allocator, conn.stream, &store, trend_limit, og_png_path, meters) catch |err| { + handleConnection(allocator, io, stream, &store, trend_limit, og_png_path, meters) catch |err| { log.debug("request error: {s}", .{@errorName(err)}); }; - conn.stream.close(); + stream.close(io); } } fn handleConnection( allocator: std.mem.Allocator, - stream: std.net.Stream, + 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; - const n = try stream.read(&buf); + 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]; @@ -74,8 +79,10 @@ fn handleConnection( const path = if (qs_sep) |i| full_path[0..i] else full_path; const query = if (qs_sep) |i| full_path[i + 1 ..] else ""; - if (std.mem.eql(u8, path, "/")) { + if (std.mem.eql(u8, path, "/") or std.mem.eql(u8, path, "/chart") or std.mem.startsWith(u8, path, "/chart/")) { try sendResponse(stream, "200 OK", "text/html", index_html); + } else if (std.mem.eql(u8, path, "/diffs-api")) { + try sendResponse(stream, "200 OK", "text/html", diffs_api_html); } else if (std.mem.eql(u8, path, "/sonify")) { try sendResponse(stream, "200 OK", "text/html", sonify_html); } else if (std.mem.eql(u8, path, "/sonify/og.svg")) { @@ -86,8 +93,37 @@ fn handleConnection( 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; @@ -96,7 +132,7 @@ fn handleConnection( } 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, limit, quorum); + try serveTrend(allocator, stream, store, query, limit, quorum); } else if (std.mem.eql(u8, path, "/api/relays")) { try serveRelays(allocator, stream, store); } else if (std.mem.eql(u8, path, "/api/relays/history")) { @@ -142,14 +178,14 @@ fn parseQuorumParam(query: []const u8, default: u16) u16 { return default; } -fn serveRunCount(stream: std.net.Stream, store: *Store) !void { +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 sendResponse(stream, "200 OK", "application/json", json); } -fn serveRuns(allocator: std.mem.Allocator, stream: std.net.Stream, store: *Store) !void { +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); @@ -171,7 +207,7 @@ fn serveRuns(allocator: std.mem.Allocator, stream: std.net.Stream, store: *Store try sendResponse(stream, "200 OK", "application/json", json.items); } -fn serveLatest(allocator: std.mem.Allocator, stream: std.net.Stream, store: *Store) !void { +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; @@ -179,7 +215,7 @@ fn serveLatest(allocator: std.mem.Allocator, stream: std.net.Stream, store: *Sto try serveRunDetail(allocator, stream, store, run_id); } -fn serveRunDetail(allocator: std.mem.Allocator, stream: std.net.Stream, store: *Store, run_id: i64) !void { +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; @@ -207,6 +243,7 @@ fn serveRunDetail(allocator: std.mem.Allocator, stream: std.net.Stream, store: * allocator.free(d.relay); allocator.free(d.did); allocator.free(d.classification); + if (d.pds) |p| allocator.free(p); } allocator.free(diff_samples); } @@ -226,24 +263,366 @@ fn serveRunDetail(allocator: std.mem.Allocator, stream: std.net.Stream, store: * 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}\",\"count\":{d}}}", .{ - dc.relay, dc.classification, dc.count, + 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}\"}}", .{ - d.relay, d.did, d.classification, + 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 sendResponse(stream, "200 OK", "application/json", json.items); } -fn serveTrend(allocator: std.mem.Allocator, stream: std.net.Stream, store: *Store, limit: u32, quorum: u16) !void { - const points = try store.getTrendData(allocator, limit); +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 points = if (use_range) + try store.getTrendDataRange(allocator, since, until) + else + try store.getTrendData(allocator, limit); defer { for (points) |p| { allocator.free(p.timestamp); @@ -253,6 +632,8 @@ fn serveTrend(allocator: std.mem.Allocator, stream: std.net.Stream, store: *Stor 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(); @@ -301,6 +682,7 @@ fn serveTrend(allocator: std.mem.Allocator, stream: std.net.Stream, store: *Stor } 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) @@ -319,14 +701,17 @@ fn serveTrend(allocator: std.mem.Allocator, stream: std.net.Stream, store: *Stor try sendResponse(stream, "200 OK", "application/json", json.items); } -fn serveOgPng(allocator: std.mem.Allocator, stream: std.net.Stream, og_png_path: []const u8) !void { - const file = std.fs.cwd().openFile(og_png_path, .{}) catch { +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(); + defer file.close(io); - const body = file.readToEndAlloc(allocator, 2 * 1024 * 1024) catch { + 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; }; @@ -335,7 +720,7 @@ fn serveOgPng(allocator: std.mem.Allocator, stream: std.net.Stream, og_png_path: try sendResponse(stream, "200 OK", "image/png", body); } -fn serveOgSvg(allocator: std.mem.Allocator, stream: std.net.Stream, store: *Store) !void { +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", \\ @@ -499,7 +884,7 @@ fn findSnapshot(list: []const Store.MonitorSnapshot, host: []const u8) ?Store.Mo } fn formatIsoNow(buf: []u8) []const u8 { - const now = std.time.timestamp(); + const now = currentEpochSeconds(); const es = std.time.epoch.EpochSeconds{ .secs = @intCast(now) }; const day = es.getEpochDay(); const yd = day.calculateYearDay(); @@ -515,6 +900,10 @@ fn formatIsoNow(buf: []u8) []const u8 { }) 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); @@ -551,7 +940,7 @@ fn writeHeadline( } } -fn serveRelays(allocator: std.mem.Allocator, stream: std.net.Stream, store: *Store) !void { +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); @@ -703,7 +1092,7 @@ fn isValidHostname(s: []const u8) bool { return true; } -fn serveRelaysHistory(allocator: std.mem.Allocator, stream: std.net.Stream, store: *Store, query: []const u8) !void { +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; @@ -822,7 +1211,7 @@ fn formatIsoEpoch(epoch: i64, buf: []u8) []const u8 { }) catch "1970-01-01T00:00:00Z"; } -fn serveRelaysEvents(allocator: std.mem.Allocator, stream: std.net.Stream, store: *Store, query: []const u8) !void { +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; @@ -835,7 +1224,7 @@ fn serveRelaysEvents(allocator: std.mem.Allocator, stream: std.net.Stream, store break :blk raw; } // default: 24h ago - break :blk formatIsoEpoch(std.time.timestamp() - 86400, &since_buf); + break :blk formatIsoEpoch(currentEpochSeconds() - 86400, &since_buf); }; const until = blk: { @@ -847,7 +1236,7 @@ fn serveRelaysEvents(allocator: std.mem.Allocator, stream: std.net.Stream, store break :blk raw; } // default: now - break :blk formatIsoEpoch(std.time.timestamp(), &until_buf); + break :blk formatIsoEpoch(currentEpochSeconds(), &until_buf); }; const name = parseQueryString(query, "name") orelse ""; @@ -900,7 +1289,7 @@ fn serveRelaysEvents(allocator: std.mem.Allocator, stream: std.net.Stream, store // 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: std.net.Stream, meters: []live_rate.LiveRateMeter) !void { +fn serveRelaysLive(allocator: std.mem.Allocator, stream: net.Stream, meters: []live_rate.LiveRateMeter) !void { var json: std.ArrayList(u8) = .empty; defer json.deinit(allocator); @@ -926,7 +1315,7 @@ fn serveRelaysLive(allocator: std.mem.Allocator, stream: std.net.Stream, meters: try sendResponse(stream, "200 OK", "application/json", json.items); } -fn sendResponse(stream: std.net.Stream, status: []const u8, content_type: []const u8, body: []const u8) !void { +fn sendResponse(stream: net.Stream, status: []const u8, content_type: []const u8, body: []const u8) !void { var header_buf: [512]u8 = undefined; const header = std.fmt.bufPrint( &header_buf, @@ -939,6 +1328,9 @@ fn sendResponse(stream: std.net.Stream, status: []const u8, content_type: []cons .{ status, content_type, body.len }, ) catch return error.HeaderTooLong; - try stream.writeAll(header); - try stream.writeAll(body); + 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(); } diff --git a/relay-eval/src/static/diffs-api.html b/relay-eval/src/static/diffs-api.html new file mode 100644 index 0000000..41af7cb --- /dev/null +++ b/relay-eval/src/static/diffs-api.html @@ -0,0 +1,153 @@ + + + + + + +missing repo API · relay-eval + + + + + + +
+
+

missing repo API

+ back to relay-eval +
+ +

+ Pull repos that were observed on at least one synchronization stream but were missing from another relay during a relay-eval run. +

+ +
+

Start With A Summary

+

The summary endpoint gives counts, classification facets, PDS host facets, and the concrete run_id for a stable snapshot.

+
GET /api/latest/diffs/summary?relay=<relay-host>
+
curl -fsS 'https://relay-eval.waow.tech/api/latest/diffs/summary?relay=relay.xero.systems' | jq
+
+ +
+

Fetch Missing Repos

+

Use /api/latest/diffs for the most recent run, or use /api/runs/:run_id/diffs after reading the summary if you want every page to refer to the same run.

+
GET /api/latest/diffs?relay=<relay-host>&limit=5000
+
curl -fsS 'https://relay-eval.waow.tech/api/latest/diffs?relay=relay.xero.systems&limit=5000' | jq
+

The response contains diffs and next_cursor. Pass next_cursor back as cursor until it is null.

+
+ +
+

Pull The Full List

+

This emits tab-separated DID, classification, and optional PDS columns for every missing repo in the selected run.

+
relay='relay.xero.systems'
+run_id="$(curl -fsS "https://relay-eval.waow.tech/api/latest/diffs/summary?relay=${relay}" | jq -r '.run_id')"
+cursor=''
+
+while true; do
+  url="https://relay-eval.waow.tech/api/runs/${run_id}/diffs?relay=${relay}&limit=5000"
+  [ -n "$cursor" ] && url="${url}&cursor=${cursor}"
+
+  page="$(curl -fsS "$url")"
+  echo "$page" | jq -r '.diffs[] | [.did, .classification, (.pds // "")] | @tsv'
+
+  cursor="$(echo "$page" | jq -r '.next_cursor // empty')"
+  [ -z "$cursor" ] && break
+done
+
+ +
+

Useful Filters

+ + + + + + + + + + + + + + +
parametermeaning
relayRelay host to inspect, for example relay.xero.systems. Required for the common single-relay workflow.
limitRows per response. Defaults to 1000 and is capped at 5000.
cursorResume token from the previous response's next_cursor.
lane=liveOnly checked rows that look live elsewhere but absent from this relay.
lane=uncheckedRows not classified yet because the resolver budget was not spent on them.
lane=reviewRows that need human review, including resolver failures and malformed documents.
lane=inactiveRows classified as inactive or without a PDS endpoint.
lane=checkedRows where relay-eval did spend classification effort.
qSubstring search over the DID.
pds_hostFilter by PDS host, for example bracket.us-west.host.bsky.network.
+
+ +
+

Examples

+

Likely live gaps only:

+
curl -fsS 'https://relay-eval.waow.tech/api/latest/diffs?relay=relay.xero.systems&lane=live&limit=5000' | jq
+ +

Search for a DID fragment:

+
curl -fsS 'https://relay-eval.waow.tech/api/latest/diffs?relay=relay.xero.systems&q=76eyp&limit=100' | jq
+ +

Use a stable run id:

+
curl -fsS 'https://relay-eval.waow.tech/api/runs/4208/diffs?relay=relay.xero.systems&limit=5000' | jq
+ +

The broad full-list query includes unchecked rows. Use lane=live when you only want checked lower-bound stream gaps.

+
+ +
+ + diff --git a/relay-eval/src/static/index.html b/relay-eval/src/static/index.html index dc185e2..c52c9e0 100644 --- a/relay-eval/src/static/index.html +++ b/relay-eval/src/static/index.html @@ -36,6 +36,17 @@ width: 100vw; height: 280px; pointer-events: none; z-index: 0; } + body.chart-route #trend { + position: relative; + display: block; + width: 100vw; + height: min(52vh, 520px); + min-height: 320px; + pointer-events: auto; + z-index: 1; + border-bottom: 1px solid rgba(48, 54, 61, 0.55); + background: var(--bg); + } /* trend tooltip */ #trend-tip { @@ -62,6 +73,12 @@ opacity: 0; transform: translateY(6px); transition: opacity 0.5s ease, transform 0.5s ease; } + body.chart-route .glass { + min-height: auto; + background: rgba(13, 17, 23, 0.94); + backdrop-filter: none; + -webkit-backdrop-filter: none; + } .glass.visible { opacity: 1; transform: translateY(0); } h1 { font-size: 1.25rem; font-weight: 600; color: var(--fg); letter-spacing: -0.02em; } @@ -100,6 +117,8 @@ /* sections */ .sec { font-size: 0.9rem; font-weight: 600; margin-top: 2rem; margin-bottom: 0.2rem; } + .sec-link { font-size: 0.68rem; font-weight: 400; margin-left: 0.45rem; color: var(--muted); text-decoration: none; border-bottom: 1px dotted var(--muted); } + .sec-link:hover { color: var(--accent); border-color: var(--accent); } .sec-desc { font-size: 0.8rem; color: var(--muted); margin-bottom: 0.4rem; line-height: 1.55; } .sec-formula { font-size: 0.78rem; color: var(--fg); margin-bottom: 0.75rem; @@ -171,9 +190,101 @@ .xi { display: inline-block; width: 1.1em; text-align: center; transition: transform 0.15s; } .xbody { display: none; } .xbody.open { display: table-row-group; } - .drow td { padding-left: 2.2rem; font-size: 0.78rem; color: var(--muted); } + .gap-open td { padding: 0; background: rgba(13, 17, 23, 0.35); } + .gap-panel { + padding: 0.8rem 0 0.95rem 2.1rem; + border-bottom: 1px solid var(--border-strong); + } + .gap-head { + display: grid; grid-template-columns: minmax(14rem, 1fr) auto; + gap: 0.75rem; align-items: start; padding-right: 0.6rem; + } + .gap-title { font-size: 0.82rem; color: var(--fg); } + .gap-sub { font-size: 0.68rem; color: var(--muted); line-height: 1.45; margin-top: 0.12rem; max-width: 52rem; } + .gap-actions { display: flex; flex-wrap: wrap; gap: 0.35rem; justify-content: flex-end; } + .mini-link, .mini-btn { + display: inline-flex; align-items: center; gap: 0.25rem; + background: rgba(22, 27, 34, 0.75); color: var(--muted); + border: 1px solid var(--border); border-radius: 4px; + padding: 0.24rem 0.42rem; font: 0.66rem var(--mono, monospace); + text-decoration: none; cursor: pointer; + } + .mini-link:hover, .mini-btn:hover { color: var(--accent); border-color: var(--accent); } + .gap-metrics { display: flex; flex-wrap: wrap; gap: 0.35rem; margin: 0.65rem 0 0.55rem; padding-right: 0.6rem; } + .metric-chip { + border: 1px solid var(--border); border-radius: 4px; + background: rgba(22, 27, 34, 0.45); + padding: 0.28rem 0.42rem; font-size: 0.68rem; color: var(--muted); + } + .metric-chip strong { color: var(--fg); font-weight: 500; } + .metric-chip.gap strong { color: var(--red); } + .metric-chip.review strong { color: var(--yellow); } + .metric-chip.dead strong { color: var(--muted); } + .gap-tools { + display: grid; grid-template-columns: minmax(15rem, 1fr) auto; + gap: 0.5rem; align-items: center; padding-right: 0.6rem; + margin-bottom: 0.45rem; + } + .gap-search { + display: flex; gap: 0.35rem; min-width: 0; + } + .gap-search input { + min-width: 0; flex: 1; background: rgba(22, 27, 34, 0.85); color: var(--fg); + border: 1px solid var(--border-strong); border-radius: 4px; + padding: 0.35rem 0.5rem; font: 0.72rem var(--mono, monospace); + } + .lane-tabs { display: inline-flex; gap: 2px; background: rgba(22, 27, 34, 0.55); padding: 2px; border-radius: 6px; } + .lane-tabs button { + border: 0; background: transparent; color: var(--muted); border-radius: 4px; + padding: 0.28rem 0.45rem; font: 0.66rem var(--mono, monospace); cursor: pointer; + } + .lane-tabs button:hover { color: var(--fg); } + .lane-tabs button.active { background: rgba(88, 166, 255, 0.12); color: var(--accent); } + .pds-facets { + display: flex; flex-wrap: wrap; gap: 0.25rem; margin: 0.25rem 0 0.55rem; padding-right: 0.6rem; + } + .pds-facet { + border: 1px solid var(--border); border-radius: 4px; background: rgba(22, 27, 34, 0.45); + color: var(--muted); padding: 0.2rem 0.38rem; font: 0.64rem var(--mono, monospace); + cursor: pointer; + } + .pds-facet:hover, .pds-facet.active { color: var(--accent); border-color: var(--accent); } + .repo-list { + display: grid; gap: 1px; margin-top: 0.25rem; padding-right: 0.6rem; + } + .repo-row { + display: grid; grid-template-columns: minmax(17rem, 1.05fr) minmax(18rem, 1.35fr) minmax(5.5rem, max-content); + gap: 0.65rem; align-items: center; + padding: 0.36rem 0.45rem; + border-top: 1px solid rgba(48, 54, 61, 0.55); + color: var(--muted); font-size: 0.72rem; + } + .repo-row > div { min-width: 0; } + .repo-row:hover { background: rgba(22, 27, 34, 0.38); } + .repo-meta { + display: flex; align-items: center; gap: 0.65rem; + min-width: 0; overflow: hidden; + } + .repo-status { + display: inline-flex; align-items: center; gap: 0.3rem; + min-width: 0; max-width: 12.5rem; flex: 0 1 auto; white-space: nowrap; + } + .repo-status .cdot { width: 6px; height: 6px; } + .repo-status strong { color: var(--fg); font-weight: 500; overflow: hidden; text-overflow: ellipsis; } + .repo-evidence { display: block; min-width: 0; flex: 1 1 auto; overflow: hidden; text-overflow: ellipsis; white-space: nowrap; } + .repo-actions { + display: flex; align-items: center; justify-content: flex-start; + gap: 0.32rem; white-space: nowrap; min-width: 0; + } + .action-sep { color: var(--muted); opacity: 0.45; font-size: 0.62rem; } + .pds-host { color: var(--muted); font-size: 0.7rem; } .did-link { color: var(--accent); text-decoration: none; font-size: 0.75rem; transition: color 0.1s; } .did-link:hover { text-decoration: underline; } + .status-link { color: var(--muted); text-decoration: none; font-size: 0.66rem; font-family: var(--mono, monospace); border-bottom: 1px dotted var(--muted); } + .status-link:hover { color: var(--accent); border-bottom-color: var(--accent); } + .legend-note { color: var(--muted); font-size: 0.72rem; margin: -0.25rem 0 0.75rem; line-height: 1.5; } + .mini-btn:disabled { opacity: 0.45; cursor: default; border-color: var(--border); color: var(--muted); } + .diff-count-note { color: var(--muted); font-size: 0.7rem; } /* tooltip — JS-positioned floating div (css ::after gets clipped by overflow) */ .tip { cursor: help; border-bottom: 1px dotted var(--muted); } @@ -249,6 +360,7 @@ .trend-toggle:hover { opacity: 0.9; color: var(--fg); } .trend-toggle.active { opacity: 0.9; color: var(--accent); background: rgba(22, 27, 34, 0.85); } .glass.hidden { opacity: 0; pointer-events: none; transform: translateY(20px); } + body.chart-route .trend-toggle { display: none !important; } /* summary items — wrap on narrow screens */ .summary-item { white-space: nowrap; } @@ -274,12 +386,25 @@ .sym { font-size: 0.75em; } .bar { width: 48px; } .class-legend { flex-wrap: wrap; gap: 0.5rem 1rem; font-size: 0.72rem; } - .drow td { padding-left: 1.2rem; font-size: 0.72rem; } + .gap-open td { display: block; width: calc(100vw - 2rem); max-width: calc(100vw - 2rem); } + .gap-panel { padding-left: 0.45rem; padding-right: 0.45rem; width: calc(100vw - 2rem); max-width: calc(100vw - 2rem); } + .gap-head, .gap-tools { grid-template-columns: 1fr; padding-right: 0; max-width: 100%; } + .gap-actions { justify-content: flex-start; } + .gap-metrics, .pds-facets, .repo-list { padding-right: 0; max-width: 100%; } + .lane-tabs { max-width: 100%; flex-wrap: wrap; } + .lane-tabs button { flex: 0 0 auto; } + .gap-search { max-width: 100%; } + .repo-row { grid-template-columns: 1fr; gap: 0.2rem; padding: 0.5rem 0.35rem; } + .repo-meta { display: grid; grid-template-columns: 1fr; gap: 0.16rem; } + .repo-status { max-width: 100%; } + .repo-actions { justify-content: flex-start; flex-wrap: wrap; min-width: 0; } + .repo-evidence, .pds-host { white-space: normal; overflow-wrap: anywhere; } .did-link { font-size: 0.68rem; } .nav a { font-size: 0.7rem; padding: 0.15rem 0.35rem; } .nav .time-detail { display: none; } #ftip { display: none !important; } #trend { height: 200px; } + body.chart-route #trend { height: 320px; min-height: 320px; } .trend-zoom { top: 8px; right: 8px; font-size: 0.65rem; padding: 6px 14px; @@ -376,6 +501,16 @@ function pct(n, t) { return t === 0 ? '\u2014' : (n / t * 100).toFixed(2) + '%'; function tip(text, t) { return `${text}`; } +function esc(s) { + return String(s).replace(/[&<>"']/g, c => ({ + '&': '&', + '<': '<', + '>': '>', + '"': '"', + "'": ''', + }[c])); +} + function rn(host) { const o = op(host); return `${o.sym} ${host}`; @@ -383,7 +518,25 @@ function rn(host) { function did(d) { const s = d.length > 32 ? d.slice(0, 16) + '\u2026' + d.slice(-12) : d; - return `${s}`; + return `${s}`; +} + +function statusLinks(x, relay, connected) { + if (!x.pds) return ''; + const encodedDid = encodeURIComponent(x.did); + let h = `repo`; + const hostname = x.pds.replace(/^https?:\/\//, '').replace(/\/.*$/, ''); + if (connected && relay && hostname) { + h += `\u00b7host`; + } + return h; +} + +function pdsHost(pds) { + if (!pds) return ''; + return pds.replace(/^https?:\/\//, '').replace(/\/.*$/, ''); } function bar(seen, total) { @@ -442,6 +595,7 @@ function drawTrend(hover) { const canvas = document.getElementById('trend'); if (!canvas || !_tc) return; const { runs, series, hosts, sorted } = _tc; + const chartMode = document.body.classList.contains('chart-route'); const ctx = canvas.getContext('2d'); const dpr = window.devicePixelRatio || 1; const W = window.innerWidth, H = canvas.clientHeight || 280; @@ -449,20 +603,38 @@ function drawTrend(hover) { ctx.scale(dpr, dpr); ctx.clearRect(0, 0, W, H); - const pad = { t: 20, r: 20, b: 30, l: 20 }; + const pad = chartMode + ? { t: W < 640 ? 86 : 78, r: W < 640 ? 18 : 30, b: 54, l: W < 640 ? 48 : 64 } + : { t: 20, r: 20, b: 30, l: 20 }; const cw = W - pad.l - pad.r, ch = H - pad.t - pad.b; - // rank-based (quantile) Y mapping: space values by distribution rank, - // not absolute distance. the dense cluster near 100% gets most of the - // canvas; outliers are clearly separated at the bottom. + const allVals = []; const valSet = new Set(); for (const host of hosts) { - for (const v of series[host]) { if (v != null) valSet.add(Math.round(v * 100) / 100); } + for (const v of series[host]) { + if (v == null) continue; + allVals.push(v); + valSet.add(Math.round(v * 100) / 100); + } } const ranks = [...valSet].sort((a, b) => a - b); + let yMin = 0, yMax = 100; + if (chartMode && allVals.length) { + const rawMin = Math.min(...allVals); + const rawMax = Math.max(...allVals); + yMin = Math.max(0, Math.floor((rawMin - 3) / 5) * 5); + yMax = Math.ceil((rawMax + 3) / 5) * 5; + if (yMax - yMin < 10) { + const mid = (yMax + yMin) / 2; + yMin = Math.max(0, mid - 5); + yMax = mid + 5; + } + if (yMax <= yMin) yMax = yMin + 10; + } const toX = i => pad.l + (i / (runs.length - 1)) * cw; const toY = v => { + if (chartMode) return pad.t + ch - ((v - yMin) / (yMax - yMin)) * ch; if (ranks.length <= 1) return pad.t + ch / 2; const vr = Math.round(v * 100) / 100; // binary search for position in sorted unique values @@ -478,10 +650,56 @@ function drawTrend(hover) { const norm = (lo + frac) / (ranks.length - 1); return pad.t + ch - norm * ch; }; - _tc.layout = { pad, cw, ch, toX, toY, W, H }; + _tc.layout = { pad, cw, ch, toX, toY, W, H, yMin, yMax, chartMode }; const validPts = host => series[host].map((v, i) => [i, v]).filter(p => p[1] !== null); + if (chartMode) { + ctx.fillStyle = '#0d1117'; + ctx.fillRect(0, 0, W, H); + + ctx.font = (W < 640 ? '0.82rem' : '1rem') + " 'SF Mono', 'Cascadia Code', 'Fira Code', monospace"; + ctx.fillStyle = '#c9d1d9'; + ctx.textAlign = 'left'; + ctx.fillText('coverage over time', pad.l, 28); + + const first = new Date(runs[0].ts); + const last = new Date(runs[runs.length - 1].ts); + const rangeLabel = `${runs.length} measurements · ${first.toLocaleString(undefined, { month: 'short', day: 'numeric', hour: 'numeric', minute: '2-digit' })} → ${last.toLocaleString(undefined, { month: 'short', day: 'numeric', hour: 'numeric', minute: '2-digit' })}`; + ctx.font = (W < 640 ? '0.58rem' : '0.68rem') + " 'SF Mono', 'Cascadia Code', 'Fira Code', monospace"; + ctx.fillStyle = '#8b949e'; + ctx.fillText(rangeLabel, pad.l, 48); + + const legendHosts = [...hosts].sort((a, b) => (series[b].filter(v => v != null).at(-1) || 0) - (series[a].filter(v => v != null).at(-1) || 0)); + let lx = pad.l; + let ly = 66; + ctx.font = (W < 640 ? '0.56rem' : '0.64rem') + " 'SF Mono', 'Cascadia Code', 'Fira Code', monospace"; + for (const host of legendHosts.slice(0, W < 640 ? 3 : 8)) { + const o = op(host); + const label = `${o.sym} ${host}`; + const tw = ctx.measureText(label).width + 28; + if (lx + tw > W - pad.r) { lx = pad.l; ly += 16; } + ctx.fillStyle = o.color; + ctx.fillText(label, lx, ly); + lx += tw; + } + + ctx.font = (W < 640 ? '0.52rem' : '0.58rem') + " 'SF Mono', 'Cascadia Code', 'Fira Code', monospace"; + ctx.textAlign = 'right'; + for (let i = 0; i <= 4; i++) { + const v = yMin + ((yMax - yMin) * i / 4); + const y = toY(v); + ctx.beginPath(); ctx.moveTo(pad.l, y); ctx.lineTo(W - pad.r, y); + ctx.strokeStyle = 'rgba(201, 209, 217, 0.09)'; ctx.lineWidth = 1; + ctx.stroke(); + ctx.fillStyle = '#8b949e'; + ctx.globalAlpha = 0.75; + ctx.fillText(v.toFixed(yMax - yMin <= 12 ? 1 : 0) + '%', pad.l - 10, y + 3); + ctx.globalAlpha = 1; + } + ctx.textAlign = 'left'; + } + for (const host of sorted) { const vp = validPts(host); if (vp.length < 2) continue; @@ -502,26 +720,38 @@ function drawTrend(hover) { } const b = isHovered ? 2.5 : 1; + const chartBoost = chartMode ? 1.8 : 1; // layer 0: wide halo ctx.save(); trace(); ctx.shadowColor = color; ctx.shadowBlur = 20 * b; - ctx.strokeStyle = color; ctx.lineWidth = 2; ctx.globalAlpha = 0.04 * b; + ctx.strokeStyle = color; ctx.lineWidth = 2.5 * chartBoost; ctx.globalAlpha = chartMode ? 0.10 * b : 0.04 * b; ctx.stroke(); ctx.restore(); // layer 1: medium glow ctx.save(); trace(); ctx.shadowColor = color; ctx.shadowBlur = 6 * b; - ctx.strokeStyle = color; ctx.lineWidth = 1; ctx.globalAlpha = 0.12 * b; + ctx.strokeStyle = color; ctx.lineWidth = 1.5 * chartBoost; ctx.globalAlpha = chartMode ? 0.20 * b : 0.12 * b; ctx.stroke(); ctx.restore(); // layer 2: core trace(); ctx.strokeStyle = color; - ctx.lineWidth = isHovered ? 1.5 : 0.7; - ctx.globalAlpha = isHovered ? 0.65 : 0.35; + ctx.lineWidth = chartMode ? (isHovered ? 3 : 2) : (isHovered ? 1.5 : 0.7); + ctx.globalAlpha = chartMode ? (isHovered ? 0.95 : 0.78) : (isHovered ? 0.65 : 0.35); ctx.stroke(); ctx.globalAlpha = 1; + if (chartMode) { + for (const p of vp) { + ctx.beginPath(); + ctx.arc(toX(p[0]), toY(p[1]), isHovered ? 3.5 : 2.4, 0, Math.PI * 2); + ctx.fillStyle = color; + ctx.globalAlpha = isHovered ? 1 : 0.85; + ctx.fill(); + } + ctx.globalAlpha = 1; + } + // layer 3: bright center on hover if (isHovered) { trace(); @@ -614,6 +844,43 @@ function drawTrend(hover) { } } +function drawTrendEmpty(raw) { + const canvas = document.getElementById('trend'); + if (!canvas || !document.body.classList.contains('chart-route')) return; + + const ctx = canvas.getContext('2d'); + const dpr = window.devicePixelRatio || 1; + const W = window.innerWidth, H = canvas.clientHeight || 280; + canvas.width = W * dpr; canvas.height = H * dpr; + ctx.scale(dpr, dpr); + ctx.clearRect(0, 0, W, H); + ctx.fillStyle = '#0d1117'; + ctx.fillRect(0, 0, W, H); + + const points = Array.isArray(raw) ? raw.length : 0; + const runs = [...new Set((raw || []).map(p => p.ts))].length; + const relays = [...new Set((raw || []).map(p => p.host))]; + const selected = (_trendRoute && _trendRoute.relays && _trendRoute.relays.length) + ? _trendRoute.relays.join(', ') + : (relays.length ? relays.join(', ') : 'selected relays'); + + ctx.textAlign = 'left'; + ctx.font = (W < 640 ? '0.82rem' : '1rem') + " 'SF Mono', 'Cascadia Code', 'Fira Code', monospace"; + ctx.fillStyle = '#c9d1d9'; + ctx.fillText('coverage over time', W < 640 ? 28 : 64, 44); + + ctx.font = (W < 640 ? '0.62rem' : '0.72rem') + " 'SF Mono', 'Cascadia Code', 'Fira Code', monospace"; + ctx.fillStyle = '#8b949e'; + ctx.fillText('need at least 2 measurements for a plotted coverage line', W < 640 ? 28 : 64, H / 2 - 10); + + ctx.font = (W < 640 ? '0.54rem' : '0.62rem') + " 'SF Mono', 'Cascadia Code', 'Fira Code', monospace"; + ctx.globalAlpha = 0.75; + ctx.fillText(`found ${runs} measurement${runs === 1 ? '' : 's'} across ${points} relay point${points === 1 ? '' : 's'}`, W < 640 ? 28 : 64, H / 2 + 16); + ctx.fillText(`range/filter: ${selected}`, W < 640 ? 28 : 64, H / 2 + 38); + ctx.fillText('widen the time range or choose relays with data in this window', W < 640 ? 28 : 64, H / 2 + 60); + ctx.globalAlpha = 1; +} + // trend hover document.addEventListener('mousemove', function(e) { if (!_tc || !_tc.layout) return; @@ -703,6 +970,272 @@ document.addEventListener('mouseout', function(e) { // --- dashboard render --- +const classificationOrder = [ + 'active_missing', + 'no_pds_endpoint', + 'malformed_did', + 'unsupported_did_method', + 'did_resolution_failed', + 'invalid_did_document', + 'classification_not_attempted', + 'coverage_gap', + 'unresolvable', + 'deactivated', +]; + +const classificationMeta = { + active_missing: { label: 'stream gap', detail: 'repo activity appeared on another stream, but not this relay during the window', cls: 'c-gap', lane: 'active' }, + no_pds_endpoint: { label: 'no PDS endpoint', detail: 'DID resolved, but no atproto PDS service was advertised', cls: 'c-dead', lane: 'inactive' }, + malformed_did: { label: 'malformed DID', detail: 'DID string did not parse', cls: 'c-unr', lane: 'uncertain' }, + unsupported_did_method: { label: 'unsupported method', detail: 'DID method is not supported by the resolver', cls: 'c-unr', lane: 'uncertain' }, + did_resolution_failed: { label: 'DID resolution failed', detail: 'resolver/network did not return a usable DID document', cls: 'c-unr', lane: 'uncertain' }, + invalid_did_document: { label: 'invalid DID document', detail: 'resolver returned a document relay-eval could not interpret safely', cls: 'c-unr', lane: 'uncertain' }, + classification_not_attempted: { label: 'not checked', detail: 'classification budget was exhausted before this DID was checked', cls: 'c-unr', lane: 'uncertain' }, + coverage_gap: { label: 'legacy active account missed', detail: 'v1 label: resolved account missed by this relay', cls: 'c-gap', lane: 'active' }, + unresolvable: { label: 'legacy unresolved', detail: 'v1 label: unresolved or not attempted', cls: 'c-unr', lane: 'uncertain' }, + deactivated: { label: 'legacy inactive', detail: 'v1 label: resolved without a PDS endpoint', cls: 'c-dead', lane: 'inactive' }, +}; + +function emptyDiffs() { + const out = { dids: [] }; + for (const key of classificationOrder) out[key] = 0; + return out; +} + +function diffTotal(d) { + return classificationOrder.reduce((sum, key) => sum + (d[key] || 0), 0); +} + +function diffActive(d) { + return (d.active_missing || 0) + (d.coverage_gap || 0); +} + +function diffActiveLabel(d) { + const live = diffActive(d); + if (!live) return '\u2014'; + return live.toLocaleString() + (diffNotChecked(d) > 0 ? '+' : ''); +} + +function diffUncertain(d) { + return (d.malformed_did || 0) + + (d.unsupported_did_method || 0) + + (d.did_resolution_failed || 0) + + (d.invalid_did_document || 0) + + (d.classification_not_attempted || 0) + + (d.unresolvable || 0); +} + +function diffInactive(d) { + return (d.no_pds_endpoint || 0) + (d.deactivated || 0); +} + +function diffNotChecked(d) { + return d.classification_not_attempted || 0; +} + +function diffChecked(d) { + return Math.max(0, diffTotal(d) - diffNotChecked(d)); +} + +function diffLaneFor(x) { + return (classificationMeta[x.classification] || { lane: 'uncertain' }).lane; +} + +function laneCell(x, lane) { + const meta = classificationMeta[x.classification] || { + label: x.classification, + detail: 'unrecognized classification', + cls: 'c-unr', + lane: 'uncertain', + }; + if (meta.lane !== lane) return '\u2014'; + const version = x.classification_version === 1 ? ' (v1)' : ''; + const host = lane === 'active' ? pdsHost(x.pds) : ''; + const label = host ? `${host}` : `${meta.label}${version}`; + return `${label}`; +} + +let _diffState = {}; +let _diffSearchTimers = {}; + +function diffStateKey(runId, relay) { + return `${runId || 'latest'}|${relay}`; +} + +function laneApiName(lane) { + return lane === 'active' ? 'live' : lane; +} + +function laneDisplay(lane) { + return { + all: 'all', + active: 'stream gap', + unchecked: 'unverified', + review: 'needs review', + inactive: 'no PDS', + checked: 'checked', + }[lane] || lane; +} + +function statusDot(meta) { + const color = meta.cls === 'c-gap' ? 'var(--red)' : meta.cls === 'c-dead' ? 'var(--muted)' : 'var(--yellow)'; + return ``; +} + +function repoStatus(x) { + const meta = classificationMeta[x.classification] || { + label: x.classification, + detail: 'unrecognized classification', + cls: 'c-unr', + lane: 'uncertain', + }; + const version = x.classification_version === 1 ? ' v1' : ''; + return `${statusDot(meta)}${meta.label}${version}`; +} + +function evidenceCell(x) { + const meta = classificationMeta[x.classification] || {}; + if (x.pds) { + return `${pdsHost(x.pds)}`; + } + if (x.classification === 'classification_not_attempted') return 'resolver budget not spent'; + return `${esc(meta.detail || 'no host evidence')}`; +} + +function rowActions(x, relay, connected) { + let h = `pds.ls`; + const links = statusLinks(x, relay, connected); + if (links) h += `\u00b7${links}`; + return h; +} + +function renderRepoRow(x, state) { + return `
` + + `
${did(x.did)}
` + + `
${repoStatus(x)}${evidenceCell(x)}
` + + `
${rowActions(x, state.relay, state.connected)}
` + + `
`; +} + +function diffApiUrl(state, summary) { + const params = [`relay=${state.relay}`]; + if (state.lane && state.lane !== 'all') params.push(`lane=${laneApiName(state.lane)}`); + if (state.q) params.push(`q=${state.q}`); + if (state.pdsHost) params.push(`pds_host=${state.pdsHost}`); + if (!summary) params.push('limit=1000'); + return `/api/runs/${state.runId}/diffs${summary ? '/summary' : ''}?${params.join('&')}`; +} + +function renderPdsFacets(id, state) { + const facets = state.summary && state.summary.facets ? (state.summary.facets.pds || []) : []; + if (!facets.length && !state.pdsHost) return ''; + let h = `
`; + if (state.pdsHost) h += ``; + for (const p of facets.slice(0, 10)) { + const host = p.pds_host || pdsHost(p.pds); + const active = state.pdsHost === host; + h += ``; + } + h += `
`; + return h; +} + +function renderDiffMetrics(state) { + const shown = state.rows.length; + const filteredTotal = state.summary ? state.summary.total : state.total; + const checked = state.summary ? state.summary.checked : diffChecked(state.counts); + const unchecked = state.summary ? state.summary.unchecked : diffNotChecked(state.counts); + let h = ''; + h += `stream gap ${diffActiveLabel(state.counts)}`; + h += `unverified ${diffUncertain(state.counts).toLocaleString()}`; + if (diffInactive(state.counts) > 0) h += `no PDS ${diffInactive(state.counts).toLocaleString()}`; + h += `checked ${checked.toLocaleString()} / ${filteredTotal.toLocaleString()}`; + if (unchecked > 0) h += `resolver budget left ${unchecked.toLocaleString()} unknown`; + return h; +} + +function renderRepoList(state) { + const shown = state.rows.length; + let h = ''; + if (state.loading && shown === 0) { + h += `
loading gap rows...
`; + } else if (shown === 0) { + h += `
no rows match this filter
`; + } else { + for (const x of state.rows) h += renderRepoRow(x, state); + } + return h; +} + +function renderDiffFooter(id, state) { + const shown = state.rows.length; + const filteredTotal = state.summary ? state.summary.total : state.total; + const canLoad = state.nextCursor !== null || !state.loadedFromApi; + const note = `${shown.toLocaleString()} / ${filteredTotal.toLocaleString()} shown`; + let h = ''; + h += ``; + h += `${note}${state.loading ? ' · loading' : ''}${state.nextCursor !== null ? ' · more available' : ''}`; + return h; +} + +function updateDiffPanelParts(id, state) { + const metrics = document.getElementById('gap-metrics-' + id); + const facets = document.getElementById('pds-facets-' + id); + const list = document.getElementById('repo-list-' + id); + const footer = document.getElementById('gap-footer-' + id); + const clear = document.getElementById('clear-' + id); + const api = document.getElementById('diff-api-' + id); + const summary = document.getElementById('diff-summary-' + id); + if (metrics) metrics.innerHTML = renderDiffMetrics(state); + if (facets) facets.innerHTML = renderPdsFacets(id, state); + if (list) list.innerHTML = renderRepoList(state); + if (footer) footer.innerHTML = renderDiffFooter(id, state); + if (clear) clear.hidden = !(state.q || state.pdsHost); + if (api) api.href = diffApiUrl(state, false); + if (summary) summary.href = diffApiUrl(state, true); +} + +function renderDiffBody(id, state) { + const apiHref = diffApiUrl(state, false); + const summaryHref = diffApiUrl(state, true); + + let h = `
`; + h += `
`; + h += `
${rn(state.relay)}
`; + h += `
repos below emitted on another observed synchronization stream during this window, but not on this relay. hosting status is evidence for investigation, not the denominator.
`; + h += `
`; + h += `diffs API`; + h += `summary API`; + h += `
`; + + h += `
`; + h += renderDiffMetrics(state); + h += `
`; + + h += `
`; + h += ``; + h += `
`; + for (const lane of ['all', 'active', 'unchecked', 'review', 'inactive', 'checked']) { + h += ``; + } + h += `
`; + + h += `
${renderPdsFacets(id, state)}
`; + + h += `
`; + h += renderRepoList(state); + h += `
`; + + h += ``; + h += `
`; + return h; +} + function render(data) { if (!data) return '

no runs yet

'; @@ -730,7 +1263,7 @@ function render(data) { h += ` \u00b7 `; h += `${tip(winLabel(data.window_seconds), data.window_seconds + 's collection window')} window`; h += ` \u00b7 `; - h += `${effectiveUnion.toLocaleString()} active accounts`; + h += `${effectiveUnion.toLocaleString()} observed repos`; h += ``; // operator legend @@ -748,11 +1281,11 @@ function render(data) { // diffs by relay (aggregated counts + sample DIDs) const rd = {}; for (const c of (data.diff_counts || [])) { - if (!rd[c.relay]) rd[c.relay] = { coverage_gap: 0, unresolvable: 0, deactivated: 0, dids: [] }; + if (!rd[c.relay]) rd[c.relay] = emptyDiffs(); rd[c.relay][c.classification] = c.count; } for (const d of (data.diff_samples || [])) { - if (!rd[d.relay]) rd[d.relay] = { coverage_gap: 0, unresolvable: 0, deactivated: 0, dids: [] }; + if (!rd[d.relay]) rd[d.relay] = emptyDiffs(); rd[d.relay].dids.push(d); } @@ -763,7 +1296,7 @@ function render(data) { h += ``; h += `

`; h += `

${tip('measured', 'subscribes to every relay\'s firehose simultaneously, records which accounts (DIDs) emit events on each, then compares against the union')} over a ${tip(winLabel(data.window_seconds), data.window_seconds + 's collection window')} window.

`; - h += `

coverage = accounts seen on this relay / accounts seen on any relay

`; + h += `

coverage = repos observed on this relay / repos observed on any relay

`; h += `
`; h += `
`; @@ -779,8 +1312,8 @@ function render(data) { const ranked = [...data.stats].sort((a, b) => b.unique_dids - a.unique_dids); for (const s of ranked) { - const d = rd[s.host] || { coverage_gap: 0, unresolvable: 0, deactivated: 0 }; - const missed = d.coverage_gap + d.unresolvable + d.deactivated; + const d = rd[s.host] || emptyDiffs(); + const missed = diffTotal(d); const o = op(s.host); const byLink = o.url ? `${o.name}` @@ -809,50 +1342,67 @@ function render(data) { // breakdown const withMisses = data.stats.filter(s => { const d = rd[s.host]; - return d && (d.coverage_gap + d.unresolvable + d.deactivated) > 0; + return d && diffTotal(d) > 0; }).sort((a, b) => { const da = rd[a.host], db = rd[b.host]; - return (db.coverage_gap + db.unresolvable + db.deactivated) - (da.coverage_gap + da.unresolvable + da.deactivated); + return diffTotal(db) - diffTotal(da); }); if (withMisses.length > 0) { - h += `

why accounts were missed

`; + const showNoPds = withMisses.some(s => diffInactive(rd[s.host]) > 0); + const detailColspan = 6; + h += `

stream coverage gaps api examples

`; h += `
`; - h += ` ${tip('missed (active)', 'account has a working PDS but this relay didn\u2019t see it')}`; - h += ` ${tip('unresolvable', 'DID lookup failed \u2014 DNS issue, PLC outage, or brand-new account')}`; - h += ` ${tip('deactivated', 'account is deactivated or deleted')}`; + h += ` ${tip('stream gap', 'checked repo emitted on another stream, but not this relay')}`; + h += ` ${tip('unverified', 'status could not be checked, or the row was outside the resolver budget')}`; + if (showNoPds) h += ` ${tip('no PDS', 'DID resolved but no PDS endpoint was found')}`; h += `
`; + h += `

open a relay for a compact investigation view: DID search, status lanes, PDS host facets, pds.ls links, and live getRepoStatus/getHostStatus queries.

`; h += `
`; - h += ``; - h += ``; - h += ``; - h += ``; - h += ``; + h += ``; + h += ``; + h += ``; + h += ``; + h += ``; + h += ``; h += ``; for (const s of withMisses) { const d = rd[s.host]; - const total = d.coverage_gap + d.unresolvable + d.deactivated; + const total = diffTotal(d); const id = s.host.replace(/[^a-z0-9]/g, '_'); + const stateKey = diffStateKey(data.id, s.host); + _diffState[id] = { + runId: data.id, + relay: s.host, + connected: s.connected, + counts: d, + rows: d.dids.slice(), + total, + detailColspan, + q: '', + lane: 'all', + pdsHost: '', + summary: null, + nextCursor: null, + loadedFromApi: false, + searching: false, + loading: false, + key: stateKey, + }; h += ``; - h += ``; - h += ``; - h += ``; - h += ``; + h += ``; + h += ``; + h += ``; + h += ``; + h += ``; h += ``; h += ``; h += ``; - for (const x of d.dids) { - const cls = x.classification === 'coverage_gap' ? 'c-gap' : x.classification === 'unresolvable' ? 'c-unr' : 'c-dead'; - const label = x.classification === 'coverage_gap' ? 'active' : x.classification; - h += ``; - } - if (total > d.dids.length) { - h += ``; - } + h += renderDiffBody(id, _diffState[id]); h += ``; } h += `
${tip('relay', 'the relay that missed these accounts')}${tip('active', 'account has a working PDS but this relay didn\u2019t see it')}${tip('unresolvable', 'DID lookup failed')}${tip('deactivated', 'account is deactivated or deleted')}${tip('total', 'total accounts missed by this relay')}${tip('relay', 'the relay whose stream missed these repos')}${tip('gap lower bound', 'checked repos with current host evidence; plus means unchecked rows may also belong here')}${tip('unverified', 'rows outside resolver budget or needing review')}${tip('no PDS', 'resolved DID document has no PDS endpoint')}${tip('checked', 'resolver-classified rows')}${tip('total', 'all repos this relay missed in the run')}
\u25b8${rn(s.host)}${d.coverage_gap || '\u2014'}${d.unresolvable || '\u2014'}${d.deactivated || '\u2014'}\u25b8${rn(s.host)}${diffActiveLabel(d)}${diffUncertain(d) || '\u2014'}${diffInactive(d) || '\u2014'}${diffChecked(d).toLocaleString()}${total.toLocaleString()}
${did(x.did)}${label}
\u2026 and ${total - d.dids.length} more
`; @@ -869,6 +1419,155 @@ function toggle(id) { if (!b) return; b.classList.toggle('open'); if (i) i.textContent = b.classList.contains('open') ? '\u25be' : '\u25b8'; + if (b.classList.contains('open')) refreshDiffPanel(id, { loadSummary: true, loadRows: false }); +} + +async function fetchDiffPage(state, reset) { + const params = [`relay=${state.relay}`, 'limit=100']; + if (state.q) params.push(`q=${state.q}`); + if (state.lane && state.lane !== 'all') params.push(`lane=${laneApiName(state.lane)}`); + if (state.pdsHost) params.push(`pds_host=${state.pdsHost}`); + if (!reset && state.nextCursor !== null) params.push(`cursor=${state.nextCursor}`); + + const url = `/api/runs/${state.runId}/diffs?${params.join('&')}`; + const page = await fetch(url).then(r => r.json()); + if (page.error) throw new Error(page.error); + return page; +} + +async function fetchDiffSummary(state) { + const url = diffApiUrl(state, true) + '&limit=12'; + const summary = await fetch(url).then(r => r.json()); + if (summary.error) throw new Error(summary.error); + return summary; +} + +async function refreshDiffPanel(id, opts) { + const state = _diffState[id]; + const body = document.getElementById('xb-' + id); + if (!state || !body || state.loading) return; + + state.loading = true; + body.innerHTML = renderDiffBody(id, state); + try { + if (opts.loadSummary) state.summary = await fetchDiffSummary(state); + if (opts.loadRows) { + const page = await fetchDiffPage(state, true); + state.rows = page.diffs; + state.nextCursor = page.next_cursor; + state.loadedFromApi = true; + state.searching = Boolean(state.q || state.pdsHost || (state.lane && state.lane !== 'all')); + } + } finally { + state.loading = false; + body.innerHTML = renderDiffBody(id, state); + } +} + +async function loadMoreDiffs(id) { + const state = _diffState[id]; + const body = document.getElementById('xb-' + id); + if (!state || !body || state.loading) return; + + const reset = !state.loadedFromApi; + state.loading = true; + body.innerHTML = renderDiffBody(id, state); + try { + const page = await fetchDiffPage(state, reset); + state.rows = reset ? page.diffs : state.rows.concat(page.diffs); + state.nextCursor = page.next_cursor; + state.loadedFromApi = true; + state.searching = Boolean(state.q || state.pdsHost || (state.lane && state.lane !== 'all')); + if (!state.summary) state.summary = await fetchDiffSummary(state); + } catch (err) { + state.rows.push({ + did: 'error', + classification: 'did_resolution_failed', + classification_version: 2, + }); + } finally { + state.loading = false; + body.innerHTML = renderDiffBody(id, state); + } +} + +async function searchDiffs(id) { + const state = _diffState[id]; + const input = document.getElementById('q-' + id); + if (!state || !input) return; + + const nextQ = input.value.trim(); + if (state.q === nextQ && state.loadedFromApi) return; + const seq = (state.searchSeq || 0) + 1; + state.searchSeq = seq; + state.q = nextQ; + state.nextCursor = null; + state.loadedFromApi = true; + state.searching = Boolean(state.q || state.pdsHost || (state.lane && state.lane !== 'all')); + state.rows = []; + state.loading = true; + updateDiffPanelParts(id, state); + try { + const summary = await fetchDiffSummary(state); + if (state.searchSeq !== seq) return; + state.summary = summary; + const page = await fetchDiffPage(state, true); + if (state.searchSeq !== seq) return; + state.rows = page.diffs; + state.nextCursor = page.next_cursor; + } finally { + if (state.searchSeq === seq) { + state.loading = false; + updateDiffPanelParts(id, state); + } + } +} + +function queueDiffSearch(id) { + clearTimeout(_diffSearchTimers[id]); + _diffSearchTimers[id] = setTimeout(() => { + searchDiffs(id).catch(() => {}); + }, 300); +} + +function clearDiffSearch(id) { + const state = _diffState[id]; + if (!state) return; + clearTimeout(_diffSearchTimers[id]); + state.searchSeq = (state.searchSeq || 0) + 1; + state.q = ''; + state.pdsHost = ''; + state.lane = 'all'; + state.summary = null; + state.searching = false; + state.loadedFromApi = false; + state.nextCursor = null; + state.rows = state.counts.dids.slice(); + const body = document.getElementById('xb-' + id); + if (body) body.innerHTML = renderDiffBody(id, state); + refreshDiffPanel(id, { loadSummary: true, loadRows: false }); +} + +function setDiffLane(id, lane) { + const state = _diffState[id]; + if (!state || state.lane === lane) return; + state.searchSeq = (state.searchSeq || 0) + 1; + state.lane = lane; + state.nextCursor = null; + state.rows = []; + state.loadedFromApi = true; + refreshDiffPanel(id, { loadSummary: true, loadRows: true }); +} + +function setPdsFilter(id, host) { + const state = _diffState[id]; + if (!state) return; + state.searchSeq = (state.searchSeq || 0) + 1; + state.pdsHost = host; + state.nextCursor = null; + state.rows = []; + state.loadedFromApi = true; + refreshDiffPanel(id, { loadSummary: true, loadRows: true }); } function toggleTrendFocus() { @@ -974,8 +1673,29 @@ let _zoomDebounce = null; let _totalRuns = 0; let _zoomLevel = 48; let _zoomFlashTimer = null; +let _trendRoute = parseTrendRoute(); +if (_trendRoute) document.body.classList.add('chart-route'); + +function parseTrendRoute() { + const parts = window.location.pathname.split('/').filter(Boolean); + if (parts[0] !== 'chart') return null; + const route = {}; + if (parts[1]) route.since = decodeURIComponent(parts[1]); + if (parts[2]) route.until = decodeURIComponent(parts[2]); + if (parts[3]) route.relays = decodeURIComponent(parts[3]).split(',').filter(Boolean); + return route; +} + +function trendPathFromRuns(runs, relays) { + if (!runs || runs.length < 2) return '/chart'; + const since = Math.floor(new Date(runs[0].ts).getTime() / 1000); + const until = Math.floor(new Date(runs[runs.length - 1].ts).getTime() / 1000); + const relayPart = relays && relays.length ? '/' + encodeURIComponent(relays.join(',')) : ''; + return `/chart/${since}/${until}${relayPart}`; +} function initZoom(total) { + if (_trendRoute && (_trendRoute.since || _trendRoute.until)) return; _totalRuns = total; _zoomLevel = Math.min(48, total); const pill = document.getElementById('trend-zoom'); @@ -1049,11 +1769,27 @@ function initZoom(total) { } function fetchTrend(limit) { - const url = limit ? `/api/trend?limit=${limit}` : '/api/trend'; + let url = limit ? `/api/trend?limit=${limit}` : '/api/trend'; + if (_trendRoute && (_trendRoute.since || _trendRoute.until || (_trendRoute.relays && _trendRoute.relays.length))) { + const params = new URLSearchParams(); + if (_trendRoute.since) params.set('since', _trendRoute.since); + if (_trendRoute.until) params.set('until', _trendRoute.until); + if (_trendRoute.relays && _trendRoute.relays.length) params.set('relays', _trendRoute.relays.join(',')); + url = '/api/trend?' + params.toString(); + } return fetch(url).then(r => r.json()).then(d => { _trendData = d; _tc = buildTrendCache(d); - drawTrend(null); + if (_trendRoute && _tc && _tc.runs && _tc.runs.length >= 2 && (!_trendRoute.since || !_trendRoute.until)) { + const path = trendPathFromRuns(_tc.runs, _trendRoute ? _trendRoute.relays : null); + history.replaceState(null, '', path); + _trendRoute = parseTrendRoute(); + } + if (_tc) { + drawTrend(null); + } else { + drawTrendEmpty(d); + } return d; }).catch(() => {}); } @@ -1190,6 +1926,7 @@ async function init() { nav.innerHTML = nh; await loadRun('latest'); + if (_trendRoute) setCovTab('alltime'); // data loaded — stop pulsing, draw final trend, reveal glass stopLoading(); diff --git a/relay-eval/src/store.zig b/relay-eval/src/store.zig index 09c83bf..89b9006 100644 --- a/relay-eval/src/store.zig +++ b/relay-eval/src/store.zig @@ -52,7 +52,8 @@ pub const Store = struct { \\ run_id INTEGER REFERENCES runs(id), \\ relay TEXT NOT NULL, \\ did TEXT NOT NULL, - \\ classification TEXT NOT NULL + \\ classification TEXT NOT NULL, + \\ classification_version INTEGER NOT NULL DEFAULT 1 \\); \\CREATE INDEX IF NOT EXISTS idx_relay_stats_run_id ON relay_stats(run_id); \\CREATE INDEX IF NOT EXISTS idx_diffs_run_id ON diffs(run_id); @@ -81,6 +82,12 @@ pub const Store = struct { if (!self.hasColumn("relay_stats", "overlap")) { try self.exec("ALTER TABLE relay_stats ADD COLUMN overlap TEXT"); } + if (!self.hasColumn("diffs", "classification_version")) { + try self.exec("ALTER TABLE diffs ADD COLUMN classification_version INTEGER NOT NULL DEFAULT 1"); + } + if (!self.hasColumn("diffs", "pds")) { + try self.exec("ALTER TABLE diffs ADD COLUMN pds TEXT"); + } } /// true if `table` has a column named `col`. used to make ADD COLUMN @@ -111,6 +118,45 @@ pub const Store = struct { relay: []const u8, did: []const u8, classification: []const u8, + classification_version: u16 = 1, + pds: ?[]const u8 = null, + }; + + pub const DiffPageItem = struct { + id: i64, + relay: []const u8, + did: []const u8, + classification: []const u8, + classification_version: u16, + pds: ?[]const u8 = null, + }; + + pub const DiffSummary = struct { + total: u32, + checked: u32, + }; + + pub const DiffRelayFacet = struct { + relay: []const u8, + total: u32, + checked: u32, + live: u32, + review: u32, + inactive: u32, + }; + + pub const DiffClassificationFacet = struct { + classification: []const u8, + classification_version: u16, + count: u32, + }; + + pub const DiffPdsFacet = struct { + pds: []const u8, + total: u32, + live: u32, + review: u32, + inactive: u32, }; pub fn insertRun(self: *Store, timestamp: []const u8, window_seconds: u32, union_dids: u32) !i64 { @@ -157,7 +203,7 @@ pub const Store = struct { } pub fn insertDiff(self: *Store, run_id: i64, diff: Diff) !void { - const sql = "INSERT INTO diffs (run_id, relay, did, classification) VALUES (?, ?, ?, ?)"; + const sql = "INSERT INTO diffs (run_id, relay, did, classification, classification_version, pds) VALUES (?, ?, ?, ?, ?, ?)"; var stmt: ?*c.sqlite3_stmt = null; if (c.sqlite3_prepare_v2(self.db, sql, -1, &stmt, null) != c.SQLITE_OK) { return self.sqlError(); @@ -168,6 +214,12 @@ pub const Store = struct { _ = c.sqlite3_bind_text(stmt, 2, diff.relay.ptr, @intCast(diff.relay.len), SQLITE_STATIC); _ = c.sqlite3_bind_text(stmt, 3, diff.did.ptr, @intCast(diff.did.len), SQLITE_STATIC); _ = c.sqlite3_bind_text(stmt, 4, diff.classification.ptr, @intCast(diff.classification.len), SQLITE_STATIC); + _ = c.sqlite3_bind_int(stmt, 5, @intCast(diff.classification_version)); + if (diff.pds) |pds| { + _ = c.sqlite3_bind_text(stmt, 6, pds.ptr, @intCast(pds.len), SQLITE_STATIC); + } else { + _ = c.sqlite3_bind_null(stmt, 6); + } if (c.sqlite3_step(stmt) != c.SQLITE_DONE) { return self.sqlError(); @@ -261,12 +313,13 @@ pub const Store = struct { pub const DiffCount = struct { relay: []const u8, classification: []const u8, + classification_version: u16, count: u32, }; /// get aggregated diff counts per relay per classification pub fn getDiffCounts(self: *Store, allocator: std.mem.Allocator, run_id: i64) ![]DiffCount { - const sql = "SELECT relay, classification, COUNT(*) FROM diffs WHERE run_id = ? GROUP BY relay, classification"; + const sql = "SELECT relay, classification, classification_version, COUNT(*) FROM diffs WHERE run_id = ? GROUP BY relay, classification_version, classification"; var stmt: ?*c.sqlite3_stmt = null; if (c.sqlite3_prepare_v2(self.db, sql, -1, &stmt, null) != c.SQLITE_OK) { return self.sqlError(); @@ -280,7 +333,8 @@ pub const Store = struct { try counts.append(allocator, .{ .relay = try self.colText(allocator, stmt, 0), .classification = try self.colText(allocator, stmt, 1), - .count = @intCast(c.sqlite3_column_int(stmt, 2)), + .classification_version = @intCast(c.sqlite3_column_int(stmt, 2)), + .count = @intCast(c.sqlite3_column_int(stmt, 3)), }); } return counts.toOwnedSlice(allocator); @@ -288,7 +342,7 @@ pub const Store = struct { /// get sample diffs (up to per_relay_limit per relay) for expandable detail pub fn getDiffSamples(self: *Store, allocator: std.mem.Allocator, run_id: i64, per_relay_limit: u32) ![]Diff { - const sql = "SELECT relay, did, classification FROM diffs WHERE run_id = ? ORDER BY relay, id"; + const sql = "SELECT relay, did, classification, classification_version, pds FROM diffs WHERE run_id = ? ORDER BY relay, id"; var stmt: ?*c.sqlite3_stmt = null; if (c.sqlite3_prepare_v2(self.db, sql, -1, &stmt, null) != c.SQLITE_OK) { return self.sqlError(); @@ -321,12 +375,302 @@ pub const Store = struct { .relay = try self.colText(allocator, stmt, 0), .did = try self.colText(allocator, stmt, 1), .classification = try self.colText(allocator, stmt, 2), + .classification_version = @intCast(c.sqlite3_column_int(stmt, 3)), + .pds = try self.colTextOpt(allocator, stmt, 4), }); } } return samples.toOwnedSlice(allocator); } + /// get a cursor page of diffs for API consumers. cursor is the last seen + /// `diffs.id`; filters are optional when passed as empty strings. + pub fn getDiffPage( + self: *Store, + allocator: std.mem.Allocator, + run_id: i64, + relay: []const u8, + classification: []const u8, + lane: []const u8, + did_search: []const u8, + pds_host: []const u8, + after_id: i64, + limit: u32, + ) ![]DiffPageItem { + const sql = + \\SELECT id, relay, did, classification, classification_version, pds + \\ FROM diffs + \\ WHERE run_id = ? + \\ AND id > ? + \\ AND (? = '' OR relay = ?) + \\ AND (? = '' OR classification = ?) + \\ AND (? = '' OR did LIKE '%' || ? || '%') + \\ AND (? = '' OR pds = ? OR pds = 'https://' || ? OR pds = 'http://' || ?) + \\ AND (? = '' + \\ OR (? = 'live' AND classification IN ('active_missing', 'coverage_gap')) + \\ OR (? = 'review' AND classification IN ('malformed_did', 'unsupported_did_method', 'did_resolution_failed', 'invalid_did_document', 'classification_not_attempted', 'unresolvable')) + \\ OR (? = 'inactive' AND classification IN ('no_pds_endpoint', 'deactivated')) + \\ OR (? = 'unchecked' AND classification = 'classification_not_attempted') + \\ OR (? = 'checked' AND classification != 'classification_not_attempted')) + \\ ORDER BY id ASC + \\ LIMIT ? + ; + var stmt: ?*c.sqlite3_stmt = null; + if (c.sqlite3_prepare_v2(self.db, sql, -1, &stmt, null) != c.SQLITE_OK) { + return self.sqlError(); + } + defer _ = c.sqlite3_finalize(stmt); + + _ = c.sqlite3_bind_int64(stmt, 1, run_id); + _ = c.sqlite3_bind_int64(stmt, 2, after_id); + _ = c.sqlite3_bind_text(stmt, 3, relay.ptr, @intCast(relay.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 4, relay.ptr, @intCast(relay.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 5, classification.ptr, @intCast(classification.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 6, classification.ptr, @intCast(classification.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 7, did_search.ptr, @intCast(did_search.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 8, did_search.ptr, @intCast(did_search.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 9, pds_host.ptr, @intCast(pds_host.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 10, pds_host.ptr, @intCast(pds_host.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 11, pds_host.ptr, @intCast(pds_host.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 12, pds_host.ptr, @intCast(pds_host.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 13, lane.ptr, @intCast(lane.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 14, lane.ptr, @intCast(lane.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 15, lane.ptr, @intCast(lane.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 16, lane.ptr, @intCast(lane.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 17, lane.ptr, @intCast(lane.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 18, lane.ptr, @intCast(lane.len), SQLITE_STATIC); + _ = c.sqlite3_bind_int(stmt, 19, @intCast(limit)); + + var out: std.ArrayList(DiffPageItem) = .empty; + while (c.sqlite3_step(stmt) == c.SQLITE_ROW) { + try out.append(allocator, .{ + .id = c.sqlite3_column_int64(stmt, 0), + .relay = try self.colText(allocator, stmt, 1), + .did = try self.colText(allocator, stmt, 2), + .classification = try self.colText(allocator, stmt, 3), + .classification_version = @intCast(c.sqlite3_column_int(stmt, 4)), + .pds = try self.colTextOpt(allocator, stmt, 5), + }); + } + return out.toOwnedSlice(allocator); + } + + pub fn getDiffSummary( + self: *Store, + run_id: i64, + relay: []const u8, + classification: []const u8, + lane: []const u8, + did_search: []const u8, + pds_host: []const u8, + ) !DiffSummary { + const sql = + \\SELECT COUNT(*), + \\ SUM(CASE WHEN classification != 'classification_not_attempted' THEN 1 ELSE 0 END) + \\ FROM diffs + \\ WHERE run_id = ? + \\ AND (? = '' OR relay = ?) + \\ AND (? = '' OR classification = ?) + \\ AND (? = '' OR did LIKE '%' || ? || '%') + \\ AND (? = '' OR pds = ? OR pds = 'https://' || ? OR pds = 'http://' || ?) + \\ AND (? = '' + \\ OR (? = 'live' AND classification IN ('active_missing', 'coverage_gap')) + \\ OR (? = 'review' AND classification IN ('malformed_did', 'unsupported_did_method', 'did_resolution_failed', 'invalid_did_document', 'classification_not_attempted', 'unresolvable')) + \\ OR (? = 'inactive' AND classification IN ('no_pds_endpoint', 'deactivated')) + \\ OR (? = 'unchecked' AND classification = 'classification_not_attempted') + \\ OR (? = 'checked' AND classification != 'classification_not_attempted')) + ; + var stmt: ?*c.sqlite3_stmt = null; + if (c.sqlite3_prepare_v2(self.db, sql, -1, &stmt, null) != c.SQLITE_OK) return self.sqlError(); + defer _ = c.sqlite3_finalize(stmt); + bindDiffFilters(stmt, run_id, relay, classification, did_search, pds_host, lane, 0, null); + + if (c.sqlite3_step(stmt) == c.SQLITE_ROW) { + return .{ + .total = @intCast(c.sqlite3_column_int(stmt, 0)), + .checked = @intCast(c.sqlite3_column_int(stmt, 1)), + }; + } + return .{ .total = 0, .checked = 0 }; + } + + pub fn getDiffRelayFacets( + self: *Store, + allocator: std.mem.Allocator, + run_id: i64, + relay: []const u8, + classification: []const u8, + lane: []const u8, + did_search: []const u8, + pds_host: []const u8, + ) ![]DiffRelayFacet { + const sql = + \\SELECT relay, + \\ COUNT(*), + \\ SUM(CASE WHEN classification != 'classification_not_attempted' THEN 1 ELSE 0 END), + \\ SUM(CASE WHEN classification IN ('active_missing', 'coverage_gap') THEN 1 ELSE 0 END), + \\ SUM(CASE WHEN classification IN ('malformed_did', 'unsupported_did_method', 'did_resolution_failed', 'invalid_did_document', 'classification_not_attempted', 'unresolvable') THEN 1 ELSE 0 END), + \\ SUM(CASE WHEN classification IN ('no_pds_endpoint', 'deactivated') THEN 1 ELSE 0 END) + \\ FROM diffs + \\ WHERE run_id = ? + \\ AND (? = '' OR relay = ?) + \\ AND (? = '' OR classification = ?) + \\ AND (? = '' OR did LIKE '%' || ? || '%') + \\ AND (? = '' OR pds = ? OR pds = 'https://' || ? OR pds = 'http://' || ?) + \\ AND (? = '' + \\ OR (? = 'live' AND classification IN ('active_missing', 'coverage_gap')) + \\ OR (? = 'review' AND classification IN ('malformed_did', 'unsupported_did_method', 'did_resolution_failed', 'invalid_did_document', 'classification_not_attempted', 'unresolvable')) + \\ OR (? = 'inactive' AND classification IN ('no_pds_endpoint', 'deactivated')) + \\ OR (? = 'unchecked' AND classification = 'classification_not_attempted') + \\ OR (? = 'checked' AND classification != 'classification_not_attempted')) + \\ GROUP BY relay + \\ ORDER BY COUNT(*) DESC, relay ASC + ; + var stmt: ?*c.sqlite3_stmt = null; + if (c.sqlite3_prepare_v2(self.db, sql, -1, &stmt, null) != c.SQLITE_OK) return self.sqlError(); + defer _ = c.sqlite3_finalize(stmt); + bindDiffFilters(stmt, run_id, relay, classification, did_search, pds_host, lane, 0, null); + + var out: std.ArrayList(DiffRelayFacet) = .empty; + while (c.sqlite3_step(stmt) == c.SQLITE_ROW) { + try out.append(allocator, .{ + .relay = try self.colText(allocator, stmt, 0), + .total = @intCast(c.sqlite3_column_int(stmt, 1)), + .checked = @intCast(c.sqlite3_column_int(stmt, 2)), + .live = @intCast(c.sqlite3_column_int(stmt, 3)), + .review = @intCast(c.sqlite3_column_int(stmt, 4)), + .inactive = @intCast(c.sqlite3_column_int(stmt, 5)), + }); + } + return out.toOwnedSlice(allocator); + } + + pub fn getDiffClassificationFacets( + self: *Store, + allocator: std.mem.Allocator, + run_id: i64, + relay: []const u8, + classification: []const u8, + lane: []const u8, + did_search: []const u8, + pds_host: []const u8, + ) ![]DiffClassificationFacet { + const sql = + \\SELECT classification, classification_version, COUNT(*) + \\ FROM diffs + \\ WHERE run_id = ? + \\ AND (? = '' OR relay = ?) + \\ AND (? = '' OR classification = ?) + \\ AND (? = '' OR did LIKE '%' || ? || '%') + \\ AND (? = '' OR pds = ? OR pds = 'https://' || ? OR pds = 'http://' || ?) + \\ AND (? = '' + \\ OR (? = 'live' AND classification IN ('active_missing', 'coverage_gap')) + \\ OR (? = 'review' AND classification IN ('malformed_did', 'unsupported_did_method', 'did_resolution_failed', 'invalid_did_document', 'classification_not_attempted', 'unresolvable')) + \\ OR (? = 'inactive' AND classification IN ('no_pds_endpoint', 'deactivated')) + \\ OR (? = 'unchecked' AND classification = 'classification_not_attempted') + \\ OR (? = 'checked' AND classification != 'classification_not_attempted')) + \\ GROUP BY classification, classification_version + \\ ORDER BY COUNT(*) DESC, classification ASC + ; + var stmt: ?*c.sqlite3_stmt = null; + if (c.sqlite3_prepare_v2(self.db, sql, -1, &stmt, null) != c.SQLITE_OK) return self.sqlError(); + defer _ = c.sqlite3_finalize(stmt); + bindDiffFilters(stmt, run_id, relay, classification, did_search, pds_host, lane, 0, null); + + var out: std.ArrayList(DiffClassificationFacet) = .empty; + while (c.sqlite3_step(stmt) == c.SQLITE_ROW) { + try out.append(allocator, .{ + .classification = try self.colText(allocator, stmt, 0), + .classification_version = @intCast(c.sqlite3_column_int(stmt, 1)), + .count = @intCast(c.sqlite3_column_int(stmt, 2)), + }); + } + return out.toOwnedSlice(allocator); + } + + pub fn getDiffPdsFacets( + self: *Store, + allocator: std.mem.Allocator, + run_id: i64, + relay: []const u8, + classification: []const u8, + lane: []const u8, + did_search: []const u8, + pds_host: []const u8, + limit: u32, + ) ![]DiffPdsFacet { + const sql = + \\SELECT pds, + \\ COUNT(*), + \\ SUM(CASE WHEN classification IN ('active_missing', 'coverage_gap') THEN 1 ELSE 0 END), + \\ SUM(CASE WHEN classification IN ('malformed_did', 'unsupported_did_method', 'did_resolution_failed', 'invalid_did_document', 'classification_not_attempted', 'unresolvable') THEN 1 ELSE 0 END), + \\ SUM(CASE WHEN classification IN ('no_pds_endpoint', 'deactivated') THEN 1 ELSE 0 END) + \\ FROM diffs + \\ WHERE run_id = ? + \\ AND pds IS NOT NULL + \\ AND (? = '' OR relay = ?) + \\ AND (? = '' OR classification = ?) + \\ AND (? = '' OR did LIKE '%' || ? || '%') + \\ AND (? = '' OR pds = ? OR pds = 'https://' || ? OR pds = 'http://' || ?) + \\ AND (? = '' + \\ OR (? = 'live' AND classification IN ('active_missing', 'coverage_gap')) + \\ OR (? = 'review' AND classification IN ('malformed_did', 'unsupported_did_method', 'did_resolution_failed', 'invalid_did_document', 'classification_not_attempted', 'unresolvable')) + \\ OR (? = 'inactive' AND classification IN ('no_pds_endpoint', 'deactivated')) + \\ OR (? = 'unchecked' AND classification = 'classification_not_attempted') + \\ OR (? = 'checked' AND classification != 'classification_not_attempted')) + \\ GROUP BY pds + \\ ORDER BY COUNT(*) DESC, pds ASC + \\ LIMIT ? + ; + var stmt: ?*c.sqlite3_stmt = null; + if (c.sqlite3_prepare_v2(self.db, sql, -1, &stmt, null) != c.SQLITE_OK) return self.sqlError(); + defer _ = c.sqlite3_finalize(stmt); + bindDiffFilters(stmt, run_id, relay, classification, did_search, pds_host, lane, limit, 18); + + var out: std.ArrayList(DiffPdsFacet) = .empty; + while (c.sqlite3_step(stmt) == c.SQLITE_ROW) { + try out.append(allocator, .{ + .pds = try self.colText(allocator, stmt, 0), + .total = @intCast(c.sqlite3_column_int(stmt, 1)), + .live = @intCast(c.sqlite3_column_int(stmt, 2)), + .review = @intCast(c.sqlite3_column_int(stmt, 3)), + .inactive = @intCast(c.sqlite3_column_int(stmt, 4)), + }); + } + return out.toOwnedSlice(allocator); + } + + fn bindDiffFilters( + stmt: ?*c.sqlite3_stmt, + run_id: i64, + relay: []const u8, + classification: []const u8, + did_search: []const u8, + pds_host: []const u8, + lane: []const u8, + limit: u32, + limit_idx: ?i32, + ) void { + _ = c.sqlite3_bind_int64(stmt, 1, run_id); + _ = c.sqlite3_bind_text(stmt, 2, relay.ptr, @intCast(relay.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 3, relay.ptr, @intCast(relay.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 4, classification.ptr, @intCast(classification.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 5, classification.ptr, @intCast(classification.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 6, did_search.ptr, @intCast(did_search.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 7, did_search.ptr, @intCast(did_search.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 8, pds_host.ptr, @intCast(pds_host.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 9, pds_host.ptr, @intCast(pds_host.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 10, pds_host.ptr, @intCast(pds_host.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 11, pds_host.ptr, @intCast(pds_host.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 12, lane.ptr, @intCast(lane.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 13, lane.ptr, @intCast(lane.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 14, lane.ptr, @intCast(lane.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 15, lane.ptr, @intCast(lane.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 16, lane.ptr, @intCast(lane.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 17, lane.ptr, @intCast(lane.len), SQLITE_STATIC); + if (limit_idx) |idx| _ = c.sqlite3_bind_int(stmt, idx, @intCast(limit)); + } + /// get the most recent valid run ID pub fn getLatestRunId(self: *Store) !?i64 { const sql = "SELECT r.id FROM runs r WHERE 1=1" ++ valid_run_filter ++ " ORDER BY r.id DESC LIMIT 1"; @@ -415,6 +759,44 @@ pub const Store = struct { return points.toOwnedSlice(allocator); } + /// get coverage data across valid runs in an inclusive timestamp range. + pub fn getTrendDataRange(self: *Store, allocator: std.mem.Allocator, since: []const u8, until: []const u8) ![]TrendPoint { + const sql = + \\SELECT r.id, r.timestamp, r.union_dids, rs.host, rs.unique_dids, rs.overlap + \\FROM runs r + \\JOIN relay_stats rs ON rs.run_id = r.id + \\WHERE r.timestamp >= ? + \\ AND r.timestamp <= ? + \\ AND r.id IN (SELECT run_id FROM relay_stats GROUP BY run_id HAVING SUM(connected) > COUNT(*) / 2) + \\ORDER BY r.id ASC, rs.host + ; + var stmt: ?*c.sqlite3_stmt = null; + if (c.sqlite3_prepare_v2(self.db, sql, -1, &stmt, null) != c.SQLITE_OK) { + return self.sqlError(); + } + defer _ = c.sqlite3_finalize(stmt); + + _ = c.sqlite3_bind_text(stmt, 1, since.ptr, @intCast(since.len), SQLITE_STATIC); + _ = c.sqlite3_bind_text(stmt, 2, until.ptr, @intCast(until.len), SQLITE_STATIC); + + var points: std.ArrayList(TrendPoint) = .empty; + while (c.sqlite3_step(stmt) == c.SQLITE_ROW) { + const overlap = if (c.sqlite3_column_type(stmt, 5) == c.SQLITE_NULL) + try allocator.dupe(u8, "") + else + try self.colText(allocator, stmt, 5); + try points.append(allocator, .{ + .run_id = c.sqlite3_column_int64(stmt, 0), + .timestamp = try self.colText(allocator, stmt, 1), + .union_dids = c.sqlite3_column_int(stmt, 2), + .host = try self.colText(allocator, stmt, 3), + .unique_dids = c.sqlite3_column_int(stmt, 4), + .overlap = overlap, + }); + } + return points.toOwnedSlice(allocator); + } + // --- monitor support (/api/phi/monitors) --- pub const MonitorSnapshot = struct { @@ -723,6 +1105,14 @@ pub const Store = struct { return try allocator.dupe(u8, ptr[0..len]); } + fn colTextOpt(self: *Store, allocator: std.mem.Allocator, stmt: ?*c.sqlite3_stmt, col: c_int) !?[]const u8 { + _ = self; + if (c.sqlite3_column_type(stmt, col) == c.SQLITE_NULL) return null; + const ptr = c.sqlite3_column_text(stmt, col); + const len: usize = @intCast(c.sqlite3_column_bytes(stmt, col)); + return try allocator.dupe(u8, ptr[0..len]); + } + pub fn begin(self: *Store) !void { try self.exec("BEGIN"); } -- 2.51.2