From 10d31ae13c74404f25ce2ebd39b30f0ab59eb77b Mon Sep 17 00:00:00 2001 From: archipelago Date: Mon, 5 Oct 2026 22:27:43 -0400 Subject: [PATCH] Persist encrypted peer approval delivery and bind discovery invite identities --- .../src/api/rpc/federation/handlers.rs | 178 ++++++++++------ .../src/api/rpc/federation/handshake_tests.rs | 54 ++++- core/archipelago/src/api/rpc/handshake.rs | 17 +- .../src/federation/handshake_delivery.rs | 199 ++++++++++++++++++ core/archipelago/src/federation/invites.rs | 58 +++++ core/archipelago/src/federation/mod.rs | 2 + core/archipelago/src/federation/pending.rs | 199 +++++++++++++++++- core/archipelago/src/server.rs | 6 +- 8 files changed, 636 insertions(+), 77 deletions(-) create mode 100644 core/archipelago/src/federation/handshake_delivery.rs diff --git a/core/archipelago/src/api/rpc/federation/handlers.rs b/core/archipelago/src/api/rpc/federation/handlers.rs index 36ccee04..c314baa7 100644 --- a/core/archipelago/src/api/rpc/federation/handlers.rs +++ b/core/archipelago/src/api/rpc/federation/handlers.rs @@ -575,9 +575,11 @@ impl RpcHandler { // Reuse the minute collector instead of running expensive probes for // every peer. An absent/stalled collector is unknown, never zero load. let now = chrono::Utc::now().timestamp(); - let latest = self.metrics_store.latest().await.filter(|sample| { - (0..=180).contains(&now.saturating_sub(sample.timestamp)) - }); + let latest = self + .metrics_store + .latest() + .await + .filter(|sample| (0..=180).contains(&now.saturating_sub(sample.timestamp))); let metrics = latest.as_ref().map(|sample| &sample.system); let uptime = tokio::fs::read_to_string("/proc/uptime") .await @@ -587,11 +589,21 @@ impl RpcHandler { .map(|v| v as u64); let state = federation::build_local_state( apps, - metrics.map(|m| m.cpu_percent).filter(|v| v.is_finite() && (0.0..=100.0).contains(v)), - metrics.filter(|m| m.mem_total_bytes > 0).map(|m| m.mem_used_bytes), - metrics.filter(|m| m.mem_total_bytes > 0).map(|m| m.mem_total_bytes), - metrics.filter(|m| m.disk_total_bytes > 0).map(|m| m.disk_used_bytes), - metrics.filter(|m| m.disk_total_bytes > 0).map(|m| m.disk_total_bytes), + metrics + .map(|m| m.cpu_percent) + .filter(|v| v.is_finite() && (0.0..=100.0).contains(v)), + metrics + .filter(|m| m.mem_total_bytes > 0) + .map(|m| m.mem_used_bytes), + metrics + .filter(|m| m.mem_total_bytes > 0) + .map(|m| m.mem_total_bytes), + metrics + .filter(|m| m.disk_total_bytes > 0) + .map(|m| m.disk_used_bytes), + metrics + .filter(|m| m.disk_total_bytes > 0) + .map(|m| m.disk_total_bytes), uptime, tor_active, server_name, @@ -1237,76 +1249,120 @@ impl RpcHandler { ); } + let reply = self.prepare_peer_approval_reply(&req).await?; + // Persist the operator decision before transport. A relay outage must + // not require another approval or lose the already-authorized reply. + pending::decide(&self.config.data_dir, id, pending::PendingState::Approved).await?; + let delivered = self.deliver_peer_approval_reply(&reply).await?; + Ok(serde_json::json!({ "approved": true, "id": id, "delivery_pending": !delivered })) + } + + async fn prepare_peer_approval_reply( + &self, + req: &pending::PendingPeerRequest, + ) -> Result { + use federation::handshake_delivery::{self, ApprovalReply}; + if let Some(reply) = handshake_delivery::find(&self.config.data_dir, &req.id).await? { + anyhow::ensure!( + reply.recipient == req.from_nostr_pubkey && reply.expected_did == req.from_did, + "Approval recipient changed" + ); + return Ok(reply); + } let (data, _) = self.state_manager.get_snapshot().await; let local_did = identity::did_key_from_pubkey_hex(&data.server_info.pubkey)?; let local_onion = data .server_info .tor_address - .clone() + .as_deref() .ok_or_else(|| anyhow::anyhow!("Tor address not available"))?; - let local_pubkey = data.server_info.pubkey.clone(); - - // Generate a one-shot federation invite. The code embeds OUR onion - // and OUR pubkey, but it leaves this box only inside the NIP-44 - // ciphertext below. - let identity_dir = self.config.data_dir.join("identity"); - let local_fips_npub = identity::fips_npub(&identity_dir).await.unwrap_or(None); - // Discovery/connection-request approvals admit the requester as - // Observer — the invite itself now carries that level, so both - // sides converge on Observer without post-hoc demotion. + let local_fips_npub = identity::fips_npub(&self.config.data_dir.join("identity")) + .await + .unwrap_or(None); let invite_code = federation::create_invite( &self.config.data_dir, &local_did, - &local_onion, - &local_pubkey, + local_onion, + &data.server_info.pubkey, local_fips_npub.as_deref(), TrustLevel::Observer, ) .await?; + handshake_delivery::stage( + &self.config.data_dir, + ApprovalReply { + request_id: req.id.clone(), + recipient: req.from_nostr_pubkey.clone(), + expected_did: req.from_did.clone(), + invite_code, + attempts: 0, + next_attempt: 0, + }, + ) + .await + } - // Pre-add the requester to OUR federation list as Observer so that - // when their `federation.peer-joined` callback arrives over Tor we - // already trust their pubkey enough to accept the join. Their DID - // and pubkey come from the request — we'll cross-check the pubkey - // against the eventual peer-joined signature in the existing - // verification path (handlers.rs line ~365). - if !req.from_did.is_empty() { - // We don't know the requester's onion or ed25519 pubkey yet — - // they'll send those in the federation.peer-joined callback - // after they apply our invite. Until then we can't add a real - // FederatedNode entry. We just store the pending row as - // Approved so the UI shows progress, and trust the existing - // peer-joined handler to admit them as Observer when they call. - // - // Caveat: peer-joined currently hardcodes TrustLevel::Trusted. - // We override that below by demoting on success. - debug!( - requester_did = %req.from_did, - "Approval pending — waiting for federation.peer-joined callback over Tor" - ); - } - - // Encrypt + send the invite over NIP-44 to the requester. - let identity_dir = self.config.data_dir.join("identity"); - nostr_handshake::send_peer_invite( - &identity_dir, - &req.from_nostr_pubkey, - &invite_code, + async fn deliver_peer_approval_reply( + &self, + reply: &federation::handshake_delivery::ApprovalReply, + ) -> Result { + let Some(claimed) = federation::handshake_delivery::claim( + &self.config.data_dir, + &reply.request_id, + chrono::Utc::now().timestamp(), + ) + .await? + else { + return Ok(false); + }; + let result = nostr_handshake::send_peer_invite( + &self.config.data_dir.join("identity"), + &claimed.recipient, + &claimed.invite_code, &self.handshake_relays().await, self.config.nostr_tor_proxy.as_deref(), ) - .await?; + .await; + if result.is_err() { + warn!(request_id = %reply.request_id, "Peer approval delivery deferred; durable retry scheduled"); + } + Ok(result.is_ok()) + } - pending::set_state(&self.config.data_dir, id, pending::PendingState::Approved).await?; - info!( - id = %id, - from = %req.from_nostr_pubkey, - "Approved peer request and shipped invite over NIP-44" - ); - Ok(serde_json::json!({ - "approved": true, - "id": id, - })) + /// Recover relay loss and legacy approvals without changing trust or + /// resurrecting a node that the operator explicitly removed. + pub(in crate::api::rpc) async fn retry_peer_approval_replies(&self) -> Result<()> { + let requests = pending::load_pending(&self.config.data_dir).await?; + let nodes = federation::load_nodes(&self.config.data_dir).await?; + let removed = federation::load_removed_dids(&self.config.data_dir).await?; + let cutoff = chrono::Utc::now() - chrono::Duration::days(30); + let mut sent = 0; + for req in requests { + if req.outbound || req.state != pending::PendingState::Approved { + federation::handshake_delivery::remove(&self.config.data_dir, &req.id).await?; + continue; + } + let expired = chrono::DateTime::parse_from_rfc3339(&req.received_at) + .map(|time| time < cutoff) + .unwrap_or(true); + if expired + || removed.contains(&req.from_did) + || nodes.iter().any(|node| node.did == req.from_did) + { + federation::handshake_delivery::remove(&self.config.data_dir, &req.id).await?; + continue; + } + if req.from_did.is_empty() || sent >= 4 { + continue; + } + let reply = self.prepare_peer_approval_reply(&req).await?; + if reply.next_attempt > chrono::Utc::now().timestamp() { + continue; + } + self.deliver_peer_approval_reply(&reply).await?; + sent += 1; + } + Ok(()) } /// federation.reject-request — drop a pending request and, if requested, @@ -1336,6 +1392,7 @@ impl RpcHandler { ); } + pending::decide(&self.config.data_dir, id, pending::PendingState::Rejected).await?; if notify { let identity_dir = self.config.data_dir.join("identity"); let _ = nostr_handshake::send_peer_reject( @@ -1348,7 +1405,6 @@ impl RpcHandler { .await; } - pending::set_state(&self.config.data_dir, id, pending::PendingState::Rejected).await?; info!(id = %id, from = %req.from_nostr_pubkey, "Rejected peer request"); Ok(serde_json::json!({ "rejected": true, "id": id })) } diff --git a/core/archipelago/src/api/rpc/federation/handshake_tests.rs b/core/archipelago/src/api/rpc/federation/handshake_tests.rs index 35b58e4e..c4407d2b 100644 --- a/core/archipelago/src/api/rpc/federation/handshake_tests.rs +++ b/core/archipelago/src/api/rpc/federation/handshake_tests.rs @@ -12,6 +12,7 @@ async fn managed_relay_receives_approval_rejection_and_cancellation() { ("reject", true), ("cancel", true), ("approve", false), + ("retry", true), ] { let tmp = tempfile::tempdir().unwrap(); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); @@ -68,6 +69,9 @@ async fn managed_relay_receives_approval_rejection_and_cancellation() { tokio::fs::write(identity_dir.join("nostr_secret"), "11".repeat(32)) .await .unwrap(); + tokio::fs::write(identity_dir.join("node_key"), [0x33; 32]) + .await + .unwrap(); let state = Arc::new(crate::state::StateManager::new()); state .mutate_data(|data| { @@ -90,7 +94,7 @@ async fn managed_relay_receives_approval_rejection_and_cancellation() { tmp.path(), recipient.public_key().to_hex(), String::new(), - String::new(), + crate::identity::did_key_from_pubkey_hex(&"44".repeat(32)).unwrap(), None, None, ) @@ -101,7 +105,7 @@ async fn managed_relay_receives_approval_rejection_and_cancellation() { tmp.path(), recipient.public_key().to_hex(), String::new(), - String::new(), + crate::identity::did_key_from_pubkey_hex(&"44".repeat(32)).unwrap(), None, None, ) @@ -109,18 +113,27 @@ async fn managed_relay_receives_approval_rejection_and_cancellation() { .unwrap() .unwrap() }; + if operation == "retry" { + pending::decide(tmp.path(), &row.id, PendingState::Approved) + .await + .unwrap(); + } let params = Some(serde_json::json!({"id": row.id, "notify": true})); let action = async { match operation { "approve" => handler.handle_federation_approve_request(params).await, "reject" => handler.handle_federation_reject_request(params).await, + "retry" => handler + .retry_peer_approval_replies() + .await + .map(|_| serde_json::json!({"ok": true})), _ => handler.handle_federation_cancel_request(params).await, } }; let outcome = tokio::time::timeout(Duration::from_secs(20), action) .await .unwrap(); - assert_eq!(outcome.is_ok(), accepted); + assert_eq!(outcome.is_ok(), accepted || operation == "approve"); let event = tokio::time::timeout(Duration::from_secs(5), relay) .await .unwrap() @@ -130,14 +143,30 @@ async fn managed_relay_receives_approval_rejection_and_cancellation() { nip44::decrypt(recipient.secret_key(), &event.pubkey, &event.content).unwrap(); let message: serde_json::Value = serde_json::from_str(&plaintext).unwrap(); let expected = match operation { - "approve" => "peer-invite", + "approve" | "retry" => "peer-invite", "reject" => "peer-reject", _ => "peer-cancel", }; assert_eq!(message["type"], expected); let saved = pending::find_by_id(tmp.path(), &row.id).await.unwrap(); if !accepted { - assert_eq!(saved.unwrap().state, PendingState::Pending); + assert_eq!( + saved.unwrap().state, + if operation == "approve" { + PendingState::Approved + } else { + PendingState::Pending + } + ); + if operation == "approve" { + assert_eq!(outcome.unwrap()["delivery_pending"], true); + let durable = crate::federation::handshake_delivery::find(tmp.path(), &row.id) + .await + .unwrap() + .unwrap(); + assert_eq!(durable.recipient, recipient.public_key().to_hex()); + assert_eq!(durable.attempts, 1); + } assert!(crate::federation::load_nodes(tmp.path()) .await .unwrap() @@ -145,8 +174,21 @@ async fn managed_relay_receives_approval_rejection_and_cancellation() { continue; } match operation { - "approve" => { + "approve" | "retry" => { assert_eq!(saved.unwrap().state, PendingState::Approved); + let first = crate::federation::handshake_delivery::find(tmp.path(), &row.id) + .await + .unwrap() + .unwrap(); + assert_eq!(first.attempts, 1); + // Concurrent/background polling honors the persisted backoff. + handler.retry_peer_approval_replies().await.unwrap(); + let second = crate::federation::handshake_delivery::find(tmp.path(), &row.id) + .await + .unwrap() + .unwrap(); + assert_eq!(second.attempts, 1); + assert_eq!(first.invite_code, second.invite_code); let invite = crate::federation::parse_invite(message["invite_code"].as_str().unwrap()) .unwrap(); diff --git a/core/archipelago/src/api/rpc/handshake.rs b/core/archipelago/src/api/rpc/handshake.rs index 6399f7a1..2373256b 100644 --- a/core/archipelago/src/api/rpc/handshake.rs +++ b/core/archipelago/src/api/rpc/handshake.rs @@ -276,6 +276,11 @@ impl RpcHandler { } Err(e) => tracing::debug!("background handshake poll failed: {e:#}"), } + if load_discovery_state(&self.config.data_dir).await.enabled { + if let Err(error) = self.retry_peer_approval_replies().await { + tracing::warn!("Peer approval retry could not complete: {error:#}"); + } + } } pub(super) async fn handle_handshake_poll(&self) -> Result { @@ -356,6 +361,16 @@ impl RpcHandler { ); continue; }; + let scoped_invite = match crate::federation::restrict_discovery_invite( + invite_code, + &row.from_did, + ) { + Ok(code) => code, + Err(_) => { + tracing::warn!("Rejected peer invite with mismatched identity"); + continue; + } + }; let row_id = row.id.clone(); let (data, _) = self.state_manager.get_snapshot().await; let local_did = @@ -373,7 +388,7 @@ impl RpcHandler { let local_name = data.server_info.name.clone(); match crate::federation::accept_invite( &self.config.data_dir, - invite_code, + &scoped_invite, &local_did, &local_onion, &local_pubkey, diff --git a/core/archipelago/src/federation/handshake_delivery.rs b/core/archipelago/src/federation/handshake_delivery.rs new file mode 100644 index 00000000..9be1d309 --- /dev/null +++ b/core/archipelago/src/federation/handshake_delivery.rs @@ -0,0 +1,199 @@ +//! Durable, node-encrypted approval replies. Relay acknowledgement is not peer +//! acceptance: keep retrying the same invite until reciprocal membership exists. +use anyhow::{Context, Result}; +use serde::{Deserialize, Serialize}; +use std::path::Path; +use tokio::{fs, io::AsyncWriteExt}; + +const FILE: &str = "federation/handshake-delivery.enc"; +const DOMAIN: &[u8] = b"archipelago-handshake-delivery-v1"; +static LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); + +#[derive(Clone, Serialize, Deserialize)] +pub(crate) struct ApprovalReply { + pub request_id: String, + pub recipient: String, + pub expected_did: String, + pub invite_code: String, + pub attempts: u32, + pub next_attempt: i64, +} + +async fn load(data_dir: &Path) -> Result> { + let bytes = match fs::read(data_dir.join(FILE)).await { + Ok(bytes) => bytes, + Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()), + Err(e) => return Err(e.into()), + }; + let key = crate::storage_crypto::derive_key(data_dir, DOMAIN).await?; + let plaintext = crate::storage_crypto::open(&bytes, &key)?; + serde_json::from_slice(&plaintext) + .context("Invalid handshake delivery store; preserved for recovery") +} + +async fn save(data_dir: &Path, entries: &[ApprovalReply]) -> Result<()> { + let key = crate::storage_crypto::derive_key(data_dir, DOMAIN).await?; + let bytes = crate::storage_crypto::seal(&serde_json::to_vec(entries)?, &key)?; + let path = data_dir.join(FILE); + let parent = path.parent().context("Delivery parent missing")?; + fs::create_dir_all(parent).await?; + let temporary = parent.join(format!(".delivery-{}.tmp", uuid::Uuid::new_v4())); + let result = async { + let mut file = fs::OpenOptions::new() + .create_new(true) + .write(true) + .mode(0o600) + .open(&temporary) + .await?; + file.write_all(&bytes).await?; + file.sync_all().await?; + drop(file); + fs::rename(&temporary, &path).await?; + fs::File::open(parent).await?.sync_all().await?; + Ok::<_, anyhow::Error>(()) + } + .await; + if result.is_err() { + let _ = fs::remove_file(temporary).await; + } + result +} + +pub(crate) async fn find(data_dir: &Path, request_id: &str) -> Result> { + let _guard = LOCK.lock().await; + Ok(load(data_dir) + .await? + .into_iter() + .find(|entry| entry.request_id == request_id)) +} + +pub(crate) async fn stage(data_dir: &Path, reply: ApprovalReply) -> Result { + let _guard = LOCK.lock().await; + let mut entries = load(data_dir).await?; + if let Some(existing) = entries + .iter() + .find(|entry| entry.request_id == reply.request_id) + { + anyhow::ensure!( + existing.recipient == reply.recipient && existing.expected_did == reply.expected_did, + "Approval recipient changed; refusing delivery" + ); + return Ok(existing.clone()); + } + anyhow::ensure!(entries.len() < 1024, "Handshake delivery queue is full"); + entries.push(reply.clone()); + save(data_dir, &entries).await?; + Ok(reply) +} + +/// Claim before sending, including failed sends. Concurrent polls cannot create +/// retry storms; a crash after this write delays but never loses the reply. +pub(crate) async fn claim( + data_dir: &Path, + request_id: &str, + now: i64, +) -> Result> { + let _guard = LOCK.lock().await; + let mut entries = load(data_dir).await?; + let Some(entry) = entries + .iter_mut() + .find(|entry| entry.request_id == request_id) + else { + return Ok(None); + }; + if entry.next_attempt > now { + return Ok(None); + } + entry.attempts = entry.attempts.saturating_add(1); + let delay = 30_i64 + .saturating_mul(1_i64 << entry.attempts.min(7)) + .min(3600); + entry.next_attempt = now.saturating_add(delay); + let claimed = entry.clone(); + save(data_dir, &entries).await?; + Ok(Some(claimed)) +} + +pub(crate) async fn remove(data_dir: &Path, request_id: &str) -> Result<()> { + let _guard = LOCK.lock().await; + let mut entries = load(data_dir).await?; + let before = entries.len(); + entries.retain(|entry| entry.request_id != request_id); + if entries.len() != before { + save(data_dir, &entries).await?; + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + async fn fixture() -> tempfile::TempDir { + let dir = tempfile::tempdir().unwrap(); + fs::create_dir_all(dir.path().join("identity")) + .await + .unwrap(); + fs::write(dir.path().join("identity/node_key"), [7; 32]) + .await + .unwrap(); + dir + } + fn reply() -> ApprovalReply { + ApprovalReply { + request_id: "request-1".into(), + recipient: "recipient".into(), + expected_did: "did:key:peer".into(), + invite_code: "secret-invite".into(), + attempts: 0, + next_attempt: 0, + } + } + #[tokio::test] + async fn encrypted_reply_survives_reload_and_retry_claim_is_exclusive() { + let dir = fixture().await; + stage(dir.path(), reply()).await.unwrap(); + let raw = fs::read(dir.path().join(FILE)).await.unwrap(); + assert!(!raw.windows(13).any(|bytes| bytes == b"secret-invite")); + assert_eq!( + find(dir.path(), "request-1") + .await + .unwrap() + .unwrap() + .invite_code, + "secret-invite" + ); + let (a, b) = tokio::join!( + claim(dir.path(), "request-1", 100), + claim(dir.path(), "request-1", 100) + ); + assert_eq!( + usize::from(a.unwrap().is_some()) + usize::from(b.unwrap().is_some()), + 1 + ); + assert!(claim(dir.path(), "request-1", 159).await.unwrap().is_none()); + assert_eq!( + claim(dir.path(), "request-1", 160) + .await + .unwrap() + .unwrap() + .attempts, + 2 + ); + remove(dir.path(), "request-1").await.unwrap(); + assert!(find(dir.path(), "request-1").await.unwrap().is_none()); + } + #[tokio::test] + async fn corrupt_store_is_not_overwritten_and_recipient_cannot_change() { + let dir = fixture().await; + stage(dir.path(), reply()).await.unwrap(); + let mut other = reply(); + other.recipient = "different-recipient".into(); + assert!(stage(dir.path(), other).await.is_err()); + let path = dir.path().join(FILE); + let mut raw = fs::read(&path).await.unwrap(); + raw[20] ^= 1; + fs::write(&path, &raw).await.unwrap(); + assert!(stage(dir.path(), reply()).await.is_err()); + assert_eq!(fs::read(path).await.unwrap(), raw); + } +} diff --git a/core/archipelago/src/federation/invites.rs b/core/archipelago/src/federation/invites.rs index 63f4bf71..97b7382d 100644 --- a/core/archipelago/src/federation/invites.rs +++ b/core/archipelago/src/federation/invites.rs @@ -134,6 +134,32 @@ pub fn parse_invite(code: &str) -> Result { }) } +/// Bind a Nostr-discovery reply to the node the operator requested, and cap +/// its grant before any local node entry or callback is written. Legacy invites +/// default to Trusted, which must never transiently authorize discovery peers. +pub(crate) fn restrict_discovery_invite(code: &str, expected_did: &str) -> Result { + use base64::Engine; + let parsed = parse_invite(code)?; + anyhow::ensure!( + !expected_did.is_empty() && parsed.did == expected_did, + "Peer invite does not match the requested node" + ); + anyhow::ensure!( + crate::identity::did_key_from_pubkey_hex(&parsed.pubkey)? == parsed.did, + "Peer invite DID does not match its identity key" + ); + let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD.decode( + code.strip_prefix("fed1:") + .context("Invalid invite prefix")?, + )?; + let mut payload: serde_json::Value = serde_json::from_slice(&bytes)?; + payload["trust"] = serde_json::json!("observer"); + Ok(format!( + "fed1:{}", + base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(serde_json::to_vec(&payload)?) + )) +} + /// Accept an invite: parse code, verify the remote node, add to federation. pub async fn accept_invite( data_dir: &Path, @@ -621,3 +647,35 @@ mod tests { assert_eq!(nodes.len(), 1, "re-accept should not duplicate"); } } + +#[cfg(test)] +mod discovery_invite_scope_tests { + use super::*; + use base64::Engine; + #[test] + fn discovery_reply_binds_identity_and_caps_legacy_trust_before_acceptance() { + let key = "33".repeat(32); + let did = crate::identity::did_key_from_pubkey_hex(&key).unwrap(); + let payload = + serde_json::json!({"did":did,"pubkey":key,"onion":"test.onion","token":"test-token"}); + let code = format!( + "fed1:{}", + base64::engine::general_purpose::URL_SAFE_NO_PAD + .encode(serde_json::to_vec(&payload).unwrap()) + ); + let restricted = restrict_discovery_invite(&code, &did).unwrap(); + let parsed = parse_invite(&restricted).unwrap(); + assert_eq!(parsed.trust_level, TrustLevel::Observer); + assert_eq!(parsed.token, "test-token"); + assert!(restrict_discovery_invite(&code, "did:key:someone-else").is_err()); + assert!(restrict_discovery_invite(&code, "").is_err()); + let mut forged = payload; + forged["pubkey"] = serde_json::json!("44".repeat(32)); + let forged = format!( + "fed1:{}", + base64::engine::general_purpose::URL_SAFE_NO_PAD + .encode(serde_json::to_vec(&forged).unwrap()) + ); + assert!(restrict_discovery_invite(&forged, &did).is_err()); + } +} diff --git a/core/archipelago/src/federation/mod.rs b/core/archipelago/src/federation/mod.rs index 44b781e2..c762a457 100644 --- a/core/archipelago/src/federation/mod.rs +++ b/core/archipelago/src/federation/mod.rs @@ -6,12 +6,14 @@ mod invites; pub mod pending; +pub(crate) mod handshake_delivery; mod storage; mod sync; mod types; // Re-export all public items so `crate::federation::*` continues to work. pub use invites::{accept_invite, create_invite, parse_invite}; +pub(crate) use invites::restrict_discovery_invite; // Crate-internal: used by the periodic federation auto-sync to re-assert // membership to peers that don't list us back (asymmetry self-heal). pub(crate) use invites::notify_join; diff --git a/core/archipelago/src/federation/pending.rs b/core/archipelago/src/federation/pending.rs index dc0d0e96..e10d0963 100644 --- a/core/archipelago/src/federation/pending.rs +++ b/core/archipelago/src/federation/pending.rs @@ -14,6 +14,9 @@ use anyhow::{Context, Result}; use serde::{Deserialize, Serialize}; use std::path::Path; use tokio::fs; +use tokio::io::AsyncWriteExt; + +static PENDING_STORE_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); const PENDING_FILE: &str = "federation/pending_requests.json"; const MAX_PENDING_PER_PUBKEY: usize = 5; @@ -76,11 +79,12 @@ pub async fn load_pending(data_dir: &Path) -> Result> { let content = fs::read_to_string(&path) .await .context("Failed to read pending requests file")?; - let file: PendingRequestsFile = serde_json::from_str(&content).unwrap_or_default(); + let file: PendingRequestsFile = serde_json::from_str(&content) + .context("Invalid pending requests file; preserving existing data")?; Ok(file.requests) } -pub async fn save_pending(data_dir: &Path, requests: &[PendingPeerRequest]) -> Result<()> { +async fn save_pending(data_dir: &Path, requests: &[PendingPeerRequest]) -> Result<()> { let path = data_dir.join(PENDING_FILE); if let Some(parent) = path.parent() { fs::create_dir_all(parent) @@ -92,9 +96,27 @@ pub async fn save_pending(data_dir: &Path, requests: &[PendingPeerRequest]) -> R }; let content = serde_json::to_string_pretty(&file).context("Failed to serialize pending requests")?; - fs::write(&path, content) - .await - .context("Failed to write pending requests file")?; + let parent = path.parent().context("Pending requests parent missing")?; + let temporary = parent.join(format!(".pending-{}.tmp", uuid::Uuid::new_v4())); + let result = async { + let mut file = fs::OpenOptions::new() + .write(true) + .create_new(true) + .mode(0o600) + .open(&temporary) + .await?; + file.write_all(content.as_bytes()).await?; + file.sync_all().await?; + drop(file); + fs::rename(&temporary, &path).await?; + fs::File::open(parent).await?.sync_all().await?; + Ok::<_, anyhow::Error>(()) + } + .await; + if result.is_err() { + let _ = fs::remove_file(&temporary).await; + } + result.context("Failed to atomically save pending requests")?; Ok(()) } @@ -102,7 +124,10 @@ pub async fn save_pending(data_dir: &Path, requests: &[PendingPeerRequest]) -> R fn expire_stale(requests: &mut Vec) { let cutoff = chrono::Utc::now() - chrono::Duration::days(PENDING_EXPIRY_DAYS); for r in requests.iter_mut() { - if !matches!(r.state, PendingState::Pending | PendingState::Sent) { + if !matches!( + r.state, + PendingState::Pending | PendingState::Sent | PendingState::Approved + ) { continue; } if let Ok(ts) = chrono::DateTime::parse_from_rfc3339(&r.received_at) { @@ -131,6 +156,7 @@ pub async fn insert_inbound( from_name: Option, message: Option, ) -> Result> { + let _guard = PENDING_STORE_LOCK.lock().await; let mut requests = load_pending(data_dir).await?; expire_stale(&mut requests); @@ -189,6 +215,7 @@ pub async fn insert_outbound( to_name: Option, message: Option, ) -> Result { + let _guard = PENDING_STORE_LOCK.lock().await; let mut requests = load_pending(data_dir).await?; expire_stale(&mut requests); requests.retain(|r| { @@ -218,6 +245,7 @@ pub async fn find_by_id(data_dir: &Path, id: &str) -> Result Result<()> { + let _guard = PENDING_STORE_LOCK.lock().await; let mut requests = load_pending(data_dir).await?; if let Some(r) = requests.iter_mut().find(|r| r.id == id) { r.state = state; @@ -228,10 +256,32 @@ pub async fn set_state(data_dir: &Path, id: &str, state: PendingState) -> Result Ok(()) } +/// 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<()> { + anyhow::ensure!( + matches!(decision, PendingState::Approved | PendingState::Rejected), + "Invalid pending decision" + ); + 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::Pending, + "Request has already been decided" + ); + row.state = decision; + save_pending(data_dir, &requests).await +} + /// Remove a pending request entirely. Used when the sender cancels an /// outbound request they initiated and we want it gone (not just marked /// Rejected/Cancelled — those states fill up the UI audit trail). pub async fn delete(data_dir: &Path, id: &str) -> Result<()> { + let _guard = PENDING_STORE_LOCK.lock().await; let mut requests = load_pending(data_dir).await?; let before = requests.len(); requests.retain(|r| r.id != id); @@ -372,3 +422,140 @@ mod tests { assert_eq!(reloaded.state, PendingState::Approved); } } + +#[cfg(test)] +mod persistence_regressions { + use super::*; + + #[tokio::test] + async fn concurrent_requests_are_not_lost() { + let dir = tempfile::tempdir().unwrap(); + let mut tasks = Vec::new(); + for i in 0..24 { + let path = dir.path().to_path_buf(); + tasks.push(tokio::spawn(async move { + insert_inbound( + &path, + format!("key-{i}"), + format!("npub-{i}"), + format!("did:key:{i}"), + None, + None, + ) + .await + .unwrap() + })); + } + for task in tasks { + assert!(task.await.unwrap().is_some()); + } + let rows = load_pending(dir.path()).await.unwrap(); + assert_eq!(rows.len(), 24); + assert_eq!( + rows.iter() + .map(|r| &r.id) + .collect::>() + .len(), + 24 + ); + } + + #[tokio::test] + async fn malformed_store_is_preserved_instead_of_replaced_with_one_request() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join(PENDING_FILE); + fs::create_dir_all(path.parent().unwrap()).await.unwrap(); + let damaged = b"{incomplete existing requests"; + fs::write(&path, damaged).await.unwrap(); + assert!(insert_inbound( + dir.path(), + "key".into(), + "npub".into(), + "did:key:test".into(), + None, + None + ) + .await + .is_err()); + assert_eq!(fs::read(path).await.unwrap(), damaged); + } +} + +#[cfg(test)] +mod decision_regressions { + use super::*; + #[tokio::test] + async fn only_one_concurrent_operator_decision_wins() { + let dir = tempfile::tempdir().unwrap(); + let row = insert_inbound( + dir.path(), + "key".into(), + "npub".into(), + "did:key:peer".into(), + None, + None, + ) + .await + .unwrap() + .unwrap(); + let (approve, reject) = tokio::join!( + decide(dir.path(), &row.id, PendingState::Approved), + decide(dir.path(), &row.id, PendingState::Rejected) + ); + assert_eq!( + usize::from(approve.is_ok()) + usize::from(reject.is_ok()), + 1 + ); + let saved = find_by_id(dir.path(), &row.id).await.unwrap().unwrap(); + assert_eq!( + saved.state, + if approve.is_ok() { + PendingState::Approved + } else { + PendingState::Rejected + } + ); + } + #[tokio::test] + async fn an_expired_approval_does_not_block_a_new_request_forever() { + let dir = tempfile::tempdir().unwrap(); + let mut row = insert_inbound( + dir.path(), + "key".into(), + "npub".into(), + "did:key:peer".into(), + None, + None, + ) + .await + .unwrap() + .unwrap(); + row.state = PendingState::Approved; + row.received_at = (chrono::Utc::now() - chrono::Duration::days(31)).to_rfc3339(); + save_pending(dir.path(), &[row]).await.unwrap(); + let renewed = insert_inbound( + dir.path(), + "key".into(), + "npub".into(), + "did:key:peer".into(), + None, + None, + ) + .await + .unwrap(); + assert!(renewed.is_some()); + let rows = load_pending(dir.path()).await.unwrap(); + assert_eq!( + rows.iter() + .filter(|r| r.state == PendingState::Expired) + .count(), + 1 + ); + assert_eq!( + rows.iter() + .filter(|r| r.state == PendingState::Pending) + .count(), + 1 + ); + } +} diff --git a/core/archipelago/src/server.rs b/core/archipelago/src/server.rs index 15632e48..f720ba2b 100644 --- a/core/archipelago/src/server.rs +++ b/core/archipelago/src/server.rs @@ -291,15 +291,15 @@ impl Server { ); // Background handshake poll: fetch inbound nostr peer requests every - // 5 minutes instead of only when a user presses the Federation Poll + // 30 seconds instead of only when a user presses the Federation Poll // button (requests used to sit on relays unseen — 2026-07-22). The // handler's own discoverability gate makes this a no-op until the // user opts in. { let rpc = api_handler.rpc_handler().clone(); tokio::spawn(async move { - let mut tick = tokio::time::interval(std::time::Duration::from_secs(300)); - tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + let mut tick = tokio::time::interval(std::time::Duration::from_secs(30)); + tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); loop { tick.tick().await; rpc.background_handshake_poll().await;