diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 892082d96c..6d935520fd 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -139,7 +139,7 @@ The `kind` integer is the only dispatch switch. The relay routes, stores, and fa | 46001–46012 | KIND_WORKFLOW_* | Workflow execution events | | 20001 | KIND_PRESENCE_UPDATE | Ephemeral presence heartbeat | -`buzz-core` defines each event kind as a `pub const u32` and exports the full registry as `ALL_KINDS: &[u32]` (127 kinds at the time of writing); `crates/buzz-core/src/kind.rs` is the source of truth for the current list. Kinds are `u32` (NIP-01 specifies unsigned integer; `u32` covers the full range). Buzz uses both standard Nostr kinds (e.g., kind 7 for reactions) and custom ranges (40000+). +`buzz-core` defines each event kind as a `pub const u32` and exports the full registry as `ALL_KINDS: &[u32]` (129 kinds at the time of writing); `crates/buzz-core/src/kind.rs` is the source of truth for the current list. Kinds are `u32` (NIP-01 specifies unsigned integer; `u32` covers the full range). Buzz uses both standard Nostr kinds (e.g., kind 7 for reactions) and custom ranges (40000+). Note: `KIND_AUTH` (22242) is `pub const KIND_AUTH: u32` in `buzz-core/src/kind.rs` and imported by `buzz-relay/src/handlers/event.rs`. `KIND_CANVAS` (40100) is likewise `pub const KIND_CANVAS: u32` in `buzz-core/src/kind.rs`. @@ -158,7 +158,7 @@ Note: `KIND_AUTH` (22242) is `pub const KIND_AUTH: u32` in `buzz-core/src/kind.r | Relay → Client | `["NOTICE", "message"]` | Informational message | | Relay → Client | `["AUTH", ]` | Authentication challenge | -Max frame size: 65,536 bytes. Max subscriptions per connection: 1024. Max historical results per filter: 500. +Max frame size: 524,288 bytes (512 KiB, `DEFAULT_MAX_FRAME_BYTES`, env-overridable via `MAX_FRAME_BYTES`). Max subscriptions per connection: 1024. Max historical results per filter: 1,000 (`DEFAULT_MAX_PAGE_LIMIT`, advertised as NIP-11 `max_limit`). Note: NIP-11 `max_content_len` (65,536) is a separate value describing maximum accepted event `content` length — distinct from the WebSocket frame bound. --- @@ -230,14 +230,14 @@ When the relay receives `["EVENT", ]`, the handler in `handlers/event.rs` 5. VERIFY — spawn_blocking(verify_event) — Schnorr sig + ID hash 6. MEMBERSHIP — channel_id in event tags? → check_channel_membership 7. DB INSERT — db.insert_event (ON CONFLICT DO NOTHING — idempotent) -8. REDIS PUBLISH — pubsub.publish_event (if channel-scoped) -9. FAN-OUT — sub_registry.fan_out → conn_manager.send_to -10. SEARCH INDEX — search_index_tx.send (bounded worker queue, non-blocking) -11. AUDIT LOG — audit.log (spawned async, non-blocking) -12. WORKFLOW TRIGGER — wf.on_event (spawned async, excludes kinds 46001–46012) + 8. REDIS PUBLISH — pubsub.publish_event (if channel-scoped) + 9. FAN-OUT — sub_registry.fan_out → conn_manager.send_to + 10. SEARCH INDEX — (no separate step; FTS via `events.search_tsv` generated column, populated by the DB insert itself) + 11. AUDIT LOG — audit.log (spawned async, non-blocking) + 12. WORKFLOW TRIGGER — wf.on_event (spawned async, excludes kinds 46001–46012) ``` -Steps 10–12 are fire-and-forget. Search indexing is sent to a bounded worker queue (`search_index_tx`, capacity 1000); audit and workflow triggers are spawned as independent async tasks. A failure in any of these does not fail the event submission. The client receives `["OK", , true, ""]` at the end of the pipeline, not immediately after DB insert. +Steps 8–9 and 11–12 run inside a spawned task (`dispatch_persistent_event` in `crates/buzz-relay/src/handlers/event.rs`), fire-and-forget relative to the client-facing `OK`. Search indexing is no longer a separate worker step: under Postgres FTS the searchable row IS the persisted event row (the `insert_event` write populates the FTS column via a generated `tsvector`), so there is no out-of-band index to feed. The old `search_index_tx` mpsc is gone. A failure in any of these does not fail the event submission. The client receives `["OK", , true, ""]` after the DB insert + audit enqueue, not after all side effects complete. Step 9 (fan-out) explicitly **excludes** global subscriptions (no `channel_id` constraint) from channel-scoped events — global subscriptions do NOT receive events from private channels, regardless of filter match. This is a deliberate security boundary: only subscriptions scoped to an accessible `channel_id` receive those events. @@ -323,7 +323,7 @@ This prevents a race where a non-member receives live fan-out events from a priv ### Historical Query (EOSE) -After registering, the REQ handler queries Postgres for stored events matching the filters (up to 500 per filter, hard cap). These are sent as `["EVENT", sub_id, event]` frames before `["EOSE", sub_id]`. New events arriving after EOSE are delivered via the fan-out path. +After registering, the REQ handler queries Postgres for stored events matching the filters (up to 1,000 per filter, hard cap — `DEFAULT_MAX_PAGE_LIMIT`, advertised as NIP-11 `max_limit`). These are sent as `["EVENT", sub_id, event]` frames before `["EOSE", sub_id]`. New events arriving after EOSE are delivered via the fan-out path. --- @@ -343,7 +343,7 @@ pub struct StoredEvent { verified: bool, // private — use is_verified() } -pub const ALL_KINDS: &[u32] // 80 entries (KIND_AUTH excluded — never stored) + pub const ALL_KINDS: &[u32] // 129 entries (KIND_AUTH excluded — never stored) ``` **Key functions:** @@ -387,7 +387,7 @@ pub trait RateLimiter: Send + Sync { ... } - NIP-42 timestamp tolerance: ±60 seconds. - Dev-only key derivation: `SHA-256("buzz-test-key:{username}")` — gated behind `#[cfg(any(test, feature = "dev"))]`. The `dev` feature must not be enabled in production relay deployments. -**Does NOT:** implement `RateLimiter` beyond a test stub (`AlwaysAllowRateLimiter`, gated behind `#[cfg(any(test, feature = "test-utils"))]`). No Redis-backed rate limiter exists anywhere in the codebase — rate limiting is not currently enforced. `RateLimitConfig` defines 4 tiers (human, agent-standard, agent-elevated, agent-platform) as a design target. +**Does NOT:** implement `RateLimiter` directly. The `RateLimiter` trait lives here (`AlwaysAllowRateLimiter` is the test stub, gated behind `#[cfg(any(test, feature = "test-utils"))]`), but the production implementation (`RedisRateLimiter`) lives in `buzz-pubsub` (see its section below). `RateLimitConfig` defines 4 tiers (human, agent-standard, agent-elevated, agent-platform) that the enforcement sites in `buzz-relay` consult. --- @@ -457,7 +457,9 @@ EXPIRE buzz:typing:{channel_id} 60 ``` 5-second activity window. 60-second key TTL prevents orphaned empty sets. -**Does NOT:** implement the rate limiter. Does NOT store events. `PubSubManager` is not `Clone` — callers use `Arc`. +**Does:** implement the production rate limiter (`RedisRateLimiter` in `crates/buzz-pubsub/src/rate_limiter.rs`), shared across relay instances via Redis. The relay wires it as `admission_rate_limiter` and enforces it on WS connect, per-message, and the HTTP bridge (see `buzz-relay`). + +**Does NOT:** store events. `PubSubManager` is not `Clone` — callers use `Arc`. --- @@ -466,8 +468,13 @@ EXPIRE buzz:typing:{channel_id} 60 Full-text search via Postgres FTS. Events are searchable through the `events.search_tsv` generated `tsvector` column (populated on insert, indexed by a GIN index) — there is no separate search service or out-of-band indexer. -Privacy-sensitive kinds are excluded at the storage level (the `search_tsv` -`CASE WHEN kind IN (...)` yields `NULL`, which never matches `@@`). In +Privacy-sensitive kinds are excluded at the storage level. Fresh installs +(`migrations/0008_fresh_install_search_allowlist.sql`) use a positive allowlist: +`kind IN (0, 9, 40002, 45001, 45003)` — only these kinds populate `search_tsv`. +Existing/upgraded installs retain the original exclusion list from +`migrations/0001_initial_schema.sql` (`kind IN (1059, 30300, 30622, 44100, +44101)` yields `NULL::tsvector`) until an operator runs the sized out-of-band +maintenance script (`scripts/maintenance/nip_rs_search_allowlist.sql`). In multi-community mode every query filter includes `community_id`, so the shared `events` table is infrastructure, not a cross-community result space; the relay re-authorizes every candidate hit before returning it. @@ -589,10 +596,10 @@ pub struct AppState { pub workflow_engine: Arc, pub conn_semaphore: Arc, // connection limit pub handler_semaphore: Arc, // 1024 concurrent handlers - pub relay_keypair: nostr::Keys, // relay identity - pub local_event_ids: moka::sync::Cache, // local-echo dedup - pub search_index_tx: mpsc::Sender, // bounded search worker queue - // + config, redis_pool, membership_cache, media_storage, shutdown state + pub relay_keypair: nostr::Keys, // relay identity + pub local_event_ids: moka::sync::Cache, // local-echo dedup + pub audit_tx: Option, // bounded audit worker queue + // + config, redis_pool, membership_cache, media_storage, shutdown state } ``` @@ -632,9 +639,9 @@ pub enum AuthState { Pending { challenge: String }, Authenticated(AuthContext), | Constant | Value | Purpose | |----------|-------|---------| -| `MAX_FRAME_BYTES` | 65,536 | Max WebSocket frame size | +| `MAX_FRAME_BYTES` | 524,288 (512 KiB) | Max WebSocket frame size | | `MAX_SUBSCRIPTIONS` | 1024 | Per-connection subscription limit | -| `MAX_HISTORICAL_LIMIT` | 500 | Per-filter historical query cap | +| `MAX_HISTORICAL_LIMIT` | 1,000 | Per-filter historical query cap (`DEFAULT_MAX_PAGE_LIMIT`) | | `handler_semaphore` capacity | 1024 | Concurrent EVENT/REQ handlers | **Does NOT:** implement business logic — delegates to the appropriate crate for every operation. @@ -657,13 +664,15 @@ Buzz Relay ──WS──→ buzz-acp ──stdio (ACP/JSON-RPC)──→ Agent | Module | LOC | Responsibility | |--------|-----|---------------| -| `relay.rs` | 3,143 | WebSocket + REST relay connection, NIP-42 auth | -| `queue.rs` | 2,565 | Per-channel event queue, batching, dedup | -| `main.rs` | 2,457 | Event loop, pool orchestration, heartbeat | -| `pool.rs` | 2,253 | N-agent pool, claim/return lifecycle | -| `config.rs` | 1,903 | CLI/env/TOML configuration | -| `acp.rs` | 1,785 | ACP client, stdio JSON-RPC, timeouts | -| `filter.rs` | 814 | Subscription rules, evalexpr filtering | +| `relay.rs` | ~6,200 | WebSocket + REST relay connection, NIP-42 auth | +| `queue.rs` | ~4,800 | Per-channel event queue, batching, dedup | +| `lib.rs` | ~6,700 | Event loop, pool orchestration, heartbeat (formerly `main.rs`) | +| `pool.rs` | ~6,900 | N-agent pool, claim/return lifecycle | +| `config.rs` | ~2,900 | CLI/env/TOML configuration | +| `acp.rs` | ~4,500 | ACP client, stdio JSON-RPC, timeouts | +| `filter.rs` | ~800 | Subscription rules, evalexpr filtering | + +> LOC counts are approximate and current as of this revision — see the source files for exact counts. **Key behaviors:** - Pool of 1–32 agent subprocesses with claim/return lifecycle. @@ -699,12 +708,12 @@ The `buzz-admin` binary is shipped in the relay Docker image (`/usr/local/bin/bu | File | Tests | Scope | |------|-------|-------| -| `tests/e2e_relay.rs` | 27 | WebSocket protocol (auth, subscriptions, filters, limits, NIP-11) | -| `tests/e2e_media.rs` | 7 | Media upload/download (Blossom) | -| `tests/e2e_media_extended.rs` | 18 | Extended media scenarios | -| `tests/e2e_nostr_interop.rs` | 15 | Nostr interoperability: NIP-50 search, NIP-10 threads, NIP-17 gift wraps, DM discovery | +| `tests/e2e_relay.rs` | ~54 | WebSocket protocol (auth, subscriptions, filters, limits, NIP-11) | +| `tests/e2e_media.rs` | ~7 | Media upload/download (Blossom) | +| `tests/e2e_media_extended.rs` | ~24 | Extended media scenarios | +| `tests/e2e_nostr_interop.rs` | ~34 | Nostr interoperability: NIP-50 search, NIP-10 threads, NIP-17 gift wraps, DM discovery | -All e2e tests are `#[ignore]` — require a running relay. Total: **134 e2e tests**. +> E2e tests live in `crates/buzz-test-client/tests/` and are `#[ignore]` — they require a running relay. Counts are approximate; the directory also contains additional e2e files (persona, team, managed-agent, etc.). `src/main.rs` is a manual testing CLI (`buzz-test-cli`) with `--send`, `--subscribe`, `--channel`, `--url`, `--kind` flags. @@ -730,7 +739,7 @@ Every security-sensitive operation uses an explicit, verified pattern. No implic |---------|-----------| | Schnorr signatures | `verify_event()` in `buzz-core` — every event verified before storage | | Event ID | SHA-256 of canonical serialization verified independently of signature | -| Frame size | `MAX_FRAME_BYTES = 65,536` — oversized frames rejected, connection closed | +| Frame size | `MAX_FRAME_BYTES = 524,288` (512 KiB default, env-overridable) — oversized frames rejected, connection closed | | Search event IDs | 64-char hex validation before URL construction — prevents path injection | | Workflow step IDs | Alphanumeric + underscore only — prevents evalexpr variable injection | | Partition names | Allowlist of table names + strict suffix/date validators — prevents DDL injection | @@ -820,7 +829,7 @@ These are verified gaps in the current implementation — not design aspirations | # | Limitation | Detail | |---|-----------|--------| | 1 | **No sqlx offline query cache** | Uses `sqlx::query()` (runtime) not `sqlx::query!()` (compile-time). No `.sqlx/` directory. Queries are not validated at compile time. | -| 2 | **No rate limiting implementation** | `RateLimiter` trait exists in `buzz-auth`. Only implementation is `AlwaysAllowRateLimiter` (test stub, gated behind `#[cfg(any(test, feature = "test-utils"))]`). `RateLimitConfig` defines 4 tiers (human, agent-standard, agent-elevated, agent-platform) but none are enforced. | +| 2 | **Rate limiting enforced via Redis** | `RedisRateLimiter` (`crates/buzz-pubsub/src/rate_limiter.rs`) is the production implementation of the `RateLimiter` trait (defined in `buzz-auth`, where `AlwaysAllowRateLimiter` remains as a test stub). The relay wires it as `admission_rate_limiter` (`crates/buzz-relay/src/state.rs`) and enforces it on WS connect, per-message (`crates/buzz-relay/src/connection.rs`), and the HTTP bridge (`crates/buzz-relay/src/api/bridge.rs`). `RateLimitConfig` defines 4 tiers (human, agent-standard, agent-elevated, agent-platform) that these sites consult. Tuning requires code/config changes; there is no dynamic admin API. | | 3 | **No dedicated typing REST endpoint** | Typing indicators (kind 20002) are delivered via both local fan-out and Redis pub/sub (cross-node). There is no REST endpoint to query current typers — `/api/presence` returns online/away status only, not typing state. | | 4 | **Huddle recording/tracks not built** | Voice, room lifecycle, and join/leave/end events are wired (see Huddle Audio above). Recording and per-track publishing have reserved kinds but no producer yet. | | 5 | **Approval gates not wired end-to-end** | The executor returns `StepResult::Suspended` and the relay has grant/deny API endpoints with DB CRUD, but the engine intercepts before creating `WaitingApproval` rows — runs that hit an approval gate are marked as Failed (🚧 WF-08). |