From 9268930c103952dd2f49fb692a3c3056aab1edd2 Mon Sep 17 00:00:00 2001 From: archipelago Date: Wed, 7 Oct 2026 20:34:53 -0400 Subject: [PATCH] Prevent unsigned Fleet collectors from claiming federation provenance --- core/archipelago/src/api/rpc/analytics.rs | 277 ++++++++++++++++++++- core/archipelago/src/federation/mod.rs | 8 +- core/archipelago/src/federation/storage.rs | 17 ++ 3 files changed, 290 insertions(+), 12 deletions(-) diff --git a/core/archipelago/src/api/rpc/analytics.rs b/core/archipelago/src/api/rpc/analytics.rs index 8e6a7e97..48be6579 100644 --- a/core/archipelago/src/api/rpc/analytics.rs +++ b/core/archipelago/src/api/rpc/analytics.rs @@ -5,10 +5,47 @@ use super::RpcHandler; use anyhow::{Context, Result}; +use std::collections::HashSet; use tracing::{debug, info, warn}; const ANALYTICS_FILE: &str = "analytics-config.json"; +/// Collector reports are unsigned claims, including historical on-disk reports. +/// Provenance is assigned here, never accepted from their JSON payload. +fn collector_report( + mut report: serde_json::Value, + expected_id: Option<&str>, +) -> Result { + let id = report + .get("node_id") + .and_then(|v| v.as_str()) + .context("Missing collector identity")?; + anyhow::ensure!( + !id.is_empty() + && id.len() <= 64 + && !id.contains('/') + && !id.contains('\\') + && !id.contains("..") + && !id.ends_with("-history") + && !id.chars().any(char::is_control), + "Invalid collector identity" + ); + anyhow::ensure!( + expected_id.is_none_or(|expected| expected == id), + "Collector identity differs from record" + ); + let object = report.as_object_mut().context("Invalid collector report")?; + object.insert("source".into(), serde_json::json!("collector")); + object.insert("trust_level".into(), serde_json::json!("unverified")); + object.insert("identity_authenticated".into(), serde_json::json!(false)); + Ok(report) +} + +async fn federation_ids(data_dir: &std::path::Path) -> Result> { + // A damaged authoritative store must not promote unsigned collector data. + crate::federation::load_node_identities(data_dir).await +} + impl RpcHandler { /// Check if analytics are enabled. pub(super) async fn handle_analytics_get_status(&self) -> Result { @@ -262,7 +299,7 @@ impl RpcHandler { &self, params: Option, ) -> Result { - let report = params.context("Missing telemetry report payload")?; + let report = collector_report(params.context("Missing telemetry report payload")?, None)?; // Validate required fields let node_id = report @@ -285,6 +322,13 @@ impl RpcHandler { .and_then(|v| v.as_str()) .context("Missing required field: reported_at")?; + anyhow::ensure!( + !federation_ids(&self.config.data_dir) + .await? + .contains(node_id), + "Unsigned collector report conflicts with a known node identity" + ); + let fleet_dir = self.config.data_dir.join("telemetry-fleet"); tokio::fs::create_dir_all(&fleet_dir) .await @@ -339,9 +383,9 @@ impl RpcHandler { let mut nodes: Vec = Vec::new(); // ── Trusted federation nodes ───────────────────────────────────── - let fed_nodes = crate::federation::load_nodes(&self.config.data_dir) - .await - .unwrap_or_default(); + let fed_nodes = crate::federation::load_nodes(&self.config.data_dir).await?; + let mut authoritative = federation_ids(&self.config.data_dir).await?; + authoritative.extend(fed_nodes.iter().map(|n| n.did.clone())); for n in fed_nodes .iter() .filter(|n| n.trust_level == crate::federation::TrustLevel::Trusted) @@ -378,6 +422,7 @@ impl RpcHandler { "reported_at": reported_at, "trust_level": n.trust_level.to_string(), "source": "federation", + "identity_authenticated": true, }); annotate_fleet_report(&mut report); nodes.push(report); @@ -400,9 +445,20 @@ impl RpcHandler { match tokio::fs::read_to_string(entry.path()).await { Ok(data) => match serde_json::from_str::(&data) { - Ok(mut report) => { - annotate_fleet_report(&mut report); - nodes.push(report); + Ok(report) => { + if let Ok(mut report) = + collector_report(report, name.strip_suffix(".json")) + { + let id = report["node_id"] + .as_str() + .expect("validated collector identity"); + // Include every relationship in this exclusion, so an unsigned + // claim cannot undo observer/untrusted Fleet exclusions either. + if !authoritative.contains(id) { + annotate_fleet_report(&mut report); + nodes.push(report); + } + } } Err(e) => { warn!(file = %name, error = %e, "Skipping corrupt fleet report"); @@ -449,6 +505,14 @@ impl RpcHandler { anyhow::bail!("Invalid node_id"); } + if federation_ids(&self.config.data_dir) + .await? + .contains(node_id) + { + return Ok(serde_json::json!({"node_id":node_id,"entries":[],"count":0, + "history_available":false,"reason":"collector_history_not_authoritative"})); + } + let history_path = self .config .data_dir @@ -460,8 +524,15 @@ impl RpcHandler { Err(_) => Vec::new(), }; + let history: Vec<_> = history + .into_iter() + .filter_map(|report| collector_report(report, Some(node_id)).ok()) + .collect(); Ok(serde_json::json!({ "node_id": node_id, + "history_available": true, + "source": "collector", + "identity_authenticated": false, "entries": history, "count": history.len(), })) @@ -475,6 +546,7 @@ impl RpcHandler { return Ok(serde_json::json!({ "alerts": [] })); } + let authoritative = federation_ids(&self.config.data_dir).await?; let mut all_alerts: Vec = Vec::new(); let mut entries = tokio::fs::read_dir(&fleet_dir) .await @@ -497,19 +569,30 @@ impl RpcHandler { Err(_) => continue, }; + let report = match collector_report(report, name.strip_suffix(".json")) { + Ok(report) => report, + Err(_) => continue, + }; let node_id = report .get("node_id") .and_then(|v| v.as_str()) .unwrap_or("unknown") .to_string(); + if authoritative.contains(&node_id) { + continue; + } + if let Some(alerts) = report.get("recent_alerts").and_then(|v| v.as_array()) { for alert in alerts { let mut enriched = alert.clone(); if let Some(obj) = enriched.as_object_mut() { obj.insert("node_id".to_string(), serde_json::json!(node_id)); + obj.insert("source".into(), serde_json::json!("collector")); + obj.insert("trust_level".into(), serde_json::json!("unverified")); + obj.insert("identity_authenticated".into(), serde_json::json!(false)); + all_alerts.push(enriched); } - all_alerts.push(enriched); } } } @@ -587,3 +670,181 @@ fn annotate_fleet_report(report: &mut serde_json::Value) { obj.insert("last_seen".to_string(), serde_json::json!(last_seen)); } } + +#[cfg(test)] +mod collector_provenance_tests { + use super::*; + use serde_json::json; + + #[test] + fn unsigned_and_legacy_reports_cannot_assign_their_own_trust() { + let report = collector_report( + json!({"node_id":"claimed-node", "source":"federation", + "trust_level":"trusted", "identity_authenticated":true, "cpu_pct":0}), + Some("claimed-node"), + ) + .unwrap(); + assert_eq!(report["source"], "collector"); + assert_eq!(report["trust_level"], "unverified"); + assert_eq!(report["identity_authenticated"], false); + assert_eq!(report["cpu_pct"], 0); + assert_eq!( + collector_report(report.clone(), Some("claimed-node")).unwrap(), + report + ); + } + + #[test] + fn collector_cannot_claim_another_record_identity_or_malformed_id() { + for value in [ + json!(null), + json!([]), + json!({"node_id":"other"}), + json!({"node_id":"../escape"}), + json!({"node_id":"bad\nidentity"}), + json!({"node_id":"claimed-node-history"}), + ] { + assert!(collector_report(value, Some("claimed-node")).is_err()); + } + } + + #[test] + fn collector_history_suffix_cannot_overwrite_another_nodes_history() { + assert!(collector_report(json!({"node_id":"node-history"}), Some("node-history")).is_err()); + assert!(collector_report(json!({"node_id":"node-history"}), None).is_err()); + } + + #[tokio::test] + async fn fleet_reads_do_not_promote_old_collector_spoofs_or_history() { + let data = tempfile::tempdir().unwrap(); + let peer = serde_json::from_value(json!({"did":"known-node", "pubkey":"00".repeat(32), + "onion":format!("{}.onion", "a".repeat(56)), "trust_level":"trusted", "added_at":"now"})).unwrap(); + let observer = + serde_json::from_value(json!({"did":"observer-node", "pubkey":"11".repeat(32), + "onion":format!("{}.onion", "a".repeat(56)), "trust_level":"observer", "added_at":"now"})) + .unwrap(); + crate::federation::save_nodes(data.path(), &[peer, observer]) + .await + .unwrap(); + let mut config = crate::config::Config::default(); + config.data_dir = data.path().to_path_buf(); + let handler = RpcHandler::new( + config, + std::sync::Arc::new(crate::state::StateManager::new()), + std::sync::Arc::new(crate::monitoring::MetricsStore::new()), + crate::session::SessionStore::new_for_tests(data.path().join("sessions.json")), + None, + None, + ) + .await + .unwrap(); + let root = data.path().join("telemetry-fleet"); + tokio::fs::create_dir_all(&root).await.unwrap(); + let spoof = json!({"node_id":"known-node", "source":"federation", "trust_level":"trusted", + "identity_authenticated":true, "version":"spoof", "reported_at":"2026-10-07T00:00:00Z", + "recent_alerts":[{"message":"spoofed alert", "source":"federation", "identity_authenticated":true}]}); + tokio::fs::write(root.join("known-node.json"), spoof.to_string()) + .await + .unwrap(); + tokio::fs::write( + root.join("known-node-history.json"), + json!([spoof.clone()]).to_string(), + ) + .await + .unwrap(); + assert!(handler + .handle_telemetry_ingest(Some(spoof.clone())) + .await + .is_err()); + let status = handler.handle_telemetry_fleet_status().await.unwrap(); + let nodes = status["nodes"].as_array().unwrap(); + assert_eq!(nodes.len(), 1); + assert_eq!(nodes[0]["source"], "federation"); + assert_eq!(nodes[0]["identity_authenticated"], true); + assert_ne!(nodes[0]["version"], "spoof"); + let history = handler + .handle_telemetry_fleet_node_history(Some(json!({"node_id":"known-node"}))) + .await + .unwrap(); + assert_eq!(history["history_available"], false); + assert_eq!(history["count"], 0); + assert!( + handler.handle_telemetry_fleet_alerts().await.unwrap()["alerts"] + .as_array() + .unwrap() + .is_empty() + ); + // Display dedup collapses the shared onion, but unsigned reports must + // still be excluded for BOTH raw identities, including the observer. + let ids = federation_ids(data.path()).await.unwrap(); + assert!(ids.contains("known-node") && ids.contains("observer-node")); + let mut observer_spoof = spoof.clone(); + observer_spoof["node_id"] = json!("observer-node"); + tokio::fs::write(root.join("observer-node.json"), observer_spoof.to_string()) + .await + .unwrap(); + assert!(handler + .handle_telemetry_ingest(Some(observer_spoof)) + .await + .is_err()); + let status = handler.handle_telemetry_fleet_status().await.unwrap(); + assert!(status["nodes"] + .as_array() + .unwrap() + .iter() + .all(|node| node["source"] == "federation")); + assert!( + handler.handle_telemetry_fleet_alerts().await.unwrap()["alerts"] + .as_array() + .unwrap() + .is_empty() + ); + let observer_history = handler + .handle_telemetry_fleet_node_history(Some(json!({"node_id":"observer-node"}))) + .await + .unwrap(); + assert_eq!(observer_history["history_available"], false); + let mut independent = spoof; + independent["node_id"] = json!("independent"); + handler + .handle_telemetry_ingest(Some(independent)) + .await + .unwrap(); + let status = handler.handle_telemetry_fleet_status().await.unwrap(); + let collector = status["nodes"] + .as_array() + .unwrap() + .iter() + .find(|v| v["node_id"] == "independent") + .unwrap(); + assert_eq!(collector["source"], "collector"); + assert_eq!(collector["identity_authenticated"], false); + let alerts = handler.handle_telemetry_fleet_alerts().await.unwrap(); + assert_eq!(alerts["alerts"][0]["source"], "collector"); + assert_eq!(alerts["alerts"][0]["identity_authenticated"], false); + // Corrupt authoritative state must fail closed, not turn collector + // claims into the fallback representation of known relationships. + let nodes_path = data.path().join("federation").join("nodes.json"); + tokio::fs::write(&nodes_path, b"{broken-authoritative-store") + .await + .unwrap(); + assert!(federation_ids(data.path()).await.is_err()); + assert!(handler.handle_telemetry_fleet_status().await.is_err()); + assert!(handler.handle_telemetry_fleet_alerts().await.is_err()); + assert!(handler + .handle_telemetry_fleet_node_history(Some(json!({"node_id":"independent"}))) + .await + .is_err()); + assert!(handler + .handle_telemetry_ingest(Some( + json!({"node_id":"new-claim", "version":"spoof", "reported_at":"now"}) + )) + .await + .is_err()); + assert!(!root.join("new-claim.json").exists()); + assert_eq!( + tokio::fs::read(&nodes_path).await.unwrap(), + b"{broken-authoritative-store" + ); + } +} diff --git a/core/archipelago/src/federation/mod.rs b/core/archipelago/src/federation/mod.rs index 59ceed05..69226e58 100644 --- a/core/archipelago/src/federation/mod.rs +++ b/core/archipelago/src/federation/mod.rs @@ -4,27 +4,27 @@ //! Trust is bilateral — both sides must agree. Federated nodes periodically //! sync container status, health metrics, and availability. +pub(crate) mod handshake_delivery; 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; +pub use invites::{accept_invite, create_invite, parse_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; // Crate-internal: peer-joined resolves the granted trust level by matching // the acceptor's invite_token against our stored outgoing invites. -pub(crate) use storage::load_invites; -pub(crate) use storage::{load_unique_payment_peer, load_unique_payment_peer_by_did}; #[allow(unused_imports)] pub use storage::{ add_node, fips_npub_for_onion, load_nodes, load_removed_dids, record_peer_transport, record_sync_result, remove_node, save_nodes, set_trust_level, update_node, }; +pub(crate) use storage::{load_invites, load_node_identities}; +pub(crate) use storage::{load_unique_payment_peer, load_unique_payment_peer_by_did}; pub use sync::{build_local_state, deploy_to_peer, sync_with_peer, sync_with_peer_by_did}; pub use types::{AppStatus, FederatedNode, NodeStateSnapshot, TrustLevel, TrustSource}; diff --git a/core/archipelago/src/federation/storage.rs b/core/archipelago/src/federation/storage.rs index cf65a5e0..7f6e1239 100644 --- a/core/archipelago/src/federation/storage.rs +++ b/core/archipelago/src/federation/storage.rs @@ -71,6 +71,23 @@ pub async fn load_nodes(data_dir: &Path) -> Result> { load_nodes_inner(data_dir).await } +/// Every persisted relationship identity, before display-only onion deduplication. +/// Callers excluding unsigned claims must not lose an observer/untrusted DID +/// merely because another record currently shares its route. +pub(crate) async fn load_node_identities( + data_dir: &Path, +) -> Result> { + let _guard = FEDERATION_STORE_LOCK.lock().await; + let content = match fs::read(data_dir.join(FEDERATION_DIR).join(NODES_FILE)).await { + Ok(content) => content, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Default::default()), + Err(error) => return Err(error).context("Could not read federation identities"), + }; + let file: NodesFile = serde_json::from_slice(&content) + .context("Invalid federation identities; unsigned claims cannot replace them")?; + Ok(file.nodes.into_iter().map(|node| node.did).collect()) +} + /// Resolve payment identity from persisted records before display deduplication /// can merge fields from different identities sharing an address. pub(crate) async fn load_unique_payment_peer(