diff --git a/Cargo.lock b/Cargo.lock index 9a3f91671d..1ad797bd4d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -930,6 +930,7 @@ dependencies = [ "clap", "diffy", "dirs", + "futures-util", "hex", "infer", "nostr", @@ -942,6 +943,7 @@ dependencies = [ "tempfile", "thiserror 2.0.18", "tokio", + "tokio-tungstenite 0.29.0", "url", "uuid", ] @@ -1348,12 +1350,17 @@ name = "buzz-ws-client" version = "0.1.0" dependencies = [ "futures-util", + "log", "nostr", + "rand 0.10.1", + "serde", "serde_json", "thiserror 2.0.18", "tokio", "tokio-tungstenite 0.29.0", "tracing", + "tracing-log", + "tracing-subscriber", "url", ] diff --git a/Cargo.toml b/Cargo.toml index 3268cfaf8d..c0e2ecf0a3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -72,6 +72,7 @@ serde_yaml = "0.9" evalexpr = "11" cron = "0.16" # Observability +log = "0.4" tracing = "0.1" tracing-subscriber = { version = "0.3", features = ["env-filter", "json"] } tracing-opentelemetry = { version = "0.33" } diff --git a/crates/buzz-cli/Cargo.toml b/crates/buzz-cli/Cargo.toml index 1476e60bfd..b8d28f71e9 100644 --- a/crates/buzz-cli/Cargo.toml +++ b/crates/buzz-cli/Cargo.toml @@ -91,3 +91,6 @@ rand = { workspace = true } tempfile = "3" # Minimal HTTP test server for retry/policy integration tests axum = { workspace = true } +# Local WebSocket relay fixtures for exact-event command integration tests +futures-util = { workspace = true } +tokio-tungstenite = { workspace = true } diff --git a/crates/buzz-cli/README.md b/crates/buzz-cli/README.md index a2dcdce6d2..d3487f7881 100644 --- a/crates/buzz-cli/README.md +++ b/crates/buzz-cli/README.md @@ -20,14 +20,22 @@ export BUZZ_PRIVATE_KEY="nsec1..." buzz channels list ``` +`events get-verified` can read from a public Nostr relay without a key. If the +relay requires NIP-42, it uses `BUZZ_PRIVATE_KEY` and the optional +`BUZZ_AUTH_TAG` for authentication only; event trust still comes from local ID +recomputation and signature verification. + ## Usage All output is JSON on stdout. Errors are JSON on stderr. Exit codes: 0=ok, 1=user error, 2=network, 3=auth, 4=other, 5=write conflict. ```bash -# Set relay URL (defaults to http://localhost:3000) +# Set relay URL (defaults to http://localhost:3000 for ordinary commands) export BUZZ_RELAY_URL="https://relay.example.com" +# Exact raw event read (command-local relay is mandatory; no env/default fallback) +buzz events get-verified --relay wss://relay.example.com --event <64hex> + # Messages buzz messages send --channel --content "Hello" buzz messages send --channel --content "Reply" --reply-to --broadcast @@ -146,6 +154,7 @@ stored rules in `validation_error` so an owner can remove and repair them. | | `runs` | Get workflow run history | | | `approve` | Approve/deny a workflow step | | `feed` | `get` | Get your activity feed | +| `events` | `get-verified` | Fetch one exact raw event and verify its NIP-01 ID and Schnorr signature locally | | `social` | `publish` | Publish a NIP-01 note | | | `set-contacts` | Set NIP-02 contact list | | | `event` | Get a Nostr event | @@ -174,6 +183,7 @@ buzz [flags] │ ├─ main.rs ──▶ commands/*.rs ──▶ client.rs ──▶ Buzz Relay REST API │ (clap) (handlers) (reqwest) + │ └──────────▶ buzz-ws-client (exact Nostr reads) │ ├─ validate.rs (UUID, hex, content size, percent-encode) └─ error.rs (CliError → JSON stderr + exit code) diff --git a/crates/buzz-cli/TESTING.md b/crates/buzz-cli/TESTING.md index 77234b7faa..da395f7753 100644 --- a/crates/buzz-cli/TESTING.md +++ b/crates/buzz-cli/TESTING.md @@ -482,6 +482,25 @@ buzz notes get --name dco-check # exits non-zero: not found buzz notes rm --name does-not-exist # exits non-zero ``` +### 6.13 Exact Verified Events + +This is a raw NIP-01 WebSocket read. Its `--relay` is mandatory and does not +inherit `BUZZ_RELAY_URL` or the global `--relay`. The command waits for EOSE, +requires exactly one result, recomputes the event ID, and verifies the Schnorr +signature locally before writing anything to stdout. + +```bash +# EVENT_ID comes from the message created in section 6.3. +buzz events get-verified \ + --relay ws://localhost:3000 \ + --event "$EVENT_ID" | jq . +# Expected: {id,pubkey,created_at,kind,tags,content,sig} + +# A missing command-local relay is a user error even when BUZZ_RELAY_URL is set. +buzz events get-verified --event "$EVENT_ID" +# Expected: exit 1, no stdout +``` + --- ## 7. Error Path Testing @@ -621,3 +640,4 @@ buzz channels delete --channel "$FORUM_ID" | jq . | 60 | `notes ls` | ☐ | Own, --author all, --tag, --limit | | 61 | `notes rm` | ☐ | Delete→get 404, double-delete idempotent, missing slug → NotFound | | 62 | `users set-status` | ☐ | Text+emoji, text only, emoji-only (`--text ""`), `--clear`, `--clear` + `--text` → exit 1 | +| 63 | `events get-verified` | ☐ | Exact ID, raw seven-field event, local ID/signature verification, explicit WS relay | diff --git a/crates/buzz-cli/src/commands/events.rs b/crates/buzz-cli/src/commands/events.rs new file mode 100644 index 0000000000..5247e2618d --- /dev/null +++ b/crates/buzz-cli/src/commands/events.rs @@ -0,0 +1,733 @@ +//! Exact-event relay reads with mandatory local NIP-01 verification. + +use std::time::Duration; + +use buzz_ws_client::{NostrWsConnection, RelayMessage, WsClientError}; +use nostr::{Event, EventId, Keys, Tag}; +use serde::Deserialize; + +use crate::error::CliError; +use crate::validate::validate_hex64; + +const FETCH_TIMEOUT: Duration = Duration::from_secs(15); +const INITIAL_SUBSCRIPTION_ID: &str = "buzz-get-verified-0"; +const AUTHENTICATED_SUBSCRIPTION_ID: &str = "buzz-get-verified-1"; + +fn transport_error(error: WsClientError) -> CliError { + match error { + WsClientError::Url(message) => CliError::Usage(format!("invalid relay URL: {message}")), + WsClientError::AuthFailed(message) => CliError::Auth(message), + WsClientError::AmbiguousAuthChallenge => CliError::Ambiguous(error.to_string()), + WsClientError::WebSocket(_) + | WsClientError::Timeout + | WsClientError::ConnectionClosed + | WsClientError::AuthTransportPoisoned => CliError::Transport(error.to_string()), + WsClientError::Json(_) + | WsClientError::InvalidEvent { .. } + | WsClientError::EventBuilder(_) + | WsClientError::UnexpectedMessage(_) + | WsClientError::EventRejected(_) + | WsClientError::NoAuthChallenge + | WsClientError::AuthChallengeTooLarge { .. } + | WsClientError::AuthFrameTooLarge { .. } + | WsClientError::ReflectedAuthMaterial => CliError::RelayProtocol(error.to_string()), + } +} + +fn exact_event_id(value: &str) -> Result { + validate_hex64(value)?; + if !value + .chars() + .all(|character| character.is_ascii_digit() || matches!(character, 'a'..='f')) + { + return Err(CliError::Usage( + "event ID must be canonical lowercase hex".into(), + )); + } + EventId::parse(value).map_err(|error| CliError::Usage(format!("invalid event ID: {error}"))) +} + +fn exact_relay_url(value: &str) -> Result { + if value.trim() != value { + return Err(CliError::Usage( + "relay URL must not contain leading or trailing whitespace".into(), + )); + } + let url = url::Url::parse(value) + .map_err(|error| CliError::Usage(format!("invalid relay URL: {error}")))?; + if !matches!(url.scheme(), "ws" | "wss") { + return Err(CliError::Usage("relay URL scheme must be ws or wss".into())); + } + if !url.username().is_empty() || url.password().is_some() { + return Err(CliError::Usage( + "relay URL must not contain credentials".into(), + )); + } + if url.fragment().is_some() { + return Err(CliError::Usage( + "relay URL must not contain a fragment".into(), + )); + } + if url.host().is_none() { + return Err(CliError::Usage("relay URL must contain a host".into())); + } + Ok(value.to_string()) +} + +fn request(subscription_id: &str, event_id: &EventId) -> serde_json::Value { + serde_json::json!([ + "REQ", + subscription_id, + { + "ids": [event_id.to_hex()], + "limit": 2 + } + ]) +} + +// NIP-42 documents `auth-required` as a machine-readable CLOSED prefix. Match +// only that exact token or its colon-delimited form; human prose and lookalike +// strings must never authorize signing with ambient credentials. +fn is_auth_required_reason(message: &str) -> bool { + message == "auth-required" || message.starts_with("auth-required:") +} + +#[derive(Debug)] +struct ReceivedEvent { + event: Event, + raw_json: Box, +} + +struct CollectedEvents { + events: Vec, + private_auth_started: bool, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct RawEventEnvelope { + id: String, + pubkey: String, + created_at: u64, + kind: u16, + tags: Vec>, + content: String, + sig: String, +} + +fn is_canonical_hex(value: &str, length: usize) -> bool { + value.len() == length + && value + .bytes() + .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) +} + +fn validate_raw_event(requested_id: &EventId, received: &ReceivedEvent) -> Result<(), CliError> { + let raw: RawEventEnvelope = serde_json::from_str(&received.raw_json).map_err(|error| { + CliError::RelayMismatch(format!( + "relay returned a malformed signed-event envelope: {error}" + )) + })?; + + if !is_canonical_hex(&raw.id, 64) { + return Err(CliError::RelayMismatch( + "relay returned a non-canonical event id".into(), + )); + } + if raw.id != requested_id.to_hex() { + return Err(CliError::RelayMismatch(format!( + "requested event {} but raw relay payload declared {}", + requested_id.to_hex(), + raw.id + ))); + } + if !is_canonical_hex(&raw.pubkey, 64) { + return Err(CliError::RelayMismatch( + "relay returned a non-canonical event pubkey".into(), + )); + } + if !is_canonical_hex(&raw.sig, 128) { + return Err(CliError::SignatureInvalid( + "relay returned a malformed Schnorr signature".into(), + )); + } + + let typed_tags: Vec> = received + .event + .tags + .iter() + .map(|tag| tag.as_slice().to_vec()) + .collect(); + if raw.id != received.event.id.to_hex() + || raw.pubkey != received.event.pubkey.to_hex() + || raw.created_at != received.event.created_at.as_secs() + || raw.kind != received.event.kind.as_u16() + || raw.tags != typed_tags + || raw.content != received.event.content + || raw.sig != received.event.sig.to_string() + { + return Err(CliError::RelayMismatch( + "raw signed-event fields did not match their typed representation".into(), + )); + } + + Ok(()) +} + +fn malformed_event_error(requested_id: &EventId, raw_event_json: &str, message: &str) -> CliError { + let value = match serde_json::from_str::(raw_event_json) { + Ok(serde_json::Value::Object(value)) => value, + _ => { + return CliError::RelayMismatch(format!( + "relay returned a malformed signed-event envelope: {message}" + )); + } + }; + + match value.get("id") { + Some(serde_json::Value::String(id)) + if is_canonical_hex(id, 64) && id == &requested_id.to_hex() => {} + Some(serde_json::Value::String(id)) if is_canonical_hex(id, 64) => { + return CliError::RelayMismatch(format!( + "requested event {} but raw relay payload declared {id}", + requested_id.to_hex() + )); + } + _ => { + return CliError::RelayMismatch( + "relay returned a malformed or non-canonical event id".into(), + ); + } + } + + if !matches!(value.get("sig"), Some(serde_json::Value::String(sig)) if is_canonical_hex(sig, 128)) + { + return CliError::SignatureInvalid("relay returned a malformed Schnorr signature".into()); + } + + CliError::RelayMismatch(format!( + "relay returned a malformed signed-event envelope: {message}" + )) +} + +fn receive_error(error: WsClientError, requested_id: &EventId, authenticated: bool) -> CliError { + if authenticated { + return match error { + WsClientError::WebSocket(_) + | WsClientError::Timeout + | WsClientError::ConnectionClosed + | WsClientError::AuthTransportPoisoned => CliError::Transport( + "authenticated relay connection failed during exact-event fetch".into(), + ), + WsClientError::AmbiguousAuthChallenge => { + CliError::Ambiguous("relay sent another authentication challenge".into()) + } + WsClientError::AuthFailed(_) => { + CliError::Auth("relay rejected NIP-42 authentication".into()) + } + _ => CliError::RelayProtocol( + "relay sent an invalid response after authentication".into(), + ), + }; + } + + match error { + WsClientError::InvalidEvent { + raw_event_json, + message, + } => malformed_event_error(requested_id, &raw_event_json, &message), + other => transport_error(other), + } +} + +// Once private AUTH may have crossed the wire, no relay-derived value may be +// rendered into CLI output. This static boundary covers typed +// fields, serde diagnostics, WebSocket errors, and validation failures without +// attempting to recognize transformed or split secrets. +fn private_auth_error(error: CliError) -> CliError { + match error { + CliError::Usage(_) => CliError::Usage( + "exact-event fetch could not continue after private authentication".into(), + ), + CliError::Relay { .. } | CliError::RelayProtocol(_) => CliError::RelayProtocol( + "relay sent an invalid response after private authentication".into(), + ), + CliError::Network(_) | CliError::Transport(_) => { + CliError::Transport("relay connection failed after private authentication".into()) + } + CliError::Auth(_) | CliError::Key(_) => { + CliError::Auth("private relay authentication failed".into()) + } + CliError::Conflict(_) => { + CliError::Conflict("relay reported a conflict after private authentication".into()) + } + CliError::NotFound(_) => { + CliError::NotFound("event not found after private authentication".into()) + } + CliError::Ambiguous(_) => CliError::Ambiguous( + "relay returned an ambiguous result after private authentication".into(), + ), + CliError::RelayMismatch(_) => CliError::RelayMismatch( + "relay response did not match the exact-event request after private authentication" + .into(), + ), + CliError::IdMismatch(_) => CliError::IdMismatch( + "relay event failed local ID verification after private authentication".into(), + ), + CliError::SignatureInvalid(_) => CliError::SignatureInvalid( + "relay event failed local signature verification after private authentication".into(), + ), + CliError::DeliveryUnknown(_) => CliError::DeliveryUnknown( + "relay delivery state is unknown after private authentication".into(), + ), + CliError::Other(_) => { + CliError::Other("exact-event fetch failed after private authentication".into()) + } + } +} + +fn verify_exact_result( + requested_id: &EventId, + relay: &str, + mut events: Vec, +) -> Result { + if events.is_empty() { + return Err(CliError::NotFound(format!( + "event {} not found at {relay}", + requested_id.to_hex() + ))); + } + if events.len() != 1 { + return Err(CliError::Ambiguous(format!( + "relay returned {} events for exact event ID {}", + events.len(), + requested_id.to_hex() + ))); + } + + let received = events.remove(0); + validate_raw_event(requested_id, &received)?; + if received.event.id != *requested_id { + return Err(CliError::RelayMismatch(format!( + "requested event {} but relay returned {}", + requested_id.to_hex(), + received.event.id.to_hex() + ))); + } + if !received.event.verify_id() { + return Err(CliError::IdMismatch(format!( + "event {} does not match its local NIP-01 ID recomputation", + received.event.id.to_hex() + ))); + } + if !received.event.verify_signature() { + return Err(CliError::SignatureInvalid(format!( + "event {} has an invalid Schnorr signature", + received.event.id.to_hex() + ))); + } + + Ok(received) +} + +async fn collect_exact_event( + relay: &str, + event_id: &EventId, + keys: Option<&Keys>, + auth_tag: Option<&Tag>, +) -> Result { + let deadline = tokio::time::Instant::now() + FETCH_TIMEOUT; + let mut private_auth_started = false; + let result = tokio::time::timeout_at( + deadline, + collect_exact_event_inner(relay, event_id, keys, auth_tag, &mut private_auth_started), + ) + .await + .map_err(|_| { + CliError::Transport(format!( + "exact-event fetch timed out after {} seconds for {relay}", + FETCH_TIMEOUT.as_secs() + )) + }) + .and_then(|result| result); + + let events = result.map_err(|error| { + if private_auth_started { + private_auth_error(error) + } else { + error + } + })?; + Ok(CollectedEvents { + events, + private_auth_started, + }) +} + +async fn collect_exact_event_inner( + relay: &str, + event_id: &EventId, + keys: Option<&Keys>, + auth_tag: Option<&Tag>, + private_auth_started: &mut bool, +) -> Result, CliError> { + let mut connection = NostrWsConnection::connect(relay) + .await + .map_err(transport_error)?; + let mut subscription_id = INITIAL_SUBSCRIPTION_ID; + let mut retired_subscription_id = None; + let mut authenticated = false; + let mut auth_challenge_received = false; + let mut authentication_required = false; + let mut events = Vec::new(); + + let result = async { + connection + .send_raw(&request(subscription_id, event_id)) + .await + .map_err(transport_error)?; + + loop { + let next_message = if authenticated { + connection + .next_event_for_exact_id(FETCH_TIMEOUT, event_id) + .await + } else { + connection.next_event(FETCH_TIMEOUT).await + }; + match next_message.map_err(|error| receive_error(error, event_id, authenticated))? + { + RelayMessage::Event { + subscription_id: received, + event, + raw_event_json, + } if received == subscription_id && !authentication_required => { + if !events.is_empty() { + return Err(CliError::Ambiguous(format!( + "relay returned more than one event for exact event ID {}", + event_id.to_hex() + ))); + } + events.push(ReceivedEvent { + event: *event, + raw_json: raw_event_json, + }); + } + RelayMessage::Eose { + subscription_id: received, + } if received == subscription_id && !authentication_required => break, + RelayMessage::Event { + subscription_id: received, + .. + } + | RelayMessage::Eose { + subscription_id: received, + } if received == subscription_id && authentication_required => {} + RelayMessage::Closed { + subscription_id: received, + message, + } if received == subscription_id => { + if authenticated { + return Err(CliError::RelayProtocol( + "relay closed authenticated exact-event subscription".into(), + )); + } + if is_auth_required_reason(&message) && keys.is_none() { + return Err(CliError::Auth( + "relay requires NIP-42 authentication; set BUZZ_PRIVATE_KEY or pass --private-key" + .into(), + )); + } + if is_auth_required_reason(&message) && !authenticated { + authentication_required = true; + } else { + return Err(CliError::RelayProtocol(format!( + "relay closed exact-event subscription: {message}" + ))); + } + } + RelayMessage::Event { + subscription_id: received, + .. + } + | RelayMessage::Eose { + subscription_id: received, + } + | RelayMessage::Closed { + subscription_id: received, + .. + } if retired_subscription_id == Some(received.as_str()) => {} + RelayMessage::Event { + subscription_id: received, + .. + } + | RelayMessage::Eose { + subscription_id: received, + } + | RelayMessage::Closed { + subscription_id: received, + .. + } => { + let message = if authenticated { + "relay responded for an unexpected authenticated subscription".into() + } else { + format!( + "relay responded for subscription {received}, expected {subscription_id}" + ) + }; + return Err(CliError::RelayMismatch(message)); + } + // NIP-42 relays may advertise a challenge even when the requested read + // is public. Keep waiting for EOSE or an explicit auth-required CLOSED + // instead of treating the challenge itself as denial. + RelayMessage::Auth { .. } if !authenticated => { + auth_challenge_received = true; + } + RelayMessage::Auth { .. } => { + return Err(CliError::RelayProtocol( + "relay sent a second authentication challenge".into(), + )); + } + RelayMessage::Notice { .. } => {} + RelayMessage::Ok(_) | RelayMessage::Count { .. } => { + return Err(CliError::RelayProtocol( + "relay sent an unexpected message during exact-event fetch".into(), + )); + } + } + + if authentication_required && auth_challenge_received && !authenticated { + let keys = keys.ok_or_else(|| { + CliError::Auth("NIP-42 authentication key is unavailable".into()) + })?; + retired_subscription_id = Some(subscription_id); + // Start conservatively before the cancellable await, then use + // the connection's precise write state on any normal return. + // A relay can reflect the private event before sending OK. + *private_auth_started = true; + let authentication = connection.authenticate(keys, auth_tag).await; + *private_auth_started = connection.private_auth_started(); + authentication.map_err(transport_error)?; + authenticated = true; + authentication_required = false; + auth_challenge_received = false; + subscription_id = AUTHENTICATED_SUBSCRIPTION_ID; + events.clear(); + connection + .send_raw(&request(subscription_id, event_id)) + .await + .map_err(transport_error)?; + } + } + + Ok(events) + } + .await; + + let _ = connection + .send_raw(&serde_json::json!(["CLOSE", subscription_id])) + .await; + let _ = connection.disconnect().await; + result +} + +/// Fetch exactly one event from one relay and emit it only after local verification. +pub async fn cmd_get_verified( + relay: &str, + event_id: &str, + keys: Option<&Keys>, + auth_tag: Option<&Tag>, +) -> Result<(), CliError> { + let relay = exact_relay_url(relay)?; + let event_id = exact_event_id(event_id)?; + let collected = collect_exact_event(&relay, &event_id, keys, auth_tag).await?; + let event = verify_exact_result(&event_id, &relay, collected.events).map_err(|error| { + if collected.private_auth_started { + private_auth_error(error) + } else { + error + } + })?; + println!("{}", event.raw_json); + Ok(()) +} + +pub async fn dispatch( + cmd: &crate::EventsCmd, + keys: Option<&Keys>, + auth_tag: Option<&Tag>, +) -> Result<(), CliError> { + match cmd { + crate::EventsCmd::GetVerified { relay, event } => { + cmd_get_verified(relay, event, keys, auth_tag).await + } + } +} + +#[cfg(test)] +mod tests { + use nostr::{EventBuilder, JsonUtil, Kind}; + + use super::*; + + fn signed_event(content: &str) -> Event { + EventBuilder::new(Kind::TextNote, content) + .sign_with_keys(&Keys::generate()) + .unwrap() + } + + fn mutate(event: &Event, field: &str, value: serde_json::Value) -> Event { + let mut json = serde_json::to_value(event).unwrap(); + json[field] = value; + Event::from_json(json.to_string()).unwrap() + } + + fn received(event: Event) -> ReceivedEvent { + let raw_json = event.as_json().into_boxed_str(); + ReceivedEvent { event, raw_json } + } + + #[test] + fn valid_exact_event_is_returned() { + let event = signed_event("valid"); + let result = verify_exact_result( + &event.id, + "wss://relay.example", + vec![received(event.clone())], + ); + assert_eq!(result.unwrap().event, event); + } + + #[test] + fn zero_results_are_not_found() { + let id = signed_event("missing").id; + let error = verify_exact_result(&id, "wss://relay.example", vec![]).unwrap_err(); + assert!(matches!(error, CliError::NotFound(_))); + } + + #[test] + fn multiple_results_are_ambiguous() { + let event = signed_event("duplicate"); + let event_id = event.id; + let error = verify_exact_result( + &event_id, + "wss://relay.example", + vec![received(event.clone()), received(event)], + ) + .unwrap_err(); + assert!(matches!(error, CliError::Ambiguous(_))); + } + + #[test] + fn returned_wrong_event_is_relay_mismatch() { + let requested = signed_event("requested"); + let returned = signed_event("returned"); + let error = verify_exact_result( + &requested.id, + "wss://relay.example", + vec![received(returned)], + ) + .unwrap_err(); + assert!(matches!(error, CliError::RelayMismatch(_))); + } + + #[test] + fn mutated_content_is_id_mismatch() { + let valid = signed_event("original"); + let mutated = mutate(&valid, "content", serde_json::json!("mutated")); + let error = verify_exact_result(&valid.id, "wss://relay.example", vec![received(mutated)]) + .unwrap_err(); + assert!(matches!(error, CliError::IdMismatch(_))); + } + + #[test] + fn invalid_signature_is_distinct_from_id_mismatch() { + let valid = signed_event("original"); + let mutated = mutate(&valid, "sig", serde_json::json!("0".repeat(128))); + let error = verify_exact_result(&valid.id, "wss://relay.example", vec![received(mutated)]) + .unwrap_err(); + assert!(matches!(error, CliError::SignatureInvalid(_))); + } + + #[test] + fn request_is_exact_id_and_cardinality_bounded() { + let event = signed_event("filter"); + assert_eq!( + request("subscription", &event.id), + serde_json::json!([ + "REQ", + "subscription", + {"ids": [event.id.to_hex()], "limit": 2} + ]) + ); + } + + #[test] + fn auth_required_reason_accepts_only_documented_machine_token() { + for accepted in ["auth-required", "auth-required:", "auth-required: policy"] { + assert!(is_auth_required_reason(accepted), "rejected {accepted:?}"); + } + for rejected in [ + "not-auth-required", + "restricted: not-auth-required", + "restricted auth-required response", + "auth-required-suffix", + "AUTH-REQUIRED", + ] { + assert!(!is_auth_required_reason(rejected), "accepted {rejected:?}"); + } + } + + #[test] + fn relay_must_be_explicit_websocket_url() { + assert!(exact_relay_url("wss://relay.example").is_ok()); + assert!(exact_relay_url("ws://127.0.0.1:3000").is_ok()); + assert!(exact_relay_url("https://relay.example").is_err()); + assert!(exact_relay_url("relay.example").is_err()); + assert!(exact_relay_url(" wss://relay.example").is_err()); + assert!(exact_relay_url("wss://relay.example ").is_err()); + } + + #[test] + fn relay_connection_host_is_not_rewritten() { + assert_eq!( + exact_relay_url("wss://localhost:443/community").unwrap(), + "wss://localhost:443/community" + ); + } + + #[test] + fn event_id_must_be_canonical_lowercase_hex() { + let lowercase = "abcdef0123456789".repeat(4); + assert!(exact_event_id(&lowercase).is_ok()); + assert!(exact_event_id(&lowercase.to_ascii_uppercase()).is_err()); + } + + #[test] + fn raw_event_requires_exact_signed_field_set() { + let event = signed_event("extra field"); + let mut raw = serde_json::to_value(&event).unwrap(); + raw["unexpected"] = serde_json::json!(true); + let received = ReceivedEvent { + event: event.clone(), + raw_json: raw.to_string().into_boxed_str(), + }; + + let error = validate_raw_event(&event.id, &received).unwrap_err(); + assert!(matches!(error, CliError::RelayMismatch(_))); + } + + #[test] + fn malformed_raw_signature_has_signature_category() { + let event = signed_event("malformed signature"); + let mut raw = serde_json::to_value(&event).unwrap(); + raw["sig"] = serde_json::json!("ABC"); + let received = ReceivedEvent { + event: event.clone(), + raw_json: raw.to_string().into_boxed_str(), + }; + + let error = validate_raw_event(&event.id, &received).unwrap_err(); + assert!(matches!(error, CliError::SignatureInvalid(_))); + } +} diff --git a/crates/buzz-cli/src/commands/mod.rs b/crates/buzz-cli/src/commands/mod.rs index 8691590636..6b28624770 100644 --- a/crates/buzz-cli/src/commands/mod.rs +++ b/crates/buzz-cli/src/commands/mod.rs @@ -3,6 +3,7 @@ pub mod channel_templates; pub mod channels; pub mod dms; pub mod emoji; +pub mod events; pub mod feed; pub mod issues; pub mod mem; diff --git a/crates/buzz-cli/src/error.rs b/crates/buzz-cli/src/error.rs index 2edcd6aa9d..0ec9ea9eb0 100644 --- a/crates/buzz-cli/src/error.rs +++ b/crates/buzz-cli/src/error.rs @@ -32,6 +32,30 @@ pub enum CliError { #[error("{0}")] NotFound(String), + /// An exact query returned more than one event. + #[error("ambiguous result: {0}")] + Ambiguous(String), + + /// A relay returned an event or subscription other than the requested one. + #[error("relay mismatch: {0}")] + RelayMismatch(String), + + /// A declared event ID did not match a local NIP-01 recomputation. + #[error("event ID mismatch: {0}")] + IdMismatch(String), + + /// An event failed local BIP-340 Schnorr verification. + #[error("signature invalid: {0}")] + SignatureInvalid(String), + + /// A WebSocket connection, timeout, or framing failure. + #[error("transport error: {0}")] + Transport(String), + + /// A relay returned a well-formed but invalid response sequence. + #[error("relay protocol error: {0}")] + RelayProtocol(String), + /// A non-idempotent command's outcome is unknown: the request may have /// reached the relay, but the response was lost. Never auto-retried and /// never labeled retryable — the relay executes these commands before any @@ -80,6 +104,7 @@ pub fn is_retryable_error(e: &CliError) -> bool { } CliError::Relay { status, .. } => matches!(status, 429 | 502 | 503 | 504), CliError::DeliveryUnknown(_) => false, + CliError::Transport(_) => true, _ => false, } } @@ -102,6 +127,11 @@ pub fn exit_code(e: &CliError) -> i32 { CliError::Key(_) => 3, CliError::Conflict(_) => 5, CliError::NotFound(_) => 1, + CliError::Ambiguous(_) + | CliError::RelayMismatch(_) + | CliError::IdMismatch(_) + | CliError::SignatureInvalid(_) => 4, + CliError::Transport(_) | CliError::RelayProtocol(_) => 2, CliError::DeliveryUnknown(_) => 2, CliError::Other(_) => 4, } @@ -124,6 +154,12 @@ pub fn print_error(e: &CliError) { CliError::Key(_) => "key_error", CliError::Conflict(_) => "conflict", CliError::NotFound(_) => "not_found", + CliError::Ambiguous(_) => "ambiguous", + CliError::RelayMismatch(_) => "relay_mismatch", + CliError::IdMismatch(_) => "id_mismatch", + CliError::SignatureInvalid(_) => "signature_invalid", + CliError::Transport(_) => "transport_error", + CliError::RelayProtocol(_) => "relay_error", CliError::DeliveryUnknown(_) => "delivery_unknown", CliError::Other(_) => "error", }; diff --git a/crates/buzz-cli/src/lib.rs b/crates/buzz-cli/src/lib.rs index 0726406d29..a3f071bc37 100644 --- a/crates/buzz-cli/src/lib.rs +++ b/crates/buzz-cli/src/lib.rs @@ -68,7 +68,7 @@ Buzz CLI — interact with a Buzz relay Configuration (flags override env vars): BUZZ_RELAY_URL Relay base URL [default: http://localhost:3000] - BUZZ_PRIVATE_KEY Nostr private key (hex or nsec) [required] + BUZZ_PRIVATE_KEY Nostr private key (hex or nsec) [required except public events get-verified reads] BUZZ_AUTH_TAG NIP-OA auth tag JSON [optional] The 'pack' subcommand runs locally and does not require a relay connection. @@ -81,7 +81,7 @@ struct Cli { #[arg(long, env = "BUZZ_RELAY_URL", default_value = "http://localhost:3000")] relay: String, - /// Nostr private key (hex or nsec). This is the CLI's identity. + /// Nostr private key (hex or nsec). Optional only for public events get-verified reads. #[arg(long, env = "BUZZ_PRIVATE_KEY", hide_env_values = true)] private_key: Option, @@ -206,6 +206,9 @@ enum Cmd { /// Publish notes and manage the social graph (NIP-01/02) #[command(subcommand)] Social(SocialCmd), + /// Fetch raw Nostr events with mandatory local verification + #[command(subcommand)] + Events(EventsCmd), /// Publish and edit long-form NIP-23 notes — team knowledge base #[command(subcommand)] Notes(NotesCmd), @@ -954,6 +957,23 @@ pub enum FeedCmd { }, } +#[derive(Subcommand)] +pub enum EventsCmd { + /// Fetch one exact event and verify its NIP-01 ID and Schnorr signature locally + #[command( + name = "get-verified", + after_help = "Example:\n buzz events get-verified --relay wss://relay.example --event <64hex>\n\nThe relay is required here and never inherited from BUZZ_RELAY_URL or the global --relay flag. No event is emitted unless exactly one result passes local NIP-01 ID recomputation and Schnorr verification." + )] + GetVerified { + /// Single WebSocket relay URL (ws:// or wss://); no default or fallback + #[arg(long)] + relay: String, + /// Exact 64-character hex event ID + #[arg(long)] + event: String, + }, +} + #[derive(Subcommand)] pub enum SocialCmd { /// Publish a text note (NIP-01 kind:1) @@ -1769,8 +1789,6 @@ pub enum ModerationCmd { } async fn run(cli: Cli) -> Result<(), CliError> { - let relay_url = client::normalize_relay_url(&cli.relay); - // Pack commands are local-only — no relay connection needed. if let Cmd::Pack(ref sub) = cli.command { return match sub { @@ -1779,7 +1797,17 @@ async fn run(cli: Cli) -> Result<(), CliError> { }; } - // Auth: private key is required for all relay operations. + // Verified event reads deliberately use their command-local relay. The global + // relay value (including BUZZ_RELAY_URL and its localhost default) is ignored. + if let Cmd::Events(ref sub) = cli.command { + let (keys, auth_tag) = + parse_optional_identity(cli.private_key.as_deref(), cli.auth_tag.as_deref())?; + return commands::events::dispatch(sub, keys.as_ref(), auth_tag.as_ref()).await; + } + + let relay_url = client::normalize_relay_url(&cli.relay); + + // Auth: private key is required for all other relay operations. // The keypair IS the identity — no tokens, no other auth. let private_key_str = cli.private_key.ok_or_else(|| { CliError::Auth("BUZZ_PRIVATE_KEY is required (use --private-key or set env var)".into()) @@ -1826,10 +1854,40 @@ async fn run(cli: Cli) -> Result<(), CliError> { Cmd::Upload(sub) => commands::upload::dispatch(sub, &client).await, Cmd::Mem(sub) => commands::mem::dispatch(sub, &client).await, Cmd::Moderation(sub) => commands::moderation::dispatch(sub, &client, &cli.format).await, - Cmd::Pack(_) => unreachable!("handled above"), + Cmd::Events(_) | Cmd::Pack(_) => unreachable!("handled above"), } } +fn parse_optional_identity( + private_key: Option<&str>, + auth_tag_json: Option<&str>, +) -> Result<(Option, Option), CliError> { + let keys = private_key + .map(|value| { + Keys::parse(value) + .map_err(|error| CliError::Key(format!("invalid BUZZ_PRIVATE_KEY: {error}"))) + }) + .transpose()?; + let auth_tag = match auth_tag_json { + Some(json) if !json.is_empty() => { + let keys = keys + .as_ref() + .ok_or_else(|| CliError::Auth("BUZZ_AUTH_TAG requires BUZZ_PRIVATE_KEY".into()))?; + let tag = buzz_sdk::nip_oa::parse_auth_tag(json) + .map_err(|error| CliError::Auth(format!("BUZZ_AUTH_TAG is malformed: {error}")))?; + buzz_sdk::nip_oa::verify_auth_tag(json, &keys.public_key()).map_err(|error| { + CliError::Auth(format!( + "BUZZ_AUTH_TAG verification failed for pubkey {}: {error}", + keys.public_key().to_hex() + )) + })?; + Some(tag) + } + _ => None, + }; + Ok((keys, auth_tag)) +} + #[cfg(test)] mod tests { use super::*; @@ -1873,6 +1931,7 @@ mod tests { "channels", "dms", "emoji", + "events", "feed", "issues", "media", @@ -1999,6 +2058,7 @@ mod tests { vec!["approve", "create", "delete", "get", "list", "runs", "trigger", "update"] ); assert_eq!(names(&cmd, "feed"), vec!["get"]); + assert_eq!(names(&cmd, "events"), vec!["get-verified"]); assert_eq!( names(&cmd, "social"), vec![ @@ -2069,6 +2129,7 @@ mod tests { ("dms", 4), ("emoji", 5), ("feed", 1), + ("events", 1), ("issues", 4), ("media", 1), ("messages", 8), diff --git a/crates/buzz-cli/tests/events_get_verified.rs b/crates/buzz-cli/tests/events_get_verified.rs new file mode 100644 index 0000000000..9ccfea4098 --- /dev/null +++ b/crates/buzz-cli/tests/events_get_verified.rs @@ -0,0 +1,1280 @@ +use std::collections::BTreeSet; +use std::process::Output; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; + +use futures_util::{SinkExt, StreamExt}; +use nostr::{Event, EventBuilder, JsonUtil, Keys, Kind}; +use serde_json::{json, Value}; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::TcpListener; +use tokio::process::Command; +use tokio::time::{timeout, Duration, Instant}; +use tokio_tungstenite::tungstenite::Message; + +fn signed_event(content: &str) -> Event { + EventBuilder::new(Kind::TextNote, content) + .sign_with_keys(&Keys::generate()) + .unwrap() +} + +fn mutate(event: &Event, field: &str, value: Value) -> Event { + let mut json = serde_json::to_value(event).unwrap(); + json[field] = value; + Event::from_json(json.to_string()).unwrap() +} + +fn raw_event_with_layout(event: &Event, id: &str) -> String { + let mut value = serde_json::to_value(event).unwrap(); + value["id"] = json!(id); + format!( + concat!( + "{{\n", + " \"sig\" : {},\n", + " \"content\" : {},\n", + " \"tags\" : {},\n", + " \"kind\" : {},\n", + " \"created_at\" : {},\n", + " \"pubkey\" : {},\n", + " \"id\" : {}\n", + "}}" + ), + value["sig"], + value["content"], + value["tags"], + value["kind"], + value["created_at"], + value["pubkey"], + value["id"], + ) +} + +async fn relay_sending(frames: Vec) -> (String, tokio::task::JoinHandle) { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let handle = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + let request = websocket.next().await.unwrap().unwrap(); + let request: Value = serde_json::from_str(request.to_text().unwrap()).unwrap(); + let subscription_id = request[1].as_str().unwrap(); + for frame in frames { + let frame = match frame { + Value::Array(mut parts) if parts.get(1) == Some(&json!("$subscription")) => { + parts[1] = json!(subscription_id); + Value::Array(parts) + } + other => other, + }; + websocket + .send(Message::Text(frame.to_string().into())) + .await + .unwrap(); + } + request + }); + (format!("ws://{address}"), handle) +} + +async fn relay_sending_raw_event(raw_event: String) -> (String, tokio::task::JoinHandle) { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let handle = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + let request = websocket.next().await.unwrap().unwrap(); + let request: Value = serde_json::from_str(request.to_text().unwrap()).unwrap(); + let subscription_id = request[1].as_str().unwrap(); + websocket + .send(Message::Text( + format!(r#"["EVENT","{subscription_id}",{raw_event}]"#).into(), + )) + .await + .unwrap(); + websocket + .send(Message::Text( + json!(["EOSE", subscription_id]).to_string().into(), + )) + .await + .unwrap(); + request + }); + (format!("ws://{address}"), handle) +} + +async fn run_buzz(relay: &str, event_id: &str) -> Output { + let mut command = Command::new(env!("CARGO_BIN_EXE_buzz")); + command + .kill_on_drop(true) + .env_remove("BUZZ_RELAY_URL") + .env_remove("BUZZ_PRIVATE_KEY") + .env_remove("BUZZ_AUTH_TAG") + .args([ + "events", + "get-verified", + "--relay", + relay, + "--event", + event_id, + ]); + command.output().await.unwrap() +} + +async fn run_buzz_with_private_key(relay: &str, event_id: &str, private_key: &str) -> Output { + run_buzz_with_identity(relay, event_id, private_key, None).await +} + +async fn run_buzz_with_identity( + relay: &str, + event_id: &str, + private_key: &str, + auth_tag: Option<&str>, +) -> Output { + let mut command = Command::new(env!("CARGO_BIN_EXE_buzz")); + command + .kill_on_drop(true) + .env_remove("BUZZ_RELAY_URL") + .env_remove("BUZZ_PRIVATE_KEY") + .env_remove("BUZZ_AUTH_TAG") + .args([ + "--private-key", + private_key, + "events", + "get-verified", + "--relay", + relay, + "--event", + event_id, + ]); + if let Some(auth_tag) = auth_tag { + command.env("BUZZ_AUTH_TAG", auth_tag); + } + command.output().await.unwrap() +} + +fn error_category(output: &Output) -> String { + let error: Value = serde_json::from_slice(&output.stderr).unwrap(); + error["error"].as_str().unwrap().to_string() +} + +fn assert_failure(output: &Output, exit_code: i32, category: &str) { + assert_eq!(output.status.code(), Some(exit_code)); + assert!(output.stdout.is_empty(), "failure wrote to stdout"); + assert_eq!(error_category(output), category); +} + +#[tokio::test] +async fn emits_only_raw_signed_fields_after_exact_verified_fetch() { + let event = signed_event("verified output"); + let raw_event = raw_event_with_layout(&event, &event.id.to_hex()); + let (relay, server) = relay_sending_raw_event(raw_event.clone()).await; + + let output = run_buzz(&relay, &event.id.to_hex()).await; + let request = server.await.unwrap(); + + assert!(output.status.success()); + assert!(output.stderr.is_empty()); + assert_eq!(output.stdout, format!("{raw_event}\n").as_bytes()); + let emitted: Value = serde_json::from_slice(&output.stdout).unwrap(); + assert_eq!(emitted, serde_json::to_value(&event).unwrap()); + let keys: BTreeSet<&str> = emitted + .as_object() + .unwrap() + .keys() + .map(String::as_str) + .collect(); + assert_eq!( + keys, + BTreeSet::from([ + "content", + "created_at", + "id", + "kind", + "pubkey", + "sig", + "tags" + ]) + ); + assert_eq!( + request, + json!([ + "REQ", + "buzz-get-verified-0", + {"ids": [event.id.to_hex()], "limit": 2} + ]) + ); +} + +#[tokio::test] +async fn uppercase_raw_event_id_is_rejected_instead_of_normalized() { + let event = signed_event("uppercase raw ID"); + let raw_event = raw_event_with_layout(&event, &event.id.to_hex().to_ascii_uppercase()); + let (relay, server) = relay_sending_raw_event(raw_event).await; + + let output = run_buzz(&relay, &event.id.to_hex()).await; + server.await.unwrap(); + + assert_failure(&output, 4, "relay_mismatch"); +} + +#[tokio::test] +async fn malformed_raw_signature_is_not_a_generic_protocol_error() { + let event = signed_event("malformed raw signature"); + let mut raw_event = serde_json::to_value(&event).unwrap(); + raw_event["sig"] = json!("not-a-signature"); + let (relay, server) = relay_sending_raw_event(raw_event.to_string()).await; + + let output = run_buzz(&relay, &event.id.to_hex()).await; + server.await.unwrap(); + + assert_failure(&output, 4, "signature_invalid"); +} + +#[tokio::test] +async fn optional_nip42_challenge_does_not_make_public_read_require_a_key() { + let event = signed_event("public read"); + let frames = vec![ + json!(["AUTH", "optional-challenge"]), + json!(["EVENT", "$subscription", event]), + json!(["EOSE", "$subscription"]), + ]; + let (relay, server) = relay_sending(frames).await; + + let output = run_buzz(&relay, &event.id.to_hex()).await; + server.await.unwrap(); + + assert!(output.status.success()); + assert!(output.stderr.is_empty()); + let emitted: Event = serde_json::from_slice(&output.stdout).unwrap(); + assert_eq!(emitted, event); +} + +#[tokio::test] +async fn optional_challenge_does_not_transmit_ambient_identity_or_auth_tag() { + let event = signed_event("public read with ambient identity"); + let event_id = event.id; + let agent_keys = Keys::generate(); + let private_key = agent_keys.secret_key().to_secret_hex(); + let owner_keys = Keys::generate(); + let auth_tag = + buzz_sdk::nip_oa::compute_auth_tag(&owner_keys, &agent_keys.public_key(), "kind=1") + .unwrap(); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let relay = format!("ws://{address}"); + + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + websocket + .send(Message::Text( + json!(["AUTH", "optional-challenge"]).to_string().into(), + )) + .await + .unwrap(); + let request = websocket.next().await.unwrap().unwrap(); + let request: Value = serde_json::from_str(request.to_text().unwrap()).unwrap(); + assert_eq!(request[0], "REQ"); + let subscription_id = request[1].as_str().unwrap(); + websocket + .send(Message::Text( + json!(["EVENT", subscription_id, event]).to_string().into(), + )) + .await + .unwrap(); + websocket + .send(Message::Text( + json!(["EOSE", subscription_id]).to_string().into(), + )) + .await + .unwrap(); + + let next = websocket.next().await.unwrap().unwrap(); + let next: Value = serde_json::from_str(next.to_text().unwrap()).unwrap(); + assert_eq!(next[0], "CLOSE", "optional challenge elicited AUTH: {next}"); + }); + + let output = + run_buzz_with_identity(&relay, &event_id.to_hex(), &private_key, Some(&auth_tag)).await; + server.await.unwrap(); + + assert!(output.status.success()); + assert!(output.stderr.is_empty()); +} + +#[tokio::test] +async fn zero_and_multiple_results_fail_without_stdout() { + let requested = signed_event("requested"); + + let (empty_relay, empty_server) = relay_sending(vec![json!(["EOSE", "$subscription"])]).await; + let empty = run_buzz(&empty_relay, &requested.id.to_hex()).await; + empty_server.await.unwrap(); + assert_failure(&empty, 1, "not_found"); + + let frames = vec![ + json!(["EVENT", "$subscription", requested]), + json!(["EVENT", "$subscription", requested]), + json!(["EOSE", "$subscription"]), + ]; + let (multiple_relay, multiple_server) = relay_sending(frames).await; + let multiple = run_buzz(&multiple_relay, &requested.id.to_hex()).await; + multiple_server.await.unwrap(); + assert_failure(&multiple, 4, "ambiguous"); +} + +#[tokio::test] +async fn second_event_fails_immediately_without_waiting_for_eose() { + let event = signed_event("immediate ambiguity"); + let event_id = event.id; + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let relay = format!("ws://{address}"); + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + let request = websocket.next().await.unwrap().unwrap(); + let request: Value = serde_json::from_str(request.to_text().unwrap()).unwrap(); + let subscription_id = request[1].as_str().unwrap(); + for _ in 0..2 { + websocket + .send(Message::Text( + json!(["EVENT", subscription_id, event]).to_string().into(), + )) + .await + .unwrap(); + } + std::future::pending::<()>().await; + }); + + let output = timeout(Duration::from_secs(2), run_buzz(&relay, &event_id.to_hex())) + .await + .expect("command waited for EOSE after the second event"); + + assert_failure(&output, 4, "ambiguous"); + server.abort(); +} + +#[tokio::test] +async fn wrong_event_mutated_content_and_invalid_signature_are_distinct() { + let requested = signed_event("requested"); + let wrong = signed_event("wrong"); + let id_mismatch = mutate(&requested, "content", json!("mutated")); + let invalid_signature = mutate(&requested, "sig", json!("0".repeat(128))); + + for (returned, category) in [ + (wrong, "relay_mismatch"), + (id_mismatch, "id_mismatch"), + (invalid_signature, "signature_invalid"), + ] { + let frames = vec![ + json!(["EVENT", "$subscription", returned]), + json!(["EOSE", "$subscription"]), + ]; + let (relay, server) = relay_sending(frames).await; + let output = run_buzz(&relay, &requested.id.to_hex()).await; + server.await.unwrap(); + assert_failure(&output, 4, category); + } +} + +#[tokio::test] +async fn subscription_mismatch_and_transport_failure_are_distinct() { + let requested = signed_event("requested"); + let (relay, server) = relay_sending(vec![json!(["EOSE", "wrong-subscription"])]).await; + let mismatch = run_buzz(&relay, &requested.id.to_hex()).await; + server.await.unwrap(); + assert_failure(&mismatch, 4, "relay_mismatch"); + + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let unavailable = format!("ws://{}", listener.local_addr().unwrap()); + drop(listener); + let transport = run_buzz(&unavailable, &requested.id.to_hex()).await; + assert_failure(&transport, 2, "transport_error"); +} + +#[tokio::test] +async fn websocket_upgrade_401_and_403_are_non_retryable_auth_errors() { + let event = signed_event("upgrade denial"); + + for (status, reason) in [(401, "Unauthorized"), (403, "Forbidden")] { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let relay = format!("ws://{address}"); + let server = tokio::spawn(async move { + let (mut stream, _) = listener.accept().await.unwrap(); + let mut request = vec![0_u8; 4096]; + let bytes = stream.read(&mut request).await.unwrap(); + assert!(bytes > 0, "client did not attempt a WebSocket upgrade"); + stream + .write_all( + format!( + "HTTP/1.1 {status} {reason}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n" + ) + .as_bytes(), + ) + .await + .unwrap(); + }); + + let output = run_buzz(&relay, &event.id.to_hex()).await; + server.await.unwrap(); + assert_failure(&output, 3, "auth_error"); + let error: Value = serde_json::from_slice(&output.stderr).unwrap(); + assert_eq!(error["retryable"], false); + } +} + +#[tokio::test] +async fn command_local_relay_is_required_even_when_default_env_is_set() { + let event = signed_event("no fallback"); + let output = Command::new(env!("CARGO_BIN_EXE_buzz")) + .env("BUZZ_RELAY_URL", "ws://127.0.0.1:1") + .env_remove("BUZZ_PRIVATE_KEY") + .env_remove("BUZZ_AUTH_TAG") + .args(["events", "get-verified", "--event", &event.id.to_hex()]) + .output() + .await + .unwrap(); + + assert_failure(&output, 1, "user_error"); +} + +#[tokio::test] +async fn uppercase_event_id_is_rejected_before_any_relay_connection() { + let event = signed_event("canonical ID"); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let relay = format!("ws://{}", listener.local_addr().unwrap()); + + let output = run_buzz(&relay, &event.id.to_hex().to_ascii_uppercase()).await; + + assert_failure(&output, 1, "user_error"); + assert!( + timeout(Duration::from_millis(100), listener.accept()) + .await + .is_err(), + "uppercase ID unexpectedly reached the relay" + ); +} + +async fn authenticated_relay( + event: Event, + challenge_first: bool, + challenge: String, + required_reason: String, +) -> (String, tokio::task::JoinHandle<()>) { + let event_id = event.id; + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let relay = format!("ws://{address}"); + + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + if challenge_first { + websocket + .send(Message::Text(json!(["AUTH", challenge]).to_string().into())) + .await + .unwrap(); + } + + let first_request = websocket.next().await.unwrap().unwrap(); + let first_request: Value = serde_json::from_str(first_request.to_text().unwrap()).unwrap(); + let first_subscription = first_request[1].as_str().unwrap(); + + websocket + .send(Message::Text( + json!(["CLOSED", first_subscription, required_reason]) + .to_string() + .into(), + )) + .await + .unwrap(); + if !challenge_first { + websocket + .send(Message::Text(json!(["AUTH", challenge]).to_string().into())) + .await + .unwrap(); + } + + let auth = websocket.next().await.unwrap().unwrap(); + let auth: Value = serde_json::from_str(auth.to_text().unwrap()).unwrap(); + assert_eq!(auth[0], "AUTH"); + let auth_event: Event = serde_json::from_value(auth[1].clone()).unwrap(); + auth_event.verify().unwrap(); + assert!(auth_event + .tags + .iter() + .any(|tag| tag.as_slice() == ["challenge", challenge.as_str()])); + + websocket + .send(Message::Text( + json!(["OK", auth_event.id.to_hex(), true, ""]) + .to_string() + .into(), + )) + .await + .unwrap(); + + let second_request = websocket.next().await.unwrap().unwrap(); + let second_request: Value = + serde_json::from_str(second_request.to_text().unwrap()).unwrap(); + assert_eq!(second_request[1], "buzz-get-verified-1"); + assert_eq!(second_request[2]["ids"], json!([event_id.to_hex()])); + let second_subscription = second_request[1].as_str().unwrap(); + websocket + .send(Message::Text( + json!(["EVENT", second_subscription, event]) + .to_string() + .into(), + )) + .await + .unwrap(); + websocket + .send(Message::Text( + json!(["EOSE", second_subscription]).to_string().into(), + )) + .await + .unwrap(); + }); + + (relay, server) +} + +#[tokio::test] +async fn nip42_buffers_challenge_and_auth_required_closure_in_either_order() { + let event = signed_event("authenticated read"); + let event_id = event.id; + let keys = Keys::generate(); + let private_key = keys.secret_key().to_secret_hex(); + + for required_reason in ["auth-required", "auth-required: not authenticated"] { + for challenge_first in [true, false] { + let (relay, server) = authenticated_relay( + event.clone(), + challenge_first, + "challenge".into(), + required_reason.into(), + ) + .await; + let output = run_buzz_with_private_key(&relay, &event_id.to_hex(), &private_key).await; + server.await.unwrap(); + + assert!(output.status.success()); + assert!(output.stderr.is_empty()); + let emitted: Event = serde_json::from_slice(&output.stdout).unwrap(); + assert_eq!(emitted.id, event_id); + } + } +} + +#[tokio::test] +async fn authenticated_exact_event_may_carry_the_connections_nip_oa_authority() { + let agent_keys = Keys::generate(); + let private_key = agent_keys.secret_key().to_secret_hex(); + let owner_keys = Keys::generate(); + let auth_tag_json = + buzz_sdk::nip_oa::compute_auth_tag(&owner_keys, &agent_keys.public_key(), "kind=1") + .unwrap(); + let auth_tag = buzz_sdk::nip_oa::parse_auth_tag(&auth_tag_json).unwrap(); + let event = EventBuilder::new(Kind::TextNote, "delegated exact event") + .tags([auth_tag.clone()]) + .sign_with_keys(&agent_keys) + .unwrap(); + let event_id = event.id; + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let relay = format!("ws://{address}"); + + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + websocket + .send(Message::Text( + json!(["AUTH", "production-shape-challenge"]) + .to_string() + .into(), + )) + .await + .unwrap(); + + let first_request = websocket.next().await.unwrap().unwrap(); + let first_request: Value = serde_json::from_str(first_request.to_text().unwrap()).unwrap(); + let first_subscription = first_request[1].clone(); + websocket + .send(Message::Text( + json!(["NOTICE", "auth-required: authenticate before subscribing"]) + .to_string() + .into(), + )) + .await + .unwrap(); + websocket + .send(Message::Text( + json!([ + "CLOSED", + first_subscription, + "auth-required: not authenticated" + ]) + .to_string() + .into(), + )) + .await + .unwrap(); + + let auth_message = websocket.next().await.unwrap().unwrap(); + let auth: Value = serde_json::from_str(auth_message.to_text().unwrap()).unwrap(); + let auth_event: Event = serde_json::from_value(auth[1].clone()).unwrap(); + assert!(auth_event.tags.iter().any(|tag| tag == &auth_tag)); + websocket + .send(Message::Text( + json!(["OK", auth_event.id.to_hex(), true, ""]) + .to_string() + .into(), + )) + .await + .unwrap(); + + let second_request = websocket.next().await.unwrap().unwrap(); + let second_request: Value = + serde_json::from_str(second_request.to_text().unwrap()).unwrap(); + assert_eq!(second_request[1], "buzz-get-verified-1"); + websocket + .send(Message::Text( + json!(["EVENT", second_request[1], event]) + .to_string() + .into(), + )) + .await + .unwrap(); + websocket + .send(Message::Text( + json!(["EOSE", second_request[1]]).to_string().into(), + )) + .await + .unwrap(); + }); + + let output = run_buzz_with_identity( + &relay, + &event_id.to_hex(), + &private_key, + Some(&auth_tag_json), + ) + .await; + server.await.unwrap(); + + assert!(output.status.success(), "{}", error_category(&output)); + assert!(output.stderr.is_empty()); + let emitted: Event = serde_json::from_slice(&output.stdout).unwrap(); + assert_eq!(emitted.id, event_id); +} + +async fn auth_sequence_relay( + sequence: Vec, +) -> (String, tokio::task::JoinHandle>) { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + let request = websocket.next().await.unwrap().unwrap(); + let request: Value = serde_json::from_str(request.to_text().unwrap()).unwrap(); + let subscription_id = request[1].clone(); + + for frame in sequence { + let frame = match frame { + Value::Array(mut parts) if parts.get(1) == Some(&json!("$subscription")) => { + parts[1] = subscription_id.clone(); + Value::Array(parts) + } + other => other, + }; + websocket + .send(Message::Text(frame.to_string().into())) + .await + .unwrap(); + } + + let mut received = Vec::new(); + while let Ok(Some(Ok(message))) = timeout(Duration::from_secs(2), websocket.next()).await { + match message { + Message::Text(text) => { + received.push(serde_json::from_str(&text).unwrap()); + } + Message::Close(_) => break, + _ => {} + } + } + received + }); + (format!("ws://{address}"), server) +} + +#[tokio::test] +async fn duplicate_pre_auth_challenges_are_ambiguous_in_either_order() { + let event = signed_event("duplicate auth challenge"); + let keys = Keys::generate(); + let private_key = keys.secret_key().to_secret_hex(); + + for (sequence, expected_auth_count) in [ + ( + vec![ + json!(["AUTH", "first-challenge"]), + json!(["AUTH", "second-challenge"]), + json!(["CLOSED", "$subscription", "auth-required: duplicate"]), + ], + 0, + ), + ( + vec![ + json!(["CLOSED", "$subscription", "auth-required: duplicate"]), + json!(["AUTH", "first-challenge"]), + json!(["AUTH", "second-challenge"]), + ], + 1, + ), + ] { + let (relay, server) = auth_sequence_relay(sequence).await; + let output = run_buzz_with_private_key(&relay, &event.id.to_hex(), &private_key).await; + let received = server.await.unwrap(); + + assert_failure(&output, 4, "ambiguous"); + let error: Value = serde_json::from_slice(&output.stderr).unwrap(); + if expected_auth_count == 0 { + assert!(error["message"] + .as_str() + .unwrap() + .contains("multiple AUTH challenges")); + } else { + // The second challenge arrives after the signed AUTH event but + // before relay OK. The output boundary starts at transmission. + assert_eq!( + error["message"], + "ambiguous result: relay returned an ambiguous result after private authentication" + ); + } + let auth_messages: Vec<&Value> = received + .iter() + .filter(|message| message.get(0) == Some(&json!("AUTH"))) + .collect(); + assert_eq!(auth_messages.len(), expected_auth_count, "{received:?}"); + if let Some(auth) = auth_messages.first() { + let event: Event = serde_json::from_value(auth[1].clone()).unwrap(); + assert!(event + .tags + .iter() + .any(|tag| tag.as_slice() == ["challenge", "first-challenge"])); + } + } +} + +#[tokio::test] +async fn auth_reason_lookalikes_never_trigger_credential_signing() { + let event = signed_event("auth-required lookalikes"); + let keys = Keys::generate(); + let private_key = keys.secret_key().to_secret_hex(); + + for reason in [ + "not-auth-required", + "restricted: not-auth-required", + "restricted auth-required response", + "auth-required-suffix", + "AUTH-REQUIRED", + ] { + let (relay, server) = auth_sequence_relay(vec![ + json!(["AUTH", "untrusted-reason-challenge"]), + json!(["CLOSED", "$subscription", reason]), + ]) + .await; + let output = run_buzz_with_private_key(&relay, &event.id.to_hex(), &private_key).await; + let received = server.await.unwrap(); + + assert_failure(&output, 2, "relay_error"); + assert!( + received + .iter() + .all(|message| message.get(0) != Some(&json!("AUTH"))), + "lookalike reason {reason:?} elicited AUTH: {received:?}" + ); + } +} + +#[tokio::test] +async fn auth_challenge_limit_accepts_1024_bytes_and_refuses_1025_before_signing() { + let event = signed_event("challenge size"); + let keys = Keys::generate(); + let private_key = keys.secret_key().to_secret_hex(); + let exact = "x".repeat(1024); + let (relay, server) = authenticated_relay( + event.clone(), + false, + exact, + "auth-required: oversized-boundary-control".into(), + ) + .await; + + let output = run_buzz_with_private_key(&relay, &event.id.to_hex(), &private_key).await; + server.await.unwrap(); + assert!(output.status.success()); + assert!(output.stderr.is_empty()); + + let oversized = "x".repeat(1025); + let (relay, server) = auth_sequence_relay(vec![ + json!(["CLOSED", "$subscription", "auth-required: oversized"]), + json!(["AUTH", oversized]), + ]) + .await; + let output = run_buzz_with_private_key(&relay, &event.id.to_hex(), &private_key).await; + let received = server.await.unwrap(); + + assert_failure(&output, 2, "relay_error"); + let error: Value = serde_json::from_slice(&output.stderr).unwrap(); + assert!(error["message"].as_str().unwrap().contains("1025 bytes")); + assert!( + received + .iter() + .all(|message| message.get(0) != Some(&json!("AUTH"))), + "oversized challenge was signed before refusal: {received:?}" + ); +} + +#[tokio::test] +async fn reflection_before_auth_ok_cannot_escape_through_stderr() { + let event = signed_event("malformed auth echo"); + let event_id = event.id; + let agent_keys = Keys::generate(); + let private_key = agent_keys.secret_key().to_secret_hex(); + let owner_keys = Keys::generate(); + let auth_tag = + buzz_sdk::nip_oa::compute_auth_tag(&owner_keys, &agent_keys.public_key(), "kind=1") + .unwrap(); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let relay = format!("ws://{address}"); + + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + let request = websocket.next().await.unwrap().unwrap(); + let request: Value = serde_json::from_str(request.to_text().unwrap()).unwrap(); + let subscription_id = request[1].clone(); + websocket + .send(Message::Text( + json!(["CLOSED", subscription_id, "auth-required"]) + .to_string() + .into(), + )) + .await + .unwrap(); + websocket + .send(Message::Text( + json!(["AUTH", "malformed-echo-challenge"]) + .to_string() + .into(), + )) + .await + .unwrap(); + + let auth_message = websocket.next().await.unwrap().unwrap(); + let auth_text = auth_message.to_text().unwrap().to_string(); + let auth: Value = serde_json::from_str(&auth_text).unwrap(); + assert_eq!(auth[0], "AUTH"); + + // Deliberately invalid JSON that reflects the complete signed event. + // The client must report only a static protocol category. + websocket + .send(Message::Text(format!("['AUTH',{}]", auth[1]).into())) + .await + .unwrap(); + auth_text + }); + + let output = + run_buzz_with_identity(&relay, &event_id.to_hex(), &private_key, Some(&auth_tag)).await; + let auth_text = server.await.unwrap(); + + assert_failure(&output, 2, "relay_error"); + let stderr = String::from_utf8(output.stderr).unwrap(); + let auth: Value = serde_json::from_str(&auth_text).unwrap(); + let auth_event_json = auth[1].to_string(); + let signature = auth[1]["sig"].as_str().unwrap(); + for private_bytes in [ + auth_tag.as_str(), + auth_text.as_str(), + auth_event_json.as_str(), + signature, + ] { + assert!( + !stderr.contains(private_bytes), + "post-AUTH error leaked private bytes {private_bytes:?}: {stderr}" + ); + } + let error: Value = serde_json::from_str(&stderr).unwrap(); + assert_eq!( + error["message"], + "relay protocol error: relay sent an invalid response after private authentication" + ); +} + +#[tokio::test] +async fn authenticated_closed_reflection_cannot_escape_through_cli_sinks() { + let event = signed_event("authenticated CLOSED reflection"); + let event_id = event.id; + let agent_keys = Keys::generate(); + let private_key = agent_keys.secret_key().to_secret_hex(); + let owner_keys = Keys::generate(); + let auth_tag = + buzz_sdk::nip_oa::compute_auth_tag(&owner_keys, &agent_keys.public_key(), "kind=1") + .unwrap(); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let relay = format!("ws://{address}"); + + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + let request = websocket.next().await.unwrap().unwrap(); + let request: Value = serde_json::from_str(request.to_text().unwrap()).unwrap(); + let subscription_id = request[1].clone(); + websocket + .send(Message::Text( + json!(["CLOSED", subscription_id, "auth-required"]) + .to_string() + .into(), + )) + .await + .unwrap(); + websocket + .send(Message::Text( + json!(["AUTH", "closed-reflection-challenge"]) + .to_string() + .into(), + )) + .await + .unwrap(); + + let auth_message = websocket.next().await.unwrap().unwrap(); + let auth_text = auth_message.to_text().unwrap().to_string(); + let auth: Value = serde_json::from_str(&auth_text).unwrap(); + let auth_event: Event = serde_json::from_value(auth[1].clone()).unwrap(); + websocket + .send(Message::Text( + json!(["OK", auth_event.id.to_hex(), true, ""]) + .to_string() + .into(), + )) + .await + .unwrap(); + + let second_request = websocket.next().await.unwrap().unwrap(); + let second_request: Value = + serde_json::from_str(second_request.to_text().unwrap()).unwrap(); + assert_eq!(second_request[1], "buzz-get-verified-1"); + websocket + .send(Message::Text( + json!(["CLOSED", second_request[1], format!("echo:{auth_text}")]) + .to_string() + .into(), + )) + .await + .unwrap(); + auth_text + }); + + let output = + run_buzz_with_identity(&relay, &event_id.to_hex(), &private_key, Some(&auth_tag)).await; + let auth_text = server.await.unwrap(); + + assert_failure(&output, 2, "relay_error"); + let stdout = String::from_utf8(output.stdout).unwrap(); + let stderr = String::from_utf8(output.stderr).unwrap(); + let auth: Value = serde_json::from_str(&auth_text).unwrap(); + let event_signature = auth[1]["sig"].as_str().unwrap(); + let auth_tag_json: Value = serde_json::from_str(&auth_tag).unwrap(); + let authority_signature = auth_tag_json[3].as_str().unwrap(); + for private_bytes in [ + auth_tag.as_str(), + auth_text.as_str(), + event_signature, + authority_signature, + ] { + assert!( + !stdout.contains(private_bytes), + "authenticated CLOSED reflected private bytes to stdout {private_bytes:?}: {stdout}" + ); + assert!( + !stderr.contains(private_bytes), + "authenticated CLOSED reflected private bytes to stderr {private_bytes:?}: {stderr}" + ); + } +} + +#[derive(Clone, Copy)] +enum AuthFieldReflection { + UppercaseEventSignature, + SplitAuthoritySignature, +} + +async fn authenticated_unknown_field_relay( + event: Event, + reflection: AuthFieldReflection, +) -> (String, tokio::task::JoinHandle) { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let relay = format!("ws://{address}"); + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + let request = websocket.next().await.unwrap().unwrap(); + let request: Value = serde_json::from_str(request.to_text().unwrap()).unwrap(); + let subscription_id = request[1].clone(); + websocket + .send(Message::Text( + json!(["CLOSED", subscription_id, "auth-required"]) + .to_string() + .into(), + )) + .await + .unwrap(); + websocket + .send(Message::Text( + json!(["AUTH", "unknown-field-reflection-challenge"]) + .to_string() + .into(), + )) + .await + .unwrap(); + + let auth_message = websocket.next().await.unwrap().unwrap(); + let auth: Value = serde_json::from_str(auth_message.to_text().unwrap()).unwrap(); + let auth_event: Event = serde_json::from_value(auth[1].clone()).unwrap(); + websocket + .send(Message::Text( + json!(["OK", auth_event.id.to_hex(), true, ""]) + .to_string() + .into(), + )) + .await + .unwrap(); + + let second_request = websocket.next().await.unwrap().unwrap(); + let second_request: Value = + serde_json::from_str(second_request.to_text().unwrap()).unwrap(); + let second_subscription = second_request[1].clone(); + let authority_signature = auth_event + .tags + .iter() + .find_map(|tag| { + let parts = tag.as_slice(); + (parts.first().map(String::as_str) == Some("auth")) + .then(|| parts.last().cloned()) + .flatten() + }) + .unwrap(); + let (field_name, private_signature) = match reflection { + AuthFieldReflection::UppercaseEventSignature => ( + auth_event.sig.to_string().to_ascii_uppercase(), + auth_event.sig.to_string(), + ), + AuthFieldReflection::SplitAuthoritySignature => { + let midpoint = authority_signature.len() / 2; + ( + format!( + "{}--{}", + &authority_signature[..midpoint], + &authority_signature[midpoint..] + ), + authority_signature, + ) + } + }; + let mut reflected_event = serde_json::to_value(&event).unwrap(); + reflected_event + .as_object_mut() + .unwrap() + .insert(field_name, json!("relay-controlled")); + websocket + .send(Message::Text( + json!(["EVENT", second_subscription, reflected_event]) + .to_string() + .into(), + )) + .await + .unwrap(); + websocket + .send(Message::Text( + json!(["EOSE", second_subscription]).to_string().into(), + )) + .await + .unwrap(); + private_signature + }); + (relay, server) +} + +async fn assert_post_auth_field_is_static(reflection: AuthFieldReflection) { + let event = signed_event("transformed auth reflection"); + let event_id = event.id; + let agent_keys = Keys::generate(); + let private_key = agent_keys.secret_key().to_secret_hex(); + let owner_keys = Keys::generate(); + let auth_tag = + buzz_sdk::nip_oa::compute_auth_tag(&owner_keys, &agent_keys.public_key(), "kind=1") + .unwrap(); + let (relay, server) = authenticated_unknown_field_relay(event, reflection).await; + + let output = + run_buzz_with_identity(&relay, &event_id.to_hex(), &private_key, Some(&auth_tag)).await; + let private_signature = server.await.unwrap(); + + assert_failure(&output, 4, "relay_mismatch"); + assert!(output.stdout.is_empty()); + let stderr = String::from_utf8(output.stderr).unwrap(); + let error: Value = serde_json::from_str(&stderr).unwrap(); + assert_eq!( + error["message"], + "relay mismatch: relay response did not match the exact-event request after private authentication" + ); + let normalized_stderr = stderr.to_ascii_lowercase().replace("--", ""); + assert!( + !normalized_stderr.contains(&private_signature.to_ascii_lowercase()), + "post-AUTH relay-controlled field leaked transformed private bytes: {stderr}" + ); +} + +#[tokio::test] +async fn uppercased_post_auth_field_cannot_escape_through_cli_errors() { + assert_post_auth_field_is_static(AuthFieldReflection::UppercaseEventSignature).await; +} + +#[tokio::test] +async fn split_post_auth_field_cannot_escape_through_cli_errors() { + assert_post_auth_field_is_static(AuthFieldReflection::SplitAuthoritySignature).await; +} + +async fn ping_stall_relay() -> (String, Arc, tokio::task::JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let request_seen = Arc::new(AtomicBool::new(false)); + let request_seen_by_server = Arc::clone(&request_seen); + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + let request = websocket.next().await.unwrap().unwrap(); + let request: Value = serde_json::from_str(request.to_text().unwrap()).unwrap(); + assert_eq!(request[0], "REQ"); + request_seen_by_server.store(true, Ordering::SeqCst); + + loop { + if websocket + .send(Message::Ping(Vec::new().into())) + .await + .is_err() + { + break; + } + tokio::time::sleep(Duration::from_millis(100)).await; + } + }); + (format!("ws://{address}"), request_seen, server) +} + +async fn handshake_stall_relay() -> (String, Arc, tokio::task::JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let connected = Arc::new(AtomicBool::new(false)); + let connected_by_server = Arc::clone(&connected); + let server = tokio::spawn(async move { + let (_stream, _) = listener.accept().await.unwrap(); + connected_by_server.store(true, Ordering::SeqCst); + std::future::pending::<()>().await; + }); + (format!("ws://{address}"), connected, server) +} + +async fn auth_stall_relay() -> (String, Arc, tokio::task::JoinHandle<()>) { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let auth_seen = Arc::new(AtomicBool::new(false)); + let auth_seen_by_server = Arc::clone(&auth_seen); + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + websocket + .send(Message::Text( + json!(["AUTH", "deadline-challenge"]).to_string().into(), + )) + .await + .unwrap(); + + let first = websocket.next().await.unwrap().unwrap(); + let first: Value = serde_json::from_str(first.to_text().unwrap()).unwrap(); + assert_eq!(first[0], "REQ"); + let subscription_id = first[1].clone(); + websocket + .send(Message::Text( + json!(["CLOSED", subscription_id, "auth-required: deadline"]) + .to_string() + .into(), + )) + .await + .unwrap(); + let second = websocket.next().await.unwrap().unwrap(); + let second: Value = serde_json::from_str(second.to_text().unwrap()).unwrap(); + assert_eq!(second[0], "AUTH"); + auth_seen_by_server.store(true, Ordering::SeqCst); + std::future::pending::<()>().await; + }); + (format!("ws://{address}"), auth_seen, server) +} + +async fn deadline_bounded_output( + output: impl std::future::Future, +) -> (Duration, Output) { + const TEST_CEILING: Duration = Duration::from_millis(16_500); + let started = Instant::now(); + let output = timeout(TEST_CEILING, output) + .await + .expect("command exceeded the 15-second transaction deadline"); + (started.elapsed(), output) +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn one_deadline_bounds_ping_connect_and_auth_stalls() { + let event = signed_event("deadline"); + let event_id = event.id.to_hex(); + let keys = Keys::generate(); + let private_key = keys.secret_key().to_secret_hex(); + let (ping_relay, ping_request_seen, ping_server) = ping_stall_relay().await; + let (connect_relay, connect_seen, connect_server) = handshake_stall_relay().await; + let (auth_relay, auth_seen, auth_server) = auth_stall_relay().await; + + let (ping, connect, auth) = tokio::join!( + deadline_bounded_output(run_buzz(&ping_relay, &event_id)), + deadline_bounded_output(run_buzz(&connect_relay, &event_id)), + deadline_bounded_output(run_buzz_with_private_key( + &auth_relay, + &event_id, + &private_key + )), + ); + + for (name, (elapsed, output)) in [("PING", ping), ("connect", connect), ("AUTH", auth)] { + assert!( + elapsed >= Duration::from_secs(14), + "{name} stall ended vacuously after {elapsed:?}" + ); + assert!( + elapsed < Duration::from_millis(16_500), + "{name} stall exceeded deadline: {elapsed:?}" + ); + assert_failure(&output, 2, "transport_error"); + } + + assert!(ping_request_seen.load(Ordering::SeqCst)); + assert!(connect_seen.load(Ordering::SeqCst)); + assert!(auth_seen.load(Ordering::SeqCst)); + ping_server.abort(); + connect_server.abort(); + auth_server.abort(); +} diff --git a/crates/buzz-test-client/src/lib.rs b/crates/buzz-test-client/src/lib.rs index 89bff3020d..adee108838 100644 --- a/crates/buzz-test-client/src/lib.rs +++ b/crates/buzz-test-client/src/lib.rs @@ -62,6 +62,9 @@ impl From for TestClientError { match e { WsClientError::WebSocket(e) => TestClientError::WebSocket(e), WsClientError::Json(e) => TestClientError::Json(e), + WsClientError::InvalidEvent { message, .. } => { + TestClientError::UnexpectedMessage(format!("invalid EVENT payload: {message}")) + } WsClientError::EventBuilder(s) => TestClientError::EventBuilder(s), WsClientError::Url(s) => TestClientError::Url(s), WsClientError::Timeout => TestClientError::Timeout, @@ -70,6 +73,13 @@ impl From for TestClientError { WsClientError::AuthFailed(s) => TestClientError::AuthFailed(s), WsClientError::EventRejected(s) => TestClientError::EventRejected(s), WsClientError::NoAuthChallenge => TestClientError::NoAuthChallenge, + WsClientError::AuthChallengeTooLarge { .. } + | WsClientError::AmbiguousAuthChallenge + | WsClientError::AuthFrameTooLarge { .. } + | WsClientError::AuthTransportPoisoned + | WsClientError::ReflectedAuthMaterial => { + TestClientError::UnexpectedMessage(e.to_string()) + } } } } @@ -197,6 +207,7 @@ impl BuzzTestClient { RelayMessage::Event { subscription_id, event, + .. } if subscription_id == sub_id => { events.push(*event); } diff --git a/crates/buzz-test-client/src/main.rs b/crates/buzz-test-client/src/main.rs index 858f795c83..976fbd3446 100644 --- a/crates/buzz-test-client/src/main.rs +++ b/crates/buzz-test-client/src/main.rs @@ -120,6 +120,7 @@ async fn run_subscribe(url: &str, keys: &Keys, channel: &str, kind: u16) { Ok(RelayMessage::Event { subscription_id: _, event, + .. }) => { println!( "[{}] kind={} pubkey={} content={}", diff --git a/crates/buzz-test-client/tests/conformance_multitenant.rs b/crates/buzz-test-client/tests/conformance_multitenant.rs index 15002142e4..91b4b100ff 100644 --- a/crates/buzz-test-client/tests/conformance_multitenant.rs +++ b/crates/buzz-test-client/tests/conformance_multitenant.rs @@ -2443,6 +2443,7 @@ mod pubsub_presence_typing { Ok(RelayMessage::Event { subscription_id, event, + .. }) if subscription_id == sub_id => events.push(*event), Ok(_) => {} Err(_) => return events, diff --git a/crates/buzz-test-client/tests/e2e_event_reminder.rs b/crates/buzz-test-client/tests/e2e_event_reminder.rs index 4c17eb0f05..34cd531d03 100644 --- a/crates/buzz-test-client/tests/e2e_event_reminder.rs +++ b/crates/buzz-test-client/tests/e2e_event_reminder.rs @@ -1119,6 +1119,7 @@ async fn await_scheduler_push( Ok(RelayMessage::Event { subscription_id, event, + .. }) if subscription_id == sub_id && has_d_tag(&event, d_tag) => return Ok(()), Ok(_) => {} Err(e) => return Err(format!("recv failed: {e:?}")), diff --git a/crates/buzz-test-client/tests/e2e_nostr_interop.rs b/crates/buzz-test-client/tests/e2e_nostr_interop.rs index fce7877676..af7e64ae9c 100644 --- a/crates/buzz-test-client/tests/e2e_nostr_interop.rs +++ b/crates/buzz-test-client/tests/e2e_nostr_interop.rs @@ -732,6 +732,7 @@ async fn test_nip17_gift_wrap_recipient_receives() { RelayMessage::Event { subscription_id, event, + .. } => { assert_eq!( subscription_id, sid_b, diff --git a/crates/buzz-ws-client/Cargo.toml b/crates/buzz-ws-client/Cargo.toml index 5cec925677..2df21944b8 100644 --- a/crates/buzz-ws-client/Cargo.toml +++ b/crates/buzz-ws-client/Cargo.toml @@ -11,7 +11,18 @@ nostr = { workspace = true } tokio = { workspace = true } tokio-tungstenite = { workspace = true } futures-util = { workspace = true } -serde_json = { workspace = true } +serde = { workspace = true, features = ["derive"] } +serde_json = { workspace = true, features = ["raw_value"] } thiserror = { workspace = true } url = { workspace = true } tracing = { workspace = true } +rand = { workspace = true } +# Security boundary: Tungstenite 0.29 logs complete data frames at trace and +# relay-controlled Close reasons at debug. This graph-wide log-facade cap +# compiles those dependency sinks out of every buzz-ws-client consumer while +# leaving info/warn/error and `tracing` diagnostics available. +log = { workspace = true, features = ["max_level_info"] } + +[dev-dependencies] +tracing-subscriber = { workspace = true } +tracing-log = "0.2" diff --git a/crates/buzz-ws-client/src/connection.rs b/crates/buzz-ws-client/src/connection.rs index bec5b56bb4..1d2b66c54d 100644 --- a/crates/buzz-ws-client/src/connection.rs +++ b/crates/buzz-ws-client/src/connection.rs @@ -2,17 +2,263 @@ use std::collections::VecDeque; use std::time::Duration; use futures_util::{SinkExt, StreamExt}; -use nostr::{Event, Keys, Tag}; +use nostr::{Event, EventId, Keys, Tag}; +use serde::Deserialize; +use serde_json::value::RawValue; use serde_json::{json, Value}; +use tokio::io::{AsyncWrite, AsyncWriteExt}; use tokio::time::timeout; use tokio_tungstenite::{connect_async, tungstenite::Message, MaybeTlsStream, WebSocketStream}; use tracing::debug; use crate::error::WsClientError; -use crate::message::{build_auth_event, parse_relay_message, OkResponse, RelayMessage}; +use crate::message::{ + build_auth_event, parse_relay_message, validate_auth_challenge, OkResponse, RelayMessage, +}; type WsStream = WebSocketStream>; +// AUTH events can contain reusable NIP-OA authority. Tungstenite 0.29 logs +// complete message/frame payloads at `trace` and relay-controlled Close reasons +// at `debug`, with no payload-redaction feature. `buzz-ws-client` therefore +// activates log/max_level_info as a compile-time dependency boundary. Keep the +// outbound private event below Tungstenite as defense in depth, with an explicit +// cap that still exercises every RFC 6455 payload-length encoding. +const MAX_PRIVATE_AUTH_FRAME_BYTES: u64 = 1024 * 1024; + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum AuthTransportState { + Ready, + Writing, + Poisoned, +} + +struct AuthWriteGuard<'a> { + state: &'a mut AuthTransportState, + completed: bool, +} + +impl<'a> AuthWriteGuard<'a> { + fn begin(state: &'a mut AuthTransportState) -> Result { + if *state != AuthTransportState::Ready { + return Err(WsClientError::AuthTransportPoisoned); + } + *state = AuthTransportState::Writing; + Ok(Self { + state, + completed: false, + }) + } + + fn complete(mut self) { + *self.state = AuthTransportState::Ready; + self.completed = true; + } +} + +impl Drop for AuthWriteGuard<'_> { + fn drop(&mut self) { + if !self.completed { + *self.state = AuthTransportState::Poisoned; + } + } +} + +#[derive(Debug)] +enum AuthChallengeState { + Empty, + Pending(String), + InFlight, +} + +fn encode_masked_client_text_frame( + payload: &[u8], + mask: [u8; 4], +) -> Result, WsClientError> { + let size = u64::try_from(payload.len()).map_err(|_| WsClientError::AuthFrameTooLarge { + size: u64::MAX, + max: MAX_PRIVATE_AUTH_FRAME_BYTES, + })?; + if size > MAX_PRIVATE_AUTH_FRAME_BYTES { + return Err(WsClientError::AuthFrameTooLarge { + size, + max: MAX_PRIVATE_AUTH_FRAME_BYTES, + }); + } + + // FIN plus text opcode. AUTH is never fragmented by this boundary, and no + // WebSocket extensions are requested. + let mut frame = Vec::with_capacity(payload.len() + 14); + frame.push(0x81); + match size { + 0..=125 => frame.push(0x80 | size as u8), + 126..=65_535 => { + frame.push(0x80 | 126); + frame.extend_from_slice(&(size as u16).to_be_bytes()); + } + _ => { + frame.push(0x80 | 127); + frame.extend_from_slice(&size.to_be_bytes()); + } + } + frame.extend_from_slice(&mask); + frame.extend( + payload + .iter() + .enumerate() + .map(|(index, byte)| byte ^ mask[index % mask.len()]), + ); + Ok(frame) +} + +async fn write_private_text_frame( + stream: &mut S, + state: &mut AuthTransportState, + payload: &[u8], +) -> Result<(), WsClientError> +where + S: AsyncWrite + Unpin, +{ + // `rand::random` uses the thread-local CSPRNG; every client frame gets a + // fresh unpredictable masking key as required by RFC 6455 section 5.3. + let frame = encode_masked_client_text_frame(payload, rand::random())?; + // Cancellation or I/O failure after the first byte may leave an incomplete + // frame on the wire. The guard permanently poisons this connection unless + // both the write and flush complete, preventing a later Tungstenite write + // from being appended to a partial private frame. + let guard = AuthWriteGuard::begin(state)?; + stream.write_all(&frame).await.map_err(|error| { + WsClientError::WebSocket(tokio_tungstenite::tungstenite::Error::Io(error)) + })?; + stream.flush().await.map_err(|error| { + WsClientError::WebSocket(tokio_tungstenite::tungstenite::Error::Io(error)) + })?; + guard.complete(); + Ok(()) +} + +fn log_outbound(value: &Value, text: &str) { + if is_auth_message(value) { + debug!("→ relay: AUTH "); + } else { + debug!("→ relay: {text}"); + } +} + +fn is_auth_message(value: &Value) -> bool { + value + .as_array() + .and_then(|parts| parts.first()) + .and_then(Value::as_str) + == Some("AUTH") +} + +fn contains_private_marker(text: &str, marker: &str) -> bool { + if text.contains(marker) { + return true; + } + + // Canonicalize valid JSON once more so escape sequences cannot disguise a + // reflected signature from the raw-byte check. Malformed frames are still + // returned only through the static post-auth parse error below. + serde_json::from_str::(text) + .ok() + .map(|value| value.to_string()) + .is_some_and(|normalized| normalized.contains(marker)) +} + +#[cfg(test)] +fn contains_private_auth_marker(text: &str, markers: &[String]) -> bool { + markers + .iter() + .any(|marker| contains_private_marker(text, marker)) +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct StrictPrivateEventEnvelope { + id: String, + pubkey: String, + created_at: u64, + kind: u16, + tags: Vec>, + content: String, + sig: String, +} + +fn has_exact_event_envelope(text: &str, subscription_id: &str, raw_event_json: &str) -> bool { + let Ok(parts) = serde_json::from_str::>>(text) else { + return false; + }; + parts.len() == 3 + && parts + .first() + .and_then(|part| serde_json::from_str::(part.get()).ok()) + .as_deref() + == Some("EVENT") + && parts + .get(1) + .and_then(|part| serde_json::from_str::(part.get()).ok()) + .as_deref() + == Some(subscription_id) + && parts + .get(2) + .is_some_and(|part| part.get() == raw_event_json) +} + +fn strict_raw_event_matches_typed(raw_event_json: &str, event: &Event) -> bool { + let Ok(raw) = serde_json::from_str::(raw_event_json) else { + return false; + }; + let typed_tags: Vec> = event + .tags + .iter() + .map(|tag| tag.as_slice().to_vec()) + .collect(); + + raw.id == event.id.to_hex() + && raw.pubkey == event.pubkey.to_hex() + && raw.created_at == event.created_at.as_secs() + && raw.kind == event.kind.as_u16() + && raw.tags == typed_tags + && raw.content == event.content + && raw.sig == event.sig.to_string() +} + +fn is_expected_verified_authority_event( + parsed: &Result, + text: &str, + expected_event_id: Option<&EventId>, + authority_markers: &[String], +) -> bool { + let ( + Ok(RelayMessage::Event { + subscription_id, + event, + raw_event_json, + }), + Some(expected_event_id), + ) = (parsed, expected_event_id) + else { + return false; + }; + + let reflected_markers_stay_inside_event = authority_markers.iter().all(|marker| { + !contains_private_marker(text, marker) + || (contains_private_marker(raw_event_json, marker) + && !contains_private_marker(subscription_id, marker)) + }); + let exact_event_envelope = has_exact_event_envelope(text, subscription_id, raw_event_json); + let raw_matches_typed = strict_raw_event_matches_typed(raw_event_json, event); + + reflected_markers_stay_inside_event + && exact_event_envelope + && raw_matches_typed + && event.id == *expected_event_id + && event.verify_id() + && event.verify_signature() +} + /// Seconds to wait for the relay to send the NIP-42 AUTH challenge after connecting. pub const AUTH_CHALLENGE_TIMEOUT_SECS: u64 = 20; @@ -26,7 +272,10 @@ pub const PUBLISH_OK_TIMEOUT_SECS: u64 = 30; pub struct NostrWsConnection { ws: WsStream, buffer: VecDeque, - pending_challenge: Option, + auth_challenge: AuthChallengeState, + auth_transport: AuthTransportState, + private_auth_started: bool, + private_auth_markers: Vec, relay_url: String, } @@ -50,16 +299,30 @@ impl NostrWsConnection { .parse::() .map_err(|e| WsClientError::Url(e.to_string()))?; - let (ws, _response) = connect_async(parsed.as_str()) - .await - .map_err(WsClientError::WebSocket)?; + let (ws, _response) = + connect_async(parsed.as_str()) + .await + .map_err(|error| match error { + tokio_tungstenite::tungstenite::Error::Http(ref response) + if matches!(response.status().as_u16(), 401 | 403) => + { + WsClientError::AuthFailed(format!( + "relay rejected WebSocket upgrade with HTTP {}", + response.status() + )) + } + other => WsClientError::WebSocket(other), + })?; debug!("connected to relay at {url}"); Ok(Self { ws, buffer: VecDeque::new(), - pending_challenge: None, + auth_challenge: AuthChallengeState::Empty, + auth_transport: AuthTransportState::Ready, + private_auth_started: false, + private_auth_markers: Vec::new(), relay_url: url.to_string(), }) } @@ -72,6 +335,7 @@ impl NostrWsConnection { keys: &Keys, auth_tag: Option<&Tag>, ) -> Result<(), WsClientError> { + self.ensure_auth_transport_ready()?; let challenge = self .wait_for_auth_challenge(Duration::from_secs(AUTH_CHALLENGE_TIMEOUT_SECS)) .await?; @@ -79,13 +343,33 @@ impl NostrWsConnection { let auth_event = build_auth_event(&challenge, &self.relay_url, keys, auth_tag)?; let event_id = auth_event.id.to_hex(); - self.send_raw(&json!(["AUTH", auth_event])).await?; + // Retain only the two signatures that uniquely identify reflected + // private authority: the one-time AUTH event signature and, when + // present, the reusable NIP-OA authority signature. Checking these + // before returning any parsed message also prevents an exact reflected + // AUTH event from reaching a caller's output path. + self.private_auth_markers.clear(); + self.private_auth_markers.push(auth_event.sig.to_string()); + if let Some(authority_signature) = auth_tag + .and_then(|tag| tag.as_slice().last()) + .filter(|value| !value.is_empty()) + { + self.private_auth_markers + .push(authority_signature.to_string()); + } + + self.send_private_auth(&auth_event).await?; let ok = self - .wait_for_ok(&event_id, Duration::from_secs(AUTH_OK_TIMEOUT_SECS)) + .wait_for_auth_ok(&event_id, Duration::from_secs(AUTH_OK_TIMEOUT_SECS)) .await?; + self.auth_challenge = AuthChallengeState::Empty; if !ok.accepted { - return Err(WsClientError::AuthFailed(ok.message)); + // Relay-controlled text must not be allowed to echo the signed + // AUTH event or its reusable authority into caller logs/stderr. + return Err(WsClientError::AuthFailed( + "relay rejected NIP-42 authentication".into(), + )); } debug!("NIP-42 authentication successful"); @@ -100,32 +384,195 @@ impl NostrWsConnection { .await } + /// Reports whether this connection has begun a private AUTH write. + /// + /// This becomes `true` before awaiting the write and does not depend on a + /// relay `OK`, so callers can select non-reflective error reporting even if + /// authentication fails or is cancelled after transmission starts. + pub fn private_auth_started(&self) -> bool { + self.private_auth_started + } + /// Receives the next relay message, waiting up to `timeout_dur`. pub async fn next_event( &mut self, timeout_dur: Duration, ) -> Result { + self.ensure_auth_transport_ready()?; if let Some(msg) = self.buffer.pop_front() { + // Messages enter the buffer only through `buffer_message`, which + // records an AUTH challenge before storing it. return Ok(msg); } - self.recv_one(timeout_dur).await + self.recv_one(timeout_dur, None).await + } + + /// Receives the next relay message for an exact-event lookup. + /// + /// A stored event signed by a delegated agent legitimately carries the + /// same reusable NIP-OA authority tag used during connection AUTH. This + /// method permits that authority signature only inside an `EVENT` whose + /// raw object is semantically identical to its locally ID-and-signature- + /// verified typed form and matches `expected_event_id`. The one-time AUTH + /// event signature remains forbidden in every frame, and all other + /// authority reflections fail closed. + pub async fn next_event_for_exact_id( + &mut self, + timeout_dur: Duration, + expected_event_id: &EventId, + ) -> Result { + self.ensure_auth_transport_ready()?; + if let Some(msg) = self.buffer.pop_front() { + return Ok(msg); + } + self.recv_one(timeout_dur, Some(expected_event_id)).await } /// Closes the WebSocket connection gracefully. pub async fn disconnect(mut self) -> Result<(), WsClientError> { + self.ensure_auth_transport_ready()?; self.ws.close(None).await?; Ok(()) } /// Sends a raw JSON value as a WebSocket text frame. + /// + /// Raw `AUTH` envelopes are rejected. Use [`Self::authenticate`] so the + /// signed event is registered as private state and written below + /// Tungstenite's payload-bearing trace boundary. pub async fn send_raw(&mut self, value: &Value) -> Result<(), WsClientError> { + self.ensure_auth_transport_ready()?; + if is_auth_message(value) { + return Err(WsClientError::AuthFailed( + "raw AUTH messages are not accepted; use authenticate".into(), + )); + } let text = serde_json::to_string(value)?; - debug!("→ relay: {text}"); + log_outbound(value, &text); self.ws.send(Message::Text(text.into())).await?; Ok(()) } - async fn recv_one(&mut self, timeout_dur: Duration) -> Result { + async fn send_private_auth(&mut self, event: &Event) -> Result<(), WsClientError> { + let value = json!(["AUTH", event]); + let text = serde_json::to_string(&value)?; + log_outbound(&value, &text); + // Tungstenite 0.29 traces plaintext frame payloads. Drain all prior + // writes, then keep the registered private AUTH event below that + // logger. This path is intentionally unavailable through `send_raw`. + self.ws.flush().await?; + self.private_auth_started = true; + write_private_text_frame(self.ws.get_mut(), &mut self.auth_transport, text.as_bytes()).await + } + + fn ensure_auth_transport_ready(&self) -> Result<(), WsClientError> { + if self.auth_transport == AuthTransportState::Ready { + Ok(()) + } else { + Err(WsClientError::AuthTransportPoisoned) + } + } + + fn observe_auth_challenge(&mut self, challenge: &str) -> Result<(), WsClientError> { + validate_auth_challenge(challenge)?; + if !matches!(self.auth_challenge, AuthChallengeState::Empty) { + return Err(WsClientError::AmbiguousAuthChallenge); + } + self.auth_challenge = AuthChallengeState::Pending(challenge.to_string()); + Ok(()) + } + + fn parse_inbound_message( + &self, + text: &str, + authenticating: bool, + expected_private_event_id: Option<&EventId>, + ) -> Result { + let auth_event_signature_reflected = self + .private_auth_markers + .first() + .is_some_and(|marker| contains_private_marker(text, marker)); + let authority_signature_reflected = self + .private_auth_markers + .iter() + .skip(1) + .any(|marker| contains_private_marker(text, marker)); + let parsed = parse_relay_message(text); + + if auth_event_signature_reflected { + return Err(WsClientError::ReflectedAuthMaterial); + } + if authority_signature_reflected + && !is_expected_verified_authority_event( + &parsed, + text, + expected_private_event_id, + &self.private_auth_markers[1..], + ) + { + return Err(WsClientError::ReflectedAuthMaterial); + } + + match parsed { + Ok(message) => Ok(message), + Err(WsClientError::AuthChallengeTooLarge { .. }) if authenticating => { + Err(WsClientError::AmbiguousAuthChallenge) + } + Err(_) if !self.private_auth_markers.is_empty() => { + let phase = if authenticating { + "during authentication" + } else { + "after authentication" + }; + Err(WsClientError::UnexpectedMessage(format!( + "malformed relay response {phase}" + ))) + } + Err(error) => Err(error), + } + } + + fn begin_auth_transaction(&mut self) -> Result, WsClientError> { + match std::mem::replace(&mut self.auth_challenge, AuthChallengeState::Empty) { + AuthChallengeState::Empty => Ok(None), + AuthChallengeState::Pending(challenge) => { + self.auth_challenge = AuthChallengeState::InFlight; + if let Some(index) = self + .buffer + .iter() + .position(|message| matches!(message, RelayMessage::Auth { .. })) + { + let _ = self.buffer.remove(index); + } + Ok(Some(challenge)) + } + AuthChallengeState::InFlight => { + self.auth_challenge = AuthChallengeState::InFlight; + Err(WsClientError::AmbiguousAuthChallenge) + } + } + } + + fn buffer_message(&mut self, message: RelayMessage) -> Result<(), WsClientError> { + if let RelayMessage::Auth { ref challenge } = message { + self.observe_auth_challenge(challenge)?; + } + self.buffer.push_back(message); + Ok(()) + } + + fn deliver_message(&mut self, message: RelayMessage) -> Result { + if let RelayMessage::Auth { ref challenge } = message { + self.observe_auth_challenge(challenge)?; + } + Ok(message) + } + + async fn recv_one( + &mut self, + timeout_dur: Duration, + expected_private_event_id: Option<&EventId>, + ) -> Result { if let Some(msg) = self.buffer.pop_front() { return Ok(msg); } @@ -139,11 +586,9 @@ impl NostrWsConnection { match raw { Message::Text(text) => { - let msg = parse_relay_message(&text)?; - if let RelayMessage::Auth { ref challenge } = msg { - self.pending_challenge = Some(challenge.clone()); - } - return Ok(msg); + let msg = + self.parse_inbound_message(&text, false, expected_private_event_id)?; + return self.deliver_message(msg); } Message::Ping(data) => { self.ws.send(Message::Pong(data)).await?; @@ -158,21 +603,10 @@ impl NostrWsConnection { &mut self, timeout_dur: Duration, ) -> Result { - if let Some(challenge) = self.pending_challenge.take() { + if let Some(challenge) = self.begin_auth_transaction()? { return Ok(challenge); } - if let Some(idx) = self - .buffer - .iter() - .position(|m| matches!(m, RelayMessage::Auth { .. })) - { - match self.buffer.remove(idx).unwrap() { - RelayMessage::Auth { challenge } => return Ok(challenge), - _ => unreachable!(), - } - } - let deadline = tokio::time::Instant::now() + timeout_dur; loop { @@ -192,17 +626,15 @@ impl NostrWsConnection { match raw { Message::Text(text) => { - let msg = parse_relay_message(&text)?; + let msg = self.parse_inbound_message(&text, false, None)?; match msg { RelayMessage::Auth { challenge } => { - if challenge.len() > 1024 { - return Err(WsClientError::AuthFailed( - "challenge exceeds 1024 bytes".into(), - )); + self.observe_auth_challenge(&challenge)?; + if let Some(challenge) = self.begin_auth_transaction()? { + return Ok(challenge); } - return Ok(challenge); } - other => self.buffer.push_back(other), + other => self.buffer_message(other)?, } } Message::Ping(data) => { @@ -214,10 +646,27 @@ impl NostrWsConnection { } } + async fn wait_for_auth_ok( + &mut self, + event_id: &str, + timeout_dur: Duration, + ) -> Result { + self.wait_for_ok_inner(event_id, timeout_dur, true).await + } + async fn wait_for_ok( &mut self, event_id: &str, timeout_dur: Duration, + ) -> Result { + self.wait_for_ok_inner(event_id, timeout_dur, false).await + } + + async fn wait_for_ok_inner( + &mut self, + event_id: &str, + timeout_dur: Duration, + authenticating: bool, ) -> Result { let deadline = tokio::time::Instant::now() + timeout_dur; @@ -249,14 +698,10 @@ impl NostrWsConnection { match raw { Message::Text(text) => { - let msg = parse_relay_message(&text)?; + let msg = self.parse_inbound_message(&text, authenticating, None)?; match msg { RelayMessage::Ok(ok) if ok.event_id == event_id => return Ok(ok), - RelayMessage::Auth { ref challenge } => { - self.pending_challenge = Some(challenge.clone()); - self.buffer.push_back(msg); - } - other => self.buffer.push_back(other), + other => self.buffer_message(other)?, } } Message::Ping(data) => { @@ -295,6 +740,13 @@ pub async fn publish_event( #[cfg(test)] mod tests { + use std::io::Write; + use std::pin::Pin; + use std::sync::{Arc, Mutex}; + use std::task::{Context, Poll}; + + use nostr::{EventBuilder, Kind}; + use super::*; #[test] @@ -311,4 +763,254 @@ mod tests { fn publish_ok_timeout_meets_floor() { const { assert!(PUBLISH_OK_TIMEOUT_SECS >= 30) }; } + + #[test] + fn dependency_payload_log_levels_are_compiled_out() { + assert_eq!(log::STATIC_MAX_LEVEL, log::LevelFilter::Info); + } + + #[test] + fn json_escapes_cannot_disguise_private_auth_markers() { + let markers = vec!["signature-marker".to_string()]; + let escaped = r#"["CLOSED","subscription","signature\u002dmarker"]"#; + + assert!(!escaped.contains(&markers[0])); + assert!(contains_private_auth_marker(escaped, &markers)); + } + + #[test] + fn delegated_authority_is_allowed_only_inside_the_exact_signed_event() { + let authority_signature = "a".repeat(128); + let authority_markers = vec![authority_signature.clone()]; + let auth_tag = Tag::parse([ + "auth", + &"b".repeat(64), + "kind=1", + authority_signature.as_str(), + ]) + .unwrap(); + let event = EventBuilder::new(Kind::TextNote, "delegated") + .tags([auth_tag]) + .sign_with_keys(&Keys::generate()) + .unwrap(); + let raw_event_json = serde_json::to_string(&event).unwrap(); + let frame = serde_json::json!(["EVENT", "exact", event]).to_string(); + let parsed = parse_relay_message(&frame); + + assert!(is_expected_verified_authority_event( + &parsed, + &frame, + Some(&event.id), + &authority_markers, + )); + + let extra_envelope_value = format!( + r#"["EVENT","exact",{},{}]"#, + raw_event_json, + serde_json::to_string(&authority_markers[0]).unwrap() + ); + let parsed = parse_relay_message(&extra_envelope_value); + assert!(matches!(parsed, Ok(RelayMessage::Event { .. }))); + assert!(!is_expected_verified_authority_event( + &parsed, + &extra_envelope_value, + Some(&event.id), + &authority_markers, + )); + + let duplicate_content_raw = format!( + r#"{},"content":{}}}"#, + raw_event_json.strip_suffix('}').unwrap(), + serde_json::to_string(&event.content).unwrap() + ); + let duplicate_content_frame = format!(r#"["EVENT","exact",{duplicate_content_raw}]"#); + let parsed = parse_relay_message(&duplicate_content_frame); + assert!(!strict_raw_event_matches_typed( + &duplicate_content_raw, + &event, + )); + assert!(!is_expected_verified_authority_event( + &parsed, + &duplicate_content_frame, + Some(&event.id), + &authority_markers, + )); + + let reflected_subscription = serde_json::json!([ + "EVENT", + authority_signature, + serde_json::from_str::(&raw_event_json).unwrap() + ]) + .to_string(); + let parsed = parse_relay_message(&reflected_subscription); + assert!(!is_expected_verified_authority_event( + &parsed, + &reflected_subscription, + Some(&event.id), + &authority_markers, + )); + + let mut raw_with_unknown_field = serde_json::from_str::(&raw_event_json).unwrap(); + raw_with_unknown_field["relay-controlled"] = Value::String(authority_markers[0].clone()); + let reflected_unknown_field = + serde_json::json!(["EVENT", "exact", raw_with_unknown_field]).to_string(); + let parsed = parse_relay_message(&reflected_unknown_field); + assert!(!is_expected_verified_authority_event( + &parsed, + &reflected_unknown_field, + Some(&event.id), + &authority_markers, + )); + } + + fn assert_masked_text_frame_at_boundary(size: usize, length_marker: u8, header_len: usize) { + let payload: Vec = (0..size).map(|index| index as u8).collect(); + let mask = [0x12, 0x34, 0x56, 0x78]; + let frame = encode_masked_client_text_frame(&payload, mask).unwrap(); + + assert_eq!(frame[0], 0x81, "AUTH frame must be FIN + text"); + assert_eq!(frame[1] & 0x80, 0x80, "client MASK bit must be set"); + assert_eq!(frame[1] & 0x7f, length_marker); + let encoded_size = match length_marker { + 0..=125 => u64::from(length_marker), + 126 => u64::from(u16::from_be_bytes([frame[2], frame[3]])), + 127 => u64::from_be_bytes(frame[2..10].try_into().unwrap()), + _ => unreachable!(), + }; + assert_eq!(encoded_size, size as u64); + assert_eq!(&frame[header_len - 4..header_len], &mask); + + let unmasked: Vec = frame[header_len..] + .iter() + .enumerate() + .map(|(index, byte)| byte ^ mask[index % mask.len()]) + .collect(); + assert_eq!(unmasked, payload); + } + + #[test] + fn private_auth_frame_has_canonical_lengths_across_rfc6455_boundaries() { + assert_masked_text_frame_at_boundary(125, 125, 6); + assert_masked_text_frame_at_boundary(126, 126, 8); + assert_masked_text_frame_at_boundary(65_535, 126, 8); + assert_masked_text_frame_at_boundary(65_536, 127, 14); + } + + #[test] + fn private_auth_frame_has_an_explicit_allocation_and_length_cap() { + let oversized = vec![0_u8; MAX_PRIVATE_AUTH_FRAME_BYTES as usize + 1]; + assert!(matches!( + encode_masked_client_text_frame(&oversized, [1, 2, 3, 4]), + Err(WsClientError::AuthFrameTooLarge { + size: 1_048_577, + max: 1_048_576 + }) + )); + } + + struct PartialThenPending { + bytes_written: usize, + } + + impl AsyncWrite for PartialThenPending { + fn poll_write( + mut self: Pin<&mut Self>, + _context: &mut Context<'_>, + bytes: &[u8], + ) -> Poll> { + if self.bytes_written == 0 && !bytes.is_empty() { + self.bytes_written = 1; + Poll::Ready(Ok(1)) + } else { + Poll::Pending + } + } + + fn poll_flush( + self: Pin<&mut Self>, + _context: &mut Context<'_>, + ) -> Poll> { + Poll::Ready(Ok(())) + } + + fn poll_shutdown( + self: Pin<&mut Self>, + _context: &mut Context<'_>, + ) -> Poll> { + Poll::Ready(Ok(())) + } + } + + #[tokio::test] + async fn cancelled_partial_auth_write_permanently_poisons_transport() { + let mut state = AuthTransportState::Ready; + let mut writer = PartialThenPending { bytes_written: 0 }; + + let cancelled = timeout( + Duration::from_millis(10), + write_private_text_frame(&mut writer, &mut state, b"private-auth"), + ) + .await; + assert!( + cancelled.is_err(), + "partial AUTH write unexpectedly completed" + ); + assert_eq!(writer.bytes_written, 1); + assert_eq!(state, AuthTransportState::Poisoned); + + assert!(matches!( + write_private_text_frame(&mut writer, &mut state, b"later-frame").await, + Err(WsClientError::AuthTransportPoisoned) + )); + assert_eq!(writer.bytes_written, 1, "poisoned transport was reused"); + } + + #[test] + fn auth_event_payload_is_redacted_from_debug_logs() { + #[derive(Clone)] + struct CaptureWriter(Arc>>); + + impl Write for CaptureWriter { + fn write(&mut self, bytes: &[u8]) -> std::io::Result { + self.0.lock().unwrap().extend_from_slice(bytes); + Ok(bytes.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + + let captured = Arc::new(Mutex::new(Vec::new())); + let writer = Arc::clone(&captured); + let subscriber = tracing_subscriber::fmt() + .with_ansi(false) + .without_time() + .with_target(false) + .with_max_level(tracing::Level::DEBUG) + .with_writer(move || CaptureWriter(Arc::clone(&writer))) + .finish(); + let auth = json!([ + "AUTH", + { + "id": "private-auth-event-bytes", + "tags": [["auth", "reusable-auth-tag-secret"]], + "sig": "private-auth-signature" + } + ]); + let serialized = serde_json::to_string(&auth).unwrap(); + + tracing::subscriber::with_default(subscriber, || log_outbound(&auth, &serialized)); + + let logs = String::from_utf8(captured.lock().unwrap().clone()).unwrap(); + assert!(logs.contains("AUTH ")); + for secret in [ + "private-auth-event-bytes", + "reusable-auth-tag-secret", + "private-auth-signature", + ] { + assert!(!logs.contains(secret), "AUTH log leaked {secret}: {logs}"); + } + assert!(!logs.contains(&serialized), "AUTH log leaked its raw frame"); + } } diff --git a/crates/buzz-ws-client/src/error.rs b/crates/buzz-ws-client/src/error.rs index c8781bb591..29e1f9fe33 100644 --- a/crates/buzz-ws-client/src/error.rs +++ b/crates/buzz-ws-client/src/error.rs @@ -11,6 +11,17 @@ pub enum WsClientError { #[error("JSON error: {0}")] Json(#[from] serde_json::Error), + /// An EVENT frame had a syntactically valid raw payload that did not decode + /// as a Nostr event. The raw object is retained so callers can classify the + /// malformed signed fields without losing their lexical representation. + #[error("invalid EVENT payload: {message}")] + InvalidEvent { + /// Exact JSON bytes of the event object from the relay frame. + raw_event_json: Box, + /// Deserialization failure without the raw payload. + message: String, + }, + /// Failed to build a Nostr event. #[error("Nostr event builder error: {0}")] EventBuilder(String), @@ -42,6 +53,37 @@ pub enum WsClientError { /// No NIP-42 AUTH challenge was received from the relay. #[error("No AUTH challenge received from relay")] NoAuthChallenge, + + /// A relay supplied an AUTH challenge larger than the client will sign. + #[error("AUTH challenge is {size} bytes; maximum is {max} bytes")] + AuthChallengeTooLarge { + /// Challenge length in UTF-8 bytes. + size: usize, + /// Maximum accepted challenge length in UTF-8 bytes. + max: usize, + }, + + /// More than one AUTH challenge was received before authentication. + #[error("relay sent multiple AUTH challenges before authentication")] + AmbiguousAuthChallenge, + + /// The serialized AUTH frame exceeded the private transport boundary's cap. + #[error("serialized AUTH frame is {size} bytes; maximum is {max} bytes")] + AuthFrameTooLarge { + /// Serialized JSON payload length in bytes. + size: u64, + /// Maximum accepted private frame payload length in bytes. + max: u64, + }, + + /// A private AUTH frame write did not complete, so the WebSocket cannot be + /// safely reused without appending bytes to a potentially partial frame. + #[error("WebSocket connection is unusable after an incomplete AUTH frame write")] + AuthTransportPoisoned, + + /// The relay reflected bytes from the signed private authentication event. + #[error("relay reflected private authentication material")] + ReflectedAuthMaterial, } impl From for WsClientError { diff --git a/crates/buzz-ws-client/src/message.rs b/crates/buzz-ws-client/src/message.rs index 3c646bc8ab..9da537ea46 100644 --- a/crates/buzz-ws-client/src/message.rs +++ b/crates/buzz-ws-client/src/message.rs @@ -1,8 +1,22 @@ use nostr::{Event, EventBuilder, Keys, RelayUrl, Tag}; -use serde_json::Value; +use serde_json::value::RawValue; use crate::error::WsClientError; +/// Maximum UTF-8 byte length accepted for a relay AUTH challenge. +pub const MAX_AUTH_CHALLENGE_BYTES: usize = 1024; + +pub(crate) fn validate_auth_challenge(challenge: &str) -> Result<(), WsClientError> { + let size = challenge.len(); + if size > MAX_AUTH_CHALLENGE_BYTES { + return Err(WsClientError::AuthChallengeTooLarge { + size, + max: MAX_AUTH_CHALLENGE_BYTES, + }); + } + Ok(()) +} + /// A message received from a Nostr relay. #[derive(Debug, Clone)] pub enum RelayMessage { @@ -12,6 +26,8 @@ pub enum RelayMessage { subscription_id: String, /// The Nostr event payload. event: Box, + /// The event object's exact JSON bytes from the relay text frame. + raw_event_json: Box, }, /// Acknowledgement of a published event. Ok(OkResponse), @@ -60,42 +76,49 @@ pub struct OkResponse { /// Parse a raw relay text frame into a typed [`RelayMessage`]. #[allow(clippy::result_large_err)] pub fn parse_relay_message(text: &str) -> Result { - let arr: Vec = serde_json::from_str(text)?; + let arr: Vec> = serde_json::from_str(text)?; let msg_type = arr .first() - .and_then(|v| v.as_str()) + .and_then(|value| serde_json::from_str::(value.get()).ok()) .ok_or_else(|| WsClientError::UnexpectedMessage(text.to_string()))?; - match msg_type { + match msg_type.as_str() { "EVENT" => { let sub_id = arr .get(1) - .and_then(|v| v.as_str()) + .and_then(|value| serde_json::from_str::(value.get()).ok()) .ok_or_else(|| WsClientError::UnexpectedMessage(text.to_string()))? .to_string(); - let event: Event = serde_json::from_value( - arr.get(2) - .cloned() - .ok_or_else(|| WsClientError::UnexpectedMessage(text.to_string()))?, - )?; + let raw_event = arr + .get(2) + .ok_or_else(|| WsClientError::UnexpectedMessage(text.to_string()))?; + let event: Event = serde_json::from_str(raw_event.get()).map_err(|error| { + WsClientError::InvalidEvent { + raw_event_json: raw_event.get().to_string().into_boxed_str(), + message: error.to_string(), + } + })?; Ok(RelayMessage::Event { subscription_id: sub_id, event: Box::new(event), + raw_event_json: raw_event.get().to_string().into_boxed_str(), }) } "OK" => { let event_id = arr .get(1) - .and_then(|v| v.as_str()) + .and_then(|value| serde_json::from_str::(value.get()).ok()) .ok_or_else(|| WsClientError::UnexpectedMessage(text.to_string()))? .to_string(); - let accepted = arr.get(2).and_then(|v| v.as_bool()).unwrap_or(false); + let accepted = arr + .get(2) + .and_then(|value| serde_json::from_str::(value.get()).ok()) + .unwrap_or(false); let message = arr .get(3) - .and_then(|v| v.as_str()) - .unwrap_or("") - .to_string(); + .and_then(|value| serde_json::from_str::(value.get()).ok()) + .unwrap_or_default(); Ok(RelayMessage::Ok(OkResponse { event_id, accepted, @@ -105,7 +128,7 @@ pub fn parse_relay_message(text: &str) -> Result { "EOSE" => { let sub_id = arr .get(1) - .and_then(|v| v.as_str()) + .and_then(|value| serde_json::from_str::(value.get()).ok()) .ok_or_else(|| WsClientError::UnexpectedMessage(text.to_string()))? .to_string(); Ok(RelayMessage::Eose { @@ -115,14 +138,13 @@ pub fn parse_relay_message(text: &str) -> Result { "CLOSED" => { let sub_id = arr .get(1) - .and_then(|v| v.as_str()) + .and_then(|value| serde_json::from_str::(value.get()).ok()) .ok_or_else(|| WsClientError::UnexpectedMessage(text.to_string()))? .to_string(); let message = arr .get(2) - .and_then(|v| v.as_str()) - .unwrap_or("") - .to_string(); + .and_then(|value| serde_json::from_str::(value.get()).ok()) + .unwrap_or_default(); Ok(RelayMessage::Closed { subscription_id: sub_id, message, @@ -131,29 +153,33 @@ pub fn parse_relay_message(text: &str) -> Result { "NOTICE" => { let message = arr .get(1) - .and_then(|v| v.as_str()) - .unwrap_or("") - .to_string(); + .and_then(|value| serde_json::from_str::(value.get()).ok()) + .unwrap_or_default(); Ok(RelayMessage::Notice { message }) } "AUTH" => { let challenge = arr .get(1) - .and_then(|v| v.as_str()) + .and_then(|value| serde_json::from_str::(value.get()).ok()) .ok_or_else(|| WsClientError::UnexpectedMessage(text.to_string()))? .to_string(); + validate_auth_challenge(&challenge)?; Ok(RelayMessage::Auth { challenge }) } "COUNT" => { let sub_id = arr .get(1) - .and_then(|v| v.as_str()) + .and_then(|value| serde_json::from_str::(value.get()).ok()) .ok_or_else(|| WsClientError::UnexpectedMessage(text.to_string()))? .to_string(); - let count = arr - .get(2) - .and_then(|o| o.get("count")) - .and_then(|c| c.as_u64()) + let count_object: serde_json::Value = serde_json::from_str( + arr.get(2) + .ok_or_else(|| WsClientError::UnexpectedMessage(text.to_string()))? + .get(), + )?; + let count = count_object + .get("count") + .and_then(|count| count.as_u64()) .ok_or_else(|| WsClientError::UnexpectedMessage(text.to_string()))?; Ok(RelayMessage::Count { subscription_id: sub_id, @@ -177,6 +203,9 @@ pub fn build_auth_event( keys: &Keys, auth_tag: Option<&Tag>, ) -> Result { + // Keep this check at the signing boundary as well as relay-message ingress: + // callers may invoke this public helper without parsing an AUTH frame first. + validate_auth_challenge(challenge)?; let url = RelayUrl::parse(relay_url).map_err(|e| WsClientError::Url(e.to_string()))?; let builder = EventBuilder::auth(challenge, url); let builder = if let Some(tag) = auth_tag { @@ -188,3 +217,65 @@ pub fn build_auth_event( .sign_with_keys(keys) .map_err(|e| WsClientError::EventBuilder(e.to_string())) } + +#[cfg(test)] +mod tests { + use nostr::{EventBuilder, JsonUtil, Kind}; + + use super::*; + + #[test] + fn event_message_retains_exact_raw_object_bytes() { + let event = EventBuilder::new(Kind::TextNote, "raw relay bytes") + .sign_with_keys(&Keys::generate()) + .unwrap(); + let raw_event = event.as_json().replacen(',', ",\n ", 1); + let frame = format!(r#"["EVENT","subscription",{raw_event}]"#); + + match parse_relay_message(&frame).unwrap() { + RelayMessage::Event { + event: parsed, + raw_event_json, + .. + } => { + assert_eq!(*parsed, event); + assert_eq!(raw_event_json.as_ref(), raw_event); + } + other => panic!("expected EVENT, got {other:?}"), + } + } + + #[test] + fn auth_challenge_limit_is_measured_in_utf8_bytes_at_parse_ingress() { + let exact = "x".repeat(MAX_AUTH_CHALLENGE_BYTES); + let oversized = "x".repeat(MAX_AUTH_CHALLENGE_BYTES + 1); + + assert!(matches!( + parse_relay_message(&serde_json::json!(["AUTH", exact]).to_string()), + Ok(RelayMessage::Auth { .. }) + )); + assert!(matches!( + parse_relay_message(&serde_json::json!(["AUTH", oversized]).to_string()), + Err(WsClientError::AuthChallengeTooLarge { + size: 1025, + max: 1024 + }) + )); + } + + #[test] + fn auth_challenge_limit_is_rechecked_before_signing() { + let keys = Keys::generate(); + let exact = "x".repeat(MAX_AUTH_CHALLENGE_BYTES); + let oversized = "x".repeat(MAX_AUTH_CHALLENGE_BYTES + 1); + + assert!(build_auth_event(&exact, "wss://relay.example", &keys, None).is_ok()); + assert!(matches!( + build_auth_event(&oversized, "wss://relay.example", &keys, None), + Err(WsClientError::AuthChallengeTooLarge { + size: 1025, + max: 1024 + }) + )); + } +} diff --git a/crates/buzz-ws-client/tests/auth_trace_redaction.rs b/crates/buzz-ws-client/tests/auth_trace_redaction.rs new file mode 100644 index 0000000000..c1cb87c5bb --- /dev/null +++ b/crates/buzz-ws-client/tests/auth_trace_redaction.rs @@ -0,0 +1,366 @@ +use std::process::Stdio; +use std::time::Duration; + +use futures_util::{SinkExt, StreamExt}; +use nostr::{Event, EventBuilder, JsonUtil, Keys, Kind, RelayUrl, Tag}; +use serde_json::{json, Value}; +use tokio::net::TcpListener; +use tokio::process::Command; +use tokio::time::timeout; +use tokio_tungstenite::tungstenite::Message; + +use buzz_ws_client::{NostrWsConnection, RelayMessage, WsClientError}; + +const RELAY_ENV: &str = "BUZZ_WS_AUTH_TRACE_TEST_RELAY"; +const PRIVATE_KEY_ENV: &str = "BUZZ_WS_AUTH_TRACE_TEST_PRIVATE_KEY"; +const AUTH_TAG_ENV: &str = "BUZZ_WS_AUTH_TRACE_TEST_TAG"; + +// Runs only in the explicitly spawned child process. Keeping the trace bridge +// out of the parent means the local test relay's inbound Tungstenite trace is +// not confused with the client-side boundary under test. +#[tokio::test] +#[ignore = "subprocess helper for auth_trace_hides_private_event_from_tungstenite"] +async fn auth_trace_client_helper() { + assert_eq!(log::STATIC_MAX_LEVEL, log::LevelFilter::Info); + let relay = std::env::var(RELAY_ENV).unwrap(); + let keys = Keys::parse(&std::env::var(PRIVATE_KEY_ENV).unwrap()).unwrap(); + let tag_parts: Vec = + serde_json::from_str(&std::env::var(AUTH_TAG_ENV).unwrap()).unwrap(); + let auth_tag = Tag::parse(tag_parts).unwrap(); + + tracing_log::LogTracer::init().unwrap(); + let subscriber = tracing_subscriber::fmt() + .with_ansi(false) + .without_time() + .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) + .finish(); + tracing::subscriber::set_global_default(subscriber).unwrap(); + + let mut connection = NostrWsConnection::connect(&relay).await.unwrap(); + connection + .authenticate(&keys, Some(&auth_tag)) + .await + .unwrap(); + connection + .send_raw(&json!([ + "REQ", + "post-auth-diagnostic", + {"kinds": [1], "limit": 1} + ])) + .await + .unwrap(); + assert!(matches!( + connection.next_event(Duration::from_secs(2)).await.unwrap(), + RelayMessage::Event { .. } + )); + assert!(matches!( + connection.next_event(Duration::from_secs(2)).await.unwrap(), + RelayMessage::Eose { .. } + )); + connection.disconnect().await.unwrap(); +} + +#[tokio::test] +#[ignore = "subprocess helper for malformed_auth_echo_is_sanitized_in_error_logs"] +async fn auth_error_log_client_helper() { + assert_eq!(log::STATIC_MAX_LEVEL, log::LevelFilter::Info); + let relay = std::env::var(RELAY_ENV).unwrap(); + let keys = Keys::parse(&std::env::var(PRIVATE_KEY_ENV).unwrap()).unwrap(); + let tag_parts: Vec = + serde_json::from_str(&std::env::var(AUTH_TAG_ENV).unwrap()).unwrap(); + let auth_tag = Tag::parse(tag_parts).unwrap(); + + tracing_log::LogTracer::init().unwrap(); + let subscriber = tracing_subscriber::fmt() + .with_ansi(false) + .without_time() + .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) + .finish(); + tracing::subscriber::set_global_default(subscriber).unwrap(); + + let mut connection = NostrWsConnection::connect(&relay).await.unwrap(); + let error = connection + .authenticate(&keys, Some(&auth_tag)) + .await + .unwrap_err(); + tracing::error!(error = %error, "authentication failed"); + assert!(matches!(error, WsClientError::ReflectedAuthMaterial)); +} + +#[tokio::test] +async fn auth_trace_hides_private_event_from_tungstenite() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let relay = format!("ws://{}", listener.local_addr().unwrap()); + let keys = Keys::generate(); + let auth_tag_parts = vec![ + "auth".to_string(), + "private-owner-authority-9f8e7d6c".to_string(), + "private-condition-kind-verified".to_string(), + "private-reusable-signature-1a2b3c4d".to_string(), + ]; + let auth_tag_json = serde_json::to_string(&auth_tag_parts).unwrap(); + + let child = Command::new(std::env::current_exe().unwrap()) + .args([ + "--ignored", + "--exact", + "auth_trace_client_helper", + "--nocapture", + ]) + .env(RELAY_ENV, &relay) + .env(PRIVATE_KEY_ENV, keys.secret_key().to_secret_hex()) + .env(AUTH_TAG_ENV, &auth_tag_json) + .env( + "RUST_LOG", + "tungstenite=trace,tokio_tungstenite=trace,buzz_ws_client=trace", + ) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .kill_on_drop(true) + .spawn() + .unwrap(); + + let (stream, _) = timeout(Duration::from_secs(2), listener.accept()) + .await + .unwrap() + .unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + websocket + .send(Message::Text( + json!(["AUTH", "trace-capture-challenge"]) + .to_string() + .into(), + )) + .await + .unwrap(); + + // A real Tungstenite server proves the hand-built client frame is masked, + // valid UTF-8 JSON, and exactly the expected signed AUTH envelope. + let auth_message = timeout(Duration::from_secs(2), websocket.next()) + .await + .unwrap() + .unwrap() + .unwrap(); + let auth_text = auth_message.to_text().unwrap().to_string(); + let auth_json: Value = serde_json::from_str(&auth_text).unwrap(); + assert_eq!(auth_json[0], "AUTH"); + let auth_event: Event = serde_json::from_value(auth_json[1].clone()).unwrap(); + auth_event.verify().unwrap(); + assert!(auth_event + .tags + .iter() + .any(|tag| tag.as_slice() == auth_tag_parts.as_slice())); + + websocket + .send(Message::Text( + json!(["OK", auth_event.id.to_hex(), true, ""]) + .to_string() + .into(), + )) + .await + .unwrap(); + + // Tungstenite remains usable after the raw AUTH write: the authenticated + // REQ and its EVENT/EOSE response complete on the original stream. + let request = timeout(Duration::from_secs(2), websocket.next()) + .await + .unwrap() + .unwrap() + .unwrap(); + let request: Value = serde_json::from_str(request.to_text().unwrap()).unwrap(); + assert_eq!(request[0], "REQ"); + assert_eq!(request[1], "post-auth-diagnostic"); + let event = EventBuilder::new(Kind::TextNote, "post-auth response") + .sign_with_keys(&Keys::generate()) + .unwrap(); + websocket + .send(Message::Text( + json!(["EVENT", "post-auth-diagnostic", event]) + .to_string() + .into(), + )) + .await + .unwrap(); + websocket + .send(Message::Text( + json!(["EOSE", "post-auth-diagnostic"]).to_string().into(), + )) + .await + .unwrap(); + + let output = timeout(Duration::from_secs(5), child.wait_with_output()) + .await + .unwrap() + .unwrap(); + let stdout = String::from_utf8_lossy(&output.stdout); + let stderr = String::from_utf8_lossy(&output.stderr); + let logs = format!("{stdout}{stderr}"); + assert!(output.status.success(), "trace helper failed:\n{logs}"); + assert!( + logs.contains("connected to relay") && logs.contains("post-auth-diagnostic"), + "nonsecret tracing controls were vacuous:\n{logs}" + ); + assert!(logs.contains("AUTH ")); + + let auth_event_json = auth_event.as_json(); + for private_bytes in auth_tag_parts.iter().skip(1).map(String::as_str).chain([ + auth_tag_json.as_str(), + auth_text.as_str(), + auth_event_json.as_str(), + ]) { + assert!( + !logs.contains(private_bytes), + "Tungstenite/log bridge leaked private AUTH bytes `{private_bytes}`:\n{logs}" + ); + } +} + +#[tokio::test] +async fn malformed_auth_echo_is_sanitized_in_error_logs() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let relay = format!("ws://{}", listener.local_addr().unwrap()); + let keys = Keys::generate(); + let auth_tag_parts = vec![ + "auth".to_string(), + "echo-private-owner-authority-9f8e7d6c".to_string(), + "echo-private-condition-kind-verified".to_string(), + "echo-private-reusable-signature-1a2b3c4d".to_string(), + ]; + let auth_tag_json = serde_json::to_string(&auth_tag_parts).unwrap(); + + let child = Command::new(std::env::current_exe().unwrap()) + .args([ + "--ignored", + "--exact", + "auth_error_log_client_helper", + "--nocapture", + ]) + .env(RELAY_ENV, &relay) + .env(PRIVATE_KEY_ENV, keys.secret_key().to_secret_hex()) + .env(AUTH_TAG_ENV, &auth_tag_json) + .env( + "RUST_LOG", + "tungstenite=trace,tokio_tungstenite=trace,buzz_ws_client=trace,auth_trace_redaction=trace", + ) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .kill_on_drop(true) + .spawn() + .unwrap(); + + let (stream, _) = timeout(Duration::from_secs(2), listener.accept()) + .await + .unwrap() + .unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + websocket + .send(Message::Text( + json!(["AUTH", "malformed-log-echo-challenge"]) + .to_string() + .into(), + )) + .await + .unwrap(); + + let auth_message = timeout(Duration::from_secs(2), websocket.next()) + .await + .unwrap() + .unwrap() + .unwrap(); + let auth_text = auth_message.to_text().unwrap().to_string(); + let auth_json: Value = serde_json::from_str(&auth_text).unwrap(); + let auth_event: Event = serde_json::from_value(auth_json[1].clone()).unwrap(); + auth_event.verify().unwrap(); + websocket + .send(Message::Text(format!("['AUTH',{}]", auth_json[1]).into())) + .await + .unwrap(); + + let output = timeout(Duration::from_secs(5), child.wait_with_output()) + .await + .unwrap() + .unwrap(); + let stdout = String::from_utf8_lossy(&output.stdout); + let stderr = String::from_utf8_lossy(&output.stderr); + let logs = format!("{stdout}{stderr}"); + assert!(output.status.success(), "error-log helper failed:\n{logs}"); + assert!( + logs.contains("relay reflected private authentication material"), + "sanitized log control was vacuous:\n{logs}" + ); + assert!( + logs.contains("connected to relay") && logs.contains("AUTH "), + "nonsecret tracing controls were vacuous:\n{logs}" + ); + + let auth_event_json = auth_event.as_json(); + let signature = auth_json[1]["sig"].as_str().unwrap(); + for private_bytes in auth_tag_parts.iter().skip(1).map(String::as_str).chain([ + auth_tag_json.as_str(), + auth_text.as_str(), + auth_event_json.as_str(), + signature, + ]) { + assert!( + !stdout.contains(private_bytes), + "post-AUTH reflection reached helper stdout: `{private_bytes}`:\n{stdout}" + ); + assert!( + !stderr.contains(private_bytes), + "post-AUTH reflection reached dependency/application stderr logs: `{private_bytes}`:\n{stderr}" + ); + } +} + +#[tokio::test] +async fn direct_send_raw_auth_is_rejected_before_transmission() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let relay = format!("ws://{}", listener.local_addr().unwrap()); + let relay_url = RelayUrl::parse(&relay).unwrap(); + let auth_event = EventBuilder::auth("direct-send-raw-challenge", relay_url) + .sign_with_keys(&Keys::generate()) + .unwrap(); + let private_signature = auth_event.sig.to_string(); + + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut websocket = tokio_tungstenite::accept_async(stream).await.unwrap(); + let Ok(Some(Ok(auth_message))) = + timeout(Duration::from_millis(250), websocket.next()).await + else { + return false; + }; + let auth: Value = serde_json::from_str(auth_message.to_text().unwrap()).unwrap(); + assert_eq!(auth[0], "AUTH"); + let signature = auth[1]["sig"].as_str().unwrap(); + websocket + .send(Message::Text( + json!(["NOTICE", format!("echo:{signature}")]) + .to_string() + .into(), + )) + .await + .unwrap(); + true + }); + + let mut connection = NostrWsConnection::connect(&relay).await.unwrap(); + let result = connection.send_raw(&json!(["AUTH", auth_event])).await; + let safely_rejected = matches!( + &result, + Err(WsClientError::AuthFailed(message)) + if message == "raw AUTH messages are not accepted; use authenticate" + ); + assert!( + safely_rejected, + "direct send_raw AUTH was not rejected with the static API error" + ); + if let Err(error) = result { + assert!(!error.to_string().contains(&private_signature)); + } + assert!(!connection.private_auth_started()); + assert!( + !server.await.unwrap(), + "rejected direct send_raw AUTH still reached the relay" + ); +}