//! Shared Iroh endpoint, authenticated folder control streams, and blob transfers. use std::{ collections::{BTreeMap, BTreeSet, HashSet}, path::Path, str::FromStr, sync::{Arc, Mutex, MutexGuard}, }; use ::iroh as iroh_network; use anyhow::Context; use futures_util::{StreamExt, stream}; use iroh_blobs::{ BlobsProtocol, Hash, HashAndFormat, api::remote::{GetProgress, GetProgressItem}, protocol::GetManyRequest, provider::events::{EventMask, EventSender, ProviderMessage, RequestMode}, store::fs::FsStore, store::{GcConfig, ProtectCb, ProtectOutcome}, }; use iroh_mdns_address_lookup::{DiscoveryEvent, MdnsAddressLookup}; use iroh_network::{ Endpoint, EndpointAddr, EndpointId, SecretKey, endpoint::{Connection, presets}, protocol::Router, }; use tokio::sync::RwLock; use crate::{ domain::{ AuditEvent, FolderId, FolderRoster, Manifest, ManifestMerkleIndex, ManifestMerkleNode, }, protocol::{ControlMessage, Invite, ManifestRequestKind, ManifestSummary}, }; mod handler; use handler::AppaProtocol; /// ALPN for Appa's Iroh control streams. A new value is wire-incompatible. pub const APPA_ALPN: &[u8] = b"appa/sync/4"; /// Re-exported so callers can hold temp-tag guards across blob imports. pub use iroh_blobs::api::TempTag; const APPA_MDNS_SERVICE_NAME: &str = "appa"; const MAX_CONTROL_MESSAGE_BYTES: usize = 16 * 1024 * 1024; pub(crate) const MAX_AUDIT_EVENTS_PER_RESPONSE: usize = 256; const MAX_CONCURRENT_ANNOUNCEMENTS: usize = 4; const ANNOUNCEMENT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); pub(super) const CONTROL_REQUEST_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10); const BLOB_TRANSFER_IDLE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5 * 60); const ONLINE_WAIT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30); const BLOB_ACTIVITY_EVENT_BUFFER: usize = 32; #[derive(Clone, Debug)] pub struct FolderSession { pub folder_id: FolderId, pub capability: String, pub active_member_ids: BTreeSet, pub roster: FolderRoster, pub peer_endpoints: Vec, pub audit_events: Vec, } #[derive(Clone, Debug)] pub struct DiscoveredPeer { pub folder_id: FolderId, pub endpoint: EndpointAddr, } #[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] pub struct AnnouncementReport { pub announced_count: usize, pub unavailable_count: usize, } #[derive(Clone, Debug)] pub struct PeerSummaryResponse { pub summary: ManifestSummary, pub roster: FolderRoster, pub peer_endpoints: Vec, pub capability: String, } #[derive(Clone, Debug)] pub(super) struct HostedFolder { pub(super) capability: String, pub(super) active_member_ids: BTreeSet, pub(super) roster: FolderRoster, pub(super) peer_endpoints: Vec, pub(super) manifest: Manifest, pub(super) merkle_index: ManifestMerkleIndex, pub(super) audit_events: Vec, } pub struct NodeHost { endpoint: Endpoint, store: FsStore, router: Router, folders: Arc>>, // These maps are read or updated without awaiting. A standard mutex avoids // holding an async lock across network work, unlike `folders` below. discovered_peers: Arc>>, announced_folders: Arc>>, lan_discovered_peers: Arc>>, lan_discovery_task: Option>, blob_activity_task: tokio::task::JoinHandle<()>, } impl NodeHost { pub async fn load(data_directory: &Path, identity: SecretKey) -> anyhow::Result { Self::load_with_lan_discovery(data_directory, identity, true).await } pub async fn load_with_lan_discovery( data_directory: &Path, identity: SecretKey, lan_discovery_enabled: bool, ) -> anyhow::Result { let blob_directory = data_directory.join("blobs"); std::fs::create_dir_all(&blob_directory).with_context(|| { format!( "could not create Appa blob store at {}", blob_directory.display() ) })?; let endpoint = Endpoint::builder(presets::N0) .secret_key(identity) .bind() .await?; let store = Self::load_store_with_gc(&blob_directory, data_directory).await?; let folders = Arc::new(RwLock::new(BTreeMap::new())); let discovered_peers = Arc::new(Mutex::new(BTreeMap::new())); let announced_folders = Arc::new(Mutex::new(BTreeMap::new())); let lan_discovered_peers = Arc::new(Mutex::new(BTreeMap::new())); let lan_discovery_task = if lan_discovery_enabled { match start_lan_discovery(&endpoint, Arc::clone(&lan_discovered_peers)).await { Ok(task) => Some(task), Err(error) => { endpoint.close().await; return Err(error); } } } else { None }; let protocol = AppaProtocol { folders: Arc::clone(&folders), discovered_peers: Arc::clone(&discovered_peers), announced_folders: Arc::clone(&announced_folders), }; let (blob_activity_events, blob_activity_task) = blob_activity_events(); let router = Router::builder(endpoint.clone()) .accept( iroh_blobs::ALPN, BlobsProtocol::new(&store, Some(blob_activity_events)), ) .accept(APPA_ALPN, protocol) .spawn(); Ok(Self { endpoint, store, router, folders, discovered_peers, announced_folders, lan_discovered_peers, lan_discovery_task, blob_activity_task, }) } pub async fn wait_until_online(&self) -> anyhow::Result<()> { self.wait_until_online_for(ONLINE_WAIT_TIMEOUT).await } pub async fn wait_until_online_for(&self, timeout: std::time::Duration) -> anyhow::Result<()> { tokio::time::timeout(timeout, self.endpoint.online()) .await .map_err(|_| { anyhow::anyhow!( "timed out waiting {} seconds for Iroh network connectivity", timeout.as_secs() ) })?; Ok(()) } pub fn endpoint_address(&self) -> EndpointAddr { self.endpoint.addr() } pub async fn register_folder( &self, session: FolderSession, manifest: Manifest, ) -> anyhow::Result<()> { let merkle_index = ManifestMerkleIndex::build(&manifest)?; self.folders.write().await.insert( session.folder_id, HostedFolder { capability: session.capability, active_member_ids: session.active_member_ids, roster: session.roster, peer_endpoints: session.peer_endpoints, manifest, merkle_index, audit_events: session.audit_events, }, ); Ok(()) } pub async fn import_file(&self, file_path: &Path) -> anyhow::Result<(Hash, TempTag)> { let temp_tag = self.store.blobs().add_path(file_path).temp_tag().await?; Ok((temp_tag.hash(), temp_tag)) } /// Creates temp tags that protect the given blobs from garbage collection until /// the returned tags are dropped. Used to keep downloaded data alive between /// receiving it and persisting the manifest that references it. pub async fn protect_blobs(&self, hashes: &[Hash]) -> anyhow::Result> { let tags = self.store.tags(); let mut temp_tags = Vec::with_capacity(hashes.len()); for &hash in hashes { temp_tags.push(tags.temp_tag(HashAndFormat::raw(hash)).await?); } Ok(temp_tags) } /// Removes every named tag from the store. Appa's source of truth for live /// blobs is the SQLite manifest database, so legacy persistent tags created by /// older Appa versions can be dropped safely before the next garbage-collection /// cycle reclaims the orphaned bytes. pub async fn sweep_legacy_tags(&self) -> anyhow::Result { Ok(self.store.tags().delete_all().await?) } pub async fn download_blob( &self, blob_hash: Hash, provider: EndpointAddr, ) -> anyhow::Result<()> { self.download_blobs(vec![blob_hash], provider).await } pub async fn download_blobs( &self, blob_hashes: Vec, provider: EndpointAddr, ) -> anyhow::Result<()> { if blob_hashes.is_empty() { return Ok(()); } let connection = self.connect_to_blob_provider(provider).await?; let request = blob_hashes.into_iter().collect::(); wait_for_blob_download(self.store.remote().execute_get_many(connection, request)).await?; Ok(()) } pub async fn export_blob(&self, blob_hash: Hash, destination: &Path) -> anyhow::Result<()> { self.store.blobs().export(blob_hash, destination).await?; Ok(()) } pub async fn publish_manifest(&self, manifest: Manifest) -> anyhow::Result<()> { let merkle_index = ManifestMerkleIndex::build(&manifest)?; let mut folders = self.folders.write().await; let Some(folder) = folders.get_mut(&manifest.folder_id) else { anyhow::bail!("folder is not registered with this Appa node"); }; folder.manifest = manifest; folder.merkle_index = merkle_index; Ok(()) } pub async fn publish_audit_events( &self, folder_id: FolderId, audit_events: 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.audit_events = audit_events; 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, folder_id: FolderId, capability: String, root_hash: String, ) -> AnnouncementReport { let outcomes = stream::iter(peers) .map(|peer| { let peer_capability = capability.clone(); let peer_root_hash = root_hash.clone(); async move { let result = tokio::time::timeout( ANNOUNCEMENT_TIMEOUT, self.announce_manifest_to_peer( peer.clone(), folder_id, peer_capability, peer_root_hash, ), ) .await; (peer, result) } }) .buffer_unordered(MAX_CONCURRENT_ANNOUNCEMENTS) .collect::>() .await; summarize_announcement_outcomes(folder_id, outcomes) } async fn announce_manifest_to_peer( &self, peer: EndpointAddr, folder_id: FolderId, capability: String, root_hash: String, ) -> anyhow::Result<()> { let connection = self.endpoint.connect(peer, APPA_ALPN).await?; let response = request_control_message( &connection, &ControlMessage::Announcement { folder_id, capability, root_hash, sender_endpoint: self.endpoint_address(), }, ) .await?; if response != ControlMessage::AnnouncementAccepted { anyhow::bail!("peer rejected manifest announcement"); } Ok(()) } pub async fn request_manifest_summary( &self, peer_address: EndpointAddr, folder_id: FolderId, capability: String, invite: Option, ) -> anyhow::Result { let connection = self.connect_to_peer(peer_address).await?; let response = request_control_message( &connection, &ControlMessage::Request { folder_id, capability, invite: invite.map(Box::new), requester_endpoint: self.endpoint_address(), request: ManifestRequestKind::Summary, }, ) .await?; let ControlMessage::Summary { summary, roster, peer_endpoints, capability, } = response else { anyhow::bail!("peer returned an unexpected control message"); }; Ok(PeerSummaryResponse { summary, roster, peer_endpoints, capability, }) } pub async fn request_manifest_delta( &self, peer_address: EndpointAddr, folder_id: FolderId, capability: String, invite: Option, local_manifest: &Manifest, ) -> anyhow::Result { let connection = self.connect_to_peer(peer_address).await?; let local_merkle_index = ManifestMerkleIndex::build(local_manifest)?; let mut remote_manifest = local_manifest.clone(); let mut pending_prefixes = vec![String::new()]; while let Some(prefix) = pending_prefixes.pop() { let node = self .request_manifest_merkle_node( &connection, folder_id, capability.clone(), invite.clone(), prefix.clone(), ) .await?; if node.prefix() != prefix { anyhow::bail!("peer returned a Merkle node for the wrong prefix"); } let local_node = local_merkle_index.node(&prefix)?; if local_node .as_ref() .is_some_and(|local| local.hash() == node.hash()) { continue; } match node { ManifestMerkleNode::Branch { prefix, children, .. } => { for (label, remote_hash) in children { let child_prefix = format!("{prefix}{label}"); let local_node = local_merkle_index.node(&child_prefix)?; let local_hash = local_node.as_ref().map(ManifestMerkleNode::hash); if local_hash != Some(remote_hash.as_str()) { pending_prefixes.push(child_prefix); } } } ManifestMerkleNode::Leaf { entries, .. } => { remote_manifest.entries.extend(entries); } } } Ok(remote_manifest) } pub async fn request_audit_events( &self, peer_address: EndpointAddr, folder_id: FolderId, capability: String, invite: Option, after_sequences: BTreeMap, ) -> anyhow::Result> { let connection = self.connect_to_peer(peer_address).await?; let response = request_control_message( &connection, &ControlMessage::Request { folder_id, capability, invite: invite.map(Box::new), requester_endpoint: self.endpoint_address(), request: ManifestRequestKind::AuditEvents { after_sequences }, }, ) .await?; let ControlMessage::AuditEvents { events } = response else { anyhow::bail!("peer returned an unexpected audit response"); }; Ok(events) } async fn request_manifest_merkle_node( &self, connection: &Connection, folder_id: FolderId, capability: String, invite: Option, prefix: String, ) -> anyhow::Result { let response = request_control_message( connection, &ControlMessage::Request { folder_id, capability, invite: invite.map(Box::new), requester_endpoint: self.endpoint_address(), request: ManifestRequestKind::MerkleNode { prefix }, }, ) .await?; let ControlMessage::MerkleNode { node } = response else { anyhow::bail!("peer returned an unexpected control message"); }; Ok(node) } pub fn take_discovered_peers(&self) -> Vec { std::mem::take(&mut *recover_mutex(&self.discovered_peers)) .into_values() .collect() } async fn connect_to_peer(&self, peer_address: EndpointAddr) -> anyhow::Result { self.connect_with_alpn(peer_address, APPA_ALPN, "Appa peer") .await } async fn connect_to_blob_provider( &self, provider_address: EndpointAddr, ) -> anyhow::Result { self.connect_with_alpn(provider_address, iroh_blobs::ALPN, "Appa blob provider") .await } async fn connect_with_alpn( &self, peer_address: EndpointAddr, alpn: &[u8], connection_name: &str, ) -> anyhow::Result { let connection = tokio::time::timeout( CONTROL_REQUEST_TIMEOUT, self.endpoint.connect(peer_address.clone(), alpn), ) .await .map_err(|_| anyhow::anyhow!("{connection_name} connection timed out"))??; verify_connected_peer(&connection, &peer_address)?; Ok(connection) } pub fn take_announced_folders(&self) -> Vec { std::mem::take(&mut *recover_mutex(&self.announced_folders)) .into_keys() .collect() } pub fn lan_discovered_peers(&self) -> Vec { recover_mutex(&self.lan_discovered_peers) .values() .cloned() .collect() } /// Loads the filesystem blob store with automatic garbage collection. /// /// GC runs on a fixed interval (configurable via `APPA_GC_INTERVAL_SECS`). /// Before each sweep the [`ProtectCb`] callback opens a fresh read-only /// connection to the SQLite state and feeds every live blob hash — those /// referenced by current manifests, retained history revisions, and pending /// materializations — into the GC's live set so they survive the sweep. async fn load_store_with_gc( blob_directory: &Path, data_directory: &Path, ) -> anyhow::Result { let database_path = data_directory.join("state.sqlite3"); let db_path = blob_directory.join("blobs.db"); let mut options = iroh_blobs::store::fs::options::Options::new(blob_directory); let interval = crate::storage::gc_interval()?; let add_protected: ProtectCb = Arc::new(move |live: &mut HashSet| { let path = database_path.clone(); Box::pin(async move { protect_live_blobs(&path, live) }) }); options.gc = Some(GcConfig { interval, add_protected: Some(add_protected), }); Ok(FsStore::load_with_opts(db_path, options).await?) } pub async fn shutdown(self) -> anyhow::Result<()> { if let Some(task) = self.lan_discovery_task { task.abort(); } self.blob_activity_task.abort(); self.router.shutdown().await?; self.endpoint.close().await; Ok(()) } } fn add_live_blobs_from_state(database_path: &Path, live: &mut HashSet) -> anyhow::Result<()> { let live_hashes = crate::storage::StateStore::live_blob_hashes_for_gc(database_path)?; if live_hashes.is_empty() { tracing::trace!("GC live set is empty; protecting nothing extra"); } for hash_str in live_hashes { match Hash::from_str(&hash_str) { Ok(hash) => { live.insert(hash); } Err(error) => { tracing::warn!(%hash_str, %error, "Could not parse a blob hash from Appa state; skipping"); } } } Ok(()) } fn protect_live_blobs(database_path: &Path, live: &mut HashSet) -> ProtectOutcome { match add_live_blobs_from_state(database_path, live) { Ok(()) => ProtectOutcome::Continue, Err(error) => { tracing::warn!(%error, "Could not compute live blob set for GC; aborting sweep"); ProtectOutcome::Abort } } } pub(super) fn recover_mutex(mutex: &Mutex) -> MutexGuard<'_, T> { mutex .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) } async fn wait_for_blob_download(progress: GetProgress) -> anyhow::Result<()> { let updates = progress.stream(); tokio::pin!(updates); loop { let update = tokio::time::timeout(BLOB_TRANSFER_IDLE_TIMEOUT, updates.next()) .await .map_err(|_| anyhow::anyhow!("blob download made no progress for 5 minutes"))? .ok_or_else(|| anyhow::anyhow!("blob download ended before completing"))?; match update { GetProgressItem::Progress(_) => {} GetProgressItem::Done(_) => return Ok(()), GetProgressItem::Error(error) => return Err(error.into()), } } } fn verify_connected_peer( connection: &Connection, peer_address: &EndpointAddr, ) -> anyhow::Result<()> { if connection.remote_id() != peer_address.id { anyhow::bail!("connected Appa peer does not match the requested endpoint"); } Ok(()) } fn summarize_announcement_outcomes( folder_id: FolderId, outcomes: Vec<( EndpointAddr, Result, tokio::time::error::Elapsed>, )>, ) -> AnnouncementReport { let mut report = AnnouncementReport::default(); for (peer, outcome) in outcomes { match outcome { Ok(Ok(())) => report.announced_count += 1, Ok(Err(error)) => { report.unavailable_count += 1; tracing::debug!(folder = %folder_id, peer = %peer.id, %error, "Could not announce folder change"); } Err(error) => { report.unavailable_count += 1; tracing::debug!(folder = %folder_id, peer = %peer.id, %error, "Folder change announcement timed out"); } } } if report.unavailable_count > 0 { tracing::warn!( folder = %folder_id, unavailable_count = report.unavailable_count, "Some peers did not receive the folder change announcement" ); } report } async fn start_lan_discovery( endpoint: &Endpoint, discovered_peers: Arc>>, ) -> anyhow::Result> { let discovery = MdnsAddressLookup::builder() .service_name(APPA_MDNS_SERVICE_NAME) .build(endpoint.id())?; endpoint.address_lookup()?.add(discovery.clone()); let mut events = discovery.subscribe().await; Ok(tokio::spawn(async move { while let Some(event) = n0_future::StreamExt::next(&mut events).await { let DiscoveryEvent::Discovered { endpoint_info, .. } = event else { continue; }; let peer_endpoint = endpoint_info.into_endpoint_addr(); recover_mutex(&discovered_peers).insert(peer_endpoint.id, peer_endpoint); } })) } async fn request_control_message( connection: &Connection, message: &ControlMessage, ) -> anyhow::Result { let payload = serde_json::to_vec(message)?; if payload.len() > MAX_CONTROL_MESSAGE_BYTES { anyhow::bail!("control message exceeds the maximum size"); } tokio::time::timeout(CONTROL_REQUEST_TIMEOUT, async { let (mut send, mut receive) = connection.open_bi().await?; send.write_all(&payload).await?; send.finish()?; Ok(serde_json::from_slice( &receive.read_to_end(MAX_CONTROL_MESSAGE_BYTES).await?, )?) }) .await .map_err(|_| anyhow::anyhow!("Appa control request timed out"))? } fn blob_activity_events() -> (EventSender, tokio::task::JoinHandle<()>) { let event_mask = EventMask { get: RequestMode::Notify, get_many: RequestMode::Notify, ..EventMask::DEFAULT }; let (sender, mut receiver) = EventSender::channel(BLOB_ACTIVITY_EVENT_BUFFER, event_mask); let task = tokio::spawn(async move { while let Some(message) = receiver.recv().await { if matches!( message, ProviderMessage::GetRequestReceivedNotify(_) | ProviderMessage::GetManyRequestReceivedNotify(_) ) { tracing::info!("Serving file data"); } } }); (sender, task) } #[cfg(test)] #[path = "iroh/tests.rs"] mod tests;