diff --git a/crates/buzz-relay/src/api/bridge.rs b/crates/buzz-relay/src/api/bridge.rs index ff4265077f..8372e49a2a 100644 --- a/crates/buzz-relay/src/api/bridge.rs +++ b/crates/buzz-relay/src/api/bridge.rs @@ -587,6 +587,28 @@ fn event_in_accessible_channel(se: &buzz_core::StoredEvent, accessible: &[uuid:: } } +/// Hard cap on the `reason` field logged for a rejected `/events` request. +/// +/// The reject message can embed event-controlled content (e.g. a submitted +/// channel's `visibility`/`channel_type` tag values, or a raw tag pubkey) — +/// attacker-controlled text that must never reach Datadog unbounded. +const REJECT_REASON_MAX_BYTES: usize = 256; + +/// Truncate `s` to at most `max_bytes`, cutting at the nearest UTF-8 character +/// boundary so a multi-byte codepoint straddling the cutoff is never split. +/// Bounds attacker-controlled text before it enters a structured log line — +/// the line's size must stay bounded regardless of the triggering input size. +fn truncate_reason(s: &str, max_bytes: usize) -> &str { + if s.len() <= max_bytes { + return s; + } + let mut end = max_bytes; + while !s.is_char_boundary(end) { + end -= 1; + } + &s[..end] +} + /// Submit a signed Nostr event via HTTP bridge (NIP-98 auth). pub async fn submit_event( State(state): State>, @@ -618,49 +640,237 @@ pub async fn submit_event( Some(&body), state.config.require_auth_token, )?; - enforce_http_admission(&state, &tenant, &pubkey).await?; - check_nip98_replay(&state, &tenant, event_id_bytes).await?; + let pubkey_hex = pubkey.to_hex(); + + // Everything after auth — admission, replay, membership, parse, ingest — + // runs inside the helper. The thin wrapper here owns the single terminal + // attribution line so it fires for every outcome, including admission/ + // replay/membership failures that previously returned before any log fired. + let outcome = + submit_event_authed(&state, &tenant, &headers, &body, pubkey, event_id_bytes).await; + + match &outcome { + SubmitOutcome::Ok { accepted, .. } => { + tracing::info!( + pubkey = %pubkey_hex, + route = "/events", + status = 200u16, + accepted, + "HTTP bridge request" + ); + } + SubmitOutcome::ParseFail { + category, + line, + column, + .. + } => { + tracing::warn!( + pubkey = %pubkey_hex, + route = "/events", + status = 400u16, + accepted = false, + category, + line, + column, + "HTTP bridge request" + ); + } + SubmitOutcome::Rejected { kind, reason, .. } => { + tracing::warn!( + pubkey = %pubkey_hex, + route = "/events", + status = 400u16, + accepted = false, + kind, + reason = %reason, + "HTTP bridge request" + ); + } + SubmitOutcome::Err { status, .. } => { + tracing::warn!( + pubkey = %pubkey_hex, + route = "/events", + status = status.as_u16(), + accepted = false, + "HTTP bridge request" + ); + } + } + + outcome.into_response() +} + +/// Log-context outcome for a single [`submit_event`] call. +/// +/// Carries enough structured data for the terminal attribution log while also +/// holding the HTTP response so the thin wrapper can return it unchanged. +enum SubmitOutcome { + /// Ingest pipeline ran and returned a result (accepted or not). + Ok { + accepted: bool, + response: Json, + }, + /// JSON parse failure before ingest — log category/line/column, not msg. + ParseFail { + category: &'static str, + line: usize, + column: usize, + response: (StatusCode, Json), + }, + /// IngestError::Rejected — log kind + truncated reason. + Rejected { + kind: u32, + reason: String, + response: (StatusCode, Json), + }, + /// Any other error (admission, replay, membership, auth, internal) — + /// only the HTTP status is logged; the response body is returned as-is. + Err { + status: StatusCode, + response: (StatusCode, Json), + }, +} + +impl SubmitOutcome { + fn into_response(self) -> Result, (StatusCode, Json)> { + match self { + SubmitOutcome::Ok { response, .. } => Ok(response), + SubmitOutcome::ParseFail { response, .. } => Err(response), + SubmitOutcome::Rejected { response, .. } => Err(response), + SubmitOutcome::Err { response, .. } => Err(response), + } + } +} + +/// Post-auth execution for [`submit_event`]: admission, replay, membership, +/// parse, and ingest. Returns a [`SubmitOutcome`] that carries both the log +/// fields and the HTTP response so the thin wrapper can emit exactly one +/// terminal attribution line covering every outcome. +async fn submit_event_authed( + state: &Arc, + tenant: &TenantContext, + headers: &HeaderMap, + body: &[u8], + pubkey: nostr::PublicKey, + event_id_bytes: [u8; 32], +) -> SubmitOutcome { + // Admission and replay checks fire before body parse — a 429 or replay + // reject on a malformed body must still be attributed. + if let Err(e) = enforce_http_admission(state, tenant, &pubkey).await { + return SubmitOutcome::Err { + status: e.0, + response: e, + }; + } + if let Err(e) = check_nip98_replay(state, tenant, event_id_bytes).await { + return SubmitOutcome::Err { + status: e.0, + response: e, + }; + } let pubkey_bytes = pubkey.to_bytes().to_vec(); - let event: nostr::Event = serde_json::from_slice(&body) - .map_err(|e| api_error(StatusCode::BAD_REQUEST, &format!("invalid event JSON: {e}")))?; + let event: nostr::Event = match serde_json::from_slice(body) { + Ok(ev) => ev, + Err(e) => { + // Never log `e`'s Display string: serde_json embeds the offending + // input verbatim in its error message, so a malformed field of + // arbitrary size (the router allows 1 MiB bodies) would otherwise + // reflect attacker-controlled text into a log line at full size. + // `category`/`line`/`column` are bounded, structured, and still + // enough to tell "which parse failure" apart at a glance. + crate::handlers::ingest::reject_with_transport("http", "invalid"); + return SubmitOutcome::ParseFail { + category: match e.classify() { + serde_json::error::Category::Io => "io", + serde_json::error::Category::Syntax => "syntax", + serde_json::error::Category::Data => "data", + serde_json::error::Category::Eof => "eof", + }, + line: e.line(), + column: e.column(), + response: api_error(StatusCode::BAD_REQUEST, &format!("invalid event JSON: {e}")), + }; + } + }; + // Enforce relay membership (with NIP-OA fallback via x-auth-tag header). let auth_tag = headers.get("x-auth-tag").and_then(|v| v.to_str().ok()); - let nip_oa_owner = super::relay_members::enforce_relay_membership( - &state, + let nip_oa_owner = match super::relay_members::enforce_relay_membership( + state, tenant.community(), &pubkey_bytes, auth_tag, ) - .await? - .or_else(|| { - if !state.config.require_relay_membership { - super::relay_members::extract_nip_oa_owner(&pubkey_bytes, auth_tag) - } else { - None + .await + { + Ok(owner) => owner.or_else(|| { + if !state.config.require_relay_membership { + super::relay_members::extract_nip_oa_owner(&pubkey_bytes, auth_tag) + } else { + None + } + }), + Err(e) => { + return SubmitOutcome::Err { + status: e.0, + response: e, + }; } - }); + }; if let Some(owner) = nip_oa_owner { - super::relay_members::materialize_nip_oa_owner(&state, &tenant, &pubkey, &owner).await; + super::relay_members::materialize_nip_oa_owner(state, tenant, &pubkey, &owner).await; } + let kind_u32 = buzz_core::kind::event_kind_u32(&event); let auth = IngestAuth::Http { pubkey, scopes: buzz_auth::Scope::all_known(), // Pure Nostr: full scopes, channel access via membership auth_method: crate::handlers::ingest::HttpAuthMethod::Nip98, }; - match crate::handlers::ingest::ingest_event(&state, &tenant, event, auth).await { - Ok(result) => Ok(Json(serde_json::json!({ - "event_id": result.event_id, - "accepted": result.accepted, - "message": result.message, - }))), - Err(e) => match e { - IngestError::Rejected(msg) => Err(api_error(StatusCode::BAD_REQUEST, &msg)), - IngestError::AuthFailed(msg) => Err(api_error(StatusCode::FORBIDDEN, &msg)), - IngestError::Internal(msg) => Err(internal_error(&msg)), - }, + match crate::handlers::ingest::ingest_event(state, tenant, event, auth).await { + Ok(result) => { + let response = Json(serde_json::json!({ + "event_id": result.event_id, + "accepted": result.accepted, + "message": result.message, + })); + SubmitOutcome::Ok { + accepted: result.accepted, + response, + } + } + Err(IngestError::Rejected(msg)) => { + // `msg` can embed event-controlled content (e.g. a channel + // create's raw `visibility`/`channel_type` tag values, or a raw + // tag pubkey) — truncate before logging, but return the full msg + // in the HTTP response body (unchanged from prior behaviour). + let reason = truncate_reason(&msg, REJECT_REASON_MAX_BYTES).to_owned(); + crate::handlers::ingest::reject_with_transport("http", "invalid"); + SubmitOutcome::Rejected { + kind: kind_u32, + reason, + response: api_error(StatusCode::BAD_REQUEST, &msg), + } + } + Err(IngestError::AuthFailed(msg)) => { + crate::handlers::ingest::reject_with_transport("http", "auth"); + let e = api_error(StatusCode::FORBIDDEN, &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); + SubmitOutcome::Err { + status: e.0, + response: e, + } + } } } @@ -698,13 +908,57 @@ pub async fn query_events( Some(&body), state.config.require_auth_token, )?; - enforce_http_admission(&state, &tenant, &pubkey).await?; - check_nip98_replay(&state, &tenant, event_id_bytes).await?; + let pubkey_hex = pubkey.to_hex(); + + // Admission, replay, membership, and filter execution all run inside the + // helper. The single terminal attribution line fires here from the Result + // so every outcome — including admission/replay/membership failures that + // previously returned before any log — is attributed. + let result = + query_events_authed(&state, &tenant, &headers, &body, pubkey, event_id_bytes).await; + match &result { + Ok(Json(Value::Array(events))) => { + tracing::info!( + pubkey = %pubkey_hex, + route = "/query", + status = 200u16, + result_count = events.len(), + "HTTP bridge request" + ); + } + Ok(_) => { + tracing::info!(pubkey = %pubkey_hex, route = "/query", status = 200u16, "HTTP bridge request"); + } + Err((status, _)) => { + tracing::warn!( + pubkey = %pubkey_hex, + route = "/query", + status = status.as_u16(), + "HTTP bridge request" + ); + } + } + result +} + +/// Filter execution for [`query_events`], run once NIP-98 auth succeeds. +/// Handles admission, replay, membership, and all filter paths so the thin +/// wrapper above can emit exactly one terminal attribution line from the Result. +async fn query_events_authed( + state: &Arc, + tenant: &TenantContext, + headers: &HeaderMap, + body: &[u8], + pubkey: nostr::PublicKey, + event_id_bytes: [u8; 32], +) -> Result, (StatusCode, Json)> { + enforce_http_admission(state, tenant, &pubkey).await?; + check_nip98_replay(state, tenant, event_id_bytes).await?; let pubkey_bytes = pubkey.to_bytes().to_vec(); let auth_tag = headers.get("x-auth-tag").and_then(|v| v.to_str().ok()); super::relay_members::enforce_relay_membership( - &state, + state, tenant.community(), &pubkey_bytes, auth_tag, @@ -713,7 +967,7 @@ pub async fn query_events( // Two-pass parse: preserve raw JSON for custom extension fields (before_id, // depth_limit, feed_types) that nostr::Filter silently drops. - let raw_filters: Vec = serde_json::from_slice(&body) + let raw_filters: Vec = serde_json::from_slice(body) .map_err(|e| api_error(StatusCode::BAD_REQUEST, &format!("invalid filters: {e}")))?; let filters: Vec = raw_filters .iter() @@ -757,18 +1011,18 @@ pub async fn query_events( )); } return handle_bridge_search( - &state, + state, &raw_filters, &filters, &accessible_channels, - &tenant, + tenant, &authed_pubkey_hex, &pubkey_bytes, ) .await; } - if let Some(presence_events) = synthesize_presence(&state, &tenant, &filters).await { + if let Some(presence_events) = synthesize_presence(state, tenant, &filters).await { return Ok(Json(Value::Array(presence_events))); } @@ -782,8 +1036,8 @@ pub async fn query_events( continue; } handle_channel_window_filter( - &state, - &tenant, + state, + tenant, raw, filter, &accessible_channels, @@ -961,7 +1215,7 @@ pub async fn query_events( let mut query = crate::handlers::req::build_event_query_from_filter( filter, &pubkey_bytes, - &state, + state, tenant.community(), ) .await; @@ -1087,20 +1341,62 @@ pub async fn count_events( Some(&body), state.config.require_auth_token, )?; - enforce_http_admission(&state, &tenant, &pubkey).await?; - check_nip98_replay(&state, &tenant, event_id_bytes).await?; + let pubkey_hex = pubkey.to_hex(); + + // Admission, replay, membership, and count execution all run inside the + // helper. The single terminal attribution line fires here from the Result + // so every outcome — including admission/replay/membership failures that + // previously returned before any log — is attributed. + let result = + count_events_authed(&state, &tenant, &headers, &body, pubkey, event_id_bytes).await; + match &result { + Ok(Json(value)) => { + let count = value.get("count").and_then(Value::as_u64); + tracing::info!( + pubkey = %pubkey_hex, + route = "/count", + status = 200u16, + result_count = count, + "HTTP bridge request" + ); + } + Err((status, _)) => { + tracing::warn!( + pubkey = %pubkey_hex, + route = "/count", + status = status.as_u16(), + "HTTP bridge request" + ); + } + } + result +} + +/// Filter execution for [`count_events`], run once NIP-98 auth succeeds. +/// Handles admission, replay, membership, and count execution so the thin +/// wrapper above can emit exactly one terminal attribution line from the Result. +async fn count_events_authed( + state: &Arc, + tenant: &TenantContext, + headers: &HeaderMap, + body: &[u8], + pubkey: nostr::PublicKey, + event_id_bytes: [u8; 32], +) -> Result, (StatusCode, Json)> { + enforce_http_admission(state, tenant, &pubkey).await?; + check_nip98_replay(state, tenant, event_id_bytes).await?; let pubkey_bytes = pubkey.to_bytes().to_vec(); let auth_tag = headers.get("x-auth-tag").and_then(|v| v.to_str().ok()); super::relay_members::enforce_relay_membership( - &state, + state, tenant.community(), &pubkey_bytes, auth_tag, ) .await?; - let filters: Vec = serde_json::from_slice(&body) + let filters: Vec = serde_json::from_slice(body) .map_err(|e| api_error(StatusCode::BAD_REQUEST, &format!("invalid filters: {e}")))?; // P-gated kinds enforcement — same as WS REQ and /query. @@ -1153,7 +1449,7 @@ pub async fn count_events( let query = crate::handlers::req::build_event_query_from_filter( filter, &pubkey_bytes, - &state, + state, tenant.community(), ) .await; @@ -1215,7 +1511,7 @@ pub async fn count_events( let mut query = crate::handlers::req::build_event_query_from_filter( filter, &pubkey_bytes, - &state, + state, tenant.community(), ) .await; @@ -1889,6 +2185,7 @@ fn ban_json(b: &buzz_db::moderation::BanRecord) -> Value { mod tests { use super::*; use nostr::{Alphabet, EventBuilder, Keys, Kind, SingleLetterTag, Tag}; + use std::sync::Mutex; fn redis_pool() -> deadpool_redis::Pool { let url = std::env::var("REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1:6379".into()); @@ -2900,4 +3197,494 @@ mod tests { "owner must still receive their own snapshot" ); } + + // ────────────────────────────────────────────────────────────────────────── + // truncate_reason regression tests + // + // Required by T1: prove a near-limit malformed input cannot produce a + // near-limit log payload. No infrastructure needed — pure unit tests. + // ────────────────────────────────────────────────────────────────────────── + + /// Near-limit ASCII input (1 MiB) must be capped at exactly `max_bytes`. + #[test] + fn truncate_reason_large_ascii_input_is_bounded() { + let big = "a".repeat(1024 * 1024); // 1 MiB + let result = truncate_reason(&big, 256); + assert_eq!(result.len(), 256); + assert!(result.is_ascii()); + } + + /// A multi-byte codepoint that straddles the byte boundary must not be + /// split: the output must end at the last complete codepoint before the + /// cap, keeping the slice valid UTF-8. + #[test] + fn truncate_reason_multibyte_codepoint_at_boundary_is_not_split() { + // Build a string of 255 ASCII bytes followed by a 3-byte codepoint (€, U+20AC). + // Naïve cutoff at byte 256 would land inside the 3-byte sequence. + let mut s = "a".repeat(255); + s.push('€'); // 3 bytes: 0xE2 0x82 0xAC — spans bytes 255..258 + assert_eq!(s.len(), 258); + + let result = truncate_reason(&s, 256); + // Must end before the multi-byte codepoint. + assert_eq!(result.len(), 255); + assert_eq!(result, "a".repeat(255)); + assert!(std::str::from_utf8(result.as_bytes()).is_ok()); + } + + /// Short input well under the cap must be returned unchanged. + #[test] + fn truncate_reason_short_input_returned_unchanged() { + let s = "invalid: kind 24620 rejected"; + let result = truncate_reason(s, 256); + assert_eq!(result, s); + } + + // ────────────────────────────────────────────────────────────────────────── + // Handler-level tests: submit_event HTTP-counter seam + // + // These tests drive real HTTP requests through the axum router to prove + // that the bridge code path (not just the shared helper) actually + // increments buzz_events_rejected_total{transport="http"}. They are + // discriminating: removing either bridge call site causes the corresponding + // test to fail. + // + // Why `#[ignore = "requires Postgres"]`: submit_event calls bind_community + // (needs a real communities row) and enforce_http_admission (needs Redis). + // Both services run locally in dev. The admission check succeeds for a + // freshly generated pubkey (first request, far below quota), so no Redis + // seeding is required beyond the default local instance. + // + // Why `#[test]` + manual runtime instead of `#[tokio::test]`: + // metrics::with_local_recorder stores the recorder in a thread-local. An + // async test uses a multi-thread scheduler by default; when submit_event + // runs, it may land on a different thread and miss the recorder entirely. + // Using a current_thread runtime with rt.block_on() inside the recorder + // closure guarantees the handler runs on the same thread as the recorder. + // ────────────────────────────────────────────────────────────────────────── + + struct AlwaysFreshReplayGuard; + + impl Nip98ReplayGuard for AlwaysFreshReplayGuard { + fn try_mark_in_scope<'a>( + &'a self, + _scope: &'a str, + _event_id: &'a nostr::EventId, + _ttl_secs: u64, + ) -> std::pin::Pin< + Box> + Send + 'a>, + > { + Box::pin(async { Ok(true) }) + } + } + + const TEST_DB_URL: &str = "postgres://buzz:buzz_dev@localhost:5432/buzz"; // sadscan:disable np.postgres.1 + + /// Build an AppState suitable for handler-level bridge tests. + /// + /// - `require_auth_token = false` → X-Pubkey dev-mode fallback active. + /// - `require_relay_membership = false` → membership check short-circuits to + /// OpenRelay without a DB lookup. + /// - `nip98_replay` replaced with an always-fresh guard → no Redis needed + /// for replay detection. + /// - Redis pool points at the local dev instance for the admission check. + /// + /// Returns `None` when local Postgres is not reachable. + async fn bridge_handler_test_state() -> Option> { + let mut config = crate::config::Config::from_env().ok()?; + config.database_url = TEST_DB_URL.to_string(); + // Use the real local Redis so enforce_http_admission can pass. + config.redis_url = + std::env::var("REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string()); + config.relay_url = "wss://bridge-test.local".to_string(); + config.require_auth_token = false; + config.require_relay_membership = false; + + let pool = sqlx::PgPool::connect(TEST_DB_URL).await.ok()?; + let db = buzz_db::Db::from_pool(pool.clone()); + let redis_pool = deadpool_redis::Config::from_url(&config.redis_url) + .create_pool(Some(deadpool_redis::Runtime::Tokio1)) + .ok()?; + let pubsub = Arc::new( + buzz_pubsub::PubSubManager::new(&config.redis_url, redis_pool.clone()) + .await + .ok()?, + ); + let audit = buzz_audit::AuditService::new(pool.clone()); + 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).ok()?; + + let (mut state, _audit_shutdown) = crate::state::AppState::new( + config, + db, + redis_pool, + audit, + pubsub, + auth, + search, + workflow_engine, + Keys::generate(), + media_storage, + ); + state.nip98_replay = Arc::new(AlwaysFreshReplayGuard); + Some(Arc::new(state)) + } + + /// Drive a single POST /events request through the router and return the + /// HTTP status code. + async fn post_events( + state: Arc, + host: &str, + pubkey_hex: &str, + body: &[u8], + ) -> axum::http::StatusCode { + use axum::body::Body; + use axum::http::{header, Request}; + use tower::ServiceExt; + + crate::router::build_router(state) + .oneshot( + Request::builder() + .method("POST") + .uri("/events") + .header(header::HOST, host) + .header("x-pubkey", pubkey_hex) + .body(Body::from(body.to_vec())) + .expect("build request"), + ) + .await + .expect("router oneshot") + .status() + } + + /// Collect buzz_events_rejected_total with (transport, reason) labels from + /// a DebuggingRecorder snapshot. + fn http_reject_counts( + snapshotter: &metrics_util::debugging::Snapshotter, + ) -> std::collections::HashMap<(String, String), u64> { + snapshotter + .snapshot() + .into_vec() + .into_iter() + .filter(|(key, ..)| key.key().name() == "buzz_events_rejected_total") + .map(|(key, _, _, value)| { + let metrics_util::debugging::DebugValue::Counter(n) = value else { + panic!("buzz_events_rejected_total must be a counter"); + }; + let labels: Vec<_> = key.key().labels().collect(); + let transport = labels + .iter() + .find(|l| l.key() == "transport") + .map(|l| l.value().to_owned()) + .unwrap_or_default(); + let reason = labels + .iter() + .find(|l| l.key() == "reason") + .map(|l| l.value().to_owned()) + .unwrap_or_default(); + ((transport, reason), n) + }) + .collect() + } + + /// T2a — pre-parse 400 arm: a POST /events with an invalid JSON body must + /// increment buzz_events_rejected_total{transport="http",reason="invalid"}. + /// + /// Discriminating: if the `reject_with_transport` call in bridge.rs's + /// `serde_json::from_slice` map_err closure is removed, this test fails. + #[test] + #[ignore = "requires Postgres"] + fn submit_event_invalid_json_body_increments_http_transport_counter() { + let rt = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("current_thread runtime"); + + let Some(state) = rt.block_on(bridge_handler_test_state()) else { + panic!("local Postgres not reachable — start Postgres on 127.0.0.1:5432 before running ignored bridge handler tests"); + }; + + // Provision a fresh community so bind_community succeeds. + let host = { + let h = format!("bridge-test-{}.local", uuid::Uuid::new_v4().simple()); + rt.block_on(state.db.ensure_configured_community(&h)) + .expect("ensure community"); + h + }; + + let pubkey_hex = Keys::generate().public_key().to_hex(); + + let recorder = metrics_util::debugging::DebuggingRecorder::new(); + let snapshotter = recorder.snapshotter(); + + metrics::with_local_recorder(&recorder, || { + let status = rt.block_on(post_events( + state.clone(), + &host, + &pubkey_hex, + b"not valid json at all", + )); + assert_eq!( + status, + axum::http::StatusCode::BAD_REQUEST, + "malformed body must yield 400" + ); + }); + + let counts = http_reject_counts(&snapshotter); + assert_eq!( + counts.get(&("http".to_owned(), "invalid".to_owned())), + Some(&1), + "pre-parse 400 arm must increment transport=http,reason=invalid" + ); + } + + /// T2b — post-parse IngestError::Rejected arm: a POST /events with a + /// valid but relay-only-kind event (kind 13534 = membership snapshot) must + /// increment buzz_events_rejected_total{transport="http",reason="invalid"}. + /// + /// Kind 13534 is rejected in ingest_event before signature verification, + /// so any properly signed Nostr event of this kind triggers the arm. + /// + /// Discriminating: if the `reject_with_transport` call in bridge.rs's + /// IngestError::Rejected match arm is removed, this test fails. + #[test] + #[ignore = "requires Postgres"] + fn submit_event_relay_only_kind_increments_http_transport_counter() { + let rt = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("current_thread runtime"); + + let Some(state) = rt.block_on(bridge_handler_test_state()) else { + panic!("local Postgres not reachable — start Postgres on 127.0.0.1:5432 before running ignored bridge handler tests"); + }; + + let host = { + let h = format!("bridge-test-{}.local", uuid::Uuid::new_v4().simple()); + rt.block_on(state.db.ensure_configured_community(&h)) + .expect("ensure community"); + h + }; + + let client_keys = Keys::generate(); + let pubkey_hex = client_keys.public_key().to_hex(); + + // Kind 13534 (membership snapshot) is relay-only; ingest_event rejects + // it before reaching signature verification. + let relay_only_event = EventBuilder::new( + Kind::Custom(buzz_core::kind::KIND_NIP43_MEMBERSHIP_LIST as u16), + "", + ) + .sign_with_keys(&client_keys) + .expect("sign relay-only event"); + let event_json = serde_json::to_vec(&relay_only_event).expect("serialize event"); + + let recorder = metrics_util::debugging::DebuggingRecorder::new(); + let snapshotter = recorder.snapshotter(); + + metrics::with_local_recorder(&recorder, || { + let status = rt.block_on(post_events(state.clone(), &host, &pubkey_hex, &event_json)); + assert_eq!( + status, + axum::http::StatusCode::BAD_REQUEST, + "relay-only-kind event must yield 400" + ); + }); + + let counts = http_reject_counts(&snapshotter); + assert_eq!( + counts.get(&("http".to_owned(), "invalid".to_owned())), + Some(&1), + "IngestError::Rejected arm must increment transport=http,reason=invalid" + ); + } + + // ────────────────────────────────────────────────────────────────────────── + // Log-capture helpers and attribution-invariant tests + // + // These tests assert that exactly ONE "HTTP bridge request" log line + // appears per authed request, pinning the W1 single-terminal-log invariant + // and the E1 no-double-log rule. + // + // Infrastructure: same `#[ignore = "requires Postgres"]` + current_thread + // runtime discipline as the HTTP-counter tests above. + // ────────────────────────────────────────────────────────────────────────── + + /// Shared buffer that collects all bytes written by the tracing fmt layer. + #[derive(Clone)] + struct CapturingMakeWriter { + buf: Arc>>, + } + + struct CapturingWriter { + buf: Arc>>, + } + + impl std::io::Write for CapturingWriter { + fn write(&mut self, data: &[u8]) -> std::io::Result { + self.buf.lock().unwrap().extend_from_slice(data); + Ok(data.len()) + } + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + + impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for CapturingMakeWriter { + type Writer = CapturingWriter; + fn make_writer(&'a self) -> Self::Writer { + CapturingWriter { + buf: Arc::clone(&self.buf), + } + } + } + + /// Run `post_events` with a capturing subscriber and return `(status_code, captured_log)`. + fn run_and_capture( + rt: &tokio::runtime::Runtime, + state: Arc, + host: &str, + pubkey_hex: &str, + body: &[u8], + ) -> (axum::http::StatusCode, String) { + let buf = Arc::new(Mutex::new(Vec::::new())); + let make_writer = CapturingMakeWriter { + buf: Arc::clone(&buf), + }; + let subscriber = tracing_subscriber::fmt() + .with_writer(make_writer) + .with_ansi(false) + .finish(); + + let status = tracing::subscriber::with_default(subscriber, || { + rt.block_on(post_events(state, host, pubkey_hex, body)) + }); + + let captured = String::from_utf8(buf.lock().unwrap().clone()).unwrap_or_default(); + (status, captured) + } + + /// Count lines that contain the terminal attribution marker. + fn count_attribution_lines(log: &str) -> usize { + log.lines() + .filter(|l| l.contains("HTTP bridge request")) + .count() + } + + /// T3a — exactly-once invariant, pre-parse 400 arm (invalid JSON). + /// + /// An authenticated client that submits a non-JSON body must produce + /// exactly ONE "HTTP bridge request" log line. This pins W1 (early exits + /// are attributed) and E1 (no double line even for the parse-fail arm). + /// + /// Discriminating: if the attribution log is removed from the ParseFail + /// arm in submit_event, this test fails. + #[test] + #[ignore = "requires Postgres"] + fn submit_event_invalid_json_emits_exactly_one_attribution_line() { + let rt = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("current_thread runtime"); + + let state = rt + .block_on(bridge_handler_test_state()) + .expect("local Postgres not reachable — start Postgres on 127.0.0.1:5432 before running ignored bridge handler tests"); + + let host = { + let h = format!("bridge-attr-{}.local", uuid::Uuid::new_v4().simple()); + rt.block_on(state.db.ensure_configured_community(&h)) + .expect("ensure community"); + h + }; + + let pubkey_hex = Keys::generate().public_key().to_hex(); + + let recorder = metrics_util::debugging::DebuggingRecorder::new(); + let (status, log) = metrics::with_local_recorder(&recorder, || { + run_and_capture(&rt, state, &host, &pubkey_hex, b"not valid json at all") + }); + + assert_eq!( + status, + axum::http::StatusCode::BAD_REQUEST, + "invalid JSON must yield 400" + ); + + let n = count_attribution_lines(&log); + assert_eq!( + n, 1, + "expected exactly 1 attribution line for invalid-JSON arm, got {n};\nlog:\n{log}" + ); + assert!( + log.contains(&pubkey_hex[..16]), + "attribution line must carry the pubkey;\nlog:\n{log}" + ); + } + + /// T3b — exactly-once invariant, post-parse IngestError::Rejected arm (relay-only kind). + /// + /// An authenticated client that submits a relay-only-kind event (kind 13534) + /// must produce exactly ONE "HTTP bridge request" log line — not two + /// (the old code emitted both an info! attribution and a warn! reason line). + /// + /// Discriminating: if two log lines are emitted (old double-log bug), this + /// test fails. + #[test] + #[ignore = "requires Postgres"] + fn submit_event_relay_only_kind_emits_exactly_one_attribution_line() { + let rt = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("current_thread runtime"); + + let state = rt + .block_on(bridge_handler_test_state()) + .expect("local Postgres not reachable — start Postgres on 127.0.0.1:5432 before running ignored bridge handler tests"); + + let host = { + let h = format!("bridge-attr-{}.local", uuid::Uuid::new_v4().simple()); + rt.block_on(state.db.ensure_configured_community(&h)) + .expect("ensure community"); + h + }; + + let client_keys = Keys::generate(); + let pubkey_hex = client_keys.public_key().to_hex(); + + let relay_only_event = EventBuilder::new( + Kind::Custom(buzz_core::kind::KIND_NIP43_MEMBERSHIP_LIST as u16), + "", + ) + .sign_with_keys(&client_keys) + .expect("sign relay-only event"); + let event_json = serde_json::to_vec(&relay_only_event).expect("serialize event"); + + let recorder = metrics_util::debugging::DebuggingRecorder::new(); + let (status, log) = metrics::with_local_recorder(&recorder, || { + run_and_capture(&rt, state, &host, &pubkey_hex, &event_json) + }); + + assert_eq!( + status, + axum::http::StatusCode::BAD_REQUEST, + "relay-only-kind event must yield 400" + ); + + let n = count_attribution_lines(&log); + assert_eq!( + n, 1, + "expected exactly 1 attribution line for IngestError::Rejected arm, got {n};\nlog:\n{log}" + ); + assert!( + log.contains(&pubkey_hex[..16]), + "attribution line must carry the pubkey;\nlog:\n{log}" + ); + } } diff --git a/crates/buzz-relay/src/handlers/event.rs b/crates/buzz-relay/src/handlers/event.rs index 3d096f934f..ecb3e41fcc 100644 --- a/crates/buzz-relay/src/handlers/event.rs +++ b/crates/buzz-relay/src/handlers/event.rs @@ -24,11 +24,11 @@ use crate::connection::{AuthState, ConnectionState}; use crate::protocol::RelayMessage; use crate::state::AppState; -use super::ingest::{IngestAuth, IngestError}; +use super::ingest::{reject_with_transport, IngestAuth, IngestError}; /// Increment the rejection counter with a bounded reason label. fn reject(reason: &'static str) { - metrics::counter!("buzz_events_rejected_total", "reason" => reason).increment(1); + reject_with_transport("ws", reason); } /// Bound the `kind` label to prevent cardinality explosion from arbitrary Nostr kinds. diff --git a/crates/buzz-relay/src/handlers/ingest.rs b/crates/buzz-relay/src/handlers/ingest.rs index 415a9e8baf..f53ea0f018 100644 --- a/crates/buzz-relay/src/handlers/ingest.rs +++ b/crates/buzz-relay/src/handlers/ingest.rs @@ -146,6 +146,22 @@ fn emit_product_feedback_success( ); } +/// Increment the rejection counter with a bounded reason and transport label. +/// +/// Shared by the WS `EVENT` handler and the HTTP `POST /events` handler so +/// both transports feed the same series — `transport` distinguishes them so +/// existing WS-only dashboards aren't silently diluted by HTTP volume. +/// `reason` is one of a small closed set ("auth", "invalid", "scope", +/// "error") — bounded, no cardinality risk. +pub fn reject_with_transport(transport: &'static str, reason: &'static str) { + metrics::counter!( + "buzz_events_rejected_total", + "transport" => transport, + "reason" => reason + ) + .increment(1); +} + /// Successful ingestion result. pub struct IngestResult { /// Hex-encoded event ID. @@ -3612,4 +3628,54 @@ mod tests { // error comes from validate_engram_nip44_content with label replaced assert!(err.contains("agent-turn-metric"), "got: {err}"); } + + /// The HTTP bridge's `submit_event` 400 arm and the WS `EVENT` handler's + /// reject path must land on the same counter, distinguished only by the + /// `transport` label — this is what lets a dashboard tell "server got + /// hammered with bad HTTP requests" apart from "a WS client is + /// misbehaving" without losing the combined total. + #[test] + fn reject_with_transport_labels_http_and_ws_as_separate_series() { + let recorder = metrics_util::debugging::DebuggingRecorder::new(); + let snapshotter = recorder.snapshotter(); + + metrics::with_local_recorder(&recorder, || { + reject_with_transport("http", "invalid"); + reject_with_transport("ws", "invalid"); + reject_with_transport("http", "invalid"); + }); + + let counts: std::collections::HashMap<(String, String), u64> = snapshotter + .snapshot() + .into_vec() + .into_iter() + .filter(|(key, ..)| key.key().name() == "buzz_events_rejected_total") + .map(|(key, _, _, value)| { + let metrics_util::debugging::DebugValue::Counter(n) = value else { + panic!("buzz_events_rejected_total must be a counter"); + }; + let labels: Vec<_> = key.key().labels().collect(); + let transport = labels + .iter() + .find(|l| l.key() == "transport") + .map(|l| l.value().to_owned()) + .unwrap_or_default(); + let reason = labels + .iter() + .find(|l| l.key() == "reason") + .map(|l| l.value().to_owned()) + .unwrap_or_default(); + ((transport, reason), n) + }) + .collect(); + + assert_eq!( + counts.get(&("http".to_owned(), "invalid".to_owned())), + Some(&2) + ); + assert_eq!( + counts.get(&("ws".to_owned(), "invalid".to_owned())), + Some(&1) + ); + } }