From 284008faceadbaa188d9acffc0d810622c2d4fed Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Thu, 9 Apr 2026 17:14:45 -0500 Subject: [PATCH] reset tree to b91382b for canary 1 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit forward-only rewind: every commit on main between b91382b and 4f3d1d4 has been superseded or is suspected of being implicated in the 2026-04-09 HTTP / delivery outage. rather than force-pushing history backward, this commit creates a new snapshot whose tree matches b91382b exactly, parented on 4f3d1d4. git pull --ff-only continues to work. the superseded commits remain in ancestry and can be referenced via the ops-changelog: - 4f3d1d4 gcLoop: disable malloc_trim, bump interval 10min→1h - 795cc41 host_authority: slot recovery + pool metrics + preload account count - bbba92c fix build: drop unused err1 capture in resolveHostAuthority - 584571a disable keep_alive on host authority resolver pool + log resolve errors - ee4e368 bump per-consumer buffer 8192→65536 + host_authority reject breakdown - 31825b2 subscriber: extract prepareFrameWork + add UAF regression test - 1eec324 fix UAF: dupe FrameWork.hostname per submit (will be re-applied on top) - 168d9f1 bump websocket.zig + zat: fix requestCrawl POST hang - fbdffbe mark DB success on did_cache hits - 3dc21b9 fix gcLoop: silently exited after one tick - e5f415f update README, CLAUDE.md, Dockerfile for current state this commit and the two following it (cherry-pick 1eec324 + pin zat alpha.21) constitute canary 1 per docs/zlay-canary-plan-2026-04-09.md. Co-Authored-By: Claude Opus 4.6 (1M context) --- CLAUDE.md | 15 +-- Dockerfile | 2 +- README.md | 20 ++- build.zig.zon | 8 +- docs/zlay-gcloop-stall-2026-04-09.md | 108 ----------------- src/broadcaster.zig | 77 +----------- src/event_log.zig | 28 +---- src/frame_worker.zig | 3 +- src/main.zig | 53 +++----- src/slurper.zig | 52 ++++---- src/subscriber.zig | 110 +++-------------- src/validator.zig | 174 +++++---------------------- 12 files changed, 107 insertions(+), 543 deletions(-) delete mode 100644 docs/zlay-gcloop-stall-2026-04-09.md diff --git a/CLAUDE.md b/CLAUDE.md index 32ef90d..c302e57 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -1,20 +1,18 @@ # zlay -AT Protocol relay in zig 0.16. Io.Evented backend (io_uring fibers, -~47 OS threads for ~2,800 PDS connections). ReleaseFast in production -(Evented + ReleaseSafe GPFs on startup — upstream zig bug, see -scripts/fiber_gpf_issue.md). +AT Protocol relay in zig 0.16. reader thread per PDS + shared frame +processing pool. ReleaseSafe in production. Io.Threaded backend (Evented +attempt shelved — see docs/evented-attempt.md). ## before pushing - `zig fmt --check .` and `zig build test` - MUST use `-Dtarget=x86_64-linux-gnu` for production (musl breaks RocksDB) -- MUST use `-Doptimize=ReleaseFast` — ReleaseSafe GPFs under Evented -- do NOT use Debug builds in production (2.5 GiB vs 1.5 GiB RSS) +- ReleaseFast has a known double-free — do not use ## deploy -configs at `../@zzstoatzz.io/relay/` — `just zlay publish-remote ReleaseFast` +configs at `../@zzstoatzz.io/relay/` — `just zlay publish-remote ReleaseSafe` KUBECONFIG is set automatically by the zlay module (`zlay/kubeconfig.yaml`). @@ -24,5 +22,4 @@ KUBECONFIG is set automatically by the zlay module (`zlay/kubeconfig.yaml`). - [docs/deployment.md](docs/deployment.md) — build flags, infra, resource usage - [docs/gotchas.md](docs/gotchas.md) — zig/pg.zig/rocksdb-zig/deploy traps - [docs/incident-2026-03-04.md](docs/incident-2026-03-04.md) — ReleaseSafe RSS analysis -- [docs/evented-attempt.md](docs/evented-attempt.md) — Evented backend migration story -- [scripts/fiber_gpf_issue.md](scripts/fiber_gpf_issue.md) — upstream zig bug report +- [docs/evented-attempt.md](docs/evented-attempt.md) — Evented backend attempt and why we reverted diff --git a/Dockerfile b/Dockerfile index 42ef3f1..3854d7a 100644 --- a/Dockerfile +++ b/Dockerfile @@ -22,7 +22,7 @@ COPY src/ src/ # contextSwitch, confirmed 2026-04-05). This is a zig codegen bug, not our code. # ReleaseFast avoids the bad optimization path. The previous production SIGSEGV # under ReleaseFast was a websocket handshake bug, now fixed (9ac64da). -RUN zig build -Doptimize=ReleaseFast -Dcpu=baseline -Dtarget=x86_64-linux-gnu +RUN zig build -Doptimize=ReleaseSafe -Dcpu=baseline -Dtarget=x86_64-linux-gnu FROM --platform=linux/amd64 debian:bookworm-slim RUN apt-get update && apt-get install -y --no-install-recommends ca-certificates && rm -rf /var/lib/apt/lists/* diff --git a/README.md b/README.md index 8d88023..571d8ec 100644 --- a/README.md +++ b/README.md @@ -8,14 +8,12 @@ an [AT Protocol](https://atproto.com/) relay in zig. subscribes to every PDS on - **direct PDS crawl** — the bootstrap relay is called once at startup for the host list via `listHosts`, then all data flows directly from each PDS. -- **Io.Evented backend** — uses zig 0.16's [`std.Io`](https://ziglang.org/documentation/master/std/#std.Io) with the [Evented](https://ziglang.org/documentation/master/std/#std.Io.Evented) backend (io_uring). ~2,800 PDS subscriber fibers run on ~47 OS threads — a 60x reduction from the 0.15 thread-per-PDS model. - -- **cross-Io architecture** — networking runs on Evented fibers, database access runs on Threaded workers. a lock-free [DbRequestQueue](https://tangled.org/zzstoatzz.io/zlay/blob/main/src/db_request_queue.zig) bridges the two using atomic spinlocks — no futex, no cross-Io boundary violations. see [devlog 008](https://zat.dev/#devlog/008-the-io-migration.md) for the full migration story. - - **optimistic signature validation** — on signing key cache miss, the frame passes through immediately and the DID is queued for background resolution. all subsequent commits are verified against the cached key. the cache caps at a configurable size and evicts the least recently used entry when full. - **inline collection index** — indexes `(DID, collection)` pairs in the event processing pipeline using RocksDB. 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). +- **reader thread per PDS + frame processing pool** — each PDS gets a lightweight reader thread (cursor tracking, rate limiting, header decode). heavy work (full CBOR decode, validation, DB persist, broadcast) runs on a shared pool of frame workers (configurable, default 16). + ## spec compliance implements the relay endpoints from the [AT Protocol sync spec](https://atproto.com/specs/sync): `subscribeRepos`, `listRepos`, `getRepoStatus`, `listHosts`, `getHostStatus`, and `requestCrawl`. @@ -26,23 +24,21 @@ also implements `getLatestCommit`, `listReposByCollection`, and `getRepo` (302 r | dependency | purpose | |---|---| -| [zat](https://tangled.org/zat.dev/zat) | AT Protocol primitives (CBOR, CAR, signatures, DID resolution) | -| [websocket.zig](https://github.com/zzstoatzz/websocket.zig) | WebSocket client/server (fork with write lock, HTTP fallback, TCP split fix) | +| [zat](https://tangled.org/zzstoatzz.io/zat) | AT Protocol primitives (CBOR, CAR, signatures, DID resolution) | +| [websocket.zig](https://github.com/zzstoatzz/websocket.zig) | WebSocket client/server (fork with HTTP fallback + TCP split fixes) | | [pg.zig](https://github.com/karlseguin/pg.zig) | PostgreSQL driver | | [rocksdb-zig](https://github.com/Syndica/rocksdb-zig) | RocksDB bindings | ## build -requires zig 0.16 and a C/C++ toolchain (for RocksDB). +requires zig 0.15 and a C/C++ toolchain (for RocksDB). ```bash -zig build # build (debug) -zig build test # run tests -zig build -Doptimize=ReleaseFast # release build (production) +zig build # build (debug) +zig build test # run tests +zig build -Doptimize=ReleaseSafe # release build (production default) ``` -note: `ReleaseSafe` GPFs on startup with the Evented backend due to a [zig stdlib bug](https://tangled.org/zzstoatzz.io/zlay/blob/main/scripts/fiber_gpf_issue.md) in `fiber.zig` context switching. use `ReleaseFast` for production. - ## configuration | variable | default | description | diff --git a/build.zig.zon b/build.zig.zon index 15e2ae9..b053db0 100644 --- a/build.zig.zon +++ b/build.zig.zon @@ -5,12 +5,12 @@ .minimum_zig_version = "0.16.0", .dependencies = .{ .zat = .{ - .url = "https://tangled.org/zat.dev/zat/archive/v0.3.0-alpha.23.tar.gz", - .hash = "zat-0.3.0-alpha.23-5PuC7k1VCACPkCoMnIstIbVu1yIrVP5Yx3l0G-lZ2Qoa", + .url = "https://tangled.org/zat.dev/zat/archive/v0.3.0-alpha.17.tar.gz", + .hash = "zat-0.3.0-alpha.15-5PuC7nVhBQBnyDEw50Zwitd8ujG7mJumQlVtqYh08QML", }, .websocket = .{ - .url = "https://github.com/zzstoatzz/websocket.zig/archive/3c6794a.tar.gz", - .hash = "websocket-0.1.0-ZPISdYv9AwB5YM5SQKd7B9kRGNcQ4O8f68yfYsCMfU5h", + .url = "https://github.com/zzstoatzz/websocket.zig/archive/9ac64da.tar.gz", + .hash = "websocket-0.1.0-ZPISdebwAwAGqC8MRLe7MFedo_K_Yy2mjYVAch4Ya4hZ", }, .pg = .{ .url = "git+https://github.com/zzstoatzz/pg.zig?ref=dev#5ce2355b1d851075523709c7d3068dcdb0224322", diff --git a/docs/zlay-gcloop-stall-2026-04-09.md b/docs/zlay-gcloop-stall-2026-04-09.md deleted file mode 100644 index 0a7c0f7..0000000 --- a/docs/zlay-gcloop-stall-2026-04-09.md +++ /dev/null @@ -1,108 +0,0 @@ -# zlay gcLoop stall — 2026-04-09 - -## tl;dr - -zlay pods were flapping on a ~10 minute cadence: ~10 min healthy, then /metrics -and /_readyz stop responding, kubelet marks NotReady, pod eventually restarts, -cycle repeats. Initial hypothesis was broadcaster writeLoop starvation -(see [zlay-broadcaster-starvation-2026-04-09.md] — now superseded as primary -cause). The actual cadence lines up precisely with `gcLoop` (main.zig), which -became functional again on 2026-04-06 via commit 3dc21b9 ("fix gcLoop: silently -exited after one tick"). Prior to that fix, gc was silently dead after one tick -per pod, masking the underlying problem. - -## causal chain (suspected) - -1. `gcLoop` fires every 10 minutes (main.zig). -2. `dp.gc()` holds `DiskPersist.mutex` for its entire duration - (event_log.zig:977-1033). That critical section includes: - - `SELECT` of expired log_file_refs - - per-file `DELETE FROM log_file_refs` + `unlink` - - `gcBySize()` — another pass of queries + `stat` + `unlink` -3. `DiskPersist.persist()` (event_log.zig:864-899) takes the same mutex on the - frame-worker hot path. For the duration of gc, every frame worker blocks on - persist → the broadcast queue dries up → consumers see ~0 events/sec. This - alone explains the "zlay only delivers 0.035 events/sec" symptom that was - previously blamed on writeLoop polling. -4. After `dp.gc()` returns, `gcLoop` called `malloc_trim(0)`. The pod runs with - `MALLOC_ARENA_MAX=4`, so glibc holds per-arena locks while walking free - lists. On a ~1.5 GiB RSS process this can stall every allocator user for - seconds. The Evented fiber serving /metrics and /_readyz would stall on its - next malloc → kubelet liveness/readiness probes time out → pod marked - NotReady → restart. - -Two separable suspects (important for isolation): -- **`dp.gc()` mutex hold** freezes ingest via the persist hot path. -- **`malloc_trim(0)`** freezes everything via arena locks. - -Either alone is sufficient to flunk probes. Don't conflate them when -investigating. - -## stabilization shipped - -Commit on top of 795cc41: - -- **Disabled `malloc_trim(0)`** in `gcLoop` (main.zig:502). Comment preserved - so future maintainers know why. If RSS growth becomes an issue, prefer - tuning `MALLOC_MMAP_THRESHOLD_` or running trim out-of-band. -- **Bumped gc interval from 10 min → 1 hour** (main.zig:473). Reduces - frequency and blast radius of the persist-mutex hold while the real fix - (mutex narrowing) is prepared. Not a cure — a stall at hour boundaries is - still possible if gc runs long. -- **Added timing log** around `dp.gc()` using `clock_gettime(.MONOTONIC)` - (plain thread — no Io available). Next incident will tell us how long - `dp.gc()` actually runs on a production dataset. - -These changes are in main.zig only. No dependency or schema changes. Safe to -revert by undoing the single commit. - -## validation plan post-deploy - -After deploying on top of 795cc41: - -1. Pod uptime should exceed 10 minutes with no NotReady flap. -2. Grep logs for `gc: dp.gc complete in` — verify gc runs and record its - duration on production data. -3. `frames_broadcast_total` should climb at ingest rate (~300/sec), not - the previously measured ~0.035/sec. -4. `tap run --relay-url https://zlay.waow.tech` for 60 seconds should - deliver thousands of events, not single digits. - -If the pod still flaps after this change, the hypothesis is wrong or -incomplete — do **not** proceed to the follow-up fixes until we re-diagnose. -Most likely remaining suspect in that case is the `dp.gc()` mutex hold on -a dataset big enough that even hourly gc takes long enough to trip probes. -Mitigation in that case: temporarily comment out the `dp.gc()` call body to -isolate. - -## follow-up work (not in this change) - -1. **Narrow `DiskPersist.mutex` scope in `gc()`**. The mutex protects - `evtbuf`/`outbuf`/`cur_seq`/`current_file_path`/`flushLocked`. Nothing in - gc's DB iteration or per-file unlink genuinely needs that lock. The only - shared state is a read of `current_file_path` to skip the active file. - Plan: do DB discovery and file discovery without the lock; acquire briefly - only to re-check `current_file_path` against each candidate before unlink. - Same treatment for `gcBySize()` and `takeDownUser()`. -2. **Broadcaster writeLoop polling** (broadcaster.zig:447-453). Real bug — - `cond.signal` at line 413 is a no-op because writeLoop polls with - `io.sleep(100ms)` instead of `cond.wait`. This caps per-consumer drain at - ~10/sec even under zero contention. Fix: use `cond.wait`, schedule pings - via a separate timer fiber or piggyback on next wakeup. - **Do NOT** move writeLoop off Evented to pool_io — commit 6674812 documents - that cross-Io path crashes via `Thread.current()` NULL deref. -3. **Consider whether `malloc_trim` should ever run on-process**. For a - steady-state relay, the answer is probably no; mmap threshold tuning is a - better choice. - -## code pointers - -- main.zig:467-510 — `gcLoop` -- event_log.zig:977-1033 — `DiskPersist.gc` -- event_log.zig:1036-1099 — `DiskPersist.gcBySize` -- event_log.zig:864-899 — `DiskPersist.persist` (hot path, same mutex) -- broadcaster.zig:439-477 — `Consumer.writeLoop` (secondary bug, not fixed here) -- commit 3dc21b9 — "fix gcLoop: silently exited after one tick" (the fix that - unmasked this) -- commit 6674812 — "fix SIGSEGV: plain threads calling Evented Io.Mutex" - (cross-Io landmine, read before touching Consumer) diff --git a/src/broadcaster.zig b/src/broadcaster.zig index 95210e1..1035825 100644 --- a/src/broadcaster.zig +++ b/src/broadcaster.zig @@ -55,31 +55,6 @@ pub const Stats = struct { host_authority_is_new: std.atomic.Value(u64) = .{ .raw = 0 }, host_authority_host_changed: std.atomic.Value(u64) = .{ .raw = 0 }, host_authority_time_us: std.atomic.Value(u64) = .{ .raw = 0 }, - // host_authority resolver pool mechanics (added 2026-04-09 per external - // review). these expose pool contention and the slot-recovery code path - // that previous reject-branch counters couldn't see. see relay - // docs/zlay-external-review-2026-04-09.md. - host_resolver_acquire_wait_us_total: std.atomic.Value(u64) = .{ .raw = 0 }, - host_resolver_in_use: std.atomic.Value(u32) = .{ .raw = 0 }, - host_resolver_resets_total: std.atomic.Value(u64) = .{ .raw = 0 }, - host_resolver_resolve_fail_total: std.atomic.Value(u64) = .{ .raw = 0 }, - // background DID resolveLoop ok/fail. previously these failures were - // log.debug + continue with no observability — we treated the loop as - // "working" without ever measuring it. baseline the rate so we can - // tell if it's silently degraded. - resolve_loop_resolve_ok_total: std.atomic.Value(u64) = .{ .raw = 0 }, - resolve_loop_resolve_fail_total: std.atomic.Value(u64) = .{ .raw = 0 }, - // per-branch reject breakdown (subsets of failed_host_authority). - // added 2026-04-08 to diagnose the 100% host_authority failure rate — - // without this breakdown we can't tell whether the DID doc lookup is - // failing, the endpoint is unparseable, the host isn't in our table, - // or the resolved host genuinely differs from the incoming host. - host_authority_reject_parse_did: std.atomic.Value(u64) = .{ .raw = 0 }, - host_authority_reject_resolve: std.atomic.Value(u64) = .{ .raw = 0 }, - host_authority_reject_no_endpoint: std.atomic.Value(u64) = .{ .raw = 0 }, - host_authority_reject_bad_url: std.atomic.Value(u64) = .{ .raw = 0 }, - host_authority_reject_unknown_host: std.atomic.Value(u64) = .{ .raw = 0 }, - host_authority_reject_host_mismatch: std.atomic.Value(u64) = .{ .raw = 0 }, // frame pool memory pressure pool_queued_bytes: std.atomic.Value(u64) = .{ .raw = 0 }, // persist/broadcast pipeline contention @@ -378,12 +353,7 @@ const ping_interval_ns: u64 = 30 * std.time.ns_per_s; const ping_timeout_ns: u64 = 5 * std.time.ns_per_s; pub const Consumer = struct { - // per-consumer ring buffer. sized to absorb transient write stalls - // without kicking the consumer with ConsumerTooSlow. at steady-state - // ~250 events/sec, 65536 entries = ~4.4 minutes of headroom. - // previously 8192 (~33s) which was short enough that pulsar's 60-min - // snapshot run accumulated repeated kicks (see ops_changelog 2026-04-01). - const BUFFER_CAP = 65536; + const BUFFER_CAP = 8192; conn: *websocket.Conn, allocator: Allocator, @@ -1079,42 +1049,12 @@ pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, attrib \\# HELP relay_broadcast_no_consumers_total frames skipped broadcast (no consumers) \\relay_broadcast_no_consumers_total {d} \\ - \\# TYPE relay_host_resolver_acquire_wait_us_total counter - \\# HELP relay_host_resolver_acquire_wait_us_total cumulative microseconds spent spinning in acquireHostResolver - \\relay_host_resolver_acquire_wait_us_total {d} - \\ - \\# TYPE relay_host_resolver_in_use gauge - \\# HELP relay_host_resolver_in_use host_authority resolver pool slots currently held by callers - \\relay_host_resolver_in_use {d} - \\ - \\# TYPE relay_host_resolver_resets_total counter - \\# HELP relay_host_resolver_resets_total slot deinit+reinit on first-attempt resolve failure (slot recovery path) - \\relay_host_resolver_resets_total {d} - \\ - \\# TYPE relay_host_resolver_resolve_fail_total counter - \\# HELP relay_host_resolver_resolve_fail_total first-attempt resolve failures in the host_authority pool (before recovery/retry) - \\relay_host_resolver_resolve_fail_total {d} - \\ - \\# TYPE relay_resolve_loop_resolve_ok_total counter - \\# HELP relay_resolve_loop_resolve_ok_total successful resolves in the background signing-key resolveLoop - \\relay_resolve_loop_resolve_ok_total {d} - \\ - \\# TYPE relay_resolve_loop_resolve_fail_total counter - \\# HELP relay_resolve_loop_resolve_fail_total failed resolves in the background signing-key resolveLoop - \\relay_resolve_loop_resolve_fail_total {d} - \\ , .{ stats.persist_order_spins.load(.acquire), stats.broadcast_queue_push_lock_spins.load(.acquire), stats.broadcast_queue_full.load(.acquire), stats.broadcast_queue_depth_hwm.load(.acquire), stats.broadcast_no_consumers.load(.acquire), - stats.host_resolver_acquire_wait_us_total.load(.acquire), - stats.host_resolver_in_use.load(.acquire), - stats.host_resolver_resets_total.load(.acquire), - stats.host_resolver_resolve_fail_total.load(.acquire), - stats.resolve_loop_resolve_ok_total.load(.acquire), - stats.resolve_loop_resolve_fail_total.load(.acquire), }) catch {}; // validation failure breakdown by reason @@ -1129,15 +1069,6 @@ pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, attrib \\relay_validation_failed{{reason="host_authority"}} {d} \\relay_validation_failed{{reason="future_rev"}} {d} \\ - \\# TYPE relay_host_authority_reject counter - \\# HELP relay_host_authority_reject host authority reject breakdown by branch - \\relay_host_authority_reject{{branch="parse_did"}} {d} - \\relay_host_authority_reject{{branch="resolve"}} {d} - \\relay_host_authority_reject{{branch="no_endpoint"}} {d} - \\relay_host_authority_reject{{branch="bad_url"}} {d} - \\relay_host_authority_reject{{branch="unknown_host"}} {d} - \\relay_host_authority_reject{{branch="host_mismatch"}} {d} - \\ , .{ stats.failed_bad_did.load(.acquire), stats.failed_bad_rev.load(.acquire), @@ -1146,12 +1077,6 @@ pub fn formatPrometheusMetrics(stats: *const Stats, cache_entries: usize, attrib stats.failed_bad_structure.load(.acquire), stats.failed_host_authority.load(.acquire), stats.failed_future_rev.load(.acquire), - stats.host_authority_reject_parse_did.load(.acquire), - stats.host_authority_reject_resolve.load(.acquire), - stats.host_authority_reject_no_endpoint.load(.acquire), - stats.host_authority_reject_bad_url.load(.acquire), - stats.host_authority_reject_unknown_host.load(.acquire), - stats.host_authority_reject_host_mismatch.load(.acquire), }) catch return w.buffered(); // memory attribution — internal capacities help identify what's consuming RSS diff --git a/src/event_log.zig b/src/event_log.zig index 861e4a8..7e65814 100644 --- a/src/event_log.zig +++ b/src/event_log.zig @@ -432,15 +432,8 @@ pub const DiskPersist = struct { /// resolve a DID to a numeric UID. creates a new account row on first encounter. /// matches indigo's Relay.DidToUid → Account.UID mapping. pub fn uidForDid(self: *DiskPersist, did: []const u8) !u64 { - // fast path: check in-memory cache. mark DB success even on hits — - // a hot cache means the DB-derived data path is functioning, which - // is what /_readyz cares about. without this, the 30s health window - // depends entirely on cache misses + the 10-min GC tick, which can - // gap during steady-state and trip k8s liveness probes. - if (self.did_cache.get(did)) |uid| { - self.markDbSuccess(); - return uid; - } + // fast path: check in-memory cache + if (self.did_cache.get(did)) |uid| return uid; // check database if (try self.db.rowUnsafe( @@ -622,10 +615,6 @@ pub const DiskPersist = struct { last_seq: u64, failed_attempts: u32, account_limit: ?u64 = null, - // computed by listActiveHosts at load time so cold-start spawn doesn't - // need a per-host DbRequest round-trip. equals account_limit when set, - // otherwise COUNT(account.uid) for this host. - effective_account_count: u64 = 0, }; const HostResult = struct { id: u64, last_seq: u64 }; @@ -722,23 +711,13 @@ pub const DiskPersist = struct { hosts.deinit(allocator); } - // batch the effective_account_count into the host load — same JOIN/COUNT - // shape as getEffectiveAccountCountImpl but folded into one query so - // spawnWorker doesn't need a per-host DbRequest round-trip during - // cold-start. h.id is the primary key so other h.* columns are - // functionally dependent for the GROUP BY. var result = try db.query( - "SELECT h.id, h.hostname, h.status, h.last_seq, h.failed_attempts, h.account_limit, " ++ - "COALESCE(h.account_limit, COUNT(a.uid)) AS effective_account_count " ++ - "FROM host h LEFT JOIN account a ON a.host_id = h.id " ++ - "WHERE h.status = 'active' " ++ - "GROUP BY h.id ORDER BY h.id ASC", + "SELECT id, hostname, status, last_seq, failed_attempts, account_limit FROM host WHERE status = 'active' ORDER BY id ASC", .{}, ); defer result.deinit(); while (result.nextUnsafe() catch null) |row| { - const eff_count_i64 = row.get(i64, 6); try hosts.append(allocator, .{ .id = @intCast(row.get(i64, 0)), .hostname = try allocator.dupe(u8, row.get([]const u8, 1)), @@ -746,7 +725,6 @@ pub const DiskPersist = struct { .last_seq = @intCast(row.get(i64, 3)), .failed_attempts = @intCast(row.get(i32, 4)), .account_limit = if (row.get(?i64, 5)) |v| @as(?u64, @intCast(v)) else null, - .effective_account_count = if (eff_count_i64 > 0) @intCast(eff_count_i64) else 0, }); } diff --git a/src/frame_worker.zig b/src/frame_worker.zig index 76f44a1..ff41796 100644 --- a/src/frame_worker.zig +++ b/src/frame_worker.zig @@ -27,7 +27,7 @@ fn microTimestamp(io: Io) i64 { pub const FrameWork = struct { data: []u8, // raw frame bytes (heap-duped by reader, freed by worker) host_id: u64, - hostname: []const u8, // owned (heap-duped at submit, freed by worker) + hostname: []const u8, // borrowed from subscriber (stable lifetime) allocator: Allocator, io: Io, // shared references (all thread-safe, all outlive the work item) @@ -41,7 +41,6 @@ pub const FrameWork = struct { pub fn processFrame(work: *FrameWork) void { _ = work.bc.stats.pool_queued_bytes.fetchSub(work.data.len, .monotonic); defer work.allocator.free(work.data); - defer work.allocator.free(work.hostname); var arena = std.heap.ArenaAllocator.init(work.allocator); defer arena.deinit(); diff --git a/src/main.zig b/src/main.zig index 52f509f..960a725 100644 --- a/src/main.zig +++ b/src/main.zig @@ -351,9 +351,7 @@ pub fn main() !void { // start GC loop on a plain thread — dp.gc() uses pool_io (Threaded) mutex // and pg.Pool. MUST NOT run as Evented fiber: Threaded futex on Evented // fiber dereferences NULL Thread.current() threadlocal → heap corruption. - // sleeps via std.Thread.sleep (NOT io.sleep) — io.sleep on pool_io from - // a non-Io thread fails on the second tick and silently kills the loop. - const gc_thread = std.Thread.spawn(.{}, gcLoop, .{&dp}) catch |err| { + const gc_thread = std.Thread.spawn(.{}, gcLoop, .{ &dp, pool_io }) catch |err| { log.err("failed to start GC thread: {s}", .{@errorName(err)}); return err; }; @@ -464,48 +462,27 @@ fn runWsServer(server: *websocket.Server(broadcaster.Handler), listener: *Io.net server.runIo(listener, bc); } -fn gcLoop(dp: *event_log_mod.DiskPersist) void { - // gc cadence: 1 hour (was 10 min before 2026-04-09 incident). DiskPersist.gc() - // currently holds DiskPersist.mutex for its entire duration — covering DB - // iteration and per-file unlinks — which blocks persist() on every frame - // worker. until the mutex scope is narrowed (follow-up), run hourly to - // bound blast radius. see docs/zlay-gcloop-stall-2026-04-09.md. - const gc_interval_s: u64 = 60 * 60; // 1 hour +fn gcLoop(dp: *event_log_mod.DiskPersist, io: Io) void { + const gc_interval: u64 = 10 * 60; // 10 minutes in seconds while (!shutdown_flag.load(.acquire)) { - // sleep in 1s ticks so shutdown is checked frequently. uses - // std.c.nanosleep directly — this is a plain OS thread, so calling - // io.sleep on pool_io would fail on the second tick and silently - // exit the loop via `catch return`. zig 0.16 has no std.Thread.sleep - // and std.posix.nanosleep was removed during the Io migration. - var elapsed: u64 = 0; - while (elapsed < gc_interval_s and !shutdown_flag.load(.acquire)) { - const ts: std.c.timespec = .{ .sec = 1, .nsec = 0 }; - _ = std.c.nanosleep(&ts, null); - elapsed += 1; + // sleep in small increments to check shutdown + var remaining: u64 = gc_interval; + while (remaining > 0 and !shutdown_flag.load(.acquire)) { + const chunk = @min(remaining, 1); + io.sleep(Io.Duration.fromSeconds(@intCast(chunk)), .awake) catch return; + remaining -= chunk; } if (shutdown_flag.load(.acquire)) return; - // time the gc call so the next incident log tells us whether gc itself - // or something around it is the stall. plain-thread context — use - // clock_gettime directly rather than Io.Timestamp. - var ts_start: std.c.timespec = undefined; - _ = std.c.clock_gettime(.MONOTONIC, &ts_start); dp.gc() catch |err| { log.warn("event log GC failed: {s}", .{@errorName(err)}); }; - var ts_end: std.c.timespec = undefined; - _ = std.c.clock_gettime(.MONOTONIC, &ts_end); - const elapsed_ns: i64 = (@as(i64, ts_end.sec) - @as(i64, ts_start.sec)) * std.time.ns_per_s + - (@as(i64, ts_end.nsec) - @as(i64, ts_start.nsec)); - log.info("gc: dp.gc complete in {d}ms", .{@divTrunc(elapsed_ns, std.time.ns_per_ms)}); - - // NOTE: malloc_trim(0) disabled 2026-04-09. with MALLOC_ARENA_MAX=4 it - // walks free lists holding per-arena locks, which on a ~1.5 GiB RSS - // process can stall every allocator user (including the Evented http - // fiber serving /metrics and /_readyz) long enough to flunk liveness - // probes. if memory reclamation becomes an issue, prefer tuning - // MALLOC_MMAP_THRESHOLD_ or calling trim from a dedicated maintenance - // window rather than on a hot process. + + // return freed pages to OS (glibc-specific, no-op on other platforms) + if (comptime malloc_trim) |trim| { + _ = trim(0); + log.info("gc: malloc_trim complete", .{}); + } } } diff --git a/src/slurper.zig b/src/slurper.zig index 8c8fee6..964de23 100644 --- a/src/slurper.zig +++ b/src/slurper.zig @@ -526,40 +526,38 @@ pub const Slurper = struct { db_queue.push(&reset_req.base); reset_req.base.wait(self.io, self.shutdown); - // phase 5: fetch effective account count (single DbRequest, addHost is - // a rare one-off path — not the cold-start hot loop). cold-start uses - // listActiveHosts which preloads this in the batch query. - var account_count: u64 = 0; - const GetCountReq = struct { - base: event_log_mod.DbRequest = .{ .callback = &execute }, - hid: u64, - count: u64 = 0, - - fn execute(b: *event_log_mod.DbRequest, dp: *event_log_mod.DiskPersist) void { - const s: *@This() = @fieldParentPtr("base", b); - s.count = dp.getEffectiveAccountCount(s.hid); - } - }; - var count_req: GetCountReq = .{ .hid = db_req.host_id }; - db_queue.push(&count_req.base); - count_req.base.wait(self.io, self.shutdown); - account_count = count_req.count; - - // phase 6: spawn worker (Evented) - try self.spawnWorker(db_req.host_id, hostname, db_req.last_seq, account_count); + // phase 5: spawn worker (Evented) + try self.spawnWorker(db_req.host_id, hostname, db_req.last_seq); log.info("added host {s} (id={d})", .{ hostname, db_req.host_id }); } - /// spawn a subscriber thread for a host. callers must pass the - /// effective_account_count up front — see comment on the addHost path - /// below for the one site that still computes it inline. - fn spawnWorker(self: *Slurper, host_id: u64, hostname: []const u8, last_seq: u64, account_count: u64) !void { + /// spawn a subscriber thread for a host + fn spawnWorker(self: *Slurper, host_id: u64, hostname: []const u8, last_seq: u64) !void { const hostname_duped = try self.allocator.dupe(u8, hostname); errdefer self.allocator.free(hostname_duped); const sub = try self.allocator.create(subscriber_mod.Subscriber); errdefer self.allocator.destroy(sub); + // get effective account count via DbRequestQueue + var account_count: u64 = 0; + if (self.db_queue) |db_queue| { + const GetCountReq = struct { + base: event_log_mod.DbRequest = .{ .callback = &execute }, + hid: u64, + count: u64 = 0, + + fn execute(b: *event_log_mod.DbRequest, dp: *event_log_mod.DiskPersist) void { + const s: *@This() = @fieldParentPtr("base", b); + s.count = dp.getEffectiveAccountCount(s.hid); + } + }; + var count_req: GetCountReq = .{ .hid = host_id }; + db_queue.push(&count_req.base); + count_req.base.wait(self.io, self.shutdown); + account_count = count_req.count; + } + sub.* = subscriber_mod.Subscriber.init( self.allocator, self.io, @@ -672,9 +670,7 @@ pub const Slurper = struct { var spawned: usize = 0; for (hosts) |host| { if (self.shutdown.load(.acquire)) break; - // host.effective_account_count was preloaded by listActiveHostsImpl - // — no per-host DbRequest round-trip during cold start. - self.spawnWorker(host.id, host.hostname, host.last_seq, host.effective_account_count) catch |err| { + self.spawnWorker(host.id, host.hostname, host.last_seq) catch |err| { log.warn("failed to spawn worker for {s}: {s}", .{ host.hostname, @errorName(err) }); }; spawned += 1; diff --git a/src/subscriber.zig b/src/subscriber.zig index e49675d..51ed362 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -256,35 +256,6 @@ pub const Subscriber = struct { return self.shutdown.load(.acquire) or self.host_shutdown.load(.acquire); } - /// build a FrameWork that owns its data + hostname, decoupled from the - /// subscriber lifetime. returns null on OOM. - /// - /// UAF-safe contract: after this returns, the caller (or the slurper) - /// may free `self.options.hostname` and the input `data` immediately - /// without affecting the work item. the worker will free both through - /// `work.allocator` when it finishes processing. - /// - /// see `test "prepareFrameWork dupes hostname and data (UAF regression)"` - fn prepareFrameWork(self: *Subscriber, data: []const u8) ?frame_worker_mod.FrameWork { - const duped = self.allocator.dupe(u8, data) catch return null; - const hostname_dup = self.allocator.dupe(u8, self.options.hostname) catch { - self.allocator.free(duped); - return null; - }; - return .{ - .data = duped, - .host_id = self.options.host_id, - .hostname = hostname_dup, - .allocator = self.allocator, - .io = self.pool_io orelse self.io, - .bc = self.bc, - .validator = self.validator, - .persist = self.persist, - .collection_index = self.collection_index, - .resyncer = self.resyncer, - }; - } - /// run the subscriber loop. reconnects with exponential backoff. /// blocks until shutdown flag is set or host is exhausted. pub fn run(self: *Subscriber) void { @@ -506,23 +477,29 @@ const FrameHandler = struct { const d = payload.getString("repo") orelse payload.getString("did"); break :blk if (d) |s| std.hash.Wyhash.hash(0, s) else sub.options.host_id; }; - // dupe data + hostname per-frame: subscriber teardown (slurper.runWorker) - // frees sub.options.hostname after sub.run() returns, but FrameWorks - // can still be queued in the pool. borrowing the slice would be a - // use-after-free (corrupt hostnames in chain-break logs, etc.). - const work = sub.prepareFrameWork(data) orelse return; + const duped = sub.allocator.dupe(u8, data) catch return; const t0 = nanoTimestamp(io); - if (pool.submit(did_key, work, sub.shutdown)) { + if (pool.submit(did_key, .{ + .data = duped, + .host_id = sub.options.host_id, + .hostname = sub.options.hostname, + .allocator = sub.allocator, + .io = sub.pool_io orelse sub.io, // pool_io (Threaded) for worker-safe ops + .bc = sub.bc, + .validator = sub.validator, + .persist = sub.persist, + .collection_index = sub.collection_index, + .resyncer = sub.resyncer, + }, sub.shutdown)) { // pool accepted — advance cursor past this frame - _ = sub.bc.stats.pool_queued_bytes.fetchAdd(work.data.len, .monotonic); + _ = sub.bc.stats.pool_queued_bytes.fetchAdd(duped.len, .monotonic); if (upstream_seq) |s| sub.last_upstream_seq = s; if (nanoTimestamp(io) - t0 > 1_000_000) { // >1ms = had to wait _ = sub.bc.stats.pool_backpressure.fetchAdd(1, .monotonic); } } else { // shutdown requested — don't advance cursor so reconnect replays this frame - sub.allocator.free(work.data); - sub.allocator.free(work.hostname); + sub.allocator.free(duped); } return; } @@ -793,63 +770,6 @@ const FrameHandler = struct { // --- tests --- -test "prepareFrameWork dupes hostname and data (UAF regression)" { - // regression test for UAF: slurper.runWorker frees sub.options.hostname - // after sub.run() returns, but FrameWorks can still be queued in the - // frame pool with that hostname slice. prepareFrameWork must heap-dupe - // both `data` and `hostname` so the work item is independent of the - // subscriber's lifetime. - // - // the test simulates subscriber teardown by freeing the source hostname - // after the FrameWork is built, then asserts the work item is still - // intact (would trip the testing allocator's use-after-free detection - // if the dupe was skipped). - const alloc = std.testing.allocator; - - const orig_hostname = try alloc.dupe(u8, "example.pds.host"); - const data_input = try alloc.dupe(u8, "raw frame bytes"); - defer alloc.free(data_input); - - var shutdown: std.atomic.Value(bool) = .{ .raw = false }; - var sub: Subscriber = .{ - .allocator = alloc, - .io = std.testing.io, - .options = .{ .hostname = orig_hostname, .host_id = 42 }, - .bc = undefined, // not dereferenced by prepareFrameWork - .validator = undefined, - .persist = null, - .collection_index = null, - .resyncer = null, - .pool = null, - .pool_io = null, - .shutdown = &shutdown, - }; - - const work = sub.prepareFrameWork(data_input) orelse return error.OutOfMemory; - defer alloc.free(work.data); - defer alloc.free(work.hostname); - - // content is correct - try std.testing.expectEqualStrings("example.pds.host", work.hostname); - try std.testing.expectEqualStrings("raw frame bytes", work.data); - - // pointers must be distinct from caller's buffers — this is the core - // UAF invariant. if the dupe was elided, these would alias. - try std.testing.expect(work.hostname.ptr != orig_hostname.ptr); - try std.testing.expect(work.data.ptr != data_input.ptr); - - // scalar fields propagate - try std.testing.expectEqual(@as(u64, 42), work.host_id); - - // simulate subscriber teardown (slurper.runWorker frees hostname) - alloc.free(orig_hostname); - - // work item must still be readable and correct — if hostname was - // borrowed instead of duped, the testing allocator would catch the - // read-after-free above when expectEqualStrings is called again here. - try std.testing.expectEqualStrings("example.pds.host", work.hostname); -} - test "decode frame via SDK and extract fields" { const cbor = zat.cbor; diff --git a/src/validator.zig b/src/validator.zig index 764a649..c1d453c 100644 --- a/src/validator.zig +++ b/src/validator.zig @@ -61,19 +61,15 @@ pub const Validator = struct { io: Io, // pool of reusable resolvers for inline host authority checks. // frame workers acquire/release via atomic flag to avoid creating - // a fresh resolver (and fresh TLS handshake) per call. heap-allocated - // in start() so the pool size can be tuned via HOST_RESOLVER_POOL_SIZE - // env var without recompiling. with keep_alive=false, pool width is a - // real startup throughput knob. - host_resolvers: []zat.DidResolver = &.{}, - host_resolver_available: []std.atomic.Value(bool) = &.{}, + // a fresh resolver (and fresh TLS handshake) per call. + host_resolvers: [host_resolver_pool_size]zat.DidResolver = undefined, + host_resolver_available: [host_resolver_pool_size]std.atomic.Value(bool) = .{std.atomic.Value(bool){ .raw = false }} ** host_resolver_pool_size, host_resolver_inited: bool = false, const max_resolver_threads = 8; const default_resolver_threads = 4; const max_queue_size: usize = 100_000; - const default_host_resolver_pool_size: usize = 4; - const max_host_resolver_pool_size: usize = 64; + const host_resolver_pool_size: usize = 4; pub fn init(allocator: Allocator, stats: *broadcaster.Stats, io: Io) Validator { return initWithConfig(allocator, stats, .{}, io); @@ -100,13 +96,9 @@ pub const Validator = struct { } if (self.host_resolver_inited) { - for (self.host_resolvers) |*r| { + for (&self.host_resolvers) |*r| { r.deinit(); } - self.allocator.free(self.host_resolvers); - self.allocator.free(self.host_resolver_available); - self.host_resolvers = &.{}; - self.host_resolver_available = &.{}; self.host_resolver_inited = false; } @@ -130,40 +122,14 @@ pub const Validator = struct { slot.* = try self.io.concurrent(resolveLoop, .{self}); } - // init host authority resolver pool (reused across calls). - // - // keep_alive = false: workaround for ~99% rejection rate observed - // 2026-04-08. root cause not yet known — local repro couldn't - // reproduce the failure, the leading hypothesis is that pool slots - // get poisoned by a transient network condition and never recover - // because there's no slot-recovery path. slot recovery added in - // resolveHostAuthority below; keep_alive can flip back to true via - // canary once resolve_loop_resolve_fail and the new sampled warn - // log give us the actual underlying error kind from zat - // v0.3.0-alpha.23. see relay docs/zlay-external-review-2026-04-09.md. - // - // pool size is HOST_RESOLVER_POOL_SIZE (default 4, max 64). with - // keep_alive=false, every check is a fresh TLS handshake (~tens of - // ms), so pool width is a real startup throughput knob — bumping - // it lets more is_new checks run concurrently during cold-start - // reconnect storms. - const requested_size = parseEnvInt(usize, "HOST_RESOLVER_POOL_SIZE", default_host_resolver_pool_size); - const pool_size = @min(requested_size, max_host_resolver_pool_size); - - self.host_resolvers = try self.allocator.alloc(zat.DidResolver, pool_size); - errdefer self.allocator.free(self.host_resolvers); - self.host_resolver_available = try self.allocator.alloc(std.atomic.Value(bool), pool_size); - errdefer self.allocator.free(self.host_resolver_available); - - for (self.host_resolvers) |*r| { - r.* = zat.DidResolver.initWithOptions(self.io, self.allocator, .{ .keep_alive = false }); + // init host authority resolver pool (reused across calls) + for (&self.host_resolvers) |*r| { + r.* = zat.DidResolver.initWithOptions(self.io, self.allocator, .{}); } - for (self.host_resolver_available) |*a| { - a.* = .{ .raw = true }; + for (&self.host_resolver_available) |*a| { + a.store(true, .release); } self.host_resolver_inited = true; - - log.info("host_authority resolver pool: size={d} keep_alive=false", .{pool_size}); } /// validate a #sync frame: signature verification only (no ops, no MST). @@ -491,12 +457,10 @@ pub const Validator = struct { // resolve DID → signing key const parsed = zat.Did.parse(d) orelse continue; var doc = resolver.resolve(parsed) catch |err| { - _ = self.stats.resolve_loop_resolve_fail_total.fetchAdd(1, .monotonic); log.debug("DID resolve failed for {s}: {s}", .{ d, @errorName(err) }); continue; }; defer doc.deinit(); - _ = self.stats.resolve_loop_resolve_ok_total.fetchAdd(1, .monotonic); // extract and decode signing key const vm = doc.signingKey() orelse continue; @@ -575,86 +539,41 @@ pub const Validator = struct { /// synchronous host authority check. called on first-seen DIDs (is_new) /// and host migrations (host_changed). resolves the DID doc to verify the - /// PDS endpoint matches the incoming host. + /// PDS endpoint matches the incoming host. retries once on failure to + /// handle transient network errors. /// /// uses a pooled resolver to avoid creating a fresh resolver (and fresh /// TLS handshake) per call. blocks briefly if all pool slots are in use. /// - /// on resolve failure, the slot is destroyed and re-initialized before - /// the retry, so any state corruption (poisoned http client, half-closed - /// connection, etc.) doesn't persist across calls. this is the leading - /// hypothesis for the 2026-04-08 ~99% rejection rate — pool slots had - /// no recovery path. see relay docs/zlay-external-review-2026-04-09.md. - /// /// returns: /// .accept — should not happen (caller should only call on new/mismatch) /// .migrate — DID doc confirms this host, caller should update host_id /// .reject — DID doc does not confirm, caller should drop the event pub fn resolveHostAuthority(self: *Validator, did: []const u8, incoming_host_id: u64) HostAuthority { const persist = self.persist orelse return .migrate; // no DB — can't check - const parsed = zat.Did.parse(did) orelse { - _ = self.stats.host_authority_reject_parse_did.fetchAdd(1, .monotonic); - self.sampleLogReject("parse_did", did, "", incoming_host_id, 0); - return .reject; - }; + const parsed = zat.Did.parse(did) orelse return .reject; const idx = self.acquireHostResolver(); defer self.releaseHostResolver(idx); - // first attempt on the existing pool slot. the first-attempt error is - // not captured: if state corruption was the cause, the kind from the - // fresh-slot retry below is what we want to log. - if (self.host_resolvers[idx].resolve(parsed)) |doc_first| { - var d = doc_first; - defer d.deinit(); - return self.checkPdsHost(&d, persist, did, incoming_host_id); - } else |_| { - // first-attempt failure: count it (steady-state signal independent - // of recovery success), then destroy + re-init the slot before the - // retry so any internal state corruption doesn't persist. - _ = self.stats.host_resolver_resolve_fail_total.fetchAdd(1, .monotonic); - self.recycleHostResolver(idx); - _ = self.stats.host_resolver_resets_total.fetchAdd(1, .monotonic); - - if (self.host_resolvers[idx].resolve(parsed)) |doc_retry| { - var d = doc_retry; - defer d.deinit(); - return self.checkPdsHost(&d, persist, did, incoming_host_id); - } else |err| { - _ = self.stats.host_authority_reject_resolve.fetchAdd(1, .monotonic); - self.sampleLogReject("resolve", did, @errorName(err), incoming_host_id, 0); - return .reject; - } - } - } + var resolver = &self.host_resolvers[idx]; - /// destroy and re-initialize a pool slot in place. caller must hold the - /// slot via acquireHostResolver — concurrent access is not safe. used by - /// the slot-recovery path on resolve failure. if the re-init alloc fails, - /// the slot is left in a degraded state and the next caller will see the - /// failure naturally; we don't crash on OOM here. - fn recycleHostResolver(self: *Validator, idx: usize) void { - self.host_resolvers[idx].deinit(); - self.host_resolvers[idx] = zat.DidResolver.initWithOptions( - self.io, - self.allocator, - .{ .keep_alive = false }, - ); + // first resolve attempt + var doc = resolver.resolve(parsed) catch { + // retry once on network failure + var doc2 = resolver.resolve(parsed) catch return .reject; + defer doc2.deinit(); + return self.checkPdsHost(&doc2, persist, incoming_host_id); + }; + defer doc.deinit(); + return self.checkPdsHost(&doc, persist, incoming_host_id); } /// acquire a resolver from the pool. spins until one is available. - /// records cumulative wait time + in_use gauge for diagnosing pool - /// contention. matches the codebase convention of timing via Io.Timestamp - /// (see frame_worker.zig microTimestamp). fn acquireHostResolver(self: *Validator) usize { - const start_us = Io.Timestamp.now(self.io, .real).toMicroseconds(); while (self.alive.load(.acquire)) { - for (0..self.host_resolvers.len) |i| { + for (0..host_resolver_pool_size) |i| { if (self.host_resolver_available[i].cmpxchgStrong(true, false, .acquire, .monotonic) == null) { - const now_us = Io.Timestamp.now(self.io, .real).toMicroseconds(); - const elapsed_us: u64 = @intCast(@max(0, now_us - start_us)); - _ = self.stats.host_resolver_acquire_wait_us_total.fetchAdd(elapsed_us, .monotonic); - _ = self.stats.host_resolver_in_use.fetchAdd(1, .monotonic); return i; } } @@ -664,52 +583,17 @@ pub const Validator = struct { } fn releaseHostResolver(self: *Validator, idx: usize) void { - _ = self.stats.host_resolver_in_use.fetchSub(1, .monotonic); self.host_resolver_available[idx].store(true, .release); } - fn checkPdsHost(self: *Validator, doc: *zat.DidDocument, persist: *event_log_mod.DiskPersist, did: []const u8, incoming_host_id: u64) HostAuthority { - const pds_endpoint = doc.pdsEndpoint() orelse { - _ = self.stats.host_authority_reject_no_endpoint.fetchAdd(1, .monotonic); - self.sampleLogReject("no_endpoint", did, "", incoming_host_id, 0); - return .reject; - }; - const pds_host = extractHostFromUrl(pds_endpoint) orelse { - _ = self.stats.host_authority_reject_bad_url.fetchAdd(1, .monotonic); - self.sampleLogReject("bad_url", did, pds_endpoint, incoming_host_id, 0); - return .reject; - }; - const pds_host_id = (persist.getHostIdForHostname(pds_host) catch null) orelse { - _ = self.stats.host_authority_reject_unknown_host.fetchAdd(1, .monotonic); - self.sampleLogReject("unknown_host", did, pds_host, incoming_host_id, 0); - return .reject; - }; + fn checkPdsHost(self: *Validator, doc: *zat.DidDocument, persist: *event_log_mod.DiskPersist, incoming_host_id: u64) HostAuthority { + _ = self; + const pds_endpoint = doc.pdsEndpoint() orelse return .reject; + const pds_host = extractHostFromUrl(pds_endpoint) orelse return .reject; + const pds_host_id = (persist.getHostIdForHostname(pds_host) catch null) orelse return .reject; if (pds_host_id == incoming_host_id) return .migrate; - _ = self.stats.host_authority_reject_host_mismatch.fetchAdd(1, .monotonic); - self.sampleLogReject("host_mismatch", did, pds_host, incoming_host_id, pds_host_id); return .reject; } - - /// log a rejection sample at 1-in-N rate to avoid drowning the log at - /// ~10 rejections/sec. total rejections per branch are available via - /// relay_host_authority_reject{branch=...} in prometheus. - fn sampleLogReject( - self: *Validator, - branch: []const u8, - did: []const u8, - detail: []const u8, - incoming_host_id: u64, - resolved_host_id: u64, - ) void { - const count = self.stats.failed_host_authority.load(.monotonic); - // sample 1 in 2048 — at 10/s that's one log line every ~3.5min. - // parens mandatory: `&` and `!=` have surprising precedence in zig. - if ((count & 0x7ff) != 0) return; - log.warn( - "host_authority reject branch={s} did={s} detail={s} incoming_host_id={d} resolved_host_id={d}", - .{ branch, did, detail, incoming_host_id, resolved_host_id }, - ); - } }; /// extract hostname from a URL like "https://pds.example.com" or "https://pds.example.com:443/path" -- 2.51.2