diff --git a/Cargo.lock b/Cargo.lock index dd15b2c2227..7f28640f86a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1321,10 +1321,12 @@ dependencies = [ name = "buzz-pair-relay" version = "0.1.0" dependencies = [ + "buzz-core", "futures-util", "http-body-util", "hyper", "hyper-util", + "nostr 0.44.7", "parking_lot", "secp256k1 0.31.1", "serde_json", @@ -1332,6 +1334,7 @@ dependencies = [ "tokio", "tokio-tungstenite 0.29.0", "tokio-util", + "zeroize", ] [[package]] diff --git a/crates/buzz-core/src/pairing/session.rs b/crates/buzz-core/src/pairing/session.rs index 0d43d4d8279..606315beadd 100644 --- a/crates/buzz-core/src/pairing/session.rs +++ b/crates/buzz-core/src/pairing/session.rs @@ -27,6 +27,14 @@ //! handle_complete(&event) → complete_event //! ``` +//! # Code-entry confirmation extension +//! +//! A target advertises `"confirmation":"desktop-code-v1"` in its encrypted +//! offer. The source generates a separate random six-digit code, never included +//! in the QR or challenge. The target submits user input through NIP-44. Only a +//! matching code releases the source proof and payload; five guesses abort the +//! entire session. Legacy offers still require explicit source-side approval. + use std::collections::HashSet; use std::time::{Duration, Instant}; @@ -104,6 +112,9 @@ pub struct PairingSession { created_at: Instant, /// Maximum session lifetime. timeout: Duration, + desktop_code_requested: bool, + desktop_code: Option>, + code_attempts: u8, } impl PairingSession { @@ -136,6 +147,9 @@ impl PairingSession { processed_ids: HashSet::new(), created_at: Instant::now(), timeout: DEFAULT_TIMEOUT, + desktop_code_requested: false, + desktop_code: None, + code_attempts: 0, }; (session, qr) @@ -147,17 +161,33 @@ impl PairingSession { /// formatted SAS code to display. After this call the session is in /// [`SessionState::Confirming`]. pub fn handle_offer(&mut self, event: &Event) -> Result { + self.handle_offer_with_confirmation(event) + .map(|(code, _)| code) + } + + /// Validate an offer and return its SAS and whether the target requests + /// code-entry confirmation. The capability is encrypted so strict relays + /// still receive only the required recipient tag. + pub fn handle_offer_with_confirmation( + &mut self, + event: &Event, + ) -> Result<(String, bool), PairingError> { self.check_expired()?; self.expect_state(SessionState::Waiting)?; self.expect_role(Role::Source)?; self.validate_event_basics(event)?; let msg = self.decrypt_message(event)?; - let (session_id_hex, version) = match &msg { + let (session_id_hex, version, code_entry) = match &msg { PairingMessage::Offer { session_id, version, - } => (session_id.clone(), *version), + confirmation, + } => ( + session_id.clone(), + *version, + confirmation.as_deref() == Some("desktop-code-v1"), + ), other => return Err(unexpected("offer", other)), }; @@ -192,7 +222,8 @@ impl PairingSession { self.state = SessionState::Confirming; self.record_event(event); - Ok(format_sas(code)) + self.desktop_code_requested = code_entry; + Ok((format_sas(code), code_entry)) } /// (Source) User confirmed the SAS codes match. Build the `sas-confirm` @@ -349,12 +380,16 @@ impl PairingSession { processed_ids: HashSet::new(), created_at: Instant::now(), timeout: DEFAULT_TIMEOUT, + desktop_code_requested: false, + desktop_code: None, + code_attempts: 0, }; // Build and return the offer event. let msg = PairingMessage::Offer { session_id: hex::encode(session_id), version: 1, + confirmation: None, }; let event = session.build_event(&msg)?; session.state = SessionState::Confirming; @@ -514,9 +549,15 @@ impl PairingSession { } } + /// Absolute protocol deadline. Transports must use this same deadline for + /// UI expiry so connection setup never adds time to an expired QR. + pub fn deadline(&self) -> Instant { + self.created_at + self.timeout + } + /// Check if the session has expired. pub fn is_expired(&self) -> bool { - self.created_at.elapsed() > self.timeout + Instant::now() >= self.deadline() } /// Current protocol state. @@ -784,6 +825,9 @@ impl Drop for PairingSession { fn unexpected(expected: &str, got: &PairingMessage) -> PairingError { let got_name = match got { PairingMessage::Offer { .. } => "offer", + PairingMessage::DesktopCode {} => "desktop-code", + PairingMessage::CodeSubmit { .. } => "code-submit", + PairingMessage::CodeRejected { .. } => "code-rejected", PairingMessage::SasConfirm { .. } => "sas-confirm", PairingMessage::Payload { .. } => "payload", PairingMessage::Complete { .. } => "complete", @@ -1423,3 +1467,10 @@ mod tests { ); } } + +#[cfg(test)] +#[path = "session_code_entry_tests.rs"] +mod code_entry_tests; + +#[path = "session_desktop_code.rs"] +mod desktop_code; diff --git a/crates/buzz-core/src/pairing/session_code_entry_tests.rs b/crates/buzz-core/src/pairing/session_code_entry_tests.rs new file mode 100644 index 00000000000..09d902cea5b --- /dev/null +++ b/crates/buzz-core/src/pairing/session_code_entry_tests.rs @@ -0,0 +1,196 @@ +use super::*; + +fn setup() -> (PairingSession, PairingSession, String) { + let (mut source, qr) = PairingSession::new_source("wss://relay.test".into()); + let (target, _) = PairingSession::new_target(&qr).expect("target"); + let offer = target + .build_event(&PairingMessage::Offer { + session_id: hex::encode(target.session_id), + version: 1, + confirmation: Some("desktop-code-v1".into()), + }) + .expect("offer"); + assert!( + source + .handle_offer_with_confirmation(&offer) + .expect("offer") + .1 + ); + let (code, challenge) = source.start_desktop_code().expect("challenge"); + assert_eq!( + target.decrypt_message(&challenge).expect("decrypt"), + PairingMessage::DesktopCode {} + ); + assert_eq!(code.len(), 6); + (source, target, code) +} + +fn submission(target: &PairingSession, code: &str, attempt: u8) -> Event { + target + .build_event(&PairingMessage::CodeSubmit { + code: code.into(), + request_id: attempt.to_string(), + }) + .expect("submission") +} + +#[test] +fn only_source_only_code_releases_identity() { + let (mut source, mut target, code) = setup(); + assert!(source + .send_payload(PayloadType::Custom, Zeroizing::new("secret".into())) + .is_err()); + let event = submission(&target, &code, 1); + let (proof, accepted) = source.handle_target_code(&event).expect("verify"); + assert!(accepted); + assert!(source.handle_target_code(&event).is_err()); + target.handle_sas_confirm(&proof).expect("source proof"); + target.confirm_target_sas().expect("user approves import"); + let payload = source + .send_payload(PayloadType::Custom, Zeroizing::new("secret".into())) + .expect("payload"); + assert_eq!( + &*target.handle_payload(&payload).expect("import").1, + "secret" + ); +} + +#[test] +fn qr_derived_transcript_cannot_authorize_release() { + let (mut source, target, _) = setup(); + let hash = derive_transcript_hash( + &target.session_id, + &source.pubkey().to_bytes(), + &target.pubkey().to_bytes(), + &target.sas_input.expect("sas"), + &target.session_secret, + ); + let event = target + .build_event(&PairingMessage::SasConfirm { + transcript_hash: hex::encode(hash), + }) + .expect("proof"); + assert!(source.handle_target_code(&event).is_err()); + assert_eq!(source.state(), SessionState::Confirming); + assert!(source + .send_payload(PayloadType::Custom, Zeroizing::new("secret".into())) + .is_err()); +} + +#[test] +fn five_wrong_guesses_abort_without_reset_or_replay_bypass() { + let (mut source, target, code) = setup(); + let wrong = if code == "000000" { "000001" } else { "000000" }; + for attempt in 1..=5 { + let event = submission(&target, wrong, attempt); + let (reply, accepted) = source.handle_target_code(&event).expect("reject"); + assert!(!accepted); + assert_eq!( + target.decrypt_message(&reply).expect("reply"), + PairingMessage::CodeRejected { + request_id: attempt.to_string(), + remaining_attempts: 5 - attempt + } + ); + assert!( + source.handle_target_code(&event).is_err(), + "duplicate submission" + ); + assert!( + source.start_desktop_code().is_err(), + "cannot regenerate code/budget" + ); + assert!(source + .send_payload(PayloadType::Custom, Zeroizing::new("secret".into())) + .is_err()); + } + assert_eq!(source.state(), SessionState::Aborted); + assert!(source + .handle_target_code(&submission(&target, &code, 6)) + .is_err()); +} + +#[test] +fn wrong_then_correct_code_succeeds() { + let (mut source, target, code) = setup(); + let wrong = if code == "000000" { "000001" } else { "000000" }; + assert!( + !source + .handle_target_code(&submission(&target, wrong, 1)) + .expect("reject") + .1 + ); + assert!( + source + .handle_target_code(&submission(&target, &code, 2)) + .expect("accept") + .1 + ); +} + +#[test] +fn wrong_peer_tampering_and_expiry_cannot_authorize() { + let (mut source, target, code) = setup(); + let (other, _) = PairingSession::new_source("wss://relay.test".into()); + let mut bad = submission(&target, &code, 1); + bad.pubkey = other.pubkey(); + assert!(source.handle_target_code(&bad).is_err()); + let mut tampered = submission(&target, &code, 2); + tampered.content.push('x'); + assert!(source.handle_target_code(&tampered).is_err()); + assert_eq!(source.code_attempts, 0); + source.created_at = Instant::now() - source.timeout - Duration::from_secs(1); + assert!(source + .handle_target_code(&submission(&target, &code, 3)) + .is_err()); +} + +#[test] +fn legacy_capabilities_never_enable_automatic_release() { + for capability in [ + None, + Some("code-entry"), + Some("unknown"), + Some("desktop-code-v1"), + ] { + let (mut source, qr) = PairingSession::new_source("wss://relay.test".into()); + let (target, _) = PairingSession::new_target(&qr).expect("target"); + let event = target + .build_event(&PairingMessage::Offer { + session_id: hex::encode(target.session_id), + version: 1, + confirmation: capability.map(str::to_string), + }) + .expect("offer"); + assert_eq!(event.tags.len(), 1); + let (_, enabled) = source + .handle_offer_with_confirmation(&event) + .expect("offer"); + assert_eq!(enabled, capability == Some("desktop-code-v1")); + assert_eq!(source.start_desktop_code().is_ok(), enabled); + } +} + +#[test] +fn readiness_delay_and_late_scan_share_the_original_deadline() { + let (mut source, qr) = PairingSession::new_source("wss://relay.test".into()); + // Simulate the maximum 35-second readiness wait without sleeping. + source.created_at = Instant::now() - Duration::from_secs(35); + let deadline = source.deadline(); + let remaining = deadline.saturating_duration_since(Instant::now()); + assert!(remaining <= Duration::from_secs(85)); + assert!(remaining > Duration::from_secs(84)); + // A scan 84 seconds after readiness still has one second to be accepted. + source.created_at -= Duration::from_secs(84); + let (_, offer) = PairingSession::new_target(&qr).unwrap(); + assert!(source.handle_offer(&offer).is_ok()); + assert!(!source.is_expired()); + // The transport deadline and protocol expiry both end the original window. + source.created_at -= Duration::from_secs(2); + assert!(source.deadline() < Instant::now()); + assert!(source.is_expired()); + assert!(matches!( + source.confirm_sas(), + Err(PairingError::SessionExpired) + )); +} diff --git a/crates/buzz-core/src/pairing/session_desktop_code.rs b/crates/buzz-core/src/pairing/session_desktop_code.rs new file mode 100644 index 00000000000..0ffd1c60d88 --- /dev/null +++ b/crates/buzz-core/src/pairing/session_desktop_code.rs @@ -0,0 +1,57 @@ +use super::*; +use subtle::ConstantTimeEq; + +impl PairingSession { + /// Start source-only code verification once, after capability negotiation. + /// Returns the code for local display and an encrypted challenge containing + /// no code or verifier. Repeated calls cannot reset the guess budget. + pub fn start_desktop_code(&mut self) -> Result<(String, Event), PairingError> { + self.check_expired()?; + self.expect_role(Role::Source)?; + self.expect_state(SessionState::Confirming)?; + if !self.desktop_code_requested || self.desktop_code.is_some() { + return Err(PairingError::SasMismatch); + } + let code = format!("{:06}", rand::random_range(0..1_000_000u32)); + let challenge = self.build_event(&PairingMessage::DesktopCode {})?; + self.desktop_code = Some(Zeroizing::new(code.clone())); + Ok((code, challenge)) + } + + /// Verify a signed submission from the locked peer. Only a separate random + /// desktop code authorizes release; the QR-derived SAS/transcript cannot. + /// The response is a source proof on success or a rejection with a remaining + /// guess budget. Five wrong guesses permanently abort this QR session. + pub fn handle_target_code(&mut self, event: &Event) -> Result<(Event, bool), PairingError> { + self.check_expired()?; + self.expect_role(Role::Source)?; + self.expect_state(SessionState::Confirming)?; + self.validate_event_from_peer(event)?; + let expected = self + .desktop_code + .as_ref() + .ok_or(PairingError::SasMismatch)?; + let (mut code, request_id) = match self.decrypt_message(event)? { + PairingMessage::CodeSubmit { code, request_id } => (code, request_id), + other => return Err(unexpected("code-submit", &other)), + }; + let correct: bool = code.as_bytes().ct_eq(expected.as_bytes()).into(); + code.zeroize(); + self.code_attempts += 1; + self.record_event(event); + if correct { + self.desktop_code = None; + return self.confirm_sas().map(|proof| (proof, true)); + } + let remaining_attempts = 5u8.saturating_sub(self.code_attempts); + let rejection = self.build_event(&PairingMessage::CodeRejected { + request_id, + remaining_attempts, + })?; + if remaining_attempts == 0 { + self.desktop_code = None; + self.state = SessionState::Aborted; + } + Ok((rejection, false)) + } +} diff --git a/crates/buzz-core/src/pairing/types.rs b/crates/buzz-core/src/pairing/types.rs index 0dcc0baf5ad..13de9ce898a 100644 --- a/crates/buzz-core/src/pairing/types.rs +++ b/crates/buzz-core/src/pairing/types.rs @@ -28,6 +28,26 @@ pub enum PairingMessage { /// Defaults to `1` when absent (backward compat with pre-versioned implementations). #[serde(default = "default_version")] version: u32, + /// Optional confirmation UX negotiated inside the encrypted offer. + #[serde(default, skip_serializing_if = "Option::is_none")] + confirmation: Option, + }, + + /// Source advertises a separate, source-only random code. Contains no code. + DesktopCode {}, + /// Target submits the user-entered code over the authenticated encrypted channel. + CodeSubmit { + /// Six ASCII digits displayed only on the source. + code: String, + /// Correlates rejection with the submitted attempt. + request_id: String, + }, + /// Source rejects an attempt without releasing any identity material. + CodeRejected { + /// Request being rejected. + request_id: String, + /// Remaining guesses before this QR session is permanently aborted. + remaining_attempts: u8, }, /// Either party → other. Confirms the Short Authentication String matches. @@ -104,6 +124,7 @@ mod tests { let msg = PairingMessage::Offer { session_id: "deadbeef".repeat(8), version: 1, + confirmation: None, }; let json = serde_json::to_string(&msg).expect("serialize"); assert!( @@ -129,6 +150,7 @@ mod tests { session_id: "deadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeef" .to_string(), version: 1, + confirmation: None, } ); } diff --git a/crates/buzz-pair-relay/Cargo.toml b/crates/buzz-pair-relay/Cargo.toml index 740f13da1d6..e90bc46044a 100644 --- a/crates/buzz-pair-relay/Cargo.toml +++ b/crates/buzz-pair-relay/Cargo.toml @@ -29,6 +29,9 @@ secp256k1 = { version = "0.31", features = ["global-context"] } sha2 = "0.11" [dev-dependencies] +buzz-core = { path = "../buzz-core" } +nostr = { workspace = true } +zeroize = { workspace = true } tokio = { workspace = true, features = ["test-util"] } tokio-tungstenite = { workspace = true } secp256k1 = { version = "0.31", features = ["global-context", "rand", "std"] } diff --git a/crates/buzz-pair-relay/src/lib.rs b/crates/buzz-pair-relay/src/lib.rs index 8d30a4b1ec4..ffed66d7106 100644 --- a/crates/buzz-pair-relay/src/lib.rs +++ b/crates/buzz-pair-relay/src/lib.rs @@ -21,7 +21,7 @@ //! NIP-01 event ID hash. Events with invalid signatures are rejected. //! - **No persistence** — events exist only in-flight between matched pub/sub. //! - **Bounded resources** — 128 max WS connections, 4 KiB max frame, 120s TTL. -//! - **Session cap** — at most 6 accepted EVENTs per connection. +//! - **Session cap** — at most 8 attempted EVENTs per connection. //! - **Freshness** — `created_at` must be within ±120 s of relay wall-clock. //! - **Deduplication** — duplicate event IDs are rejected; dedup entries expire after 300 s. @@ -70,10 +70,12 @@ const RATE_EVENT_MAX: u32 = 10; const SUB_ID_MAX: usize = 64; /// Hard session cap: at most this many attempted EVENTs (post-sig-check) per connection. -const MAX_EVENTS_PER_CONN: u32 = 6; +// Desktop-code pairing needs seven events per peer at the fifth successful +// guess (including offer/challenge, proof, payload and completion), plus abort. +const MAX_EVENTS_PER_CONN: u32 = 8; /// Per-#p delivery budget: enough for one full pairing from each direction. -const MAX_DELIVERED_PER_P: u32 = 12; +const MAX_DELIVERED_PER_P: u32 = 2 * MAX_EVENTS_PER_CONN; /// Dedup vec rejects new events when still at capacity after TTL eviction (fail closed). const DEDUP_CAP: usize = 1024; diff --git a/crates/buzz-pair-relay/tests/desktop_code.rs b/crates/buzz-pair-relay/tests/desktop_code.rs new file mode 100644 index 00000000000..ba533ded29a --- /dev/null +++ b/crates/buzz-pair-relay/tests/desktop_code.rs @@ -0,0 +1,163 @@ +//! Full desktop-code wire exchange through the production relay and source state machine. +use std::{sync::Arc, time::Duration}; + +use buzz_core::pairing::{crypto, PairingSession, PayloadType, SessionState}; +use buzz_pair_relay::{run_server, Relay}; +use futures_util::{SinkExt, StreamExt}; +use nostr::{nips::nip44, Event, EventBuilder, Keys, Kind, PublicKey, Tag, ToBech32}; +use serde_json::{json, Value}; +use tokio::net::TcpListener; +use tokio_tungstenite::{connect_async, tungstenite::Message, MaybeTlsStream, WebSocketStream}; +use zeroize::Zeroizing; + +type Socket = WebSocketStream>; + +async fn receive(socket: &mut Socket) -> Value { + let frame = tokio::time::timeout(Duration::from_secs(2), socket.next()) + .await + .expect("relay response deadline") + .expect("connection open") + .expect("frame"); + serde_json::from_str(frame.to_text().expect("text frame")).expect("relay JSON") +} + +async fn subscribe(socket: &mut Socket, pubkey: PublicKey) { + socket + .send(Message::Text( + json!(["REQ", "pair", { + "kinds": [24134], "#p": [pubkey.to_hex()] + }]) + .to_string() + .into(), + )) + .await + .unwrap(); + assert_eq!(receive(socket).await[0], "EOSE"); +} + +async fn forward(sender: &mut Socket, receiver: &mut Socket, event: Event) -> Event { + sender + .send(Message::Text(json!(["EVENT", event]).to_string().into())) + .await + .unwrap(); + let ack = receive(sender).await; + assert_eq!(ack[0], "OK"); + assert_eq!(ack[2], true, "relay must accept each protocol event: {ack}"); + let delivery = receive(receiver).await; + assert_eq!(delivery[0], "EVENT"); + let received: Event = serde_json::from_value(delivery[2].clone()).unwrap(); + received.verify().unwrap(); + received +} + +fn target_message(keys: &Keys, source: PublicKey, message: Value) -> Event { + let encrypted = nip44::encrypt( + keys.secret_key(), + &source, + message.to_string(), + nip44::Version::V2, + ) + .unwrap(); + EventBuilder::new(Kind::Custom(24134), encrypted) + .tags([Tag::public_key(source)]) + .sign_with_keys(keys) + .unwrap() +} + +fn decrypt(keys: &Keys, event: &Event) -> Value { + serde_json::from_str(&nip44::decrypt(keys.secret_key(), &event.pubkey, &event.content).unwrap()) + .unwrap() +} + +#[tokio::test] +async fn fifth_correct_guess_delivers_identity_and_completion() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("ws://{}", listener.local_addr().unwrap()); + let server = tokio::spawn(run_server(listener, Arc::new(Relay::new()))); + let (mut source, qr) = PairingSession::new_source(url.clone()); + let target = Keys::generate(); + let (mut source_socket, _) = connect_async(&url).await.unwrap(); + let (mut target_socket, _) = connect_async(&url).await.unwrap(); + subscribe(&mut source_socket, source.pubkey()).await; + subscribe(&mut target_socket, target.public_key()).await; + let session_id = crypto::derive_session_id(&qr.session_secret); + let hex = |bytes: &[u8]| bytes.iter().map(|b| format!("{b:02x}")).collect::(); + let offer = target_message( + &target, + source.pubkey(), + json!({ + "type": "offer", "session_id": hex(&session_id), "version": 1, + "confirmation": "desktop-code-v1" + }), + ); + let offer = forward(&mut target_socket, &mut source_socket, offer).await; + assert!(source.handle_offer_with_confirmation(&offer).unwrap().1); + let (code, challenge) = source.start_desktop_code().unwrap(); + let challenge = forward(&mut source_socket, &mut target_socket, challenge).await; + assert_eq!( + decrypt(&target, &challenge), + json!({"type": "desktop-code"}) + ); + let wrong = if code == "000000" { "000001" } else { "000000" }; + for attempt in 1..=5 { + let submission = target_message( + &target, + source.pubkey(), + json!({ + "type": "code-submit", "code": if attempt == 5 { &code } else { wrong }, + "request_id": attempt.to_string() + }), + ); + let submission = forward(&mut target_socket, &mut source_socket, submission).await; + let (response, accepted) = source.handle_target_code(&submission).unwrap(); + assert_eq!(accepted, attempt == 5); + let response = forward(&mut source_socket, &mut target_socket, response).await; + let message = decrypt(&target, &response); + if attempt < 5 { + assert_eq!( + message, + json!({"type":"code-rejected", "request_id":attempt.to_string(), "remaining_attempts":5-attempt}) + ); + } else { + let ecdh = + nostr::util::generate_shared_key(target.secret_key(), &source.pubkey()).unwrap(); + let (_, sas_input) = crypto::derive_sas(&ecdh, &qr.session_secret); + let proof = crypto::derive_transcript_hash( + &session_id, + &source.pubkey().to_bytes(), + &target.public_key().to_bytes(), + &sas_input, + &qr.session_secret, + ); + assert_eq!( + message, + json!({"type":"sas-confirm", "transcript_hash":hex(&proof)}) + ); + } + } + // Source event seven must reach the target and decrypt to the identity. + let identity = Keys::generate(); + let secret = identity.secret_key().to_bech32().unwrap(); + let payload = source + .send_payload(PayloadType::Nsec, Zeroizing::new(secret.clone())) + .unwrap(); + let payload = forward(&mut source_socket, &mut target_socket, payload).await; + let imported = decrypt(&target, &payload); + assert_eq!(imported["payload"], secret); + assert_eq!( + Keys::parse(imported["payload"].as_str().unwrap()) + .unwrap() + .public_key(), + identity.public_key() + ); + // Target event seven must reach source and finish the session. + let complete = target_message( + &target, + source.pubkey(), + json!({"type":"complete", "success":true}), + ); + let complete = forward(&mut target_socket, &mut source_socket, complete).await; + source.handle_complete(&complete).unwrap(); + assert_eq!(source.state(), SessionState::Completed); + server.abort(); +} diff --git a/crates/buzz-pair-relay/tests/integration.rs b/crates/buzz-pair-relay/tests/integration.rs index 77b71ba4dd5..dbd1809e2d0 100644 --- a/crates/buzz-pair-relay/tests/integration.rs +++ b/crates/buzz-pair-relay/tests/integration.rs @@ -439,16 +439,20 @@ async fn test_second_sub_same_id() { } /// 9. Connection closes after 120 s (virtual time). -#[tokio::test(start_paused = true)] +#[tokio::test] async fn test_120s_timeout() { let url = start_relay().await; let mut ws = connect(&url).await; + subscribe(&mut ws, "ready", P_A).await; + tokio::time::pause(); // Advance virtual time past the connection timeout. tokio::time::advance(Duration::from_secs(121)).await; // Yield to let the relay task run its deadline branch. tokio::task::yield_now().await; + // Real network I/O must not race an auto-advancing assertion timeout. + tokio::time::resume(); assert_closed(&mut ws).await; } @@ -605,8 +609,8 @@ async fn test_global_conn_cap() { let _new = connect(&url).await; } -/// 18. Session event cap: 6 EVENTs are accepted; 7th is rejected with -/// "session event limit reached". (The relay's hard cap is 6 per +/// 18. Session event cap: 8 EVENTs are accepted; 9th is rejected with +/// "session event limit reached". (The relay's hard cap is 8 per /// connection, which is tighter than the per-window rate limit of 10.) #[tokio::test] async fn test_event_rate_limit() { @@ -619,8 +623,8 @@ async fn test_event_rate_limit() { let mut pub_ws = connect(&url).await; - // First 6 events must be accepted (session cap = 6). - for i in 0..6u64 { + // First 8 events must be accepted (session cap = 8). + for i in 0..8u64 { send( &mut pub_ws, &json!(["EVENT", make_signed_event(&sk, &pk, P_A, i)]), @@ -637,10 +641,10 @@ async fn test_event_rate_limit() { let _ = recv(&mut sub_ws).await; } - // 7th must be rejected by the session cap. + // 9th must be rejected by the session cap. send( &mut pub_ws, - &json!(["EVENT", make_signed_event(&sk, &pk, P_A, 6)]), + &json!(["EVENT", make_signed_event(&sk, &pk, P_A, 8)]), ) .await; let resp = recv(&mut pub_ws).await; @@ -931,9 +935,9 @@ async fn test_write_timeout() { send(&mut sub_ws, &json!(["REQ", "s1", {"#p": [P_A]}])).await; // Don't call recv — leave the EOSE unread. - // Flood from a publisher — stay within the 6-event session cap. + // Flood from a publisher — stay within the 8-event session cap. let mut pub_ws = connect(&url).await; - for i in 0..6u64 { + for i in 0..8u64 { send( &mut pub_ws, &json!(["EVENT", make_signed_event(&sk, &pk, P_A, i)]), @@ -1016,9 +1020,9 @@ async fn test_control_msg_backpressure() { send(&mut sub_ws, &json!(["REQ", "s1", {"#p": [P_A]}])).await; // Leave EOSE unread to fill the channel quickly. - // Flood within the 6-event session cap. + // Flood within the 8-event session cap. let mut pub_ws = connect(&url).await; - for i in 0..6u64 { + for i in 0..8u64 { send( &mut pub_ws, &json!(["EVENT", make_signed_event(&sk, &pk, P_A, i)]), @@ -1132,9 +1136,9 @@ async fn test_fan_out_drop_doesnt_close() { let mut sub_ws = connect(&url).await; send(&mut sub_ws, &json!(["REQ", "s1", {"#p": [P_A]}])).await; - // Flood from a publisher — stay within the 6-event session cap. + // Flood from a publisher — stay within the 8-event session cap. let mut pub_ws = connect(&url).await; - for i in 0..6u64 { + for i in 0..8u64 { send( &mut pub_ws, &json!(["EVENT", make_signed_event(&sk, &pk, P_A, i)]), @@ -1164,7 +1168,7 @@ async fn test_reader_backpressure_closes() { // Publisher floods events up to the session cap. let mut pub_ws = connect(&url).await; - for i in 0..6u64 { + for i in 0..8u64 { send( &mut pub_ws, &json!(["EVENT", make_signed_event(&sk, &pk, P_A, i)]), @@ -1178,14 +1182,21 @@ async fn test_reader_backpressure_closes() { /// 42. Connection closes promptly after 120 s (virtual time). /// Explicit duplicate of test 9 with a slightly different assertion style. -#[tokio::test(start_paused = true)] +#[tokio::test] async fn test_cancellation_immediate() { let url = start_relay().await; let mut ws = connect(&url).await; + // Observe a relay response so the connection deadline is installed. + subscribe(&mut ws, "ready", P_A).await; + tokio::time::pause(); tokio::time::advance(Duration::from_secs(121)).await; tokio::task::yield_now().await; + // TCP closure is real I/O; do not auto-advance the assertion timeout + // while the OS is still delivering the close frame. + tokio::time::resume(); + // The connection must be closed — not just slow. assert_closed(&mut ws).await; } diff --git a/desktop/src-tauri/src/commands/pairing.rs b/desktop/src-tauri/src/commands/pairing.rs index aedd67854c1..37a36fdeb07 100644 --- a/desktop/src-tauri/src/commands/pairing.rs +++ b/desktop/src-tauri/src/commands/pairing.rs @@ -10,7 +10,11 @@ use futures_util::{SinkExt, StreamExt}; use nostr::ToBech32; use serde::Serialize; use tauri::{AppHandle, Emitter, Manager, State}; -use tokio::sync::mpsc; +use tokio::sync::{mpsc, oneshot}; + +#[path = "pairing_subscription.rs"] +mod pairing_subscription; +use pairing_subscription::wait_for_eose; use tokio_tungstenite::{connect_async, tungstenite::Message}; use tokio_util::sync::CancellationToken; use zeroize::Zeroizing; @@ -21,6 +25,7 @@ use crate::relay::{relay_api_base_url_with_override, relay_ws_url_with_override} #[derive(Serialize, Clone)] struct PairingSasPayload { sas: String, + code_entry: bool, } #[derive(Serialize, Clone)] @@ -41,12 +46,35 @@ enum PairingMode { #[derive(Clone)] struct PairingTaskContext { + payload: Arc>>>, mode: PairingMode, generation: Arc, generation_fence: Arc>, task_generation: u64, } +impl PairingTaskContext { + fn take_payload(&self) -> Result, String> { + let _fence = self.generation_fence.lock().map_err(|e| e.to_string())?; + ensure_pairing_task_is_current(&self.generation, self.task_generation)?; + self.payload + .lock() + .map_err(|e| e.to_string())? + .take() + .ok_or_else(|| "Pairing payload missing".into()) + } + + fn clear_payload_if_current(&self) { + let _fence = self + .generation_fence + .lock() + .unwrap_or_else(|e| e.into_inner()); + if pairing_task_is_current(&self.generation, self.task_generation) { + *self.payload.lock().unwrap_or_else(|e| e.into_inner()) = None; + } + } +} + /// Managed Tauri state for an active pairing session. pub struct PairingHandle { session: Arc>>, @@ -61,7 +89,7 @@ pub struct PairingHandle { outbound_tx: std::sync::Mutex>>, /// Pre-built payload string (contains nsec) to send after SAS confirmation. /// Wrapped in Zeroizing so the nsec is cleared from memory on drop. - payload: std::sync::Mutex>>, + payload: Arc>>>, mode: Arc>, } @@ -74,7 +102,7 @@ impl PairingHandle { start_lock: tokio::sync::Mutex::new(()), cancel: std::sync::Mutex::new(None), outbound_tx: std::sync::Mutex::new(None), - payload: std::sync::Mutex::new(None), + payload: Arc::new(std::sync::Mutex::new(None)), mode: Arc::new(std::sync::Mutex::new(PairingMode::SendIdentity)), } } @@ -161,21 +189,36 @@ async fn start_pairing_session( *pairing.outbound_tx.lock().map_err(|e| e.to_string())? = Some(outbound_tx); *pairing.cancel.lock().map_err(|e| e.to_string())? = Some(cancel.clone()); + let (ready_tx, ready_rx) = oneshot::channel(); tauri::async_runtime::spawn(pairing_ws_task( pairing_relay_url, Arc::clone(&pairing.session), PairingTaskContext { + payload: Arc::clone(&pairing.payload), mode, generation: Arc::clone(&pairing.generation), generation_fence: Arc::clone(&pairing.generation_fence), task_generation, }, - cancel, + cancel.clone(), outbound_rx, app, + ready_tx, )); - Ok(qr_uri) + // A QR must not become scannable until its ephemeral subscription is ready. + match tokio::time::timeout(Duration::from_secs(35), ready_rx).await { + Ok(Ok(Ok(()))) => { + ensure_pairing_task_is_current(&pairing.generation, task_generation)?; + Ok(qr_uri) + } + Ok(Ok(Err(error))) => Err(error), + Ok(Err(_)) => Err("Pairing connection stopped before it was ready".into()), + Err(_) => { + cancel.cancel(); + Err("Pairing connection took too long. Try again.".into()) + } + } } /// User confirmed the SAS codes match. Sends sas-confirm + payload. @@ -277,21 +320,31 @@ async fn pairing_ws_task( cancel: CancellationToken, mut outbound_rx: mpsc::Receiver, app: AppHandle, + ready: oneshot::Sender>, ) { - if let Err(e) = pairing_ws_task_inner( + let mut ready = Some(ready); + let result = tokio::select! { + biased; + _ = cancel.cancelled() => Err("Pairing was canceled".into()), + result = pairing_ws_task_inner( &relay_url, &session, &context, &cancel, &mut outbound_rx, &app, - ) - .await - { + &mut ready, + ) => result, + }; + if let Err(e) = result { + if let Some(ready) = ready.take() { + let _ = ready.send(Err(e.clone())); + } if pairing_task_is_current(&context.generation, context.task_generation) { let _ = app.emit("pairing-error", PairingErrorPayload { message: e }); } } + context.clear_payload_if_current(); clear_pairing_session_if_current(&session, &context.generation, context.task_generation).await; } @@ -302,6 +355,7 @@ async fn pairing_ws_task_inner( cancel: &CancellationToken, outbound_rx: &mut mpsc::Receiver, app: &AppHandle, + ready: &mut Option>>, ) -> Result<(), String> { let (ws, _) = connect_async(relay_url) .await @@ -323,9 +377,24 @@ async fn pairing_ws_task_inner( .await .map_err(|e| format!("subscribe failed: {e}"))?; - wait_for_eose(&mut read, "pair", Duration::from_secs(10)).await?; + let pending = wait_for_eose(&mut read, "pair", Duration::from_secs(10)).await?; + ensure_pairing_task_is_current(&context.generation, context.task_generation)?; + let hard_timeout = { + let guard = session.lock().await; + let active = guard.as_ref().ok_or("session gone")?; + if active.is_expired() { + return Err("Pairing session expired during connection setup".into()); + } + pairing_expiry_timer(active) + }; + // Readiness may have consumed part of the protocol lifetime. Expire the + // visible QR at that original deadline, never 130 seconds after readiness. + if let Some(ready) = ready.take() { + let _ = ready.send(Ok(())); + } + let mut read = futures_util::stream::iter(pending.into_iter().map(Ok)).chain(read); - let hard_timeout = tokio::time::sleep(Duration::from_secs(130)); + let mut code_entry = false; tokio::pin!(hard_timeout); loop { @@ -361,6 +430,9 @@ async fn pairing_ws_task_inner( } let mut guard = session.lock().await; + if !pairing_task_is_current(&context.generation, context.task_generation) { + break; + } let Some(s) = guard.as_mut() else { break }; if let Ok(reason) = s.handle_abort(&event) { @@ -372,13 +444,55 @@ async fn pairing_ws_task_inner( break; } - if let Ok(sas) = s.handle_offer(&event) { + if let Ok((mut sas, target_code_entry)) = s.handle_offer_with_confirmation(&event) { + code_entry = context.mode == PairingMode::SendIdentity && target_code_entry; + if code_entry { + let (code, challenge) = s.start_desktop_code().map_err(|e| e.to_string())?; + sas = code; + write.send(Message::Text(event_to_relay_json(&challenge).into())).await + .map_err(|e| format!("publish challenge failed: {e}"))?; + } if pairing_task_is_current(&context.generation, context.task_generation) { - let _ = app.emit("pairing-sas-received", PairingSasPayload { sas }); + let _ = app.emit("pairing-sas-received", PairingSasPayload { sas, code_entry }); } continue; } + if code_entry { + match s.handle_target_code(&event) { + Ok((response, false)) => { + write.send(Message::Text(event_to_relay_json(&response).into())).await + .map_err(|e| format!("publish rejection failed: {e}"))?; + if s.state() == buzz_core_pkg::pairing::SessionState::Aborted { + return Err("Too many incorrect codes. Create a new QR code.".into()); + } + continue; + } + Ok((proof, true)) => { + let identity = context.take_payload()?; + let transfer = s.send_payload(PayloadType::Custom, identity) + .map_err(|e| e.to_string())?; + ensure_pairing_task_is_current(&context.generation, context.task_generation)?; + tokio::select! { + biased; + _ = cancel.cancelled() => break, + result = async { + write.send(Message::Text(event_to_relay_json(&proof).into())).await?; + write.send(Message::Text(event_to_relay_json(&transfer).into())).await + } => result.map_err(|e| format!("publish pairing failed: {e}"))?, + } + if pairing_task_is_current(&context.generation, context.task_generation) { + let _ = app.emit("pairing-code-entered", serde_json::json!({})); + } + continue; + } + Err(buzz_core_pkg::pairing::PairingError::TranscriptMismatch) => { + return Err("Code verification failed. Try pairing again.".into()); + } + Err(_) => {} + } + } + if context.mode == PairingMode::RecoverIdentity { if let Ok((payload_type, payload)) = s.handle_return_payload(&event) { if let Err(message) = validate_recovery_payload_type(payload_type) { @@ -552,6 +666,10 @@ fn finish_recovery( Ok(()) } +fn pairing_expiry_timer(session: &PairingSession) -> tokio::time::Sleep { + tokio::time::sleep_until(session.deadline().into()) +} + fn pairing_task_is_current(generation: &AtomicU64, task_generation: u64) -> bool { generation.load(Ordering::SeqCst) == task_generation } @@ -755,35 +873,6 @@ fn parse_auth_challenge(text: &str) -> Option { None } -async fn wait_for_eose(read: &mut S, sub_id: &str, dur: Duration) -> Result<(), String> -where - S: StreamExt> + Unpin, -{ - tokio::time::timeout(dur, async { - loop { - let msg = read - .next() - .await - .ok_or_else(|| "relay closed waiting for EOSE".to_string())? - .map_err(|e| format!("WS error waiting for EOSE: {e}"))?; - if let Message::Text(text) = msg { - if let Ok(arr) = serde_json::from_str::(text.as_str()) { - if let Some(arr) = arr.as_array() { - if arr.len() >= 2 - && arr[0].as_str() == Some("EOSE") - && arr[1].as_str() == Some(sub_id) - { - return Ok(()); - } - } - } - } - } - }) - .await - .map_err(|_| "timeout waiting for EOSE".to_string())? -} - #[cfg(test)] #[path = "pairing_generation_tests.rs"] mod pairing_generation_tests; diff --git a/desktop/src-tauri/src/commands/pairing_generation_tests.rs b/desktop/src-tauri/src/commands/pairing_generation_tests.rs index 8a2291ae86f..ca5b6ee01da 100644 --- a/desktop/src-tauri/src/commands/pairing_generation_tests.rs +++ b/desktop/src-tauri/src/commands/pairing_generation_tests.rs @@ -127,3 +127,51 @@ async fn current_task_clears_its_session() { assert!(session.lock().await.is_none()); } + +fn prepared_context(pairing: &PairingHandle) -> super::PairingTaskContext { + *pairing.payload.lock().unwrap() = Some(zeroize::Zeroizing::new("test-identity".into())); + super::PairingTaskContext { + payload: Arc::clone(&pairing.payload), + mode: super::PairingMode::SendIdentity, + generation: Arc::clone(&pairing.generation), + generation_fence: Arc::clone(&pairing.generation_fence), + task_generation: pairing.generation.load(Ordering::SeqCst), + } +} + +#[test] +fn code_entry_transfer_consumes_the_managed_secret() { + let pairing = PairingHandle::new(); + let context = prepared_context(&pairing); + let identity = context.take_payload().unwrap(); + assert_eq!(identity.as_str(), "test-identity"); + assert!(pairing.payload.lock().unwrap().is_none()); + assert!(context.take_payload().is_err()); +} + +#[test] +fn terminal_worker_clears_secret_but_stale_worker_cannot_touch_replacement() { + let pairing = PairingHandle::new(); + let context = prepared_context(&pairing); + context.clear_payload_if_current(); + assert!(pairing.payload.lock().unwrap().is_none()); + let context = prepared_context(&pairing); + pairing.generation.fetch_add(1, Ordering::SeqCst); + context.clear_payload_if_current(); + assert!(context.take_payload().is_err()); + assert!(pairing.payload.lock().unwrap().is_some()); +} + +#[tokio::test(start_paused = true)] +async fn slow_readiness_does_not_extend_the_visible_qr_past_protocol_expiry() { + let (session, _) = PairingSession::new_source("wss://relay.test".into()); + let deadline = tokio::time::Instant::from_std(session.deadline()); + tokio::time::advance(Duration::from_secs(35)).await; + let expiry = super::pairing_expiry_timer(&session); + tokio::pin!(expiry); + assert_eq!(expiry.deadline(), deadline); + tokio::time::advance(Duration::from_secs(84)).await; + assert!(futures_util::poll!(expiry.as_mut()).is_pending()); + tokio::time::advance(Duration::from_secs(2)).await; + assert!(futures_util::poll!(expiry.as_mut()).is_ready()); +} diff --git a/desktop/src-tauri/src/commands/pairing_subscription.rs b/desktop/src-tauri/src/commands/pairing_subscription.rs new file mode 100644 index 00000000000..d9cfb5bbf22 --- /dev/null +++ b/desktop/src-tauri/src/commands/pairing_subscription.rs @@ -0,0 +1,119 @@ +//! Preserve pairing events delivered before the subscription's EOSE marker. +use std::time::Duration; + +use futures_util::StreamExt; +use tokio_tungstenite::tungstenite::{Error, Message}; + +pub(super) async fn wait_for_eose( + read: &mut S, + sub_id: &str, + duration: Duration, +) -> Result, String> +where + S: StreamExt> + Unpin, +{ + tokio::time::timeout(duration, async { + let mut pending = Vec::new(); + loop { + let message = read + .next() + .await + .ok_or_else(|| "relay closed waiting for EOSE".to_string())? + .map_err(|error| format!("WS error waiting for EOSE: {error}"))?; + let Message::Text(text) = &message else { + continue; + }; + let Ok(value) = serde_json::from_str::(text.as_str()) else { + continue; + }; + let Some(values) = value.as_array() else { + continue; + }; + if values.get(1).and_then(|value| value.as_str()) != Some(sub_id) { + continue; + } + match values.first().and_then(|value| value.as_str()) { + Some("EOSE") => return Ok(pending), + Some("CLOSED") => { + return Err("The relay closed the pairing subscription. Try again.".into()) + } + Some("EVENT") if values.len() >= 3 => { + // A session needs only one offer; bound untrusted setup traffic. + if pending.len() >= 32 { + return Err("Too many pairing events during connection setup".into()); + } + pending.push(message); + } + _ => {} + } + } + }) + .await + .map_err(|_| "timeout waiting for EOSE".to_string())? +} + +#[cfg(test)] +mod tests { + use super::*; + use futures_util::stream; + + fn text(value: &str) -> Result { + Ok(Message::Text(value.into())) + } + + #[tokio::test] + async fn early_scan_is_preserved_in_order_before_live_messages() { + let offer = r#"["EVENT","pair",{"content":"encrypted offer"}]"#; + let confirmation = r#"["EVENT","pair",{"content":"encrypted confirmation"}]"#; + let mut read = stream::iter(vec![ + text(offer), + text(r#"["EOSE","other"]"#), + text(r#"["EOSE","pair"]"#), + text(confirmation), + ]); + let buffered = wait_for_eose(&mut read, "pair", Duration::from_secs(1)) + .await + .expect("ready"); + assert_eq!(buffered, vec![Message::Text(offer.into())]); + let mut combined = stream::iter(buffered.into_iter().map(Ok)).chain(read); + assert_eq!( + combined.next().await.expect("offer").expect("message"), + Message::Text(offer.into()) + ); + assert_eq!( + combined + .next() + .await + .expect("confirmation") + .expect("message"), + Message::Text(confirmation.into()) + ); + } + + #[tokio::test] + async fn closed_subscription_fails_instead_of_showing_a_dead_qr() { + let mut read = stream::iter(vec![text(r#"["CLOSED","pair","not allowed"]"#)]); + assert!(wait_for_eose(&mut read, "pair", Duration::from_secs(1)) + .await + .expect_err("closed") + .contains("closed")); + } + + #[tokio::test] + async fn missing_eose_times_out() { + let mut read = stream::pending::>(); + assert!(wait_for_eose(&mut read, "pair", Duration::from_millis(1)) + .await + .expect_err("timeout") + .contains("timeout")); + } + + #[tokio::test] + async fn event_buffer_is_bounded() { + let mut read = stream::iter((0..33).map(|_| text(r#"["EVENT","pair",{}]"#))); + assert!(wait_for_eose(&mut read, "pair", Duration::from_secs(1)) + .await + .expect_err("bound") + .contains("Too many")); + } +} diff --git a/desktop/src/features/settings/ui/MobilePairingCard.tsx b/desktop/src/features/settings/ui/MobilePairingCard.tsx index 3a25f2edfaf..55c73d58eb3 100644 --- a/desktop/src/features/settings/ui/MobilePairingCard.tsx +++ b/desktop/src/features/settings/ui/MobilePairingCard.tsx @@ -103,7 +103,13 @@ function PairingStepIndicator({ ); } -function PairingSteps({ step }: { step: PairingStep }) { +function PairingSteps({ + step, + codeEntry, +}: { + step: PairingStep; + codeEntry: boolean; +}) { const hasScanned = step === "sas" || step === "transferring" || step === "done"; const hasConfirmed = step === "transferring" || step === "done"; @@ -138,13 +144,16 @@ function PairingSteps({ step }: { step: PairingStep }) { testId="mobile-pairing-confirm-step-indicator" />
-

Confirm mobile code

+

+ {codeEntry ? "Enter code on your phone" : "Confirm mobile code"} +

- Check that the six-digit code matches on both devices, then confirm - it. + {codeEntry + ? "Type the six-digit code shown here into the Buzz mobile app." + : "Check that the six-digit code matches on both devices, then confirm it."}

@@ -168,7 +177,7 @@ function PairingSteps({ step }: { step: PairingStep }) { > {isPaired ? "Your mobile app is now connected to this relay." - : "Your mobile app will connect after you confirm the code."} + : "Finish setup on your phone to connect your mobile app."}

@@ -177,10 +186,12 @@ function PairingSteps({ step }: { step: PairingStep }) { } function PairingCodeConfirmation({ + codeEntry, onConfirm, onDeny, sasCode, }: { + codeEntry: boolean; onConfirm: () => void; onDeny: () => void; sasCode: string; @@ -196,7 +207,7 @@ function PairingCodeConfirmation({ className="self-start text-base font-medium" data-testid="pairing-sas-title" > - Confirm mobile code + {codeEntry ? "Enter code on your phone" : "Confirm mobile code"}

- + {!codeEntry && ( + + )}