Skip to content
Open
Changes from 4 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
131 changes: 119 additions & 12 deletions src/p2p/iroh/transport.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
use std::collections::HashSet;
use std::future::Future;
use std::net::SocketAddr;
use std::ops::Range;
use std::path::Path;
Expand Down Expand Up @@ -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:";
Expand Down Expand Up @@ -220,32 +222,31 @@ impl IrohBlobsP2pTransport {
peers: Vec<P2pPeer>,
key: &P2pArtifactKey,
) -> Option<P2pArtifactDescriptor> {
// 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!(
peer = %peer.node_id,
error = %err,
"P2P artifact catalog lookup failed; trying remaining peers"
);
None
}
}
}
None
})
.await
}

async fn get_local(&self, key: &P2pArtifactKey) -> Option<P2pArtifactDescriptor> {
Expand Down Expand Up @@ -720,6 +721,24 @@ fn blob_hash_from_descriptor(descriptor: &P2pArtifactDescriptor) -> Result<iroh_
})
}

async fn first_some_buffered_in_order<I, F, Fut, T>(
items: impl IntoIterator<Item = I>,
concurrency: usize,
lookup: F,
) -> Option<T>
where
F: FnMut(I) -> Fut,
Fut: Future<Output = Option<T>>,
{
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;
}
}
None
}

fn resolve_remote_local_providers(
mut descriptor: P2pArtifactDescriptor,
response_peer: &P2pPeer,
Expand Down Expand Up @@ -765,6 +784,9 @@ fn gated_gc_config(interval: Duration, pending_gc: Arc<AtomicBool>) -> 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;
Expand All @@ -773,6 +795,91 @@ 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;
Comment thread
catyans marked this conversation as resolved.
Outdated
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 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::<()>
}
})
.await
}
});

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);
}

async fn collect_range_stream(mut stream: P2pByteStream) -> Result<Vec<u8>> {
let mut out = Vec::new();
while let Some(chunk) = stream.next().await {
Expand Down
Loading