From a64dd77efa342efb552802816b7a499a1147e3ca Mon Sep 17 00:00:00 2001 From: archipelago Date: Wed, 7 Oct 2026 07:25:44 -0400 Subject: [PATCH] Bound federated sync and show progress feedback --- .../src/api/rpc/federation/handlers.rs | 90 ++++++++++++------- docs/post-1.9.0-work-backlog.md | 14 +++ neode-ui/src/views/Federation.vue | 2 + .../src/views/federation/QuickActions.vue | 8 +- .../federation/__tests__/QuickActions.test.ts | 38 ++++++++ 5 files changed, 121 insertions(+), 31 deletions(-) create mode 100644 neode-ui/src/views/federation/__tests__/QuickActions.test.ts diff --git a/core/archipelago/src/api/rpc/federation/handlers.rs b/core/archipelago/src/api/rpc/federation/handlers.rs index cd821174..c901086d 100644 --- a/core/archipelago/src/api/rpc/federation/handlers.rs +++ b/core/archipelago/src/api/rpc/federation/handlers.rs @@ -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"); } } diff --git a/docs/post-1.9.0-work-backlog.md b/docs/post-1.9.0-work-backlog.md index 578d5f95..3d97dd92 100644 --- a/docs/post-1.9.0-work-backlog.md +++ b/docs/post-1.9.0-work-backlog.md @@ -284,6 +284,20 @@ ahead of that work or as an untested addition to the current release. - Keep catalog and player integration within the Yaya-only demo scope until separately approved for general release. +## Federated Nodes state sync reliability and progress + +- Federated Nodes sync currently gives too little feedback and can take a long + time when one peer is unreachable. Show an accessible spinner, the number of + eligible nodes, current phase and a clear success/partial-failure result. +- Bound concurrent peer syncs so one slow or dead peer cannot serialize the + entire federation refresh. Prefer the authenticated FIPS route when + available, preserve the existing bounded fallback, and keep per-peer errors + visible without discarding successful results. +- Keep sync state convergent across retries and reconnects; prevent duplicate + clicks, stale responses or an incomplete refresh from looking like success. + Qualify healthy, slow, offline, mixed-transport, large-federation and mobile + layouts before acceptance. + ## 14. FIPS media transport requirement - Operator explicitly requires FIPS for streaming peer files and distributed diff --git a/neode-ui/src/views/Federation.vue b/neode-ui/src/views/Federation.vue index 0eec7cba..79c59190 100644 --- a/neode-ui/src/views/Federation.vue +++ b/neode-ui/src/views/Federation.vue @@ -62,6 +62,7 @@ :invite-type="inviteType" :invite-code="inviteCode" :syncing="syncing" + :sync-node-count="syncableNodeCount" @generate-invite="handleGenerateInvite" @show-join="showJoinModal = true" @sync="syncAll" @@ -311,6 +312,7 @@ const joinSuccess = ref(false) const syncing = ref(false) const syncResults = ref([]) +const syncableNodeCount = computed(() => nodes.value.filter(node => node.trust_level !== 'untrusted').length) const deploying = ref(false) const deployResult = ref('') diff --git a/neode-ui/src/views/federation/QuickActions.vue b/neode-ui/src/views/federation/QuickActions.vue index bf37b8f5..c8546600 100644 --- a/neode-ui/src/views/federation/QuickActions.vue +++ b/neode-ui/src/views/federation/QuickActions.vue @@ -74,11 +74,16 @@ @@ -121,6 +126,7 @@ const props = defineProps<{ inviteType: 'trusted' | 'observer' inviteCode: string syncing: boolean + syncNodeCount: number }>() defineEmits<{ diff --git a/neode-ui/src/views/federation/__tests__/QuickActions.test.ts b/neode-ui/src/views/federation/__tests__/QuickActions.test.ts new file mode 100644 index 00000000..07909c4a --- /dev/null +++ b/neode-ui/src/views/federation/__tests__/QuickActions.test.ts @@ -0,0 +1,38 @@ +import { mount } from '@vue/test-utils' +import { describe, expect, it } from 'vitest' +import QuickActions from '../QuickActions.vue' + +describe('QuickActions federation sync', () => { + it('shows bounded progress context while sync is running', () => { + const wrapper = mount(QuickActions, { + props: { + generatingInvite: false, + inviteType: 'trusted', + inviteCode: '', + syncing: true, + syncNodeCount: 3, + }, + }) + + const button = wrapper.get('[data-testid="federation-sync-button"]') + expect(button.attributes('disabled')).toBeDefined() + const status = button.get('[role="status"]') + expect(status.attributes('aria-live')).toBe('polite') + expect(status.text()).toContain('Syncing 3 nodes') + expect(status.find('.animate-spin').exists()).toBe(true) + }) + + it('keeps the action concise for a single node', () => { + const wrapper = mount(QuickActions, { + props: { + generatingInvite: false, + inviteType: 'trusted', + inviteCode: '', + syncing: true, + syncNodeCount: 1, + }, + }) + + expect(wrapper.get('[data-testid="federation-sync-button"]').text()).toContain('Syncing 1 node…') + }) +})