diff --git a/README.md b/README.md index aff1faf..c6c5363 100644 --- a/README.md +++ b/README.md @@ -1,34 +1,24 @@ # zlay -an [AT Protocol](https://atproto.com/) relay in zig. crawls PDS hosts directly and rebroadcasts their firehose as a single aggregated stream. +an [AT Protocol](https://atproto.com/) relay in zig. subscribes to every PDS on the network, verifies commit signatures, and serves the merged event stream to downstream consumers via `com.atproto.sync.subscribeRepos`. -**live instance**: [zlay.waow.tech](https://zlay.waow.tech/_health) — [metrics dashboard](https://zlay-metrics.waow.tech) - -## what it does - -a relay subscribes to every PDS on the network, verifies commit signatures, and serves the merged event stream to downstream consumers via `com.atproto.sync.subscribeRepos`. it also maintains a collection index for `com.atproto.sync.listReposByCollection`. +**live instance**: [zlay.waow.tech](https://zlay.waow.tech/_health) — [metrics](https://zlay-metrics.waow.tech) ## design -- **direct PDS crawl** — no fan-out relay in between. the bootstrap relay (bsky.network) is called once at startup for the host list, then all data flows from each PDS. -- **optimistic validation** — on signing key cache miss, frames pass through immediately and the DID is queued for background resolution. >99.9% cache hit rate after warmup. -- **inline collection index** — RocksDB with two column families for bidirectional `(DID, collection)` lookups. no sidecar process. -- **one thread per PDS** — predictable memory, no GC. ~2,750 threads is fine; most are blocked on websocket reads. +- **direct PDS crawl** — the bootstrap relay (`bsky.network`) is called once at startup for the host list via `listHosts`, then all data flows directly from each PDS. no fan-out relay in between. -## dependencies +- **optimistic signature validation** — on signing key cache miss, the frame passes through immediately and the DID is queued for background resolution. the first commit from an unknown account is unvalidated; all subsequent commits are verified against the cached key. >99.9% cache hit rate after warmup. -| dependency | purpose | -|---|---| -| [zat](https://tangled.org/zzstoatzz.io/zat) | AT Protocol primitives (CBOR, CAR, signatures, DID resolution) | -| [websocket.zig](https://github.com/nicholasgasior/websocket.zig) | WebSocket client/server | -| [pg.zig](https://github.com/karlseguin/pg.zig) | PostgreSQL driver | -| [rocksdb-zig](https://github.com/Syndica/rocksdb-zig) | RocksDB bindings | +- **inline collection index** — indexes `(DID, collection)` pairs directly in the event processing pipeline using [RocksDB](https://rocksdb.org/) with two column families: `rbc` for collection-to-DID lookups and `cbr` for DID-to-collection cleanup. serves `listReposByCollection` from the relay process — no sidecar. the index design draws on [fig](https://tangled.org/microcosm.blue)'s work on [lightrail](https://tangled.org/microcosm.blue/lightrail), which uses adjacent keys from CAR slices to enumerate collections. + +- **one OS thread per PDS** — predictable memory, no garbage collector. ~2,750 threads is fine; most are blocked on websocket reads. thread stacks are set to 2 MB (zig's default is 16 MB). ## endpoints | endpoint | method | |---|---| -| `com.atproto.sync.subscribeRepos` | WebSocket (port 3000) | +| `com.atproto.sync.subscribeRepos` | WebSocket | | `com.atproto.sync.listRepos` | GET | | `com.atproto.sync.getRepoStatus` | GET | | `com.atproto.sync.getLatestCommit` | GET | @@ -36,6 +26,17 @@ a relay subscribes to every PDS on the network, verifies commit signatures, and | `com.atproto.sync.listHosts` | GET | | `com.atproto.sync.requestCrawl` | POST | +`getRepo` is not implemented — the relay does not serve full repository exports. + +## dependencies + +| dependency | purpose | +|---|---| +| [zat](https://tangled.org/zzstoatzz.io/zat) | AT Protocol primitives (CBOR, CAR, signatures, DID resolution) | +| [websocket.zig](https://github.com/nicholasgasior/websocket.zig) | WebSocket client/server | +| [pg.zig](https://github.com/karlseguin/pg.zig) | PostgreSQL driver | +| [rocksdb-zig](https://github.com/Syndica/rocksdb-zig) | [RocksDB](https://rocksdb.org/) bindings | + ## build requires zig 0.15 and a C/C++ toolchain (for RocksDB). @@ -46,13 +47,25 @@ zig build test # run tests zig build -Doptimize=ReleaseSafe # release build ``` +## configuration + +| variable | default | description | +|---|---|---| +| `RELAY_PORT` | `3000` | WebSocket firehose port | +| `RELAY_HTTP_PORT` | `3001` | HTTP API port | +| `RELAY_UPSTREAM` | `bsky.network` | bootstrap relay for initial host list | +| `RELAY_DATA_DIR` | `data/events` | event log storage | +| `RELAY_RETENTION_HOURS` | `72` | event retention window | +| `COLLECTION_INDEX_DIR` | `data/collection-index` | RocksDB collection index path | +| `DATABASE_URL` | — | PostgreSQL connection string | +| `RELAY_ADMIN_PASSWORD` | — | bearer token for admin endpoints | + see [docs/deployment.md](docs/deployment.md) for production deployment and [docs/backfill.md](docs/backfill.md) for collection index backfill. ## numbers | metric | value | |---|---| -| code | ~6,000 lines | | connected PDS hosts | ~2,750 | -| memory | ~2.9 GiB steady state | +| memory | ~2.9 GiB | | throughput | ~600 events/sec typical | diff --git a/src/event_log.zig b/src/event_log.zig index 5665a90..ff47dbe 100644 --- a/src/event_log.zig +++ b/src/event_log.zig @@ -233,11 +233,16 @@ pub const DiskPersist = struct { self.flush_thread = try std.Thread.spawn(.{ .stack_size = 2 * 1024 * 1024 }, flushLoop, .{self}); } + pub const UidResult = struct { + uid: u64, + host_changed: bool = false, + }; + /// resolve a DID to a numeric UID, associating with a host. /// on first encounter, creates account row with host_id. - /// on subsequent encounters from a different host, updates host_id. - /// Go relay: preProcessEvent → CreateAccountHost / EnsureAccountHost - pub fn uidForDidFromHost(self: *DiskPersist, did: []const u8, host_id: u64) !u64 { + /// on subsequent encounters from a different host, returns host_changed=true + /// so the caller can queue async DID migration validation. + pub fn uidForDidFromHost(self: *DiskPersist, did: []const u8, host_id: u64) !UidResult { const uid = try self.uidForDid(did); if (host_id > 0) { const current_host = self.getAccountHostId(uid) catch 0; @@ -245,14 +250,12 @@ pub const DiskPersist = struct { // first encounter: set host_id self.setAccountHostId(uid, host_id) catch {}; } else if (current_host != host_id) { - // host mismatch: account may have migrated - // Go relay re-resolves DID doc here; we log and update - // (full DID re-resolution for migration validation is TODO) - log.info("account {s} (uid={d}) host changed: {d} → {d}", .{ did, uid, current_host, host_id }); - self.setAccountHostId(uid, host_id) catch {}; + // host mismatch: don't update yet — caller should validate via DID resolution + log.info("account {s} (uid={d}) host mismatch: current={d} new={d}, queuing migration check", .{ did, uid, current_host, host_id }); + return .{ .uid = uid, .host_changed = true }; } } - return uid; + return .{ .uid = uid }; } /// resolve a DID to a numeric UID. creates a new account row on first encounter. @@ -440,6 +443,16 @@ pub const DiskPersist = struct { ); } + /// look up host ID by hostname. returns null if not found. + pub fn getHostIdForHostname(self: *DiskPersist, hostname: []const u8) !?u64 { + var row = (try self.db.rowUnsafe( + "SELECT id FROM host WHERE hostname = $1", + .{hostname}, + )) orelse return null; + defer row.deinit() catch {}; + return @intCast(row.get(i64, 0)); + } + /// update host status (active, blocked, exhausted) pub fn updateHostStatus(self: *DiskPersist, host_id: u64, status: []const u8) !void { _ = try self.db.exec( diff --git a/src/main.zig b/src/main.zig index eb26b91..7503289 100644 --- a/src/main.zig +++ b/src/main.zig @@ -12,6 +12,7 @@ //! /xrpc/com.atproto.sync.getLatestCommit — latest commit CID + rev //! /xrpc/com.atproto.sync.listReposByCollection — repos with records in a collection //! /xrpc/com.atproto.sync.listHosts — paginated active host listing +//! /xrpc/com.atproto.sync.getHostStatus — single host status //! /xrpc/com.atproto.sync.requestCrawl — request PDS crawl (POST) //! /admin/hosts — list all hosts (GET, admin) //! /admin/hosts/block — block a host (POST, admin) @@ -45,6 +46,7 @@ const HttpServer = struct { slurper: *slurper_mod.Slurper, collection_index: *collection_index_mod.CollectionIndex, backfiller: *backfill_mod.Backfiller, + bc: *broadcaster.Broadcaster, fn run(self: *HttpServer) void { while (!shutdown_flag.load(.acquire)) { @@ -53,7 +55,7 @@ const HttpServer = struct { log.debug("http accept error: {s}", .{@errorName(err)}); continue; }; - handleHttpConn(conn.stream, self.stats, self.persist, self.slurper, self.collection_index, self.backfiller); + handleHttpConn(conn.stream, self.stats, self.persist, self.slurper, self.collection_index, self.backfiller, self.bc); } } }; @@ -98,8 +100,9 @@ pub fn main() !void { // start flush thread try dp.start(); - // wire persist into broadcaster for cursor replay + // wire persist into broadcaster for cursor replay and validator for migration checks bc.persist = &dp; + val.persist = &dp; // init collection index (RocksDB — inspired by lightrail/microcosm.blue) const ci_dir = std.posix.getenv("COLLECTION_INDEX_DIR") orelse "data/collection-index"; @@ -145,6 +148,7 @@ pub fn main() !void { .slurper = &slurper, .collection_index = &ci, .backfiller = &backfiller, + .bc = &bc, }; const http_thread = try std.Thread.spawn(.{ .stack_size = default_stack_size }, HttpServer.run, .{&http_srv}); @@ -224,7 +228,7 @@ fn installSignalHandlers() void { std.posix.sigaction(std.posix.SIG.PIPE, &ignore_act, null); } -fn handleHttpConn(stream: std.net.Stream, stats: *broadcaster.Stats, persist: *event_log_mod.DiskPersist, slurper: *slurper_mod.Slurper, ci: *collection_index_mod.CollectionIndex, backfiller: *backfill_mod.Backfiller) void { +fn handleHttpConn(stream: std.net.Stream, stats: *broadcaster.Stats, persist: *event_log_mod.DiskPersist, slurper: *slurper_mod.Slurper, ci: *collection_index_mod.CollectionIndex, backfiller: *backfill_mod.Backfiller, bc: *broadcaster.Broadcaster) void { defer stream.close(); var recv_buf: [8192]u8 = undefined; @@ -244,7 +248,7 @@ fn handleHttpConn(stream: std.net.Stream, stats: *broadcaster.Stats, persist: *e if (request.head.method == .GET) { handleGet(&request, path, query, stats, persist, slurper, ci, backfiller); } else if (request.head.method == .POST) { - handlePost(&request, path, query, persist, slurper, backfiller); + handlePost(&request, path, query, persist, slurper, backfiller, bc); } else { respondText(&request, .method_not_allowed, "method not allowed"); } @@ -274,6 +278,8 @@ fn handleGet(request: *http.Server.Request, path: []const u8, query: []const u8, handleListReposByCollection(request, query, ci); } else if (std.mem.eql(u8, path, "/xrpc/com.atproto.sync.listHosts")) { handleListHosts(request, query, persist); + } else if (std.mem.eql(u8, path, "/xrpc/com.atproto.sync.getHostStatus")) { + handleGetHostStatus(request, query, persist); } else if (std.mem.eql(u8, path, "/admin/hosts")) { handleAdminListHosts(request, persist, slurper); } else if (std.mem.eql(u8, path, "/admin/backfill-collections")) { @@ -308,9 +314,9 @@ fn handleGet(request: *http.Server.Request, path: []const u8, query: []const u8, } } -fn handlePost(request: *http.Server.Request, path: []const u8, query: []const u8, persist: *event_log_mod.DiskPersist, slurper: *slurper_mod.Slurper, backfiller: *backfill_mod.Backfiller) void { +fn handlePost(request: *http.Server.Request, path: []const u8, query: []const u8, persist: *event_log_mod.DiskPersist, slurper: *slurper_mod.Slurper, backfiller: *backfill_mod.Backfiller, bc: *broadcaster.Broadcaster) void { if (std.mem.eql(u8, path, "/admin/repo/ban")) { - handleBan(request, persist); + handleBan(request, persist, bc); } else if (std.mem.eql(u8, path, "/xrpc/com.atproto.sync.requestCrawl")) { handleRequestCrawl(request, slurper); } else if (std.mem.eql(u8, path, "/admin/hosts/block")) { @@ -324,7 +330,7 @@ fn handlePost(request: *http.Server.Request, path: []const u8, query: []const u8 } } -fn handleBan(request: *http.Server.Request, persist: *event_log_mod.DiskPersist) void { +fn handleBan(request: *http.Server.Request, persist: *event_log_mod.DiskPersist, bc: *broadcaster.Broadcaster) void { if (!checkAdmin(request)) return; // read body (after checkAdmin which uses iterateHeaders) @@ -355,6 +361,18 @@ fn handleBan(request: *http.Server.Request, persist: *event_log_mod.DiskPersist) return; }; + // emit #account event so downstream consumers see the takedown + if (buildAccountFrame(persist.allocator, did)) |frame_bytes| { + if (persist.persist(.account, uid, frame_bytes)) |relay_seq| { + bc.stats.relay_seq.store(relay_seq, .release); + const broadcast_data = broadcaster.resequenceFrame(persist.allocator, frame_bytes, relay_seq) orelse frame_bytes; + bc.broadcast(relay_seq, broadcast_data); + log.info("admin: emitted #account takedown event for {s} (seq={d})", .{ did, relay_seq }); + } else |err| { + log.warn("admin: failed to persist #account takedown event: {s}", .{@errorName(err)}); + } + } + log.info("admin: banned {s} (uid={d})", .{ did, uid }); respondJson(request, .ok, "{\"success\":true}"); } @@ -897,6 +915,125 @@ fn handleListHosts(request: *http.Server.Request, query: []const u8, persist: *e respondJson(request, .ok, fbs.getWritten()); } +fn handleGetHostStatus(request: *http.Server.Request, query: []const u8, persist: *event_log_mod.DiskPersist) void { + var hostname_buf: [256]u8 = undefined; + const hostname = queryParamDecoded(query, "hostname", &hostname_buf) orelse { + respondJson(request, .bad_request, "{\"error\":\"InvalidRequest\",\"message\":\"hostname parameter required\"}"); + return; + }; + + // look up host + var row = (persist.db.rowUnsafe( + "SELECT id, hostname, status, last_seq FROM host WHERE hostname = $1", + .{hostname}, + ) catch { + respondJson(request, .internal_server_error, "{\"error\":\"DatabaseError\",\"message\":\"query failed\"}"); + return; + }) orelse { + respondJson(request, .bad_request, "{\"error\":\"HostNotFound\",\"message\":\"host not found\"}"); + return; + }; + defer row.deinit() catch {}; + + const host_id = row.get(i64, 0); + const host_name = row.get([]const u8, 1); + const raw_status = row.get([]const u8, 2); + const seq = row.get(i64, 3); + + // map internal status to lexicon hostStatus values + const status = if (std.mem.eql(u8, raw_status, "blocked")) + "banned" + else if (std.mem.eql(u8, raw_status, "exhausted")) + "offline" + else + raw_status; // active, idle pass through + + // count accounts on this host + const account_count: i64 = if (persist.db.rowUnsafe( + "SELECT COUNT(*) FROM account WHERE host_id = $1", + .{host_id}, + ) catch null) |cnt_row| blk: { + var r = cnt_row; + defer r.deinit() catch {}; + break :blk r.get(i64, 0); + } else 0; + + var buf: [4096]u8 = undefined; + var fbs = std.io.fixedBufferStream(&buf); + const w = fbs.writer(); + + w.writeAll("{\"hostname\":\"") catch return; + w.writeAll(host_name) catch return; + w.writeAll("\"") catch return; + std.fmt.format(w, ",\"seq\":{d},\"accountCount\":{d}", .{ seq, account_count }) catch return; + w.writeAll(",\"status\":\"") catch return; + w.writeAll(status) catch return; + w.writeAll("\"}") catch return; + + respondJson(request, .ok, fbs.getWritten()); +} + +/// build a CBOR #account frame for a takedown event. +/// header: {op: 1, t: "#account"}, payload: {seq: 0, did: "...", time: "...", active: false, status: "takendown"} +fn buildAccountFrame(allocator: std.mem.Allocator, did: []const u8) ?[]const u8 { + const zat = @import("zat"); + const cbor = zat.cbor; + + const header: cbor.Value = .{ .map = &.{ + .{ .key = "op", .value = .{ .unsigned = 1 } }, + .{ .key = "t", .value = .{ .text = "#account" } }, + } }; + + var time_buf: [24]u8 = undefined; + const time_str = formatTimestamp(&time_buf); + + const payload: cbor.Value = .{ .map = &.{ + .{ .key = "seq", .value = .{ .unsigned = 0 } }, + .{ .key = "did", .value = .{ .text = did } }, + .{ .key = "time", .value = .{ .text = time_str } }, + .{ .key = "active", .value = .{ .boolean = false } }, + .{ .key = "status", .value = .{ .text = "takendown" } }, + } }; + + const header_bytes = cbor.encodeAlloc(allocator, header) catch return null; + const payload_bytes = cbor.encodeAlloc(allocator, payload) catch { + allocator.free(header_bytes); + return null; + }; + + var frame = allocator.alloc(u8, header_bytes.len + payload_bytes.len) catch { + allocator.free(header_bytes); + allocator.free(payload_bytes); + return null; + }; + @memcpy(frame[0..header_bytes.len], header_bytes); + @memcpy(frame[header_bytes.len..], payload_bytes); + + allocator.free(header_bytes); + allocator.free(payload_bytes); + + return frame; +} + +/// format current UTC time as ISO 8601 (YYYY-MM-DDTHH:MM:SSZ) +fn formatTimestamp(buf: *[24]u8) []const u8 { + const ts: u64 = @intCast(std.time.timestamp()); + const es = std.time.epoch.EpochSeconds{ .secs = ts }; + const day = es.getEpochDay(); + const yd = day.calculateYearDay(); + const md = yd.calculateMonthDay(); + const ds = es.getDaySeconds(); + + return std.fmt.bufPrint(buf, "{d:0>4}-{d:0>2}-{d:0>2}T{d:0>2}:{d:0>2}:{d:0>2}Z", .{ + yd.year, + @as(u32, @intFromEnum(md.month)) + 1, + @as(u32, md.day_index) + 1, + ds.getHoursIntoDay(), + ds.getMinutesIntoHour(), + ds.getSecondsIntoMinute(), + }) catch "1970-01-01T00:00:00Z"; +} + // --- backfill handlers --- fn handleAdminBackfillTrigger(request: *http.Server.Request, query: []const u8, backfiller: *backfill_mod.Backfiller) void { diff --git a/src/subscriber.zig b/src/subscriber.zig index 077ffce..8a58fee 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -339,10 +339,13 @@ const FrameHandler = struct { // resolve DID → numeric UID for event header (host-aware) const uid: u64 = if (sub.persist) |dp| blk: { - break :blk if (did) |d| - dp.uidForDidFromHost(d, sub.options.host_id) catch 0 - else - 0; + if (did) |d| { + const result = dp.uidForDidFromHost(d, sub.options.host_id) catch break :blk @as(u64, 0); + if (result.host_changed) { + sub.validator.queueMigrationCheck(d, sub.options.host_id); + } + break :blk result.uid; + } else break :blk @as(u64, 0); } else 0; // process #account events: update upstream status diff --git a/src/validator.zig b/src/validator.zig index 93bbce3..26ce9ab 100644 --- a/src/validator.zig +++ b/src/validator.zig @@ -8,6 +8,7 @@ const std = @import("std"); const zat = @import("zat"); const broadcaster = @import("broadcaster.zig"); +const event_log_mod = @import("event_log.zig"); const Allocator = std.mem.Allocator; const log = std.log.scoped(.relay); @@ -39,15 +40,23 @@ pub const ValidatorConfig = struct { rev_clock_skew: i64 = 300, // 5 minutes }; +const MigrationCheck = struct { + did: []const u8, // duped, owned by validator + new_host_id: u64, +}; + pub const Validator = struct { allocator: Allocator, stats: *broadcaster.Stats, config: ValidatorConfig, + persist: ?*event_log_mod.DiskPersist = null, // DID → signing key cache (decoded, ready for verification) cache: std.StringHashMapUnmanaged(CachedKey) = .{}, cache_mutex: std.Thread.Mutex = .{}, // background resolve queue queue: std.ArrayListUnmanaged([]const u8) = .{}, + // migration validation queue + migration_queue: std.ArrayListUnmanaged(MigrationCheck) = .{}, queue_mutex: std.Thread.Mutex = .{}, queue_cond: std.Thread.Condition = .{}, resolver_threads: [max_resolver_threads]?std.Thread = .{null} ** max_resolver_threads, @@ -90,6 +99,12 @@ pub const Validator = struct { self.allocator.free(did); } self.queue.deinit(self.allocator); + + // free migration queue + for (self.migration_queue.items) |mc| { + self.allocator.free(mc.did); + } + self.migration_queue.deinit(self.allocator); } /// start background resolver threads @@ -383,14 +398,17 @@ pub const Validator = struct { while (self.alive.load(.acquire)) { var did: ?[]const u8 = null; + var migration: ?MigrationCheck = null; { self.queue_mutex.lock(); defer self.queue_mutex.unlock(); - while (self.queue.items.len == 0 and self.alive.load(.acquire)) { + while (self.queue.items.len == 0 and self.migration_queue.items.len == 0 and self.alive.load(.acquire)) { self.queue_cond.timedWait(&self.queue_mutex, 1 * std.time.ns_per_s) catch {}; } if (self.queue.items.len > 0) { did = self.queue.orderedRemove(0); + } else if (self.migration_queue.items.len > 0) { + migration = self.migration_queue.orderedRemove(0); } } @@ -434,10 +452,60 @@ pub const Validator = struct { self.cache.put(self.allocator, did_duped, cached) catch { self.allocator.free(did_duped); }; + } else if (migration) |mc| { + defer self.allocator.free(mc.did); + self.processMigrationCheck(&resolver, mc); } } } + /// validate a host migration by resolving the DID document and checking the PDS endpoint + fn processMigrationCheck(self: *Validator, resolver: *zat.DidResolver, mc: MigrationCheck) void { + const persist = self.persist orelse return; + + const parsed = zat.Did.parse(mc.did) orelse { + log.debug("migration check: invalid DID {s}", .{mc.did}); + return; + }; + + var doc = resolver.resolve(parsed) catch |err| { + log.debug("migration check: DID resolve failed for {s}: {s}", .{ mc.did, @errorName(err) }); + return; + }; + defer doc.deinit(); + + const pds_endpoint = doc.pdsEndpoint() orelse { + log.debug("migration check: no PDS endpoint for {s}", .{mc.did}); + return; + }; + + // extract hostname from PDS endpoint URL (strip https:// prefix) + const pds_host = extractHostFromUrl(pds_endpoint) orelse { + log.debug("migration check: cannot parse PDS URL '{s}' for {s}", .{ pds_endpoint, mc.did }); + return; + }; + + // look up the hostname → host_id + const resolved_host_id = (persist.getHostIdForHostname(pds_host) catch { + log.debug("migration check: host lookup failed for {s}", .{pds_host}); + return; + }) orelse { + log.debug("migration check: unknown host {s} for {s}", .{ pds_host, mc.did }); + return; + }; + + if (resolved_host_id == mc.new_host_id) { + // DID document confirms the new host — update + const uid = persist.uidForDid(mc.did) catch return; + persist.setAccountHostId(uid, mc.new_host_id) catch return; + log.info("migration validated: {s} → host {d} (confirmed by DID doc)", .{ mc.did, mc.new_host_id }); + } else { + log.warn("migration rejected: {s} claims host {d}, but DID doc says {s} (host {d})", .{ + mc.did, mc.new_host_id, pds_host, resolved_host_id, + }); + } + } + /// evict a DID's cached signing key (e.g. on #identity event). /// the next commit from this DID will trigger a fresh resolution. pub fn evictKey(self: *Validator, did: []const u8) void { @@ -448,6 +516,22 @@ pub const Validator = struct { } } + /// queue a DID for async migration validation (host change detected) + pub fn queueMigrationCheck(self: *Validator, did: []const u8, new_host_id: u64) void { + const duped = self.allocator.dupe(u8, did) catch return; + + self.queue_mutex.lock(); + defer self.queue_mutex.unlock(); + self.migration_queue.append(self.allocator, .{ + .did = duped, + .new_host_id = new_host_id, + }) catch { + self.allocator.free(duped); + return; + }; + self.queue_cond.signal(); + } + /// cache size (for diagnostics) pub fn cacheSize(self: *Validator) usize { self.cache_mutex.lock(); @@ -456,6 +540,27 @@ pub const Validator = struct { } }; +/// extract hostname from a URL like "https://pds.example.com" or "https://pds.example.com:443/path" +fn extractHostFromUrl(url: []const u8) ?[]const u8 { + // strip scheme + var rest = url; + if (std.mem.startsWith(u8, rest, "https://")) { + rest = rest["https://".len..]; + } else if (std.mem.startsWith(u8, rest, "http://")) { + rest = rest["http://".len..]; + } + // strip path + if (std.mem.indexOfScalar(u8, rest, '/')) |i| { + rest = rest[0..i]; + } + // strip port + if (std.mem.indexOfScalar(u8, rest, ':')) |i| { + rest = rest[0..i]; + } + if (rest.len == 0) return null; + return rest; +} + fn parseEnvInt(comptime T: type, key: []const u8, default: T) T { const val = std.posix.getenv(key) orelse return default; return std.fmt.parseInt(T, val, 10) catch default;