Bound federated sync and show progress feedback
This commit is contained in:
@@ -7,6 +7,8 @@ use crate::mesh;
|
|||||||
use crate::network::dwn_store::DwnStore;
|
use crate::network::dwn_store::DwnStore;
|
||||||
use crate::nostr_handshake;
|
use crate::nostr_handshake;
|
||||||
use anyhow::Result;
|
use anyhow::Result;
|
||||||
|
use futures_util::stream::{FuturesUnordered, StreamExt};
|
||||||
|
use std::sync::Arc;
|
||||||
use tracing::{debug, info, warn};
|
use tracing::{debug, info, warn};
|
||||||
|
|
||||||
const FEDERATION_PROTOCOL: &str = "https://archipelago.dev/protocols/federation/v1";
|
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 identity_dir = self.config.data_dir.join("identity");
|
||||||
let node_identity = identity::NodeIdentity::load_or_create(&identity_dir).await?;
|
let node_identity = identity::NodeIdentity::load_or_create(&identity_dir).await?;
|
||||||
|
|
||||||
let mut synced = 0u32;
|
// A dead Tor peer can take the bounded transport timeout. Serialising
|
||||||
let mut failed = 0u32;
|
// those waits made one bad peer block every healthy peer behind it.
|
||||||
let mut results = Vec::new();
|
// 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 node in &nodes {
|
for (index, node) in nodes
|
||||||
if node.trust_level == TrustLevel::Untrusted {
|
.into_iter()
|
||||||
continue;
|
.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 did_clone = local_did.clone();
|
let mut results = Vec::new();
|
||||||
match federation::sync_with_peer(&self.config.data_dir, node, &did_clone, |bytes| {
|
while let Some((index, did, result)) = pending.next().await {
|
||||||
node_identity.sign(bytes)
|
let row = match result {
|
||||||
})
|
Ok(state) => serde_json::json!({
|
||||||
.await
|
"index": index,
|
||||||
{
|
"did": did,
|
||||||
Ok(state) => {
|
|
||||||
synced += 1;
|
|
||||||
results.push(serde_json::json!({
|
|
||||||
"did": node.did,
|
|
||||||
"status": "ok",
|
"status": "ok",
|
||||||
"apps": state.apps.len(),
|
"apps": state.apps.len(),
|
||||||
}));
|
}),
|
||||||
}
|
Err(e) => serde_json::json!({
|
||||||
Err(e) => {
|
"index": index,
|
||||||
failed += 1;
|
"did": did,
|
||||||
results.push(serde_json::json!({
|
|
||||||
"did": node.did,
|
|
||||||
"status": "error",
|
"status": "error",
|
||||||
"error": e.to_string(),
|
"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");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -311,6 +311,20 @@ must use a qualified private delivery path; do not republish the demo image.
|
|||||||
- Keep catalog and player integration within the Yaya-only demo scope until
|
- Keep catalog and player integration within the Yaya-only demo scope until
|
||||||
separately approved for general release.
|
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
|
## 14. FIPS media transport requirement
|
||||||
|
|
||||||
- Operator explicitly requires FIPS for streaming peer files and distributed
|
- Operator explicitly requires FIPS for streaming peer files and distributed
|
||||||
|
|||||||
@@ -62,6 +62,7 @@
|
|||||||
:invite-type="inviteType"
|
:invite-type="inviteType"
|
||||||
:invite-code="inviteCode"
|
:invite-code="inviteCode"
|
||||||
:syncing="syncing"
|
:syncing="syncing"
|
||||||
|
:sync-node-count="syncableNodeCount"
|
||||||
@generate-invite="handleGenerateInvite"
|
@generate-invite="handleGenerateInvite"
|
||||||
@show-join="showJoinModal = true"
|
@show-join="showJoinModal = true"
|
||||||
@sync="syncAll"
|
@sync="syncAll"
|
||||||
@@ -313,6 +314,7 @@ const joinSuccess = ref(false)
|
|||||||
|
|
||||||
const syncing = ref(false)
|
const syncing = ref(false)
|
||||||
const syncResults = ref<SyncResult[]>([])
|
const syncResults = ref<SyncResult[]>([])
|
||||||
|
const syncableNodeCount = computed(() => nodes.value.filter(node => node.trust_level !== 'untrusted').length)
|
||||||
|
|
||||||
const deploying = ref(false)
|
const deploying = ref(false)
|
||||||
const deployResult = ref('')
|
const deployResult = ref('')
|
||||||
|
|||||||
@@ -74,11 +74,16 @@
|
|||||||
</div>
|
</div>
|
||||||
</div>
|
</div>
|
||||||
<button
|
<button
|
||||||
|
data-testid="federation-sync-button"
|
||||||
@click="$emit('sync')"
|
@click="$emit('sync')"
|
||||||
class="w-full sm:w-fit px-3 py-1.5 glass-button glass-button-sm rounded text-xs font-medium text-white/90 hover:text-white transition-colors disabled:opacity-50"
|
class="w-full sm:w-fit px-3 py-1.5 glass-button glass-button-sm rounded text-xs font-medium text-white/90 hover:text-white transition-colors disabled:opacity-50"
|
||||||
:disabled="syncing"
|
:disabled="syncing"
|
||||||
>
|
>
|
||||||
{{ syncing ? 'Syncing...' : 'Sync Now' }}
|
<span v-if="syncing" class="inline-flex items-center gap-2" role="status" aria-live="polite">
|
||||||
|
<span class="w-3 h-3 border-2 border-white/30 border-t-white rounded-full animate-spin" aria-hidden="true"></span>
|
||||||
|
Syncing {{ syncNodeCount }} node{{ syncNodeCount === 1 ? '' : 's' }}…
|
||||||
|
</span>
|
||||||
|
<span v-else>Sync Now</span>
|
||||||
</button>
|
</button>
|
||||||
</div>
|
</div>
|
||||||
</div>
|
</div>
|
||||||
@@ -121,6 +126,7 @@ const props = defineProps<{
|
|||||||
inviteType: 'trusted' | 'observer'
|
inviteType: 'trusted' | 'observer'
|
||||||
inviteCode: string
|
inviteCode: string
|
||||||
syncing: boolean
|
syncing: boolean
|
||||||
|
syncNodeCount: number
|
||||||
}>()
|
}>()
|
||||||
|
|
||||||
defineEmits<{
|
defineEmits<{
|
||||||
|
|||||||
@@ -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…')
|
||||||
|
})
|
||||||
|
})
|
||||||
Reference in New Issue
Block a user