feat(relay): log NIP-98 pubkey attribution on HTTP bridge requests (#2206)

Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
Co-authored-by: npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7 <dcfd242e557282d7a1e2cf2e6877522682f1e5c6156dc92ca7d90eaedd3b0f95@sprout-oss.stage.blox.sqprod.co>
This commit is contained in:
Will Pfleger
2026-07-21 23:28:16 +00:00
committed by GitHub
co-authored by npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7
parent dd64d721d3
commit 7e34bee62c
3 changed files with 896 additions and 43 deletions
+828 -41
View File
@@ -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<Arc<AppState>>,
@@ -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<Value>,
},
/// JSON parse failure before ingest — log category/line/column, not msg.
ParseFail {
category: &'static str,
line: usize,
column: usize,
response: (StatusCode, Json<Value>),
},
/// IngestError::Rejected — log kind + truncated reason.
Rejected {
kind: u32,
reason: String,
response: (StatusCode, Json<Value>),
},
/// 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<Value>),
},
}
impl SubmitOutcome {
fn into_response(self) -> Result<Json<Value>, (StatusCode, Json<Value>)> {
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<AppState>,
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<AppState>,
tenant: &TenantContext,
headers: &HeaderMap,
body: &[u8],
pubkey: nostr::PublicKey,
event_id_bytes: [u8; 32],
) -> Result<Json<Value>, (StatusCode, Json<Value>)> {
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<Value> = serde_json::from_slice(&body)
let raw_filters: Vec<Value> = serde_json::from_slice(body)
.map_err(|e| api_error(StatusCode::BAD_REQUEST, &format!("invalid filters: {e}")))?;
let filters: Vec<nostr::Filter> = 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<AppState>,
tenant: &TenantContext,
headers: &HeaderMap,
body: &[u8],
pubkey: nostr::PublicKey,
event_id_bytes: [u8; 32],
) -> Result<Json<Value>, (StatusCode, Json<Value>)> {
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<nostr::Filter> = serde_json::from_slice(&body)
let filters: Vec<nostr::Filter> = 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<dyn std::future::Future<Output = Result<bool, buzz_auth::AuthError>> + 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<Arc<crate::state::AppState>> {
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<crate::state::AppState>,
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<Mutex<Vec<u8>>>,
}
struct CapturingWriter {
buf: Arc<Mutex<Vec<u8>>>,
}
impl std::io::Write for CapturingWriter {
fn write(&mut self, data: &[u8]) -> std::io::Result<usize> {
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<crate::state::AppState>,
host: &str,
pubkey_hex: &str,
body: &[u8],
) -> (axum::http::StatusCode, String) {
let buf = Arc::new(Mutex::new(Vec::<u8>::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}"
);
}
}
+2 -2
View File
@@ -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.
+66
View File
@@ -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)
);
}
}