Bound federated sync and show progress feedback

This commit is contained in:
archipelago
2026-10-07 07:25:44 -04:00
parent cedfbb2b07
commit a64dd77efa
5 changed files with 121 additions and 31 deletions
@@ -7,6 +7,8 @@ use crate::mesh;
use crate::network::dwn_store::DwnStore;
use crate::nostr_handshake;
use anyhow::Result;
use futures_util::stream::{FuturesUnordered, StreamExt};
use std::sync::Arc;
use tracing::{debug, info, warn};
const FEDERATION_PROTOCOL: &str = "https://archipelago.dev/protocols/federation/v1";
@@ -460,37 +462,65 @@ impl RpcHandler {
let identity_dir = self.config.data_dir.join("identity");
let node_identity = identity::NodeIdentity::load_or_create(&identity_dir).await?;
let mut synced = 0u32;
let mut failed = 0u32;
// A dead Tor peer can take the bounded transport timeout. Serialising
// those waits made one bad peer block every healthy peer behind it.
// Keep concurrency bounded so a large federation cannot exhaust the
// node's sockets or overwhelm a peer, while allowing healthy peers to
// finish independently. Results are tagged and sorted below so the
// response remains stable for callers and tests.
const MAX_CONCURRENT_SYNCS: usize = 4;
let semaphore = Arc::new(tokio::sync::Semaphore::new(MAX_CONCURRENT_SYNCS));
let identity = Arc::new(node_identity);
let data_dir = self.config.data_dir.clone();
let mut pending = FuturesUnordered::new();
for (index, node) in nodes
.into_iter()
.filter(|node| node.trust_level != TrustLevel::Untrusted)
.enumerate()
{
let data_dir = data_dir.clone();
let local_did = local_did.clone();
let identity = identity.clone();
let semaphore = semaphore.clone();
pending.push(async move {
let permit = semaphore
.acquire_owned()
.await
.expect("sync semaphore lives for all pending syncs");
let result = federation::sync_with_peer(&data_dir, &node, &local_did, |bytes| {
identity.sign(bytes)
})
.await;
drop(permit);
(index, node.did, result)
});
}
let mut results = Vec::new();
for node in &nodes {
if node.trust_level == TrustLevel::Untrusted {
continue;
}
let did_clone = local_did.clone();
match federation::sync_with_peer(&self.config.data_dir, node, &did_clone, |bytes| {
node_identity.sign(bytes)
})
.await
{
Ok(state) => {
synced += 1;
results.push(serde_json::json!({
"did": node.did,
"status": "ok",
"apps": state.apps.len(),
}));
}
Err(e) => {
failed += 1;
results.push(serde_json::json!({
"did": node.did,
"status": "error",
"error": e.to_string(),
}));
}
while let Some((index, did, result)) = pending.next().await {
let row = match result {
Ok(state) => serde_json::json!({
"index": index,
"did": did,
"status": "ok",
"apps": state.apps.len(),
}),
Err(e) => serde_json::json!({
"index": index,
"did": did,
"status": "error",
"error": e.to_string(),
}),
};
results.push(row);
}
results.sort_by_key(|row| row["index"].as_u64().unwrap_or(u64::MAX));
let synced = results.iter().filter(|row| row["status"] == "ok").count() as u32;
let failed = results.len() as u32 - synced;
for row in &mut results {
if let Some(object) = row.as_object_mut() {
object.remove("index");
}
}