From 53fc8dc9af5ed41eef72d4332238f46fd2965581 Mon Sep 17 00:00:00 2001 From: Aly Raffauf Date: Sun, 2 Aug 2026 21:17:48 -0400 Subject: [PATCH] isolate peer synchronization workflow --- src/app.rs | 136 +++------------------------------------------ src/app/sync.rs | 144 ++++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 151 insertions(+), 129 deletions(-) create mode 100644 src/app/sync.rs diff --git a/src/app.rs b/src/app.rs index d40bafe..8c4750b 100644 --- a/src/app.rs +++ b/src/app.rs @@ -16,12 +16,13 @@ use crate::{ IgnoreRules, collect_directories, collect_files, relative_path, remove_path, safe_destination, }, - iroh::{FolderSession, NodeHost}, + iroh::NodeHost, protocol::{Invite, PROTOCOL_VERSION, decode_invite, encode_invite}, storage::{AppPaths, FolderConfig, ManifestRevision, MemberInfo, PeerInfo, StateStore}, }; mod run; +mod sync; pub struct AppaService { paths: AppPaths, @@ -346,117 +347,6 @@ impl AppaService { NodeHost::load(&self.paths.data_directory, self.paths.load_identity()?).await } - async fn sync_with_node(&self, folder: &FolderConfig, node: &NodeHost) -> AppResult { - self.save_lan_discovered_peers(folder, node)?; - self.save_discovered_peers(folder, node)?; - let manifest_build = if folder.mode.can_send() { - build_manifest(folder, &self.state_store, node).await? - } else { - ManifestBuild { - manifest: self.state_store.load_manifest(folder.id)?, - changed_paths: Vec::new(), - } - }; - report_local_changes(&manifest_build.changed_paths); - let mut local_manifest = manifest_build.manifest; - let local_roster = self.load_roster(folder.id)?; - node.register_folder( - FolderSession { - folder_id: folder.id, - capability: folder.capability.clone(), - previous_capability: self.state_store.previous_capability(folder.id)?, - active_member_ids: local_roster.active_device_ids(), - roster: local_roster, - peer_endpoints: self.state_store.peers(folder.id)?, - }, - local_manifest.clone(), - ) - .await; - self.state_store.save_manifest(&local_manifest)?; - node.publish_manifest(local_manifest.clone()).await?; - let mut synchronized_count = 0; - for peer in self.state_store.peers(folder.id)? { - let (remote, remote_roster, remote_peers, remote_capability) = match node - .request_manifest( - peer.clone(), - folder.id, - folder.capability.clone(), - self.state_store.join_invite(folder.id)?, - ) - .await - { - Ok(response) => response, - Err(error) => { - tracing::debug!(folder = %folder.name, peer = %peer.id, %error, "peer is unavailable; continuing with other peers"); - continue; - } - }; - self.state_store.save_peer(folder.id, &peer)?; - self.state_store.accept_roster(&remote_roster)?; - self.save_roster_peers(folder.id, &remote_roster, remote_peers)?; - if remote_capability != folder.capability { - self.state_store - .update_capability(folder.id, &remote_capability)?; - tracing::info!(folder = %folder.name, "Updated folder capability after member revocation"); - } - if remote.folder_id != folder.id { - tracing::warn!(folder = %folder.name, peer = %peer.id, "peer returned a manifest for another folder"); - continue; - } - if !folder.mode.can_receive() { - tracing::debug!(folder = %folder.name, "Ignoring remote changes for send-only folder"); - continue; - } - let (merged, count) = match apply_remote_manifest( - folder, - node, - &local_manifest, - &remote, - peer.clone(), - ) - .await - { - Ok(result) => result, - Err(error) => { - tracing::warn!(folder = %folder.name, peer = %peer.id, %error, "could not apply peer manifest"); - continue; - } - }; - local_manifest = merged; - synchronized_count += 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)?; - self.state_store.record_successful_sync(folder.id)?; - Ok(synchronized_count) - } - - fn save_discovered_peers(&self, folder: &FolderConfig, node: &NodeHost) -> AppResult<()> { - for peer in node.take_discovered_peers() { - if peer.folder_id == folder.id { - self.state_store.save_peer(folder.id, &peer.endpoint)?; - self.enroll_discovered_member(folder.id, &peer.endpoint.id.to_string())?; - } - } - Ok(()) - } - - fn save_lan_discovered_peers(&self, folder: &FolderConfig, node: &NodeHost) -> AppResult<()> { - let active_member_ids = self.active_member_ids(folder.id)?; - for peer in node.lan_discovered_peers() { - let device_id = peer.id.to_string(); - if active_member_ids.contains(&device_id) { - tracing::debug!(folder = %folder.name, peer = %device_id, "Refreshed trusted local peer route"); - self.state_store.save_peer(folder.id, &peer)?; - } else { - tracing::debug!(folder = %folder.name, peer = %device_id, "Ignoring untrusted local peer"); - } - } - Ok(()) - } - fn require_folder(&self, folder_path: &Path) -> AppResult { self.state_store.find_folder(folder_path)?.ok_or_else(|| { anyhow::anyhow!( @@ -535,20 +425,6 @@ impl AppaService { self.state_store.save_roster(&roster) } - fn save_roster_peers( - &self, - folder_id: crate::domain::FolderId, - roster: &FolderRoster, - peer_endpoints: Vec, - ) -> AppResult<()> { - for peer in peer_endpoints { - if roster.contains_member(&peer.id.to_string()) { - self.state_store.save_peer(folder_id, &peer)?; - } - } - Ok(()) - } - fn device_id(&self) -> AppResult { Ok(self.paths.load_identity()?.public().to_string()) } @@ -953,11 +829,13 @@ mod tests { use tempfile::TempDir; use time::{Duration, OffsetDateTime}; - use crate::config::{AppaConfig, CONFIG_VERSION, ConfiguredFolder}; + use crate::{ + config::{AppaConfig, CONFIG_VERSION, ConfiguredFolder}, + iroh::FolderSession, + }; use super::{ - AppPaths, AppaService, ConfigAuditAction, FolderSession, Invite, PROTOCOL_VERSION, - build_manifest, + AppPaths, AppaService, ConfigAuditAction, Invite, PROTOCOL_VERSION, build_manifest, }; #[test] diff --git a/src/app/sync.rs b/src/app/sync.rs new file mode 100644 index 0000000..e673844 --- /dev/null +++ b/src/app/sync.rs @@ -0,0 +1,144 @@ +use crate::{ + app::{ + AppResult, AppaService, ManifestBuild, apply_remote_manifest, build_manifest, + report_local_changes, + }, + domain::FolderRoster, + iroh::{FolderSession, NodeHost}, + storage::FolderConfig, +}; + +impl AppaService { + pub(super) async fn sync_with_node( + &self, + folder: &FolderConfig, + node: &NodeHost, + ) -> AppResult { + self.save_lan_discovered_peers(folder, node)?; + self.save_discovered_peers(folder, node)?; + let manifest_build = if folder.mode.can_send() { + build_manifest(folder, &self.state_store, node).await? + } else { + ManifestBuild { + manifest: self.state_store.load_manifest(folder.id)?, + changed_paths: Vec::new(), + } + }; + report_local_changes(&manifest_build.changed_paths); + let mut local_manifest = manifest_build.manifest; + let local_roster = self.load_roster(folder.id)?; + node.register_folder( + FolderSession { + folder_id: folder.id, + capability: folder.capability.clone(), + previous_capability: self.state_store.previous_capability(folder.id)?, + active_member_ids: local_roster.active_device_ids(), + roster: local_roster, + peer_endpoints: self.state_store.peers(folder.id)?, + }, + local_manifest.clone(), + ) + .await; + self.state_store.save_manifest(&local_manifest)?; + node.publish_manifest(local_manifest.clone()).await?; + let mut synchronized_count = 0; + for peer in self.state_store.peers(folder.id)? { + let (remote, remote_roster, remote_peers, remote_capability) = match node + .request_manifest( + peer.clone(), + folder.id, + folder.capability.clone(), + self.state_store.join_invite(folder.id)?, + ) + .await + { + Ok(response) => response, + Err(error) => { + tracing::debug!(folder = %folder.name, peer = %peer.id, %error, "peer is unavailable; continuing with other peers"); + continue; + } + }; + self.state_store.save_peer(folder.id, &peer)?; + self.state_store.accept_roster(&remote_roster)?; + self.save_roster_peers(folder.id, &remote_roster, remote_peers)?; + if remote_capability != folder.capability { + self.state_store + .update_capability(folder.id, &remote_capability)?; + tracing::info!(folder = %folder.name, "Updated folder capability after member revocation"); + } + if remote.folder_id != folder.id { + tracing::warn!(folder = %folder.name, peer = %peer.id, "peer returned a manifest for another folder"); + continue; + } + if !folder.mode.can_receive() { + tracing::debug!(folder = %folder.name, "Ignoring remote changes for send-only folder"); + continue; + } + let (merged, count) = match apply_remote_manifest( + folder, + node, + &local_manifest, + &remote, + peer.clone(), + ) + .await + { + Ok(result) => result, + Err(error) => { + tracing::warn!(folder = %folder.name, peer = %peer.id, %error, "could not apply peer manifest"); + continue; + } + }; + local_manifest = merged; + synchronized_count += 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)?; + self.state_store.record_successful_sync(folder.id)?; + Ok(synchronized_count) + } + + pub(super) fn save_discovered_peers( + &self, + folder: &FolderConfig, + node: &NodeHost, + ) -> AppResult<()> { + for peer in node.take_discovered_peers() { + if peer.folder_id == folder.id { + self.state_store.save_peer(folder.id, &peer.endpoint)?; + self.enroll_discovered_member(folder.id, &peer.endpoint.id.to_string())?; + } + } + Ok(()) + } + + fn save_lan_discovered_peers(&self, folder: &FolderConfig, node: &NodeHost) -> AppResult<()> { + let active_member_ids = self.active_member_ids(folder.id)?; + for peer in node.lan_discovered_peers() { + let device_id = peer.id.to_string(); + if active_member_ids.contains(&device_id) { + tracing::debug!(folder = %folder.name, peer = %device_id, "Refreshed trusted local peer route"); + self.state_store.save_peer(folder.id, &peer)?; + } else { + tracing::debug!(folder = %folder.name, peer = %device_id, "Ignoring untrusted local peer"); + } + } + Ok(()) + } + + fn save_roster_peers( + &self, + folder_id: crate::domain::FolderId, + roster: &FolderRoster, + peer_endpoints: Vec, + ) -> AppResult<()> { + for peer in peer_endpoints { + if roster.contains_member(&peer.id.to_string()) { + self.state_store.save_peer(folder_id, &peer)?; + } + } + Ok(()) + } +} -- 2.51.2