From f9c6bc636e73386e6d9dd9359a39f03be2c5d5e7 Mon Sep 17 00:00:00 2001 From: Aly Raffauf Date: Mon, 3 Aug 2026 09:38:01 -0400 Subject: [PATCH] Synchronize newly introduced peers immediately --- src/app/sync.rs | 78 ++++++++++++++++++++++++++++++++++-------------- src/app/tests.rs | 60 ++++++++++++++++++++++++++++++++++++- src/iroh.rs | 13 ++++++++ 3 files changed, 128 insertions(+), 23 deletions(-) diff --git a/src/app/sync.rs b/src/app/sync.rs index 42af91b..0f0f8a7 100644 --- a/src/app/sync.rs +++ b/src/app/sync.rs @@ -1,3 +1,5 @@ +use std::collections::BTreeSet; + use crate::{ app::{AppResult, AppaService, MAX_CONCURRENT_PEER_REQUESTS, SyncContext}, app::{ @@ -77,26 +79,38 @@ impl AppaService { .await; tracing::debug!(folder = %folder.name, announced_count = report.announced_count, unavailable_count = report.unavailable_count, "Announced local manifest changes"); } - let peer_summaries = request_peer_summaries( - node, - self.state_store.peers(folder.id)?, - folder.id, - folder.capability.clone(), - self.state_store.join_invite(folder.id)?, - ) - .await; let mut synchronized_count = 0; - for (peer, response) in peer_summaries { - let Some(remote_result) = self - .synchronize_peer(&context, &local_manifest, peer, response) - .await? - else { - continue; - }; - local_manifest = remote_result.manifest; - synchronized_count += remote_result.synchronized_file_count; - self.state_store.save_manifest(&local_manifest)?; - node.publish_manifest(local_manifest.clone()).await?; + let mut attempted_peer_ids = BTreeSet::new(); + loop { + let peers: Vec<_> = self + .state_store + .peers(folder.id)? + .into_iter() + .filter(|peer| attempted_peer_ids.insert(peer.id)) + .collect(); + if peers.is_empty() { + break; + } + let peer_summaries = request_peer_summaries( + node, + peers, + folder.id, + folder.capability.clone(), + self.state_store.join_invite(folder.id)?, + ) + .await; + for (peer, response) in peer_summaries { + let Some(remote_result) = self + .synchronize_peer(&context, &local_manifest, peer, response) + .await? + else { + continue; + }; + local_manifest = remote_result.manifest; + synchronized_count += remote_result.synchronized_file_count; + self.state_store.save_manifest(&local_manifest)?; + node.publish_manifest(local_manifest.clone()).await?; + } } self.save_discovered_peers(folder, node)?; self.save_lan_discovered_peers(folder, node)?; @@ -216,9 +230,10 @@ impl AppaService { Ok(()) } - fn save_peer_response( + async fn save_peer_response( &self, folder: &FolderConfig, + node: &NodeHost, peer: &iroh::EndpointAddr, roster: &FolderRoster, peer_endpoints: Vec, @@ -239,7 +254,8 @@ impl AppaService { } if roster.epoch < current_roster.epoch { self.state_store.save_peer(folder.id, peer)?; - return Ok(()); + self.save_roster_peers(folder.id, ¤t_roster, peer_endpoints)?; + return self.publish_peer_endpoints(folder.id, node).await; } if roster.epoch == current_roster.epoch && roster != ¤t_roster { anyhow::bail!("peer returned a conflicting folder roster epoch"); @@ -247,6 +263,7 @@ impl AppaService { self.state_store.save_peer(folder.id, peer)?; let accepted_new_roster = self.state_store.save_roster_if_newer(roster)?; self.save_roster_peers(folder.id, roster, peer_endpoints)?; + self.publish_peer_endpoints(folder.id, node).await?; if accepted_new_roster && capability != folder.capability { self.state_store.update_capability(folder.id, capability)?; tracing::info!(folder = %folder.name, "Updated folder capability after member revocation"); @@ -254,6 +271,15 @@ impl AppaService { Ok(()) } + async fn publish_peer_endpoints( + &self, + folder_id: crate::domain::FolderId, + node: &NodeHost, + ) -> AppResult<()> { + node.publish_peer_endpoints(folder_id, self.state_store.peers(folder_id)?) + .await + } + fn save_roster_peers( &self, folder_id: crate::domain::FolderId, @@ -292,7 +318,15 @@ impl AppaService { tracing::warn!(folder = %folder.name, peer = %peer.id, "peer returned a summary for another folder"); return Ok(None); } - self.save_peer_response(folder, &peer, &roster, peer_endpoints, &capability)?; + self.save_peer_response( + folder, + context.node, + &peer, + &roster, + peer_endpoints, + &capability, + ) + .await?; self.receive_peer_audit_events(context, &peer).await?; if summary.root_hash == manifest_root_hash(local_manifest)? { tracing::debug!(folder = %folder.name, peer = %peer.id, "Folder roots match; skipping manifest transfer"); diff --git a/src/app/tests.rs b/src/app/tests.rs index ce7c2a7..707deac 100644 --- a/src/app/tests.rs +++ b/src/app/tests.rs @@ -11,7 +11,7 @@ use tokio::time::sleep; use crate::{ app::manifest::{build_manifest, build_manifest_with_stages}, config::{AppaConfig, CONFIG_VERSION, ConfiguredFolder}, - iroh::FolderSession, + iroh::{FolderSession, NodeHost}, }; use super::{AppPaths, AppaService, INVITATION_PROTOCOL_VERSION, Invite}; @@ -301,6 +301,64 @@ async fn synchronizes_audit_events_written_before_the_author_was_revoked() -> an Ok(()) } +#[tokio::test] +async fn learns_and_contacts_a_peer_introduced_during_the_same_sync_pass() -> anyhow::Result<()> { + let source_data = TempDir::new()?; + let source_folder = TempDir::new()?; + let bridge_data = TempDir::new()?; + let third_data = TempDir::new()?; + let source = AppaService::open_at(AppPaths::from_data_directory( + source_data.path().to_owned(), + )?)?; + let folder = source.register_folder(source_folder.path())?; + let bridge_identity = iroh::SecretKey::generate(); + let third_identity = iroh::SecretKey::generate(); + source.enroll_discovered_member(folder.id, &bridge_identity.public().to_string())?; + source.enroll_discovered_member(folder.id, &third_identity.public().to_string())?; + let roster = source + .state_store + .load_roster(folder.id)? + .expect("shared roster"); + let bridge_node = + NodeHost::load_with_lan_discovery(bridge_data.path(), bridge_identity, false).await?; + let third_node = + NodeHost::load_with_lan_discovery(third_data.path(), third_identity, false).await?; + bridge_node + .register_folder( + FolderSession { + folder_id: folder.id, + capability: folder.capability.clone(), + active_member_ids: roster.active_device_ids(), + roster, + peer_endpoints: vec![third_node.endpoint_address()], + audit_events: vec![], + }, + crate::domain::Manifest::empty(folder.id), + ) + .await?; + let source_node = source.load_node().await?; + source_node.wait_until_online().await?; + bridge_node.wait_until_online().await?; + third_node.wait_until_online().await?; + source + .state_store + .save_peer(folder.id, &bridge_node.endpoint_address())?; + + source.sync_with_node(&folder, &source_node).await?; + + assert!( + source + .state_store + .peers(folder.id)? + .iter() + .any(|peer| peer.id == third_node.endpoint_address().id) + ); + source_node.shutdown().await?; + bridge_node.shutdown().await?; + third_node.shutdown().await?; + Ok(()) +} + async fn sync_until_file_materializes( appa: &AppaService, folder: &Path, diff --git a/src/iroh.rs b/src/iroh.rs index 9c504a5..424971f 100644 --- a/src/iroh.rs +++ b/src/iroh.rs @@ -259,6 +259,19 @@ impl NodeHost { Ok(()) } + pub async fn publish_peer_endpoints( + &self, + folder_id: FolderId, + peer_endpoints: Vec, + ) -> anyhow::Result<()> { + let mut folders = self.folders.write().await; + let Some(folder) = folders.get_mut(&folder_id) else { + anyhow::bail!("folder is not registered with this Appa node"); + }; + folder.peer_endpoints = peer_endpoints; + Ok(()) + } + pub async fn announce_manifest( &self, peers: Vec, -- 2.51.2