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",
\\
`;
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 += `
`;
- 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 += `
${tip('relay', 'the relay that missed these accounts')}
`;
- h += `
${tip('active', 'account has a working PDS but this relay didn\u2019t see it')}
`;
- h += `
${tip('unresolvable', 'DID lookup failed')}
`;
- h += `
${tip('deactivated', 'account is deactivated or deleted')}
`;
- h += `
${tip('total', 'total accounts missed by this relay')}
`;
+ h += `
${tip('relay', 'the relay whose stream missed these repos')}
`;
+ h += `
${tip('gap lower bound', 'checked repos with current host evidence; plus means unchecked rows may also belong here')}
`;
+ h += `
${tip('unverified', 'rows outside resolver budget or needing review')}
`;
+ h += `
${tip('no PDS', 'resolved DID document has no PDS endpoint')}
`;
+ h += `
${tip('checked', 'resolver-classified rows')}
`;
+ h += `
${tip('total', 'all repos this relay missed in the run')}