diff --git a/crates/buzz-conformance/TRACE_SCHEMA.md b/crates/buzz-conformance/TRACE_SCHEMA.md index c1dca674d55..f3c6afde9d7 100644 --- a/crates/buzz-conformance/TRACE_SCHEMA.md +++ b/crates/buzz-conformance/TRACE_SCHEMA.md @@ -1,6 +1,7 @@ # Trace Schema (`buzz-conformance`) -Schema version: **1** (`SCHEMA_VERSION` in `src/lib.rs`). +Schema version: **2** (`SCHEMA_VERSION` in `src/lib.rs`). Version 2 adds +`accept_ephemeral` to the action enum; version 1 consumers cannot decode it. This document is the contract between the relay's emitter and the independent replay checker. It is grounded in @@ -65,6 +66,10 @@ exact spec line it grounds in. via the host-community map; no `channel` field. `claimed_community` recorded for the same reason as above. +- **`accept_ephemeral { msg_id }`** + A validated ephemeral event accepted for live fan-out. No durable write + occurs. + - **`write_duplicate { msg_id, channel, claimed_community }`** spec: `WriteDuplicate` (line 612). The DB returned "already present"; no row was added. No `row_community` because no row was produced. @@ -133,7 +138,7 @@ normalized away the violation. The checker assumes you *did not*. |------|---------------| | `crates/buzz-relay/src/conformance/mod.rs` | helpers + `EmitGuard` + `sanitized_reason_for` | | `crates/buzz-relay/src/conformance/tracers.rs` | `NoopTracer` (prod default), `JsonlTracer` | -| `crates/buzz-relay/src/handlers/ingest.rs` | `AuthCheck`, `WriteInsert`, `WriteInsertGlobal`, `WriteDuplicate`, outer-wrapper `SanitizedError` | +| `crates/buzz-relay/src/handlers/ingest.rs` | `AuthCheck`, `WriteInsert`, `WriteInsertGlobal`, `AcceptEphemeral`, `WriteDuplicate`, outer-wrapper `SanitizedError` | | `crates/buzz-relay/src/handlers/req.rs` | **held back** — additive patch for integration onto Max's req.rs work | ## Where the checker lives diff --git a/crates/buzz-conformance/src/checker.rs b/crates/buzz-conformance/src/checker.rs index ceb4028cb9f..1129c708d21 100644 --- a/crates/buzz-conformance/src/checker.rs +++ b/crates/buzz-conformance/src/checker.rs @@ -258,6 +258,20 @@ mod tests { check_trace(&Scenario::unstructured(trace)).expect("deny with foreign claim is in-spec"); } + #[test] + fn accepted_ephemeral_event_is_valid_without_a_write() { + let c = cid(1); + let trace = vec![step( + TraceAction::AcceptEphemeral { + msg_id: OpaqueId("ephemeral".into()), + }, + c, + )]; + + check_trace(&Scenario::unstructured(trace)) + .expect("ephemeral acceptance is valid without a durable write"); + } + #[test] fn state_after_changing_mid_request_is_state_mismatch() { let c1 = cid(1); diff --git a/crates/buzz-conformance/src/lib.rs b/crates/buzz-conformance/src/lib.rs index b8e3f933df7..1848f306247 100644 --- a/crates/buzz-conformance/src/lib.rs +++ b/crates/buzz-conformance/src/lib.rs @@ -83,7 +83,7 @@ impl std::fmt::Display for CommunityLabel { } /// Trace schema version. Bump on any backwards-incompatible field change. -pub const SCHEMA_VERSION: u32 = 1; +pub const SCHEMA_VERSION: u32 = 2; /// An opaque ID derived from an event id or other secret material. Stable, /// no payload, no key bytes. Implementations pick a hash; the checker @@ -165,6 +165,7 @@ pub struct AbstractState { /// Action vocabulary (spec actions in parentheses): /// - [`TraceAction::WriteInsert`] (spec `WriteInsert`, lines 514–550) /// - [`TraceAction::WriteInsertGlobal`] (spec `WriteInsertGlobal`, lines 559–595) +/// - [`TraceAction::AcceptEphemeral`] (in-memory event acceptance; no modeled write) /// - [`TraceAction::WriteDuplicate`] (spec `WriteDuplicate`, lines 606–637) /// - [`TraceAction::SanitizedError`] (spec `SanitizedError`, line 778) /// - [`TraceAction::AuthCheck`] (spec `AuthCheck`, line 794) — M2/M8 target @@ -199,6 +200,11 @@ pub enum TraceAction { /// resolver but recorded for the audit trail. claimed_community: Option, }, + /// Validated ephemeral event accepted for live fan-out; no durable write occurs. + AcceptEphemeral { + /// Opaque hash of the event id. + msg_id: OpaqueId, + }, /// Channel-bearing duplicate / no-op write (spec `WriteDuplicate`, /// `ON CONFLICT (community_id, id)` returning a duplicate result). WriteDuplicate { @@ -267,6 +273,7 @@ impl TraceAction { match self { TraceAction::WriteInsert { .. } => "write_insert", TraceAction::WriteInsertGlobal { .. } => "write_insert_global", + TraceAction::AcceptEphemeral { .. } => "accept_ephemeral", TraceAction::WriteDuplicate { .. } => "write_duplicate", TraceAction::SanitizedError { .. } => "sanitized_error", TraceAction::AuthCheck { .. } => "auth_check", diff --git a/crates/buzz-conformance/src/transitions.rs b/crates/buzz-conformance/src/transitions.rs index cabd66e6919..8763d13f6f4 100644 --- a/crates/buzz-conformance/src/transitions.rs +++ b/crates/buzz-conformance/src/transitions.rs @@ -189,6 +189,8 @@ pub fn check_step( // --- Spec WriteInsertGlobal (lines 559-595) --- // resolved == HostCommunity[host]. Same shape as WriteInsert. TraceAction::WriteInsertGlobal { .. } => Ok(()), + // Ephemeral acceptance is an observation only; it mutates no durable model state. + TraceAction::AcceptEphemeral { .. } => Ok(()), // --- Spec WriteDuplicate (lines 606-637) --- // Carries the same host-axis obligation as WriteInsert: an A-host @@ -323,6 +325,7 @@ pub fn action_channel(action: &TraceAction) -> Option<&ChannelLabel> { TraceAction::ReadMessageRows { channel, .. } => channel.as_ref(), TraceAction::ReadByIdRows { channel, .. } => channel.as_ref(), TraceAction::WriteInsertGlobal { .. } + | TraceAction::AcceptEphemeral { .. } | TraceAction::ReadHostFeedRows { .. } | TraceAction::SanitizedError { .. } | TraceAction::ImplBug { .. } => None, diff --git a/crates/buzz-conformance/tests/fixtures/bad_coverage_breach.jsonl b/crates/buzz-conformance/tests/fixtures/bad_coverage_breach.jsonl index 8e165eaaafc..38cdfdbe9d3 100644 --- a/crates/buzz-conformance/tests/fixtures/bad_coverage_breach.jsonl +++ b/crates/buzz-conformance/tests/fixtures/bad_coverage_breach.jsonl @@ -1 +1 @@ -{"schema_version":1,"action":{"type":"impl_bug","kind":"ingest_exited_without_trace"},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} +{"schema_version":2,"action":{"type":"impl_bug","kind":"ingest_exited_without_trace"},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} diff --git a/crates/buzz-conformance/tests/fixtures/bad_foreign_row_leak.jsonl b/crates/buzz-conformance/tests/fixtures/bad_foreign_row_leak.jsonl index 8313562d450..857ed4ec630 100644 --- a/crates/buzz-conformance/tests/fixtures/bad_foreign_row_leak.jsonl +++ b/crates/buzz-conformance/tests/fixtures/bad_foreign_row_leak.jsonl @@ -1 +1 @@ -{"schema_version":1,"action":{"type":"read_message_rows","channel":"cafe0000-0000-0000-0000-000000000010","row_communities":["bbbb0000-0000-0000-0000-000000000002"]},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} +{"schema_version":2,"action":{"type":"read_message_rows","channel":"cafe0000-0000-0000-0000-000000000010","row_communities":["bbbb0000-0000-0000-0000-000000000002"]},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} diff --git a/crates/buzz-conformance/tests/fixtures/bad_host_channel_mismatch.jsonl b/crates/buzz-conformance/tests/fixtures/bad_host_channel_mismatch.jsonl index cf4cb26a778..b95b4f63f41 100644 --- a/crates/buzz-conformance/tests/fixtures/bad_host_channel_mismatch.jsonl +++ b/crates/buzz-conformance/tests/fixtures/bad_host_channel_mismatch.jsonl @@ -1,2 +1,2 @@ -{"schema_version":1,"action":{"type":"auth_check","channel":"dead0000-0000-0000-0000-000000000020","claimed_community":"bbbb0000-0000-0000-0000-000000000002","verdict":"allow"},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} -{"schema_version":1,"action":{"type":"write_insert","msg_id":"badbadbad0000000","channel":"dead0000-0000-0000-0000-000000000020","claimed_community":"bbbb0000-0000-0000-0000-000000000002"},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} +{"schema_version":2,"action":{"type":"auth_check","channel":"dead0000-0000-0000-0000-000000000020","claimed_community":"bbbb0000-0000-0000-0000-000000000002","verdict":"allow"},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} +{"schema_version":2,"action":{"type":"write_insert","msg_id":"badbadbad0000000","channel":"dead0000-0000-0000-0000-000000000020","claimed_community":"bbbb0000-0000-0000-0000-000000000002"},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} diff --git a/crates/buzz-conformance/tests/fixtures/good.jsonl b/crates/buzz-conformance/tests/fixtures/good.jsonl index e18be5656ce..70935f354a7 100644 --- a/crates/buzz-conformance/tests/fixtures/good.jsonl +++ b/crates/buzz-conformance/tests/fixtures/good.jsonl @@ -1,3 +1,3 @@ -{"schema_version":1,"action":{"type":"auth_check","channel":"cafe0000-0000-0000-0000-000000000010","claimed_community":"aaaa0000-0000-0000-0000-000000000001","verdict":"allow"},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} -{"schema_version":1,"action":{"type":"write_insert","msg_id":"d34db33fcafef00d","channel":"cafe0000-0000-0000-0000-000000000010","claimed_community":"aaaa0000-0000-0000-0000-000000000001"},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} -{"schema_version":1,"action":{"type":"read_message_rows","channel":"cafe0000-0000-0000-0000-000000000010","row_communities":["aaaa0000-0000-0000-0000-000000000001","aaaa0000-0000-0000-0000-000000000001"]},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} +{"schema_version":2,"action":{"type":"auth_check","channel":"cafe0000-0000-0000-0000-000000000010","claimed_community":"aaaa0000-0000-0000-0000-000000000001","verdict":"allow"},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} +{"schema_version":2,"action":{"type":"write_insert","msg_id":"d34db33fcafef00d","channel":"cafe0000-0000-0000-0000-000000000010","claimed_community":"aaaa0000-0000-0000-0000-000000000001"},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} +{"schema_version":2,"action":{"type":"read_message_rows","channel":"cafe0000-0000-0000-0000-000000000010","row_communities":["aaaa0000-0000-0000-0000-000000000001","aaaa0000-0000-0000-0000-000000000001"]},"state_after":{"resolved_community":"aaaa0000-0000-0000-0000-000000000001","bound_host":"a.example.test","actor":"0123456789abcdef"}} diff --git a/crates/buzz-relay/src/api/bridge.rs b/crates/buzz-relay/src/api/bridge.rs index 13437d9c0d8..ca504d803b9 100644 --- a/crates/buzz-relay/src/api/bridge.rs +++ b/crates/buzz-relay/src/api/bridge.rs @@ -1195,6 +1195,14 @@ async fn submit_event_authed( response: e, } } + Err(IngestError::RateLimited(msg)) => { + crate::handlers::ingest::reject_with_transport("http", "rate-limited"); + let e = api_error(StatusCode::TOO_MANY_REQUESTS, &msg); + SubmitOutcome::Err { + status: e.0, + response: e, + } + } Err(IngestError::Internal(msg)) => { crate::handlers::ingest::reject_with_transport("http", "error"); let e = internal_error(&msg); diff --git a/crates/buzz-relay/src/conformance/mod.rs b/crates/buzz-relay/src/conformance/mod.rs index 0ceb8b67b64..eee0a48cc47 100644 --- a/crates/buzz-relay/src/conformance/mod.rs +++ b/crates/buzz-relay/src/conformance/mod.rs @@ -435,6 +435,7 @@ pub fn sanitized_reason_for(err: &crate::handlers::ingest::IngestError) -> Sanit E::Rejected(_) => SanitizedReason::Invalid, E::CanvasConflict(_) => SanitizedReason::Invalid, E::AuthFailed(_) => SanitizedReason::Restricted, + E::RateLimited(_) => SanitizedReason::Restricted, E::Internal(_) => SanitizedReason::ServerError, } } diff --git a/crates/buzz-relay/src/handlers/command_executor.rs b/crates/buzz-relay/src/handlers/command_executor.rs index b1859146d86..6bb25c1ab85 100644 --- a/crates/buzz-relay/src/handlers/command_executor.rs +++ b/crates/buzz-relay/src/handlers/command_executor.rs @@ -1426,6 +1426,9 @@ mod postgres_tests { panic!("unexpected canvas conflict: {message}") } Err(IngestError::AuthFailed(message)) => panic!("unexpected auth failure: {message}"), + Err(IngestError::RateLimited(message)) => { + panic!("unexpected rate limit: {message}") + } Err(IngestError::Internal(message)) => panic!("unexpected internal failure: {message}"), Ok(_) => panic!("expected revision parsing to fail"), } diff --git a/crates/buzz-relay/src/handlers/event.rs b/crates/buzz-relay/src/handlers/event.rs index b5477ecef97..feb1c305ef2 100644 --- a/crates/buzz-relay/src/handlers/event.rs +++ b/crates/buzz-relay/src/handlers/event.rs @@ -686,7 +686,7 @@ pub async fn handle_event(event: Event, conn: Arc, state: Arc, state: Arc (message, "invalid"), IngestError::CanvasConflict(message) => (message, "invalid"), IngestError::AuthFailed(message) => (message, "auth"), + IngestError::RateLimited(message) => (message, "rate-limited"), IngestError::Internal(message) => (message, "error"), }; reject(reason); @@ -797,6 +798,7 @@ pub async fn handle_event(event: Event, conn: Arc, state: Arc (m.clone(), "invalid"), IngestError::CanvasConflict(m) => (m.clone(), "invalid"), IngestError::AuthFailed(m) => (m.clone(), "auth"), + IngestError::RateLimited(m) => (m.clone(), "rate-limited"), IngestError::Internal(_) => ("error: internal server error".to_string(), "error"), }; reject(reason); @@ -964,6 +966,9 @@ enum AgentObserverDirection { Control, } +const OBSERVER_FRESHNESS_REJECTION: &str = + "invalid: observer frame timestamp outside ±5 minute freshness window"; + #[derive(Debug, Clone, Copy)] struct AgentObserverRoute { agent: PublicKey, @@ -1004,7 +1009,6 @@ fn observer_frame_rate_limited( /// gates subscription in the REQ handler via the cleartext `p` tag. async fn handle_agent_observer_event( event: Event, - conn_id: uuid::Uuid, event_id_hex: &str, conn: Arc, state: Arc, @@ -1031,50 +1035,77 @@ async fn handle_agent_observer_event( } } - // Freshness check: reject observer frames with stale/future timestamps - let now = chrono::Utc::now().timestamp(); - let event_ts = event.created_at.as_secs() as i64; - if (event_ts - now).unsigned_abs() > 300 { - conn.send(RelayMessage::ok( - event_id_hex, - false, - "invalid: observer frame timestamp outside ±5 minute freshness window", - )); - return; + let session_owner = { + if let crate::connection::AuthState::Authenticated(ctx) = conn.auth_state_snapshot() { + ctx.agent_owner_pubkey + } else { + None + } + }; + let result = ingest_agent_observer_event(&state, &conn.tenant, &event, session_owner).await; + if let Err(error) = &result { + warn!( + conn_id = %conn.conn_id, + event_id = %event_id_hex, + error = ?error, + "Agent observer frame rejected" + ); } - - let route = match agent_observer_route(&event) { - Ok(Some(route)) => route, - Ok(None) => { - // Unknown frame value — silently drop, no error to publisher. + match result { + Ok(()) => { conn.send(RelayMessage::ok(event_id_hex, true, "")); - return; } - Err(message) => { + Err(IngestError::Rejected(message)) if message == OBSERVER_FRESHNESS_REJECTION => { + // The historic WS response rejects stale frames without counting + // them in the rejection metric. HTTP keeps its normal 400 mapping. + conn.send(RelayMessage::ok(event_id_hex, false, &message)); + } + Err(IngestError::Rejected(message)) => { reject("invalid"); conn.send(RelayMessage::ok(event_id_hex, false, &message)); - return; } - }; - - // Fast path: if this connection authenticated via NIP-OA and the verified - // owner matches the observer frame's target owner, skip the DB lookup entirely. - let session_owner_match = { - if let crate::connection::AuthState::Authenticated(ctx) = conn.auth_state_snapshot() { - ctx.agent_owner_pubkey.as_ref() == Some(&route.owner) - } else { - false + Err(IngestError::CanvasConflict(message)) => { + reject("invalid"); + conn.send(RelayMessage::ok(event_id_hex, false, &message)); + } + Err(IngestError::AuthFailed(message)) => { + reject("auth"); + conn.send(RelayMessage::ok(event_id_hex, false, &message)); + } + Err(IngestError::RateLimited(message)) => { + conn.send(RelayMessage::ok(event_id_hex, false, &message)); + } + Err(IngestError::Internal(message)) => { + conn.send(RelayMessage::ok(event_id_hex, false, &message)); } + } +} + +/// Validate and fan out one agent observer frame without persisting it. +pub(crate) async fn ingest_agent_observer_event( + state: &Arc, + tenant: &TenantContext, + event: &Event, + session_owner: Option, +) -> Result<(), IngestError> { + // Freshness check: reject observer frames with stale/future timestamps. + let now = chrono::Utc::now().timestamp(); + let event_ts = event.created_at.as_secs() as i64; + if (event_ts - now).unsigned_abs() > 300 { + return Err(IngestError::Rejected(OBSERVER_FRESHNESS_REJECTION.into())); + } + + let route = match agent_observer_route(event) { + Ok(Some(route)) => route, + Ok(None) => return Ok(()), + Err(message) => return Err(IngestError::Rejected(message)), }; + // Fast path: NIP-OA-authenticated sessions already carry the owner mapping. let agent_bytes = route.agent.to_bytes().to_vec(); let owner_bytes = route.owner.to_bytes().to_vec(); - let cache_key = ( - conn.tenant.community(), - agent_bytes.clone(), - owner_bytes.clone(), - ); - let is_owner = if session_owner_match { + let cache_key = (tenant.community(), agent_bytes.clone(), owner_bytes.clone()); + let is_owner = if session_owner.as_ref() == Some(&route.owner) { true } else { match state.observer_owner_cache.get(&cache_key) { @@ -1082,7 +1113,7 @@ async fn handle_agent_observer_event( None => { let result = state .db - .is_agent_owner(conn.tenant.community(), &agent_bytes, &owner_bytes) + .is_agent_owner(tenant.community(), &agent_bytes, &owner_bytes) .await; match result { Ok(v) => { @@ -1090,66 +1121,64 @@ async fn handle_agent_observer_event( v } Err(e) => { - warn!(conn_id = %conn_id, event_id = %event_id_hex, "agent observer owner check failed: {e}"); - conn.send(RelayMessage::ok( - event_id_hex, - false, - "error: internal server error", - )); - return; + warn!(event_id = %event.id.to_hex(), "agent observer owner check failed: {e}"); + return Err(IngestError::Internal("error: internal server error".into())); } } } } }; if !is_owner { - reject("auth"); - conn.send(RelayMessage::ok( - event_id_hex, - false, - "restricted: observer frame is not authorized for this agent owner", + return Err(IngestError::AuthFailed( + "restricted: observer frame is not authorized for this agent owner".into(), )); - return; } + // Control frames are owner commands and share the agent's budget key; keep + // telemetry bursts from consuming the budget needed for control delivery. // Rate limit telemetry frames only (100/sec per agent). - // Control frames (owner → agent) bypass the limiter — they are rare and must not - // be starved by bursty telemetry from the agent. if matches!(route.direction, AgentObserverDirection::Telemetry) { let agent_key: [u8; 32] = agent_bytes.as_slice().try_into().unwrap_or([0u8; 32]); - if observer_frame_rate_limited(&state, conn.tenant.community(), agent_key) { - conn.send(RelayMessage::ok( - event_id_hex, - false, - "rate-limited: observer frame rate exceeded (100/sec per agent)", + if observer_frame_rate_limited(state, tenant.community(), agent_key) { + return Err(IngestError::RateLimited( + "rate-limited: observer frame rate exceeded (100/sec per agent)".into(), )); - return; } } - state.mark_local_event(conn.tenant.community(), &event.id); + state.mark_local_event(tenant.community(), &event.id); if let Err(e) = state .pubsub - .publish_event(&conn.tenant, EventTopic::Global, &event) + .publish_event(tenant, EventTopic::Global, event) .await { state .local_event_ids - .invalidate(&(conn.tenant.community(), event.id.to_bytes())); - warn!(conn_id = %conn_id, event_id = %event_id_hex, "Agent observer publish failed: {e}"); + .invalidate(&(tenant.community(), event.id.to_bytes())); + warn!(event_id = %event.id.to_hex(), "Agent observer publish failed: {e}"); } let stored_event = StoredEvent::new(event.clone(), None); debug!( - event_id = %event_id_hex, + event_id = %event.id.to_hex(), agent = %route.agent.to_hex(), owner = %route.owner.to_hex(), direction = ?route.direction, "Agent observer fan-out" ); - fan_out_event_to_local_subscribers(&state, conn.tenant.community(), &stored_event).await; + fan_out_event_to_local_subscribers(state, tenant.community(), &stored_event).await; + Ok(()) +} - conn.send(RelayMessage::ok(event_id_hex, true, "")); +/// Owner pubkey named by a well-formed observer frame, if any. +/// +/// Used by the HTTP ingest path to apply the owner-to-agent ban cascade that +/// the WebSocket auth seam enforces structurally. +pub(crate) fn agent_observer_frame_owner(event: &Event) -> Option { + agent_observer_route(event) + .ok() + .flatten() + .map(|route| route.owner) } fn agent_observer_route(event: &Event) -> Result, String> { @@ -1464,14 +1493,7 @@ mod tests { grace_limit: 3, }); - super::handle_agent_observer_event( - event.clone(), - conn.conn_id, - &event.id.to_hex(), - conn, - state, - ) - .await; + super::handle_agent_observer_event(event.clone(), &event.id.to_hex(), conn, state).await; let axum::extract::ws::Message::Text(text) = send_rx.try_recv().expect("observer rejection sent") diff --git a/crates/buzz-relay/src/handlers/ingest.rs b/crates/buzz-relay/src/handlers/ingest.rs index 9b95218a026..a2aeafa5bce 100644 --- a/crates/buzz-relay/src/handlers/ingest.rs +++ b/crates/buzz-relay/src/handlers/ingest.rs @@ -12,29 +12,29 @@ use uuid::Uuid; use buzz_auth::Scope; use buzz_core::kind::{ event_kind_u32, is_identity_archive_request_kind, is_parameterized_replaceable, - is_relay_admin_kind, KIND_AGENT_ENGRAM, KIND_AGENT_PROFILE, KIND_AGENT_TURN_METRIC, - KIND_APPROVAL_DENY, KIND_APPROVAL_GRANT, KIND_AUTH, KIND_BOOKMARK_LIST, KIND_BOOKMARK_SET, - KIND_CANVAS, KIND_CONTACT_LIST, KIND_DELETION, KIND_DM_ADD_MEMBER, KIND_DM_HIDE, KIND_DM_OPEN, - KIND_EMOJI_LIST, KIND_EMOJI_SET, KIND_EVENT_REMINDER, KIND_FOLLOW_SET, KIND_FORUM_COMMENT, - KIND_FORUM_POST, KIND_FORUM_VOTE, KIND_GIFT_WRAP, KIND_GIT_ISSUE, KIND_GIT_PATCH, - KIND_GIT_PR_UPDATE, KIND_GIT_PULL_REQUEST, KIND_GIT_REPO_ANNOUNCEMENT, KIND_GIT_REPO_STATE, - KIND_GIT_STATUS_CLOSED, KIND_GIT_STATUS_DRAFT, KIND_GIT_STATUS_MERGED, KIND_GIT_STATUS_OPEN, - KIND_HUDDLE_ENDED, KIND_HUDDLE_GUIDELINES, KIND_HUDDLE_PARTICIPANT_JOINED, - KIND_HUDDLE_PARTICIPANT_LEFT, KIND_HUDDLE_STARTED, KIND_IA_ARCHIVE_REQUEST, - KIND_IA_UNARCHIVE_REQUEST, KIND_LONG_FORM, KIND_MANAGED_AGENT, KIND_MEMBER_ADDED_NOTIFICATION, - KIND_MEMBER_REMOVED_NOTIFICATION, KIND_MODERATION_BAN, KIND_MODERATION_RESOLVE_REPORT, - KIND_MODERATION_TIMEOUT, KIND_MODERATION_UNBAN, KIND_MODERATION_UNTIMEOUT, KIND_MUTE_LIST, - KIND_NIP29_CREATE_GROUP, KIND_NIP29_DELETE_EVENT, KIND_NIP29_DELETE_GROUP, - KIND_NIP29_EDIT_METADATA, KIND_NIP29_JOIN_REQUEST, KIND_NIP29_LEAVE_REQUEST, - KIND_NIP29_PUT_USER, KIND_NIP29_REMOVE_USER, KIND_NIP43_LEAVE_REQUEST, - KIND_NIP65_RELAY_LIST_METADATA, KIND_PERSONA, KIND_PIN_LIST, KIND_PRESENCE_UPDATE, - KIND_PRIVATE_MANAGED_AGENT, KIND_PRODUCT_FEEDBACK, KIND_PROFILE, KIND_PROJECT, KIND_REACTION, - KIND_READ_STATE, KIND_REPORT, KIND_STREAM_MESSAGE, KIND_STREAM_MESSAGE_BOOKMARKED, - KIND_STREAM_MESSAGE_DIFF, KIND_STREAM_MESSAGE_EDIT, KIND_STREAM_MESSAGE_PINNED, - KIND_STREAM_MESSAGE_SCHEDULED, KIND_STREAM_MESSAGE_V2, KIND_STREAM_REMINDER, KIND_TEAM, - KIND_TEAM_CATALOG, KIND_TEXT_NOTE, KIND_USER_STATUS, KIND_WORKFLOW_DEF, KIND_WORKFLOW_TRIGGER, - RELAY_ADMIN_ADD_MEMBER, RELAY_ADMIN_CHANGE_ROLE, RELAY_ADMIN_REMOVE_MEMBER, - RELAY_ADMIN_SET_WORKSPACE_PROFILE, + is_relay_admin_kind, KIND_AGENT_ENGRAM, KIND_AGENT_OBSERVER_FRAME, KIND_AGENT_PROFILE, + KIND_AGENT_TURN_METRIC, KIND_APPROVAL_DENY, KIND_APPROVAL_GRANT, KIND_AUTH, KIND_BOOKMARK_LIST, + KIND_BOOKMARK_SET, KIND_CANVAS, KIND_CONTACT_LIST, KIND_DELETION, KIND_DM_ADD_MEMBER, + KIND_DM_HIDE, KIND_DM_OPEN, KIND_EMOJI_LIST, KIND_EMOJI_SET, KIND_EVENT_REMINDER, + KIND_FOLLOW_SET, KIND_FORUM_COMMENT, KIND_FORUM_POST, KIND_FORUM_VOTE, KIND_GIFT_WRAP, + KIND_GIT_ISSUE, KIND_GIT_PATCH, KIND_GIT_PR_UPDATE, KIND_GIT_PULL_REQUEST, + KIND_GIT_REPO_ANNOUNCEMENT, KIND_GIT_REPO_STATE, KIND_GIT_STATUS_CLOSED, KIND_GIT_STATUS_DRAFT, + KIND_GIT_STATUS_MERGED, KIND_GIT_STATUS_OPEN, KIND_HUDDLE_ENDED, KIND_HUDDLE_GUIDELINES, + KIND_HUDDLE_PARTICIPANT_JOINED, KIND_HUDDLE_PARTICIPANT_LEFT, KIND_HUDDLE_STARTED, + KIND_IA_ARCHIVE_REQUEST, KIND_IA_UNARCHIVE_REQUEST, KIND_LONG_FORM, KIND_MANAGED_AGENT, + KIND_MEMBER_ADDED_NOTIFICATION, KIND_MEMBER_REMOVED_NOTIFICATION, KIND_MODERATION_BAN, + KIND_MODERATION_RESOLVE_REPORT, KIND_MODERATION_TIMEOUT, KIND_MODERATION_UNBAN, + KIND_MODERATION_UNTIMEOUT, KIND_MUTE_LIST, KIND_NIP29_CREATE_GROUP, KIND_NIP29_DELETE_EVENT, + KIND_NIP29_DELETE_GROUP, KIND_NIP29_EDIT_METADATA, KIND_NIP29_JOIN_REQUEST, + KIND_NIP29_LEAVE_REQUEST, KIND_NIP29_PUT_USER, KIND_NIP29_REMOVE_USER, + KIND_NIP43_LEAVE_REQUEST, KIND_NIP65_RELAY_LIST_METADATA, KIND_PERSONA, KIND_PIN_LIST, + KIND_PRESENCE_UPDATE, KIND_PRIVATE_MANAGED_AGENT, KIND_PRODUCT_FEEDBACK, KIND_PROFILE, + KIND_PROJECT, KIND_REACTION, KIND_READ_STATE, KIND_REPORT, KIND_STREAM_MESSAGE, + KIND_STREAM_MESSAGE_BOOKMARKED, KIND_STREAM_MESSAGE_DIFF, KIND_STREAM_MESSAGE_EDIT, + KIND_STREAM_MESSAGE_PINNED, KIND_STREAM_MESSAGE_SCHEDULED, KIND_STREAM_MESSAGE_V2, + KIND_STREAM_REMINDER, KIND_TEAM, KIND_TEAM_CATALOG, KIND_TEXT_NOTE, KIND_USER_STATUS, + KIND_WORKFLOW_DEF, KIND_WORKFLOW_TRIGGER, RELAY_ADMIN_ADD_MEMBER, RELAY_ADMIN_CHANGE_ROLE, + RELAY_ADMIN_REMOVE_MEMBER, RELAY_ADMIN_SET_WORKSPACE_PROFILE, }; use buzz_core::tenant::TenantContext; use buzz_core::verification::verify_event; @@ -449,6 +449,8 @@ pub enum IngestError { CanvasConflict(String), /// Auth/scope error — WS: OK false, HTTP: 401/403. AuthFailed(String), + /// A transient rate limit refusal; transports map this to retryable status. + RateLimited(String), /// Server error — WS: OK false, HTTP: 500. Internal(String), } @@ -507,6 +509,8 @@ fn required_scope_for_kind(kind: u32, event: &Event) -> Result Ok(Scope::MessagesWrite), + // NIP-AM: observer frames are agent-authored global events (encrypted to owner). + KIND_AGENT_OBSERVER_FRAME => Ok(Scope::MessagesWrite), // NIP-56 reports are ordinary member writes into the mod-only queue. // Ingest persists them to `moderation_reports` and suppresses public // storage/fanout; reports are signals, never enforcement triggers. @@ -760,6 +764,7 @@ pub(crate) fn is_global_only_kind(kind: u32) -> bool { // NIP-AM: agent turn metrics are owner-scoped global events. // Channel identity is encrypted inside the payload — no `h` tag. | KIND_AGENT_TURN_METRIC + | KIND_AGENT_OBSERVER_FRAME // NIP-PL leases are author-owned, addressable global state. | super::push_lease::KIND_PUSH_LEASE ) @@ -2155,6 +2160,66 @@ async fn author_type_label( } } +/// Reject when `pubkey` is banned in the community. Fails closed when the +/// restriction lookup errors. +async fn reject_if_banned( + state: &Arc, + tenant: &TenantContext, + pubkey: &nostr::PublicKey, +) -> Result<(), IngestError> { + match state + .db + .moderation_restriction_state(tenant.community(), pubkey.as_bytes()) + .await + { + Ok(r) if r.banned => Err(IngestError::AuthFailed( + "blocked: agent owner is banned from this community".to_string(), + )), + Ok(_) => Ok(()), + Err(e) => Err(IngestError::Internal(format!( + "error: internal error checking restriction state: {e}" + ))), + } +} + +/// Reject writes from a pubkey that is banned or currently timed out in the +/// community. Fails closed when the restriction lookup errors. +async fn enforce_write_restrictions( + state: &Arc, + tenant: &TenantContext, + pubkey: &nostr::PublicKey, +) -> Result<(), IngestError> { + match state + .db + .moderation_restriction_state(tenant.community(), pubkey.as_bytes()) + .await + { + Ok(r) => { + if r.banned { + return Err(IngestError::AuthFailed( + "blocked: you are banned from this community".to_string(), + )); + } + if let Some(until) = r.muted_until { + if until > chrono::Utc::now() { + return Err(IngestError::AuthFailed(format!( + "restricted: you are timed out until {}", + until.timestamp() + ))); + } + } + } + Err(e) => { + // Fail closed: a DB error must not let a banned/timed-out actor + // write. + return Err(IngestError::Internal(format!( + "error: internal error checking restriction state: {e}" + ))); + } + } + Ok(()) +} + /// Ingest a signed Nostr event through the full validation pipeline. /// /// Shared by WebSocket and HTTP transports. The caller constructs [`IngestAuth`] @@ -2178,7 +2243,8 @@ pub async fn ingest_event( // Captured before `event` moves into the inner fn: the stored-events // counter below is emitted at this shared seam so WebSocket and HTTP // transports are counted identically. - let kind_label = super::event::bounded_kind_label(event_kind_u32(&event)); + let kind_u32 = event_kind_u32(&event); + let kind_label = super::event::bounded_kind_label(kind_u32); // Classify the authenticated principal, not the event envelope signer: // NIP-59 gift wraps deliberately use an unrelated ephemeral pubkey. let author_pubkey_bytes = auth.principal_pubkey_bytes(); @@ -2197,7 +2263,7 @@ pub async fn ingest_event( // author_type is a 2-value label so it merely doubles the kind series). // Emitted here rather than per-transport so HTTP bridge ingests count too. if let Ok(r) = &result { - if r.accepted { + if should_count_as_stored(r.accepted, kind_u32) { let author_type = author_type_label(state, tenant, author_pubkey_bytes).await; metrics::counter!( "buzz_events_stored_total", @@ -2229,6 +2295,10 @@ pub async fn ingest_event( result } +fn should_count_as_stored(accepted: bool, kind: u32) -> bool { + accepted && kind != KIND_AGENT_OBSERVER_FRAME +} + /// Maximum seconds in the future a kind:40100 canvas event may be timestamped. /// Tighter than the general ±900 s drift window to prevent a ceiling-timestamped /// head from producing a write at `head + 1` that the relay would accept (the @@ -2373,6 +2443,34 @@ async fn ingest_event_inner( ))); } + if kind_u32 == KIND_AGENT_OBSERVER_FRAME { + // Observer frames return before the write-path restriction gate below, + // and NIP-98 auth never consults bans. Apply the same ban/timeout gate + // here so a restricted owner cannot keep sending control frames over + // HTTP. + enforce_write_restrictions(state, tenant, auth.pubkey()).await?; + // Agent-signed telemetry names its owner. Mirror the auth-seam cascade: + // an owner ban blocks the owner's agents (timeouts do not cascade). + if let Some(owner) = super::event::agent_observer_frame_owner(&event) { + if &owner != auth.pubkey() { + reject_if_banned(state, tenant, &owner).await?; + } + } + super::event::ingest_agent_observer_event(state, tenant, &event, None).await?; + emit( + tracer, + TraceAction::AcceptEphemeral { + msg_id: msg_id_label(event.id.as_bytes()), + }, + state_for_request(tenant, auth.pubkey()), + ); + return Ok(IngestResult { + event_id: event_id_hex, + accepted: true, + message: String::new(), + }); + } + // Command kinds are routed AFTER signature verification, timestamp check, // pubkey/auth match, and scope validation — never before. if buzz_core::kind::is_command_kind(kind_u32) { @@ -2458,34 +2556,7 @@ async fn ingest_event_inner( // the restriction-state cache (see should-fix), which can fold in owner // resolution without a per-write DB round-trip. if !buzz_core::kind::is_moderation_command_kind(kind_u32) && !is_relay_admin_kind(kind_u32) { - match state - .db - .moderation_restriction_state(tenant.community(), auth.pubkey().as_bytes()) - .await - { - Ok(r) => { - if r.banned { - return Err(IngestError::AuthFailed( - "blocked: you are banned from this community".to_string(), - )); - } - if let Some(until) = r.muted_until { - if until > chrono::Utc::now() { - return Err(IngestError::AuthFailed(format!( - "restricted: you are timed out until {}", - until.timestamp() - ))); - } - } - } - Err(e) => { - // Fail closed: a DB error must not let a banned/timed-out actor - // write. - return Err(IngestError::Internal(format!( - "error: internal error checking restriction state: {e}" - ))); - } - } + enforce_write_restrictions(state, tenant, auth.pubkey()).await?; } let mut channel_id = if kind_u32 == KIND_REACTION { @@ -3438,11 +3509,310 @@ mod postgres_tests { use super::*; use buzz_conformance::{TraceStep, Tracer}; use buzz_core::kind::{ - KIND_CANVAS, KIND_FORUM_COMMENT, KIND_FORUM_POST, KIND_FORUM_VOTE, KIND_LONG_FORM, - KIND_MANAGED_AGENT, KIND_PERSONA, KIND_PRESENCE_UPDATE, KIND_STREAM_MESSAGE, - KIND_STREAM_MESSAGE_DIFF, KIND_TEAM, KIND_USER_STATUS, + KIND_AGENT_OBSERVER_FRAME, KIND_CANVAS, KIND_FORUM_COMMENT, KIND_FORUM_POST, + KIND_FORUM_VOTE, KIND_LONG_FORM, KIND_MANAGED_AGENT, KIND_PERSONA, KIND_PRESENCE_UPDATE, + KIND_STREAM_MESSAGE, KIND_STREAM_MESSAGE_DIFF, KIND_TEAM, KIND_USER_STATUS, }; - use nostr::{EventBuilder, Kind}; + use nostr::{EventBuilder, Keys, Kind, Timestamp}; + + async fn observer_ingest_fixture() -> (Arc, TenantContext, Keys, Keys) { + let mut config = crate::config::Config::from_env().expect("load test config"); + config.require_relay_membership = false; + let pool = sqlx::PgPool::connect(&config.database_url) + .await + .expect("connect test DB"); + let db = buzz_db::Db::from_pool(pool.clone()); + if std::env::var("BUZZ_TEST_SCHEMA_MODE").as_deref() != Ok("desired") { + db.migrate().await.expect("migrate test DB"); + } + let host = format!("observer-ingest-{}.example", Uuid::new_v4().simple()); + let community = db + .ensure_configured_community(&host) + .await + .expect("community") + .id; + let redis_pool = deadpool_redis::Config::from_url(&config.redis_url) + .create_pool(Some(deadpool_redis::Runtime::Tokio1)) + .expect("redis pool"); + let pubsub = Arc::new( + buzz_pubsub::PubSubManager::new(&config.redis_url, redis_pool.clone()) + .await + .expect("pubsub manager"), + ); + let auth = buzz_auth::AuthService::new(config.auth.clone()); + let search = buzz_search::SearchService::new(pool.clone()); + let workflow_engine = Arc::new(buzz_workflow::WorkflowEngine::new( + db.clone(), + buzz_workflow::WorkflowConfig::default(), + )); + let media_storage = buzz_media::MediaStorage::new(&config.media).expect("media storage"); + let audit = buzz_audit::AuditService::new(pool); + let (state, _audit_shutdown) = AppState::new( + config, + db.clone(), + redis_pool, + audit, + pubsub, + auth, + search, + workflow_engine, + Keys::generate(), + media_storage, + ); + + let agent = Keys::generate(); + let owner = Keys::generate(); + db.ensure_user(community, &agent.public_key().to_bytes()) + .await + .expect("agent user"); + db.ensure_user(community, &owner.public_key().to_bytes()) + .await + .expect("owner user"); + db.set_agent_owner( + community, + &agent.public_key().to_bytes(), + &owner.public_key().to_bytes(), + ) + .await + .expect("agent owner"); + + ( + Arc::new(state), + TenantContext::resolved(community, host), + agent, + owner, + ) + } + + fn observer_event(agent: &Keys, owner: &Keys) -> nostr::Event { + let encrypted = buzz_core::observer::encrypt_observer_payload( + agent, + &owner.public_key(), + &serde_json::json!({"kind": "turn_started"}), + ) + .expect("encrypt observer payload"); + buzz_sdk::build_agent_observer_frame( + &owner.public_key().to_hex(), + &agent.public_key().to_hex(), + buzz_core::observer::OBSERVER_FRAME_TELEMETRY, + &encrypted, + ) + .expect("build observer frame") + .sign_with_keys(agent) + .expect("sign observer frame") + } + + async fn ingest_observer_test_event( + state: &Arc, + tenant: &TenantContext, + agent: &Keys, + event: nostr::Event, + ) -> Result { + ingest_event( + state, + tenant, + event, + IngestAuth::Http { + pubkey: agent.public_key(), + scopes: vec![Scope::MessagesWrite], + auth_method: HttpAuthMethod::Nip98, + }, + ) + .await + } + + #[tokio::test] + #[ignore = "requires Postgres and Redis"] + async fn http_observer_frame_is_accepted_and_not_stored() { + let (state, tenant, agent, owner) = observer_ingest_fixture().await; + let event = observer_event(&agent, &owner); + let event_id = event.id; + let result = ingest_observer_test_event(&state, &tenant, &agent, event) + .await + .expect("valid observer frame"); + assert!(result.accepted); + assert!(state + .db + .get_event_by_id(tenant.community(), event_id.as_bytes()) + .await + .expect("query event") + .is_none()); + } + + #[tokio::test] + #[ignore = "requires Postgres and Redis"] + async fn http_observer_frame_rejects_invalid_owner_and_envelope() { + let (state, tenant, agent, owner) = observer_ingest_fixture().await; + + let wrong_owner = Keys::generate(); + let wrong_owner_result = ingest_observer_test_event( + &state, + &tenant, + &agent, + observer_event(&agent, &wrong_owner), + ) + .await; + assert!( + matches!(wrong_owner_result, Err(IngestError::AuthFailed(message)) if message.contains("owner")) + ); + + let encrypted = buzz_core::observer::encrypt_observer_payload( + &agent, + &owner.public_key(), + &serde_json::json!({"kind": "turn_started"}), + ) + .expect("encrypt observer payload"); + let missing_agent = EventBuilder::new( + Kind::Custom(KIND_AGENT_OBSERVER_FRAME as u16), + encrypted.clone(), + ) + .tags([nostr::Tag::parse(["p", &owner.public_key().to_hex()]).expect("p tag")]) + .sign_with_keys(&agent) + .expect("sign missing-agent frame"); + assert!(matches!( + ingest_observer_test_event(&state, &tenant, &agent, missing_agent).await, + Err(IngestError::Rejected(message)) if message.contains("agent") + )); + + let p_is_agent = + EventBuilder::new(Kind::Custom(KIND_AGENT_OBSERVER_FRAME as u16), encrypted) + .tags([ + nostr::Tag::parse(["p", &agent.public_key().to_hex()]).expect("p tag"), + nostr::Tag::parse(["agent", &agent.public_key().to_hex()]).expect("agent tag"), + nostr::Tag::parse(["frame", buzz_core::observer::OBSERVER_FRAME_TELEMETRY]) + .expect("frame tag"), + ]) + .sign_with_keys(&agent) + .expect("sign self-owned frame"); + assert!(matches!( + ingest_observer_test_event(&state, &tenant, &agent, p_is_agent).await, + Err(IngestError::Rejected(_)) + )); + } + + #[tokio::test] + #[ignore = "requires Postgres and Redis"] + async fn http_observer_frame_rejects_stale_created_at() { + let (state, tenant, agent, owner) = observer_ingest_fixture().await; + let event = EventBuilder::new( + Kind::Custom(KIND_AGENT_OBSERVER_FRAME as u16), + observer_event(&agent, &owner).content, + ) + .custom_created_at(Timestamp::from(Timestamp::now().as_secs() - 301)) + .tags(observer_event(&agent, &owner).tags) + .sign_with_keys(&agent) + .expect("sign stale observer frame"); + assert!(matches!( + ingest_observer_test_event(&state, &tenant, &agent, event).await, + Err(IngestError::Rejected(message)) if message.contains("timestamp") + )); + } + + fn owner_control_event(agent: &Keys, owner: &Keys) -> nostr::Event { + let encrypted = buzz_core::observer::encrypt_observer_payload( + owner, + &agent.public_key(), + &serde_json::json!({"kind": "cancel"}), + ) + .expect("encrypt control payload"); + buzz_sdk::build_agent_observer_frame( + &agent.public_key().to_hex(), + &agent.public_key().to_hex(), + buzz_core::observer::OBSERVER_FRAME_CONTROL, + &encrypted, + ) + .expect("build control frame") + .sign_with_keys(owner) + .expect("sign control frame") + } + + #[tokio::test] + #[ignore = "requires Postgres and Redis"] + async fn http_observer_frame_rejects_banned_owner_control_frame() { + let (state, tenant, agent, owner) = observer_ingest_fixture().await; + assert!( + ingest_observer_test_event( + &state, + &tenant, + &owner, + owner_control_event(&agent, &owner) + ) + .await + .is_ok(), + "unbanned owner control frame should be accepted" + ); + + state + .db + .ban_community_member( + tenant.community(), + &owner.public_key().to_bytes(), + &Keys::generate().public_key().to_bytes(), + None, + None, + ) + .await + .expect("ban owner"); + assert!(matches!( + ingest_observer_test_event(&state, &tenant, &owner, owner_control_event(&agent, &owner)).await, + Err(IngestError::AuthFailed(message)) if message.contains("banned") + )); + } + + #[tokio::test] + #[ignore = "requires Postgres and Redis"] + async fn http_observer_frame_rejects_agent_telemetry_for_banned_owner() { + let (state, tenant, agent, owner) = observer_ingest_fixture().await; + state + .db + .ban_community_member( + tenant.community(), + &owner.public_key().to_bytes(), + &Keys::generate().public_key().to_bytes(), + None, + None, + ) + .await + .expect("ban owner"); + assert!(matches!( + ingest_observer_test_event(&state, &tenant, &agent, observer_event(&agent, &owner)).await, + Err(IngestError::AuthFailed(message)) if message.contains("owner is banned") + )); + } + + #[tokio::test] + #[ignore = "requires Postgres and Redis"] + async fn http_observer_frame_rejects_timed_out_owner_control_frame() { + let (state, tenant, agent, owner) = observer_ingest_fixture().await; + state + .db + .timeout_community_member( + tenant.community(), + &owner.public_key().to_bytes(), + &Keys::generate().public_key().to_bytes(), + chrono::Utc::now() + chrono::Duration::minutes(10), + None, + ) + .await + .expect("time out owner"); + assert!(matches!( + ingest_observer_test_event(&state, &tenant, &owner, owner_control_event(&agent, &owner)).await, + Err(IngestError::AuthFailed(message)) if message.contains("timed out") + )); + } + + #[tokio::test] + #[ignore = "requires Postgres and Redis"] + async fn http_unknown_ephemeral_kind_remains_restricted() { + let (state, tenant, agent, _owner) = observer_ingest_fixture().await; + let event = EventBuilder::new(Kind::Custom(20_099), "ephemeral") + .sign_with_keys(&agent) + .expect("sign unknown ephemeral event"); + assert!(matches!( + ingest_observer_test_event(&state, &tenant, &agent, event).await, + Err(IngestError::Rejected(message)) if message == "restricted: unknown event kind" + )); + } #[test] fn missing_huddle_backing_channel_is_a_client_rejection() { @@ -5681,6 +6051,13 @@ mod postgres_tests { ); } + #[test] + fn accepted_observer_frames_are_not_counted_as_stored() { + assert!(!should_count_as_stored(true, KIND_AGENT_OBSERVER_FRAME)); + assert!(should_count_as_stored(true, KIND_STREAM_MESSAGE)); + assert!(!should_count_as_stored(false, KIND_STREAM_MESSAGE)); + } + /// Boundary regression for the canvas-specific ingest future-timestamp guard. /// `validate_canvas_future_timestamp` is the pure seam; mutation: changing /// `CANVAS_MAX_INGEST_FUTURE_SECS` to 900 or removing the guard makes the diff --git a/crates/buzz-test-client/tests/e2e_relay.rs b/crates/buzz-test-client/tests/e2e_relay.rs index a801c1aaae6..51f6305e7a5 100644 --- a/crates/buzz-test-client/tests/e2e_relay.rs +++ b/crates/buzz-test-client/tests/e2e_relay.rs @@ -207,6 +207,244 @@ async fn create_test_channel(keys: &Keys) -> String { channel_uuid.to_string() } +/// HTTP observer frames are authorized by NIP-OA, fanned out to the owner's +/// `#p` subscription, and omitted from durable queries. +#[tokio::test] +#[ignore] +async fn test_http_agent_observer_frame_is_ephemeral_and_owner_scoped() { + let owner_keys = test_owner_keys(); + let wrong_owner_keys = Keys::generate(); + let agent_keys = Keys::generate(); + seed_relay_owner(&owner_keys).await; + + let mut owner_client = BuzzTestClient::connect(&relay_url(), &owner_keys) + .await + .expect("connect owner"); + let subscription_id = sub_id("agent-observer"); + let observer_filter = Filter::new() + .kinds(vec![Kind::Custom( + buzz_core::kind::KIND_AGENT_OBSERVER_FRAME as u16, + )]) + .custom_tag( + SingleLetterTag::lowercase(Alphabet::P), + owner_keys.public_key().to_hex(), + ); + owner_client + .subscribe(&subscription_id, vec![observer_filter]) + .await + .expect("subscribe to observer frames"); + owner_client + .collect_until_eose(&subscription_id, Duration::from_secs(5)) + .await + .expect("observer subscription EOSE"); + + let plaintext = serde_json::json!({"kind": "turn_started"}); + let encrypted = buzz_core::observer::encrypt_observer_payload( + &agent_keys, + &owner_keys.public_key(), + &plaintext, + ) + .expect("encrypt observer payload"); + let event = buzz_sdk::build_agent_observer_frame( + &owner_keys.public_key().to_hex(), + &agent_keys.public_key().to_hex(), + buzz_core::observer::OBSERVER_FRAME_TELEMETRY, + &encrypted, + ) + .expect("build observer frame") + .sign_with_keys(&agent_keys) + .expect("sign observer frame"); + let body = serde_json::to_string(&event).expect("serialize observer frame"); + let auth_tag = buzz_sdk::nip_oa::compute_auth_tag(&owner_keys, &agent_keys.public_key(), "") + .expect("compute owner auth tag"); + let http = relay_http_url(); + let response = reqwest::Client::new() + .post(format!("{http}/events")) + .header( + "Authorization", + nip98_post_header(&agent_keys, &format!("{http}/events"), &body), + ) + .header("x-auth-tag", &auth_tag) + .header("Content-Type", "application/json") + .body(body) + .send() + .await + .expect("POST observer frame"); + assert!( + response.status().is_success(), + "observer POST failed: {response:?}" + ); + let result: serde_json::Value = response.json().await.expect("observer response JSON"); + assert_eq!( + result["accepted"], true, + "observer frame rejected: {result}" + ); + + let received = owner_client + .recv_event(Duration::from_secs(5)) + .await + .expect("receive observer frame"); + let RelayMessage::Event { + event: received, .. + } = received + else { + panic!("expected observer EVENT, got {received:?}"); + }; + assert_eq!(received.id, event.id); + + let query = serde_json::json!([{ + "kinds": [buzz_core::kind::KIND_AGENT_OBSERVER_FRAME], + "authors": [agent_keys.public_key().to_hex()], + "#p": [owner_keys.public_key().to_hex()], + }]); + let query_response = reqwest::Client::new() + .post(format!("{http}/query")) + .header("X-Pubkey", owner_keys.public_key().to_hex()) + .header("Content-Type", "application/json") + .body(query.to_string()) + .send() + .await + .expect("query observer frame"); + assert!(query_response.status().is_success()); + let queried: Vec = query_response.json().await.expect("query JSON"); + assert!(queried.iter().all(|value| value["id"] != event.id.to_hex())); + + let wrong_event = buzz_sdk::build_agent_observer_frame( + &wrong_owner_keys.public_key().to_hex(), + &agent_keys.public_key().to_hex(), + buzz_core::observer::OBSERVER_FRAME_TELEMETRY, + &encrypted, + ) + .expect("build wrong-owner observer frame") + .sign_with_keys(&agent_keys) + .expect("sign wrong-owner observer frame"); + let wrong_body = serde_json::to_string(&wrong_event).expect("serialize wrong-owner frame"); + let wrong_auth_tag = + buzz_sdk::nip_oa::compute_auth_tag(&wrong_owner_keys, &agent_keys.public_key(), "") + .expect("compute wrong owner auth tag"); + let wrong_response = reqwest::Client::new() + .post(format!("{http}/events")) + .header( + "Authorization", + nip98_post_header(&agent_keys, &format!("{http}/events"), &wrong_body), + ) + .header("x-auth-tag", wrong_auth_tag) + .header("Content-Type", "application/json") + .body(wrong_body) + .send() + .await + .expect("POST wrong-owner observer frame"); + assert_eq!(wrong_response.status(), reqwest::StatusCode::FORBIDDEN); +} + +#[tokio::test] +#[ignore] +async fn test_http_observer_rate_limit_returns_429() { + let owner = test_owner_keys(); + let agent = Keys::generate(); + seed_relay_owner(&owner).await; + + let http = relay_http_url(); + let client = reqwest::Client::new(); + let auth_tag = buzz_sdk::nip_oa::compute_auth_tag(&owner, &agent.public_key(), "") + .expect("compute owner auth tag"); + + let mut bodies = Vec::with_capacity(101); + for sequence in 0..101 { + let encrypted = buzz_core::observer::encrypt_observer_payload( + &agent, + &owner.public_key(), + &serde_json::json!({"kind": "turn_started", "sequence": sequence}), + ) + .expect("encrypt observer payload"); + let event = buzz_sdk::build_agent_observer_frame( + &owner.public_key().to_hex(), + &agent.public_key().to_hex(), + buzz_core::observer::OBSERVER_FRAME_TELEMETRY, + &encrypted, + ) + .expect("build observer frame") + .sign_with_keys(&agent) + .expect("sign observer frame"); + let body = serde_json::to_string(&event).expect("serialize observer frame"); + bodies.push(body); + } + + let mut requests = tokio::task::JoinSet::new(); + for body in bodies { + let client = client.clone(); + let http = http.clone(); + let auth_tag = auth_tag.clone(); + let agent = agent.clone(); + requests.spawn(async move { + let response = client + .post(format!("{http}/events")) + .header( + "Authorization", + nip98_post_header(&agent, &format!("{http}/events"), &body), + ) + .header("x-auth-tag", auth_tag) + .header("Content-Type", "application/json") + .body(body) + .send() + .await + .expect("POST observer frame"); + response.status() + }); + } + + let mut accepted = 0; + let mut rate_limited = 0; + while let Some(result) = requests.join_next().await { + match result.expect("observer request task") { + status if status.is_success() => accepted += 1, + reqwest::StatusCode::TOO_MANY_REQUESTS => rate_limited += 1, + status => panic!("unexpected observer status: {status}"), + } + } + assert_eq!(accepted, 100); + assert_eq!(rate_limited, 1); +} + +#[tokio::test] +#[ignore] +async fn test_http_direct_relay_member_observer_frame_hits_known_owner_recording_gap() { + let agent = Keys::generate(); + let owner = Keys::generate(); + seed_relay_member(&relay_authority(), &agent, "member").await; + + let encrypted = buzz_core::observer::encrypt_observer_payload( + &agent, + &owner.public_key(), + &serde_json::json!({"kind": "turn_started"}), + ) + .expect("encrypt observer payload"); + let event = buzz_sdk::build_agent_observer_frame( + &owner.public_key().to_hex(), + &agent.public_key().to_hex(), + buzz_core::observer::OBSERVER_FRAME_TELEMETRY, + &encrypted, + ) + .expect("build observer frame") + .sign_with_keys(&agent) + .expect("sign observer frame"); + let body = serde_json::to_string(&event).expect("serialize observer frame"); + let http = relay_http_url(); + let response = reqwest::Client::new() + .post(format!("{http}/events")) + .header( + "Authorization", + nip98_post_header(&agent, &format!("{http}/events"), &body), + ) + .header("Content-Type", "application/json") + .body(body) + .send() + .await + .expect("POST observer frame without x-auth-tag"); + + assert_eq!(response.status(), reqwest::StatusCode::FORBIDDEN); +} + #[tokio::test] #[ignore] async fn test_connect_and_authenticate() { diff --git a/docs/nips/NIP-AO.md b/docs/nips/NIP-AO.md index 36adea04871..56c569259e0 100644 --- a/docs/nips/NIP-AO.md +++ b/docs/nips/NIP-AO.md @@ -154,12 +154,18 @@ subscribe attempts MUST be rejected with `AUTH required`. ## Relay Behavior -On receiving a kind 24200 event, a relay MUST: +On receiving a kind 24200 event over WebSocket or the HTTP `POST /events` +bridge, a relay MUST apply the same signature, freshness, envelope, and +agent-owner authorization checks, then: 1. Validate the event signature per NIP-01. 2. Verify authorization per the rules above. -3. Fan out to matching subscribers via in-memory pub/sub. -4. NOT invoke the normal event ingestion or persistence path. +3. Use in-memory pub/sub to fan out only to matching subscribers. +4. Do not invoke the normal event ingestion or persistence path. + +The HTTP bridge uses the same validation and authorization as WebSocket. It +returns an HTTP error for rejected frames and only accepts frames for live +fan-out; accepted frames are never persisted or made queryable. Relays SHOULD enforce a rate limit of 100 events/second per agent pubkey. Relays are RECOMMENDED to reject events whose `created_at` falls outside a ±5-minute