diff --git a/README.md b/README.md index fc8e276..8731609 100644 --- a/README.md +++ b/README.md @@ -23,7 +23,7 @@ implements the [AT Protocol sync spec](https://atproto.com/specs/sync) — `subs | 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 | +| [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 | diff --git a/docs/deployment.md b/docs/deployment.md index e2e295e..f6863b9 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -73,7 +73,9 @@ just zlay-ssh # ssh into server ## memory tuning -two changes brought steady-state memory from ~6.6 GiB down to ~2.9 GiB at 2,738 connected hosts: +three changes brought steady-state memory from ~6.6 GiB down to ~1.2 GiB at ~2,750 connected hosts: + +**shared TLS CA bundle.** the biggest single win. websocket.zig's TLS client calls `Bundle.rescan()` per connection, loading the system CA certificates into a per-connection arena. with ~2,750 PDS connections, that's ~2,750 copies of the CA bundle in memory (~800 KB each = ~2.2 GiB). fix: load the bundle once in the slurper, pass it to all subscribers via `config.ca_bundle`. memory dropped from ~3.3 GiB to ~1.2 GiB (~65% reduction). **thread stack sizes.** zig's default thread stack is 16 MB. with ~2,750 subscriber threads that maps 44 GB of virtual memory. most threads just read websockets and decode CBOR — 2 MB is generous. all `Thread.spawn` calls now pass `.{ .stack_size = 2 * 1024 * 1024 }`. the constant is defined in `main.zig` as `default_stack_size` for the threads spawned there; other modules use the literal directly. @@ -83,9 +85,10 @@ two changes brought steady-state memory from ~6.6 GiB down to ~2.9 GiB at 2,738 | metric | value | |--------|-------| -| memory | ~2.9 GiB steady state (~2,750 hosts) | +| memory | ~1.2 GiB steady state (~2,750 hosts) | | CPU | ~1.5 cores peak | -| limits | 8 GiB memory, 250m CPU request | +| requests | 1 GiB memory, 1000m CPU | +| limits | 3 GiB memory | | PVC | 20 GiB (events + RocksDB collection index) | | postgres | ~238 MiB | diff --git a/docs/design.md b/docs/design.md index ffb3d7e..f7ee0c5 100644 --- a/docs/design.md +++ b/docs/design.md @@ -34,8 +34,9 @@ Downstream consumers (WebSocket) additionally: - **collection index** (RocksDB): subscriber calls `trackCommitOps` on each validated commit; stores `(collection, did)` pairs for `listReposByCollection` -- **event log**: append-only files rotated every 10K events, 72h retention. - supports cursor replay — disk first, then in-memory ring buffer (50K frames) +- **event log**: append-only files rotated every 10K events, configurable + retention (default: 72h, env: `RELAY_RETENTION_HOURS`). supports cursor + replay — disk first, then in-memory ring buffer (50K frames) - **slurper**: orchestrates subscribers. bootstraps host list from seed relay's `listHosts` API, spawns/stops workers, processes `requestCrawl` requests @@ -68,6 +69,11 @@ arenas and `madvise`-based page return. the general-purpose allocator (GPA) is a debug allocator that never returns freed pages — unsuitable for long-running servers. +**shared TLS CA bundle**: loaded once by the slurper, passed to all ~2,750 +subscriber connections via `config.ca_bundle`. without this, each websocket.zig +TLS client calls `Bundle.rescan()` and loads its own copy (~800 KB each), +totaling ~2.2 GiB of duplicate CA certificates in memory. + **arena per frame**: each subscriber creates a `std.heap.ArenaAllocator` per WebSocket message. all CBOR decode temporaries, CAR parse buffers, and MST nodes live in this arena. freed in bulk after the frame is processed. this @@ -80,7 +86,10 @@ releases. this avoids copying frame bytes per consumer. **validator cache**: `StringHashMap(CachedKey)` — DID string → 75-byte fixed-size struct (key type + 33-byte compressed pubkey + resolve timestamp). capped at 500K entries (env: `VALIDATOR_CACHE_SIZE`), LRU-ish eviction of -oldest 10% when full. ~37 MB at capacity. +oldest 10% when full. ~37 MB at capacity. the resolve queue uses a +`StringHashMapUnmanaged(void)` as a dedupe set to prevent the same DID from +being queued multiple times. migration checks are interleaved with DID +resolutions (1 per 10) to prevent starvation. **ring buffer**: 50K-entry in-memory frame history for cursor replay when disk replay isn't available. entries are `(seq, data)` pairs with data duped @@ -127,19 +136,20 @@ populated live from firehose commits. backfill from source relay's ## scaling limits current deployment: ~2,780 PDS hosts, running on a 32 GB / 16 CPU node. -steady-state memory: ~3.5 GiB. postgres alongside at ~240 MiB. +steady-state memory: ~1.2 GiB (after shared CA bundle fix). postgres alongside +at ~240 MiB. resource limits: 3 GiB memory, 1 GiB request, 1000m CPU. | component | current (~2,750 PDS) | at 10x (~27,500 PDS) | status | |---|---|---|---| | thread stacks | ~5.5 GB virtual (2,750 × 2 MB) | ~55 GB virtual | **breaks** — exceeds 32 GB node | | pg pool | 5 connections (hardcoded) | 5 connections | **breaks** — saturates under concurrent UID lookups | -| resolver queue | unbounded `ArrayList` | unbounded | **risk** — backlog grows if resolvers can't keep up | +| resolver queue | `ArrayList` + dedupe set | same | **ok** — dedupe prevents unbounded growth from duplicate DIDs | | validator cache | 500K entries, ~37 MB | same (capped) | **degrades** — miss rate climbs with more unique DIDs | | broadcaster | O(n consumers) under mutex | same | **risk** — lock contention at high consumer count | | RocksDB | manageable write rate | ~1.4M writes/sec projected | **needs** compaction tuning | | event log | buffered, 100ms flush | fine — sequential I/O | ok | | kernel threads | ~2,800 (below 30K default) | ~28,000 (near default max) | **breaks** without `sysctl` tuning | -| RSS | ~3.5 GiB | ~15–20 GiB projected (malloc overhead scales sublinearly) | **tight** — needs larger node | +| RSS | ~1.2 GiB | ~5–8 GiB projected (shared CA bundle, malloc overhead scales sublinearly) | ok — fits 32 GB node | ### what breaks first diff --git a/src/subscriber.zig b/src/subscriber.zig index e4a00e6..86f714e 100644 --- a/src/subscriber.zig +++ b/src/subscriber.zig @@ -323,10 +323,16 @@ const FrameHandler = struct { return; } - // route by frame type + // route by frame type — unknown types are ignored per spec (forward-compat) const is_commit = std.mem.eql(u8, frame_type, "#commit"); const is_sync = std.mem.eql(u8, frame_type, "#sync"); const is_account = std.mem.eql(u8, frame_type, "#account"); + const is_identity = std.mem.eql(u8, frame_type, "#identity"); + + if (!is_commit and !is_sync and !is_account and !is_identity) { + log.debug("host {s}: unknown frame type '{s}', ignoring", .{ sub.options.hostname, frame_type }); + return; + } // extract DID: "repo" for commits, "did" for identity/account const did: ?[]const u8 = if (is_commit) @@ -424,7 +430,7 @@ const FrameHandler = struct { .sync else if (is_account) .account - else + else // is_identity (unknown types already filtered above) .identity; // persist and get relay-assigned seq, broadcast raw bytes. diff --git a/src/validator.zig b/src/validator.zig index 9d40a1e..f28cd32 100644 --- a/src/validator.zig +++ b/src/validator.zig @@ -153,7 +153,8 @@ pub const Validator = struct { }; // #sync CAR should be small (just the signed commit block) - if (blocks.len > 10 * 1024) { + // lexicon maxLength: 10000 + if (blocks.len > 10_000) { _ = self.stats.failed.fetchAdd(1, .monotonic); return .{ .valid = false, .skipped = false }; } @@ -246,8 +247,8 @@ pub const Validator = struct { // extract blocks (raw CAR bytes) from the pre-decoded payload const blocks = payload.getBytes("blocks") orelse return error.InvalidFrame; - // blocks size check - if (blocks.len > 2 * 1024 * 1024) return error.InvalidFrame; + // blocks size check — lexicon maxLength: 2000000 + if (blocks.len > 2_000_000) return error.InvalidFrame; // build public key for verification const public_key = zat.multicodec.PublicKey{