Persist encrypted peer approval delivery and bind discovery invite identities

This commit is contained in:
archipelago
2026-10-05 22:27:43 -04:00
parent e97f958f45
commit 10d31ae13c
8 changed files with 636 additions and 77 deletions
@@ -575,9 +575,11 @@ impl RpcHandler {
// Reuse the minute collector instead of running expensive probes for // Reuse the minute collector instead of running expensive probes for
// every peer. An absent/stalled collector is unknown, never zero load. // every peer. An absent/stalled collector is unknown, never zero load.
let now = chrono::Utc::now().timestamp(); let now = chrono::Utc::now().timestamp();
let latest = self.metrics_store.latest().await.filter(|sample| { let latest = self
(0..=180).contains(&now.saturating_sub(sample.timestamp)) .metrics_store
}); .latest()
.await
.filter(|sample| (0..=180).contains(&now.saturating_sub(sample.timestamp)));
let metrics = latest.as_ref().map(|sample| &sample.system); let metrics = latest.as_ref().map(|sample| &sample.system);
let uptime = tokio::fs::read_to_string("/proc/uptime") let uptime = tokio::fs::read_to_string("/proc/uptime")
.await .await
@@ -587,11 +589,21 @@ impl RpcHandler {
.map(|v| v as u64); .map(|v| v as u64);
let state = federation::build_local_state( let state = federation::build_local_state(
apps, apps,
metrics.map(|m| m.cpu_percent).filter(|v| v.is_finite() && (0.0..=100.0).contains(v)), metrics
metrics.filter(|m| m.mem_total_bytes > 0).map(|m| m.mem_used_bytes), .map(|m| m.cpu_percent)
metrics.filter(|m| m.mem_total_bytes > 0).map(|m| m.mem_total_bytes), .filter(|v| v.is_finite() && (0.0..=100.0).contains(v)),
metrics.filter(|m| m.disk_total_bytes > 0).map(|m| m.disk_used_bytes), metrics
metrics.filter(|m| m.disk_total_bytes > 0).map(|m| m.disk_total_bytes), .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, uptime,
tor_active, tor_active,
server_name, 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<federation::handshake_delivery::ApprovalReply> {
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 (data, _) = self.state_manager.get_snapshot().await;
let local_did = identity::did_key_from_pubkey_hex(&data.server_info.pubkey)?; let local_did = identity::did_key_from_pubkey_hex(&data.server_info.pubkey)?;
let local_onion = data let local_onion = data
.server_info .server_info
.tor_address .tor_address
.clone() .as_deref()
.ok_or_else(|| anyhow::anyhow!("Tor address not available"))?; .ok_or_else(|| anyhow::anyhow!("Tor address not available"))?;
let local_pubkey = data.server_info.pubkey.clone(); let local_fips_npub = identity::fips_npub(&self.config.data_dir.join("identity"))
.await
// Generate a one-shot federation invite. The code embeds OUR onion .unwrap_or(None);
// 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 invite_code = federation::create_invite( let invite_code = federation::create_invite(
&self.config.data_dir, &self.config.data_dir,
&local_did, &local_did,
&local_onion, local_onion,
&local_pubkey, &data.server_info.pubkey,
local_fips_npub.as_deref(), local_fips_npub.as_deref(),
TrustLevel::Observer, TrustLevel::Observer,
) )
.await?; .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 async fn deliver_peer_approval_reply(
// when their `federation.peer-joined` callback arrives over Tor we &self,
// already trust their pubkey enough to accept the join. Their DID reply: &federation::handshake_delivery::ApprovalReply,
// and pubkey come from the request — we'll cross-check the pubkey ) -> Result<bool> {
// against the eventual peer-joined signature in the existing let Some(claimed) = federation::handshake_delivery::claim(
// verification path (handlers.rs line ~365). &self.config.data_dir,
if !req.from_did.is_empty() { &reply.request_id,
// We don't know the requester's onion or ed25519 pubkey yet — chrono::Utc::now().timestamp(),
// they'll send those in the federation.peer-joined callback )
// after they apply our invite. Until then we can't add a real .await?
// FederatedNode entry. We just store the pending row as else {
// Approved so the UI shows progress, and trust the existing return Ok(false);
// peer-joined handler to admit them as Observer when they call. };
// let result = nostr_handshake::send_peer_invite(
// Caveat: peer-joined currently hardcodes TrustLevel::Trusted. &self.config.data_dir.join("identity"),
// We override that below by demoting on success. &claimed.recipient,
debug!( &claimed.invite_code,
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,
&self.handshake_relays().await, &self.handshake_relays().await,
self.config.nostr_tor_proxy.as_deref(), 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?; /// Recover relay loss and legacy approvals without changing trust or
info!( /// resurrecting a node that the operator explicitly removed.
id = %id, pub(in crate::api::rpc) async fn retry_peer_approval_replies(&self) -> Result<()> {
from = %req.from_nostr_pubkey, let requests = pending::load_pending(&self.config.data_dir).await?;
"Approved peer request and shipped invite over NIP-44" let nodes = federation::load_nodes(&self.config.data_dir).await?;
); let removed = federation::load_removed_dids(&self.config.data_dir).await?;
Ok(serde_json::json!({ let cutoff = chrono::Utc::now() - chrono::Duration::days(30);
"approved": true, let mut sent = 0;
"id": id, 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, /// 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 { if notify {
let identity_dir = self.config.data_dir.join("identity"); let identity_dir = self.config.data_dir.join("identity");
let _ = nostr_handshake::send_peer_reject( let _ = nostr_handshake::send_peer_reject(
@@ -1348,7 +1405,6 @@ impl RpcHandler {
.await; .await;
} }
pending::set_state(&self.config.data_dir, id, pending::PendingState::Rejected).await?;
info!(id = %id, from = %req.from_nostr_pubkey, "Rejected peer request"); info!(id = %id, from = %req.from_nostr_pubkey, "Rejected peer request");
Ok(serde_json::json!({ "rejected": true, "id": id })) Ok(serde_json::json!({ "rejected": true, "id": id }))
} }
@@ -12,6 +12,7 @@ async fn managed_relay_receives_approval_rejection_and_cancellation() {
("reject", true), ("reject", true),
("cancel", true), ("cancel", true),
("approve", false), ("approve", false),
("retry", true),
] { ] {
let tmp = tempfile::tempdir().unwrap(); let tmp = tempfile::tempdir().unwrap();
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.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)) tokio::fs::write(identity_dir.join("nostr_secret"), "11".repeat(32))
.await .await
.unwrap(); .unwrap();
tokio::fs::write(identity_dir.join("node_key"), [0x33; 32])
.await
.unwrap();
let state = Arc::new(crate::state::StateManager::new()); let state = Arc::new(crate::state::StateManager::new());
state state
.mutate_data(|data| { .mutate_data(|data| {
@@ -90,7 +94,7 @@ async fn managed_relay_receives_approval_rejection_and_cancellation() {
tmp.path(), tmp.path(),
recipient.public_key().to_hex(), recipient.public_key().to_hex(),
String::new(), String::new(),
String::new(), crate::identity::did_key_from_pubkey_hex(&"44".repeat(32)).unwrap(),
None, None,
None, None,
) )
@@ -101,7 +105,7 @@ async fn managed_relay_receives_approval_rejection_and_cancellation() {
tmp.path(), tmp.path(),
recipient.public_key().to_hex(), recipient.public_key().to_hex(),
String::new(), String::new(),
String::new(), crate::identity::did_key_from_pubkey_hex(&"44".repeat(32)).unwrap(),
None, None,
None, None,
) )
@@ -109,18 +113,27 @@ async fn managed_relay_receives_approval_rejection_and_cancellation() {
.unwrap() .unwrap()
.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 params = Some(serde_json::json!({"id": row.id, "notify": true}));
let action = async { let action = async {
match operation { match operation {
"approve" => handler.handle_federation_approve_request(params).await, "approve" => handler.handle_federation_approve_request(params).await,
"reject" => handler.handle_federation_reject_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, _ => handler.handle_federation_cancel_request(params).await,
} }
}; };
let outcome = tokio::time::timeout(Duration::from_secs(20), action) let outcome = tokio::time::timeout(Duration::from_secs(20), action)
.await .await
.unwrap(); .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) let event = tokio::time::timeout(Duration::from_secs(5), relay)
.await .await
.unwrap() .unwrap()
@@ -130,14 +143,30 @@ async fn managed_relay_receives_approval_rejection_and_cancellation() {
nip44::decrypt(recipient.secret_key(), &event.pubkey, &event.content).unwrap(); nip44::decrypt(recipient.secret_key(), &event.pubkey, &event.content).unwrap();
let message: serde_json::Value = serde_json::from_str(&plaintext).unwrap(); let message: serde_json::Value = serde_json::from_str(&plaintext).unwrap();
let expected = match operation { let expected = match operation {
"approve" => "peer-invite", "approve" | "retry" => "peer-invite",
"reject" => "peer-reject", "reject" => "peer-reject",
_ => "peer-cancel", _ => "peer-cancel",
}; };
assert_eq!(message["type"], expected); assert_eq!(message["type"], expected);
let saved = pending::find_by_id(tmp.path(), &row.id).await.unwrap(); let saved = pending::find_by_id(tmp.path(), &row.id).await.unwrap();
if !accepted { 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()) assert!(crate::federation::load_nodes(tmp.path())
.await .await
.unwrap() .unwrap()
@@ -145,8 +174,21 @@ async fn managed_relay_receives_approval_rejection_and_cancellation() {
continue; continue;
} }
match operation { match operation {
"approve" => { "approve" | "retry" => {
assert_eq!(saved.unwrap().state, PendingState::Approved); 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 = let invite =
crate::federation::parse_invite(message["invite_code"].as_str().unwrap()) crate::federation::parse_invite(message["invite_code"].as_str().unwrap())
.unwrap(); .unwrap();
+16 -1
View File
@@ -276,6 +276,11 @@ impl RpcHandler {
} }
Err(e) => tracing::debug!("background handshake poll failed: {e:#}"), 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<serde_json::Value> { pub(super) async fn handle_handshake_poll(&self) -> Result<serde_json::Value> {
@@ -356,6 +361,16 @@ impl RpcHandler {
); );
continue; 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 row_id = row.id.clone();
let (data, _) = self.state_manager.get_snapshot().await; let (data, _) = self.state_manager.get_snapshot().await;
let local_did = let local_did =
@@ -373,7 +388,7 @@ impl RpcHandler {
let local_name = data.server_info.name.clone(); let local_name = data.server_info.name.clone();
match crate::federation::accept_invite( match crate::federation::accept_invite(
&self.config.data_dir, &self.config.data_dir,
invite_code, &scoped_invite,
&local_did, &local_did,
&local_onion, &local_onion,
&local_pubkey, &local_pubkey,
@@ -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<Vec<ApprovalReply>> {
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<Option<ApprovalReply>> {
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<ApprovalReply> {
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<Option<ApprovalReply>> {
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);
}
}
@@ -134,6 +134,32 @@ pub fn parse_invite(code: &str) -> Result<ParsedInvite> {
}) })
} }
/// 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<String> {
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. /// Accept an invite: parse code, verify the remote node, add to federation.
pub async fn accept_invite( pub async fn accept_invite(
data_dir: &Path, data_dir: &Path,
@@ -621,3 +647,35 @@ mod tests {
assert_eq!(nodes.len(), 1, "re-accept should not duplicate"); 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());
}
}
+2
View File
@@ -6,12 +6,14 @@
mod invites; mod invites;
pub mod pending; pub mod pending;
pub(crate) mod handshake_delivery;
mod storage; mod storage;
mod sync; mod sync;
mod types; mod types;
// Re-export all public items so `crate::federation::*` continues to work. // Re-export all public items so `crate::federation::*` continues to work.
pub use invites::{accept_invite, create_invite, parse_invite}; 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 // Crate-internal: used by the periodic federation auto-sync to re-assert
// membership to peers that don't list us back (asymmetry self-heal). // membership to peers that don't list us back (asymmetry self-heal).
pub(crate) use invites::notify_join; pub(crate) use invites::notify_join;
+193 -6
View File
@@ -14,6 +14,9 @@ use anyhow::{Context, Result};
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use std::path::Path; use std::path::Path;
use tokio::fs; 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 PENDING_FILE: &str = "federation/pending_requests.json";
const MAX_PENDING_PER_PUBKEY: usize = 5; const MAX_PENDING_PER_PUBKEY: usize = 5;
@@ -76,11 +79,12 @@ pub async fn load_pending(data_dir: &Path) -> Result<Vec<PendingPeerRequest>> {
let content = fs::read_to_string(&path) let content = fs::read_to_string(&path)
.await .await
.context("Failed to read pending requests file")?; .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) 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); let path = data_dir.join(PENDING_FILE);
if let Some(parent) = path.parent() { if let Some(parent) = path.parent() {
fs::create_dir_all(parent) fs::create_dir_all(parent)
@@ -92,9 +96,27 @@ pub async fn save_pending(data_dir: &Path, requests: &[PendingPeerRequest]) -> R
}; };
let content = let content =
serde_json::to_string_pretty(&file).context("Failed to serialize pending requests")?; serde_json::to_string_pretty(&file).context("Failed to serialize pending requests")?;
fs::write(&path, content) let parent = path.parent().context("Pending requests parent missing")?;
.await let temporary = parent.join(format!(".pending-{}.tmp", uuid::Uuid::new_v4()));
.context("Failed to write pending requests file")?; 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(()) Ok(())
} }
@@ -102,7 +124,10 @@ pub async fn save_pending(data_dir: &Path, requests: &[PendingPeerRequest]) -> R
fn expire_stale(requests: &mut Vec<PendingPeerRequest>) { fn expire_stale(requests: &mut Vec<PendingPeerRequest>) {
let cutoff = chrono::Utc::now() - chrono::Duration::days(PENDING_EXPIRY_DAYS); let cutoff = chrono::Utc::now() - chrono::Duration::days(PENDING_EXPIRY_DAYS);
for r in requests.iter_mut() { 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; continue;
} }
if let Ok(ts) = chrono::DateTime::parse_from_rfc3339(&r.received_at) { if let Ok(ts) = chrono::DateTime::parse_from_rfc3339(&r.received_at) {
@@ -131,6 +156,7 @@ pub async fn insert_inbound(
from_name: Option<String>, from_name: Option<String>,
message: Option<String>, message: Option<String>,
) -> Result<Option<PendingPeerRequest>> { ) -> Result<Option<PendingPeerRequest>> {
let _guard = PENDING_STORE_LOCK.lock().await;
let mut requests = load_pending(data_dir).await?; let mut requests = load_pending(data_dir).await?;
expire_stale(&mut requests); expire_stale(&mut requests);
@@ -189,6 +215,7 @@ pub async fn insert_outbound(
to_name: Option<String>, to_name: Option<String>,
message: Option<String>, message: Option<String>,
) -> Result<PendingPeerRequest> { ) -> Result<PendingPeerRequest> {
let _guard = PENDING_STORE_LOCK.lock().await;
let mut requests = load_pending(data_dir).await?; let mut requests = load_pending(data_dir).await?;
expire_stale(&mut requests); expire_stale(&mut requests);
requests.retain(|r| { requests.retain(|r| {
@@ -218,6 +245,7 @@ pub async fn find_by_id(data_dir: &Path, id: &str) -> Result<Option<PendingPeerR
} }
pub async fn set_state(data_dir: &Path, id: &str, state: PendingState) -> Result<()> { pub async fn set_state(data_dir: &Path, id: &str, state: PendingState) -> Result<()> {
let _guard = PENDING_STORE_LOCK.lock().await;
let mut requests = load_pending(data_dir).await?; let mut requests = load_pending(data_dir).await?;
if let Some(r) = requests.iter_mut().find(|r| r.id == id) { if let Some(r) = requests.iter_mut().find(|r| r.id == id) {
r.state = state; r.state = state;
@@ -228,10 +256,32 @@ pub async fn set_state(data_dir: &Path, id: &str, state: PendingState) -> Result
Ok(()) 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 /// Remove a pending request entirely. Used when the sender cancels an
/// outbound request they initiated and we want it gone (not just marked /// outbound request they initiated and we want it gone (not just marked
/// Rejected/Cancelled — those states fill up the UI audit trail). /// Rejected/Cancelled — those states fill up the UI audit trail).
pub async fn delete(data_dir: &Path, id: &str) -> Result<()> { 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 mut requests = load_pending(data_dir).await?;
let before = requests.len(); let before = requests.len();
requests.retain(|r| r.id != id); requests.retain(|r| r.id != id);
@@ -372,3 +422,140 @@ mod tests {
assert_eq!(reloaded.state, PendingState::Approved); 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::<std::collections::HashSet<_>>()
.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
);
}
}
+3 -3
View File
@@ -291,15 +291,15 @@ impl Server {
); );
// Background handshake poll: fetch inbound nostr peer requests every // 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 // button (requests used to sit on relays unseen — 2026-07-22). The
// handler's own discoverability gate makes this a no-op until the // handler's own discoverability gate makes this a no-op until the
// user opts in. // user opts in.
{ {
let rpc = api_handler.rpc_handler().clone(); let rpc = api_handler.rpc_handler().clone();
tokio::spawn(async move { tokio::spawn(async move {
let mut tick = tokio::time::interval(std::time::Duration::from_secs(300)); let mut tick = tokio::time::interval(std::time::Duration::from_secs(30));
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop { loop {
tick.tick().await; tick.tick().await;
rpc.background_handshake_poll().await; rpc.background_handshake_poll().await;