From f8ca64246378b77e735605c54df82da9b29fb676 Mon Sep 17 00:00:00 2001 From: Aly Raffauf Date: Mon, 3 Aug 2026 07:52:32 -0400 Subject: [PATCH] Centralize Iroh mutex recovery --- src/iroh.rs | 41 +++++++++++++++++------------------------ src/iroh/handler.rs | 27 ++++++++++++--------------- 2 files changed, 29 insertions(+), 39 deletions(-) diff --git a/src/iroh.rs b/src/iroh.rs index 21217f3..a737b8c 100644 --- a/src/iroh.rs +++ b/src/iroh.rs @@ -3,7 +3,7 @@ use std::{ collections::{BTreeMap, BTreeSet}, path::Path, - sync::{Arc, Mutex}, + sync::{Arc, Mutex, MutexGuard}, }; use ::iroh as iroh_network; @@ -86,6 +86,8 @@ pub struct NodeHost { 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>>, @@ -407,14 +409,9 @@ impl NodeHost { } pub fn take_discovered_peers(&self) -> Vec { - std::mem::take( - &mut *self - .discovered_peers - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner), - ) - .into_values() - .collect() + std::mem::take(&mut *recover_mutex(&self.discovered_peers)) + .into_values() + .collect() } async fn connect_to_peer(&self, peer_address: EndpointAddr) -> anyhow::Result { @@ -447,20 +444,13 @@ impl NodeHost { } pub fn take_announced_folders(&self) -> Vec { - std::mem::take( - &mut *self - .announced_folders - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner), - ) - .into_keys() - .collect() + std::mem::take(&mut *recover_mutex(&self.announced_folders)) + .into_keys() + .collect() } pub fn lan_discovered_peers(&self) -> Vec { - self.lan_discovered_peers - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) + recover_mutex(&self.lan_discovered_peers) .values() .cloned() .collect() @@ -477,6 +467,12 @@ impl NodeHost { } } +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); @@ -548,10 +544,7 @@ async fn start_lan_discovery( continue; }; let peer_endpoint = endpoint_info.into_endpoint_addr(); - discovered_peers - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .insert(peer_endpoint.id, peer_endpoint); + recover_mutex(&discovered_peers).insert(peer_endpoint.id, peer_endpoint); } })) } diff --git a/src/iroh/handler.rs b/src/iroh/handler.rs index 4ad16ed..efcaa65 100644 --- a/src/iroh/handler.rs +++ b/src/iroh/handler.rs @@ -12,7 +12,10 @@ use tokio::sync::RwLock; use crate::{ domain::FolderId, - iroh::{CONTROL_REQUEST_TIMEOUT, DiscoveredPeer, HostedFolder, MAX_CONTROL_MESSAGE_BYTES}, + iroh::{ + CONTROL_REQUEST_TIMEOUT, DiscoveredPeer, HostedFolder, MAX_CONTROL_MESSAGE_BYTES, + recover_mutex, + }, protocol::{ControlMessage, ManifestRequestKind, ManifestSummary, validate_invite}, }; @@ -69,16 +72,13 @@ impl AppaProtocol { request_kind, ) .await?; - self.discovered_peers - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .insert( - (folder_id, requester_endpoint.id), - DiscoveredPeer { - folder_id, - endpoint: requester_endpoint, - }, - ); + recover_mutex(&self.discovered_peers).insert( + (folder_id, requester_endpoint.id), + DiscoveredPeer { + folder_id, + endpoint: requester_endpoint, + }, + ); response } ControlMessage::Announcement { @@ -94,10 +94,7 @@ impl AppaProtocol { } self.authorize_announcement(folder_id, &capability, &sender_endpoint) .await?; - self.announced_folders - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .insert(folder_id, root_hash); + recover_mutex(&self.announced_folders).insert(folder_id, root_hash); tracing::debug!(peer = %connection.remote_id(), folder = %folder_id, "Received folder change announcement"); ControlMessage::AnnouncementAccepted } -- 2.51.2