diff --git a/core/archipelago/src/api/rpc/federation/handlers.rs b/core/archipelago/src/api/rpc/federation/handlers.rs index c314baa7..1ec9b2b6 100644 --- a/core/archipelago/src/api/rpc/federation/handlers.rs +++ b/core/archipelago/src/api/rpc/federation/handlers.rs @@ -1430,16 +1430,7 @@ impl RpcHandler { .and_then(|v| v.as_bool()) .unwrap_or(true); - let req = pending::find_by_id(&self.config.data_dir, id) - .await? - .ok_or_else(|| anyhow::anyhow!("Pending request not found: {}", id))?; - if !req.outbound || !matches!(req.state, pending::PendingState::Sent) { - anyhow::bail!( - "Can only cancel outbound requests in Sent state (outbound={}, state={:?})", - req.outbound, - req.state - ); - } + let req = pending::cancel_outbound(&self.config.data_dir, id).await?; if notify { let identity_dir = self.config.data_dir.join("identity"); @@ -1462,7 +1453,6 @@ impl RpcHandler { } } - pending::delete(&self.config.data_dir, id).await?; info!(id = %id, to = %req.from_nostr_pubkey, notified = notify, "Cancelled outbound peer request"); Ok(serde_json::json!({ "cancelled": true, "id": id, "notified": notify })) } diff --git a/core/archipelago/src/api/rpc/federation/handshake_tests.rs b/core/archipelago/src/api/rpc/federation/handshake_tests.rs index e3812192..03745c29 100644 --- a/core/archipelago/src/api/rpc/federation/handshake_tests.rs +++ b/core/archipelago/src/api/rpc/federation/handshake_tests.rs @@ -434,6 +434,14 @@ async fn npub_only_request_accepts_bound_reply_but_rejects_other_sender_and_forg .state, PendingState::Approved ); + assert_eq!( + pending::find_by_id(dir.path(), &row.id) + .await + .unwrap() + .unwrap() + .from_did, + remote_did + ); let duplicate = handler.handle_handshake_poll().await.unwrap(); assert!(duplicate["applied_invites"].as_array().unwrap().is_empty()); assert_eq!( diff --git a/core/archipelago/src/api/rpc/handshake.rs b/core/archipelago/src/api/rpc/handshake.rs index 2373256b..8e033786 100644 --- a/core/archipelago/src/api/rpc/handshake.rs +++ b/core/archipelago/src/api/rpc/handshake.rs @@ -345,6 +345,7 @@ impl RpcHandler { } } HandshakeMessage::PeerInvite { invite_code } => { + let _decision = pending::outbound_decision_guard().await; // Match against an outbound Sent request from this nostr // pubkey. If we never sent them anything, ignore — we // don't accept unsolicited invites over Nostr. @@ -431,10 +432,11 @@ impl RpcHandler { .await; } - pending::set_state( + pending::complete_outbound( &self.config.data_dir, &row_id, - PendingState::Approved, + &hs.from_nostr_pubkey, + &node.did, ) .await?; applied_invites.push(node.did); @@ -449,6 +451,7 @@ impl RpcHandler { } } HandshakeMessage::PeerReject { reason } => { + let _decision = pending::outbound_decision_guard().await; let pendings = pending::load_pending(&self.config.data_dir).await?; if let Some(row) = pendings.iter().find(|r| { r.outbound diff --git a/core/archipelago/src/federation/pending.rs b/core/archipelago/src/federation/pending.rs index e10d0963..4929d3b8 100644 --- a/core/archipelago/src/federation/pending.rs +++ b/core/archipelago/src/federation/pending.rs @@ -18,6 +18,13 @@ use tokio::io::AsyncWriteExt; static PENDING_STORE_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); +// Serialize acceptance with cancellation/replacement without blocking inbound +// requests while the remote join notification is in flight. +static OUTBOUND_DECISION_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); +pub async fn outbound_decision_guard() -> tokio::sync::MutexGuard<'static, ()> { + OUTBOUND_DECISION_LOCK.lock().await +} + const PENDING_FILE: &str = "federation/pending_requests.json"; const MAX_PENDING_PER_PUBKEY: usize = 5; const PENDING_EXPIRY_DAYS: i64 = 30; @@ -215,6 +222,7 @@ pub async fn insert_outbound( to_name: Option, message: Option, ) -> Result { + let _decision = outbound_decision_guard().await; let _guard = PENDING_STORE_LOCK.lock().await; let mut requests = load_pending(data_dir).await?; expire_stale(&mut requests); @@ -256,6 +264,47 @@ pub async fn set_state(data_dir: &Path, id: &str, state: PendingState) -> Result Ok(()) } +/// Persist the authenticated DID learned by an npub-only connection. Callers +/// hold outbound_decision_guard across validation and membership creation. +pub async fn complete_outbound(data_dir: &Path, id: &str, sender: &str, did: &str) -> Result<()> { + let _guard = PENDING_STORE_LOCK.lock().await; + let mut requests = load_pending(data_dir).await?; + let row = requests + .iter_mut() + .find(|row| row.id == id) + .context("Pending request not found")?; + anyhow::ensure!( + row.outbound + && row.state == PendingState::Sent + && row.from_nostr_pubkey == sender + && !did.is_empty() + && (row.from_did.is_empty() || row.from_did == did), + "Outbound request no longer matches authenticated reply" + ); + row.from_did = did.to_owned(); + row.state = PendingState::Approved; + save_pending(data_dir, &requests).await +} + +/// Withdraw locally before a slow relay notification. An acceptance already in +/// progress wins; a cancellation that wins prevents later invite application. +pub async fn cancel_outbound(data_dir: &Path, id: &str) -> Result { + let _decision = outbound_decision_guard().await; + let _guard = PENDING_STORE_LOCK.lock().await; + let mut requests = load_pending(data_dir).await?; + let index = requests + .iter() + .position(|row| row.id == id) + .context("Pending request not found")?; + anyhow::ensure!( + requests[index].outbound && requests[index].state == PendingState::Sent, + "Can only cancel outbound requests in Sent state" + ); + let row = requests.remove(index); + save_pending(data_dir, &requests).await?; + Ok(row) +} + /// Resolve a pending decision once; concurrent approval/rejection cannot /// overwrite each other after a slow network request. pub async fn decide(data_dir: &Path, id: &str, decision: PendingState) -> Result<()> { @@ -296,6 +345,59 @@ pub async fn delete(data_dir: &Path, id: &str) -> Result<()> { mod tests { use super::*; + #[tokio::test] + async fn outbound_completion_binds_identity_and_serializes_cancellation() { + let dir = tempfile::tempdir().unwrap(); + let row = insert_outbound( + dir.path(), + "sender".into(), + "npub".into(), + "".into(), + None, + None, + ) + .await + .unwrap(); + assert!( + complete_outbound(dir.path(), &row.id, "stranger", "did:key:peer") + .await + .is_err() + ); + let decision = outbound_decision_guard().await; + let cancellation = cancel_outbound(dir.path(), &row.id); + tokio::pin!(cancellation); + assert!( + tokio::time::timeout(std::time::Duration::from_millis(20), &mut cancellation) + .await + .is_err() + ); + complete_outbound(dir.path(), &row.id, "sender", "did:key:peer") + .await + .unwrap(); + drop(decision); + assert!(cancellation.await.is_err()); + let saved = find_by_id(dir.path(), &row.id).await.unwrap().unwrap(); + assert_eq!(saved.from_did, "did:key:peer"); + assert_eq!(saved.state, PendingState::Approved); + let second = insert_outbound( + dir.path(), + "other".into(), + "npub".into(), + "".into(), + None, + None, + ) + .await + .unwrap(); + cancel_outbound(dir.path(), &second.id).await.unwrap(); + assert!( + complete_outbound(dir.path(), &second.id, "other", "did:key:other") + .await + .is_err() + ); + assert!(find_by_id(dir.path(), &second.id).await.unwrap().is_none()); + } + #[tokio::test] async fn test_insert_inbound_then_dedupes() { let dir = tempfile::tempdir().unwrap();