Prevent unsigned Fleet collectors from claiming federation provenance
This commit is contained in:
@@ -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<serde_json::Value> {
|
||||
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<HashSet<String>> {
|
||||
// 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<serde_json::Value> {
|
||||
@@ -262,7 +299,7 @@ impl RpcHandler {
|
||||
&self,
|
||||
params: Option<serde_json::Value>,
|
||||
) -> Result<serde_json::Value> {
|
||||
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<serde_json::Value> = 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,10 +445,21 @@ impl RpcHandler {
|
||||
|
||||
match tokio::fs::read_to_string(entry.path()).await {
|
||||
Ok(data) => match serde_json::from_str::<serde_json::Value>(&data) {
|
||||
Ok(mut 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<serde_json::Value> = Vec::new();
|
||||
let mut entries = tokio::fs::read_dir(&fleet_dir)
|
||||
.await
|
||||
@@ -497,22 +569,33 @@ 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Sort by timestamp descending (most recent first)
|
||||
all_alerts.sort_by(|a, b| {
|
||||
@@ -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"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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};
|
||||
|
||||
@@ -71,6 +71,23 @@ pub async fn load_nodes(data_dir: &Path) -> Result<Vec<FederatedNode>> {
|
||||
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<std::collections::HashSet<String>> {
|
||||
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(
|
||||
|
||||
Reference in New Issue
Block a user