From 326592f74b09b42a78395c8e2d68ab216d0c58c8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BD=A6=E7=BA=BE?= Date: Fri, 31 Jul 2026 17:17:37 +0800 Subject: [PATCH 1/5] perf(p2p): parallelize catalog peer lookups --- src/p2p/iroh/transport.rs | 125 ++++++++++++++++++++++++++++++++++---- 1 file changed, 113 insertions(+), 12 deletions(-) diff --git a/src/p2p/iroh/transport.rs b/src/p2p/iroh/transport.rs index 3d0339588..1ba22ec70 100644 --- a/src/p2p/iroh/transport.rs +++ b/src/p2p/iroh/transport.rs @@ -1,4 +1,5 @@ use std::collections::HashSet; +use std::future::Future; use std::net::SocketAddr; use std::ops::Range; use std::path::Path; @@ -40,6 +41,7 @@ use crate::p2p::types::{ use crate::p2p::P2pByteStream; const CATALOG_DB_DIR: &str = "catalog.db"; +const MAX_CONCURRENT_CATALOG_LOOKUPS: usize = 4; const ENDPOINT_ADDR_TIMEOUT: Duration = Duration::from_secs(5); const DEFAULT_STORE_GC_INTERVAL: Duration = Duration::from_mins(5); const PUBLISH_TAG_PREFIX: &str = "agentenv:p2p:v1:"; @@ -220,21 +222,19 @@ impl IrohBlobsP2pTransport { peers: Vec, key: &P2pArtifactKey, ) -> Option { - // TODO: parallelize lookups to multiple peers. - // No shuffling or prioritization for now since the scheduler should already have done that work for us. - for peer in peers { + first_some_buffered_in_order(peers, MAX_CONCURRENT_CATALOG_LOOKUPS, |peer| async move { if peer.node_id == self.node_id || peer.endpoint == self.local_endpoint { - let Some(descriptor) = self.get_local(key).await else { - continue; - }; - trace!("P2P lookup found local descriptor matching peer discovery"); - return Some(descriptor); + let descriptor = self.get_local(key).await; + if descriptor.is_some() { + trace!("P2P lookup found local descriptor matching peer discovery"); + } + return descriptor; } match self.lookup_peer_with_timeout(&peer, key).await { - Ok(None) => continue, + Ok(None) => None, Ok(Some(descriptor)) => { trace!(peer = %peer.node_id, "P2P lookup found remote descriptor"); - return Some(descriptor); + Some(descriptor) } Err(err) => { debug!( @@ -242,10 +242,11 @@ impl IrohBlobsP2pTransport { error = %err, "P2P artifact catalog lookup failed; trying remaining peers" ); + None } } - } - None + }) + .await } async fn get_local(&self, key: &P2pArtifactKey) -> Option { @@ -720,6 +721,23 @@ fn blob_hash_from_descriptor(descriptor: &P2pArtifactDescriptor) -> Result( + items: impl IntoIterator, + concurrency: usize, + lookup: F, +) -> Option +where + F: FnMut(I) -> Fut, + Fut: Future>, +{ + stream::iter(items) + .map(lookup) + .buffered(concurrency.max(1)) + .filter_map(|result| async move { result }) + .next() + .await +} + fn resolve_remote_local_providers( mut descriptor: P2pArtifactDescriptor, response_peer: &P2pPeer, @@ -765,6 +783,9 @@ fn gated_gc_config(interval: Duration, pending_gc: Arc) -> GcConfig #[cfg(test)] mod tests { use super::*; + use std::sync::atomic::AtomicUsize; + use tokio::sync::{Barrier, Notify}; + use crate::cfg::P2pConfig; use crate::p2p::config::ResolvedP2pConfig; use crate::p2p::discovery::P2pPeerDiscovery; @@ -773,6 +794,86 @@ mod tests { const TEST_TIMEOUT: Duration = Duration::from_secs(10); + #[tokio::test] + async fn buffered_lookup_preserves_candidate_order() { + let release_first = Arc::new(Notify::new()); + let first_started = Arc::new(Notify::new()); + let second_finished = Arc::new(Notify::new()); + + let lookup = tokio::spawn({ + let release_first = release_first.clone(); + let first_started = first_started.clone(); + let second_finished = second_finished.clone(); + async move { + first_some_buffered_in_order([0, 1], 2, move |candidate| { + let release_first = release_first.clone(); + let first_started = first_started.clone(); + let second_finished = second_finished.clone(); + async move { + match candidate { + 0 => { + first_started.notify_one(); + release_first.notified().await; + Some("first") + } + 1 => { + second_finished.notify_one(); + Some("second") + } + _ => None, + } + } + }) + .await + } + }); + + first_started.notified().await; + second_finished.notified().await; + assert!( + !lookup.is_finished(), + "a lower-priority result must wait for earlier candidates" + ); + + release_first.notify_one(); + assert_eq!(lookup.await.expect("lookup task"), Some("first")); + } + + #[tokio::test] + async fn buffered_lookup_limits_in_flight_candidates() { + let current = Arc::new(AtomicUsize::new(0)); + let maximum = Arc::new(AtomicUsize::new(0)); + let first_window_ready = Arc::new(Barrier::new(5)); + + let lookup = tokio::spawn({ + let current = current.clone(); + let maximum = maximum.clone(); + let first_window_ready = first_window_ready.clone(); + async move { + first_some_buffered_in_order(0..8, 4, move |candidate| { + let current = current.clone(); + let maximum = maximum.clone(); + let first_window_ready = first_window_ready.clone(); + async move { + let active = current.fetch_add(1, Ordering::SeqCst) + 1; + maximum.fetch_max(active, Ordering::SeqCst); + if candidate < 4 { + first_window_ready.wait().await; + } + current.fetch_sub(1, Ordering::SeqCst); + None::<()> + } + }) + .await + } + }); + + first_window_ready.wait().await; + assert_eq!(current.load(Ordering::SeqCst), 4); + assert_eq!(lookup.await.expect("lookup task"), None); + assert_eq!(maximum.load(Ordering::SeqCst), 4); + } + async fn collect_range_stream(mut stream: P2pByteStream) -> Result> { let mut out = Vec::new(); while let Some(chunk) = stream.next().await { From cda8d77c6a33640aad7294075d6fc517bc1b6501 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BD=A6=E7=BA=BE?= Date: Fri, 31 Jul 2026 17:18:51 +0800 Subject: [PATCH 2/5] fix(p2p): consume buffered lookup results directly --- src/p2p/iroh/transport.rs | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/src/p2p/iroh/transport.rs b/src/p2p/iroh/transport.rs index 1ba22ec70..e71bcd41d 100644 --- a/src/p2p/iroh/transport.rs +++ b/src/p2p/iroh/transport.rs @@ -730,12 +730,15 @@ where F: FnMut(I) -> Fut, Fut: Future>, { - stream::iter(items) + let mut results = stream::iter(items) .map(lookup) - .buffered(concurrency.max(1)) - .filter_map(|result| async move { result }) - .next() - .await + .buffered(concurrency.max(1)); + while let Some(result) = results.next().await { + if result.is_some() { + return result; + } + } + None } fn resolve_remote_local_providers( From f8603dd975253d590b9c2c1a23b7ee5fb9bfa4f3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BD=A6=E7=BA=BE?= Date: Fri, 31 Jul 2026 17:19:29 +0800 Subject: [PATCH 3/5] style(p2p): apply rustfmt to lookup helper --- src/p2p/iroh/transport.rs | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/src/p2p/iroh/transport.rs b/src/p2p/iroh/transport.rs index e71bcd41d..990c3ff90 100644 --- a/src/p2p/iroh/transport.rs +++ b/src/p2p/iroh/transport.rs @@ -730,9 +730,7 @@ where F: FnMut(I) -> Fut, Fut: Future>, { - let mut results = stream::iter(items) - .map(lookup) - .buffered(concurrency.max(1)); + let mut results = stream::iter(items).map(lookup).buffered(concurrency.max(1)); while let Some(result) = results.next().await { if result.is_some() { return result; From 70ccd6d270f80802945cc615b58ea7c03bcb824d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BD=A6=E7=BA=BE?= Date: Fri, 31 Jul 2026 17:21:00 +0800 Subject: [PATCH 4/5] test(p2p): synchronize concurrency window assertion --- src/p2p/iroh/transport.rs | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/src/p2p/iroh/transport.rs b/src/p2p/iroh/transport.rs index 990c3ff90..da296d517 100644 --- a/src/p2p/iroh/transport.rs +++ b/src/p2p/iroh/transport.rs @@ -845,21 +845,25 @@ mod tests { let current = Arc::new(AtomicUsize::new(0)); let maximum = Arc::new(AtomicUsize::new(0)); let first_window_ready = Arc::new(Barrier::new(5)); + let release_first_window = Arc::new(Barrier::new(5)); let lookup = tokio::spawn({ let current = current.clone(); let maximum = maximum.clone(); let first_window_ready = first_window_ready.clone(); + let release_first_window = release_first_window.clone(); async move { first_some_buffered_in_order(0..8, 4, move |candidate| { let current = current.clone(); let maximum = maximum.clone(); let first_window_ready = first_window_ready.clone(); + let release_first_window = release_first_window.clone(); async move { let active = current.fetch_add(1, Ordering::SeqCst) + 1; maximum.fetch_max(active, Ordering::SeqCst); if candidate < 4 { first_window_ready.wait().await; + release_first_window.wait().await; } current.fetch_sub(1, Ordering::SeqCst); None::<()> @@ -871,6 +875,7 @@ mod tests { first_window_ready.wait().await; assert_eq!(current.load(Ordering::SeqCst), 4); + release_first_window.wait().await; assert_eq!(lookup.await.expect("lookup task"), None); assert_eq!(maximum.load(Ordering::SeqCst), 4); } From 951f1cb27a25706abd32b73762f772f2671c6de3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BD=A6=E7=BA=BE?= Date: Wed, 5 Aug 2026 11:41:05 +0800 Subject: [PATCH 5/5] test(p2p): bound buffered lookup synchronization --- src/p2p/iroh/transport.rs | 32 ++++++++++++++++++++++++++------ 1 file changed, 26 insertions(+), 6 deletions(-) diff --git a/src/p2p/iroh/transport.rs b/src/p2p/iroh/transport.rs index da296d517..7a3bd2268 100644 --- a/src/p2p/iroh/transport.rs +++ b/src/p2p/iroh/transport.rs @@ -829,15 +829,25 @@ mod tests { } }); - first_started.notified().await; - second_finished.notified().await; + tokio::time::timeout(TEST_TIMEOUT, async { + first_started.notified().await; + second_finished.notified().await; + }) + .await + .expect("initial buffered lookups should complete"); assert!( !lookup.is_finished(), "a lower-priority result must wait for earlier candidates" ); release_first.notify_one(); - assert_eq!(lookup.await.expect("lookup task"), Some("first")); + assert_eq!( + tokio::time::timeout(TEST_TIMEOUT, lookup) + .await + .expect("ordered lookup should complete") + .expect("lookup task"), + Some("first") + ); } #[tokio::test] @@ -873,10 +883,20 @@ mod tests { } }); - first_window_ready.wait().await; + tokio::time::timeout(TEST_TIMEOUT, first_window_ready.wait()) + .await + .expect("initial concurrency window should fill"); assert_eq!(current.load(Ordering::SeqCst), 4); - release_first_window.wait().await; - assert_eq!(lookup.await.expect("lookup task"), None); + tokio::time::timeout(TEST_TIMEOUT, release_first_window.wait()) + .await + .expect("initial concurrency window should be released"); + assert_eq!( + tokio::time::timeout(TEST_TIMEOUT, lookup) + .await + .expect("bounded lookup should complete") + .expect("lookup task"), + None + ); assert_eq!(maximum.load(Ordering::SeqCst), 4); }