From 5f6f7eb0330108455c4a53d24265e40c4eefeba9 Mon Sep 17 00:00:00 2001 From: "Toast (gastown)" Date: Fri, 19 Jun 2026 05:21:30 +0000 Subject: [PATCH 1/3] connect: use deterministic BTreeSet::first for peer/IP selection MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replace s.iter().next() with s.first() in connect_name and connect_imid to make the selection strategy explicit and documented. BTreeSet::first() returns the minimum element per Ord, which is deterministic across runs. Add test_connect_name_multi_imid to exercise name → multi-IMID resolution. Add test_connect_imid_multi_ip to exercise IMID → multi-IP resolution. --- src/connect.rs | 133 ++++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 131 insertions(+), 2 deletions(-) diff --git a/src/connect.rs b/src/connect.rs index e944469..306502e 100644 --- a/src/connect.rs +++ b/src/connect.rs @@ -97,6 +97,9 @@ impl IntermeshClient { } /// Connect to a mesh name. + /// + /// If the name resolves to multiple IMIDs, the IMID with the minimum + /// value (per `Ord`) is selected deterministically via [`BTreeSet::first`]. pub(crate) async fn connect_name( &mut self, name: &Name, @@ -107,7 +110,7 @@ impl IntermeshClient { .derivation() .name_to_imid .get(name) - .and_then(|s| s.iter().next()) + .and_then(|s| s.first()) .cloned() .ok_or_else(|| anyhow!("failed to resolve name to IMID: {name}"))?; @@ -115,6 +118,9 @@ impl IntermeshClient { } /// Connect to an IMID. + /// + /// If the IMID resolves to multiple IPs, the IP with the minimum value + /// (per `Ord`) is selected deterministically via [`BTreeSet::first`]. async fn connect_imid(&mut self, imid: &Imid, port: u16) -> Result> { // Trust engine is authoritative; bootstrap hint is fallback for bootstrap. let te_ip = self @@ -122,7 +128,7 @@ impl IntermeshClient { .derivation() .imid_to_ip .get(imid) - .and_then(|s| s.iter().next()) + .and_then(|s| s.first()) .copied(); let hint_ip = self .bootstrap_hint @@ -474,6 +480,129 @@ mod tests { server_task.await.assert(); } + #[tokio::test] + async fn test_connect_name_multi_imid() { + let mut fix = TestFixture::new(); + + let server_keypair = fix.keypair("server").clone(); + let server_imid = server_keypair.to_imid(); + + let client_keypair = fix.keypair("client").clone(); + let client_imid = client_keypair.to_imid(); + + let listener = TcpListener::bind("127.0.0.1:0").await.assert(); + let server_addr = listener.local_addr().assert(); + let (shutdown_tx, mut shutdown_rx) = mpsc::channel::<()>(1); + + let server_task = tokio::spawn(async move { + let server_verifier = Arc::new(IntermeshVerifier::new_permissive()); + let server_stream = + intermesh_server_stream(&server_keypair, server_verifier, listener); + tokio::pin!(server_stream); + loop { + tokio::select! { + result = server_stream.next() => { + let mut tls_stream = result.assert().assert(); + let mut buf = [0u8; 1024]; + let n = tls_stream.read(&mut buf).await.assert(); + tls_stream.write_all(&buf[..n]).await.assert(); + } + _ = shutdown_rx.recv() => break, + } + } + }); + + let trust_engine = Arc::new(TrustEngine::new(client_imid)); + let other_imid = fix.imid("other"); + trust_engine.set_ip_map(BTreeMap::from([ + (server_imid.clone(), vec![server_addr.ip()]), + (other_imid.clone(), vec![server_addr.ip()]), + ])); + trust_engine.set_name_map(BTreeMap::from([( + "multi-imid.test".parse().assert(), + vec![server_imid.clone(), other_imid], + )])); + + let client_verifier = Arc::new(IntermeshVerifier::new_permissive()); + let mut client = + IntermeshClient::new(&client_keypair, client_verifier, trust_engine); + + for _ in 0..10 { + let mut stream = client + .connect("multi-imid.test", server_addr.port()) + .await + .assert(); + stream.write_all(b"multi-imid").await.assert(); + let mut buf = [0u8; 1024]; + let n = stream.read(&mut buf).await.assert(); + assert_eq!(&buf[..n], b"multi-imid"); + } + + shutdown_tx.send(()).await.assert(); + server_task.await.assert(); + } + + #[tokio::test] + async fn test_connect_imid_multi_ip() { + let mut fix = TestFixture::new(); + + let server_keypair = fix.keypair("server").clone(); + let server_imid = server_keypair.to_imid(); + + let client_keypair = fix.keypair("client").clone(); + let client_imid = client_keypair.to_imid(); + + let listener = TcpListener::bind("127.0.0.1:0").await.assert(); + let server_addr = listener.local_addr().assert(); + let (shutdown_tx, mut shutdown_rx) = mpsc::channel::<()>(1); + + let server_task = tokio::spawn(async move { + let server_verifier = Arc::new(IntermeshVerifier::new_permissive()); + let server_stream = + intermesh_server_stream(&server_keypair, server_verifier, listener); + tokio::pin!(server_stream); + loop { + tokio::select! { + result = server_stream.next() => { + let mut tls_stream = result.assert().assert(); + let mut buf = [0u8; 1024]; + let n = tls_stream.read(&mut buf).await.assert(); + tls_stream.write_all(&buf[..n]).await.assert(); + } + _ = shutdown_rx.recv() => break, + } + } + }); + + let trust_engine = Arc::new(TrustEngine::new(client_imid)); + // Include a second IP that sorts after 127.0.0.1 in BTreeSet order + // to exercise the multi-IP selection path. BTreeSet::first() + // deterministically returns the minimum (127.0.0.1) which matches + // the server. + trust_engine.set_ip_map(BTreeMap::from([( + server_imid.clone(), + vec![server_addr.ip(), "127.0.0.2".parse().assert()], + )])); + + let client_verifier = Arc::new(IntermeshVerifier::new_permissive()); + let mut client = + IntermeshClient::new(&client_keypair, client_verifier, trust_engine); + + for _ in 0..10 { + let mut stream = client + .connect_imid(&server_imid, server_addr.port()) + .await + .assert(); + stream.write_all(b"multi-ip").await.assert(); + let mut buf = [0u8; 1024]; + let n = stream.read(&mut buf).await.assert(); + assert_eq!(&buf[..n], b"multi-ip"); + } + + shutdown_tx.send(()).await.assert(); + server_task.await.assert(); + } + #[tokio::test(flavor = "multi_thread")] async fn test_tonic_integration() { timeout(Duration::from_secs(30), async { From 805cf967d5fc27c59945c62fd92e34995a22c1fb Mon Sep 17 00:00:00 2001 From: Tim Anglade Date: Fri, 19 Jun 2026 00:46:46 -0700 Subject: [PATCH 2/3] connect: fix intra-doc link and flaky multi-imid test --- src/connect.rs | 26 +++++++++++--------------- 1 file changed, 11 insertions(+), 15 deletions(-) diff --git a/src/connect.rs b/src/connect.rs index 306502e..688a98c 100644 --- a/src/connect.rs +++ b/src/connect.rs @@ -99,7 +99,7 @@ impl IntermeshClient { /// Connect to a mesh name. /// /// If the name resolves to multiple IMIDs, the IMID with the minimum - /// value (per `Ord`) is selected deterministically via [`BTreeSet::first`]. + /// value (per `Ord`) is selected deterministically. pub(crate) async fn connect_name( &mut self, name: &Name, @@ -120,7 +120,7 @@ impl IntermeshClient { /// Connect to an IMID. /// /// If the IMID resolves to multiple IPs, the IP with the minimum value - /// (per `Ord`) is selected deterministically via [`BTreeSet::first`]. + /// (per `Ord`) is selected deterministically. async fn connect_imid(&mut self, imid: &Imid, port: u16) -> Result> { // Trust engine is authoritative; bootstrap hint is fallback for bootstrap. let te_ip = self @@ -482,12 +482,12 @@ mod tests { #[tokio::test] async fn test_connect_name_multi_imid() { - let mut fix = TestFixture::new(); - - let server_keypair = fix.keypair("server").clone(); + // Use deterministic keypairs so "server" sorts before "zzz", + // guaranteeing server_imid < other_imid in BTreeSet order. + let server_keypair = ImidKeypair::test_keypair("server"); let server_imid = server_keypair.to_imid(); - let client_keypair = fix.keypair("client").clone(); + let client_keypair = ImidKeypair::test_keypair("client"); let client_imid = client_keypair.to_imid(); let listener = TcpListener::bind("127.0.0.1:0").await.assert(); @@ -496,8 +496,7 @@ mod tests { let server_task = tokio::spawn(async move { let server_verifier = Arc::new(IntermeshVerifier::new_permissive()); - let server_stream = - intermesh_server_stream(&server_keypair, server_verifier, listener); + let server_stream = intermesh_server_stream(&server_keypair, server_verifier, listener); tokio::pin!(server_stream); loop { tokio::select! { @@ -513,7 +512,7 @@ mod tests { }); let trust_engine = Arc::new(TrustEngine::new(client_imid)); - let other_imid = fix.imid("other"); + let other_imid = ImidKeypair::test_keypair("zzz").to_imid(); trust_engine.set_ip_map(BTreeMap::from([ (server_imid.clone(), vec![server_addr.ip()]), (other_imid.clone(), vec![server_addr.ip()]), @@ -524,8 +523,7 @@ mod tests { )])); let client_verifier = Arc::new(IntermeshVerifier::new_permissive()); - let mut client = - IntermeshClient::new(&client_keypair, client_verifier, trust_engine); + let mut client = IntermeshClient::new(&client_keypair, client_verifier, trust_engine); for _ in 0..10 { let mut stream = client @@ -558,8 +556,7 @@ mod tests { let server_task = tokio::spawn(async move { let server_verifier = Arc::new(IntermeshVerifier::new_permissive()); - let server_stream = - intermesh_server_stream(&server_keypair, server_verifier, listener); + let server_stream = intermesh_server_stream(&server_keypair, server_verifier, listener); tokio::pin!(server_stream); loop { tokio::select! { @@ -585,8 +582,7 @@ mod tests { )])); let client_verifier = Arc::new(IntermeshVerifier::new_permissive()); - let mut client = - IntermeshClient::new(&client_keypair, client_verifier, trust_engine); + let mut client = IntermeshClient::new(&client_keypair, client_verifier, trust_engine); for _ in 0..10 { let mut stream = client From d653b763ae504d4f0f732c6e84a7048b1f83afdd Mon Sep 17 00:00:00 2001 From: Tim Anglade Date: Sat, 20 Jun 2026 16:00:26 -0700 Subject: [PATCH 3/3] connect: extract pure resolve_name_to_imid and resolve_imid_to_ip helpers Replace the heavy integration tests for multi-IMID and multi-IP selection with simple unit tests on the pure helper functions. --- src/connect.rs | 204 +++++++++++++++++++------------------------------ 1 file changed, 79 insertions(+), 125 deletions(-) diff --git a/src/connect.rs b/src/connect.rs index 688a98c..6a4d792 100644 --- a/src/connect.rs +++ b/src/connect.rs @@ -8,6 +8,7 @@ use futures::{future::BoxFuture, Stream}; use hyper_util::rt::TokioIo; use rustls::pki_types::PrivateKeyDer; use rustls::{ClientConfig, ServerConfig}; +use std::collections::{BTreeMap, BTreeSet}; use std::io; use std::net::IpAddr; use std::pin::Pin; @@ -105,15 +106,7 @@ impl IntermeshClient { name: &Name, port: u16, ) -> Result> { - let imid = self - .te - .derivation() - .name_to_imid - .get(name) - .and_then(|s| s.first()) - .cloned() - .ok_or_else(|| anyhow!("failed to resolve name to IMID: {name}"))?; - + let imid = resolve_name_to_imid(name, &self.te.derivation().name_to_imid)?; self.connect_imid(&imid, port).await } @@ -122,14 +115,7 @@ impl IntermeshClient { /// If the IMID resolves to multiple IPs, the IP with the minimum value /// (per `Ord`) is selected deterministically. async fn connect_imid(&mut self, imid: &Imid, port: u16) -> Result> { - // Trust engine is authoritative; bootstrap hint is fallback for bootstrap. - let te_ip = self - .te - .derivation() - .imid_to_ip - .get(imid) - .and_then(|s| s.first()) - .copied(); + let te_ip = resolve_imid_to_ip(imid, &self.te.derivation().imid_to_ip).ok(); let hint_ip = self .bootstrap_hint .take_if(|(hint_imid, _)| hint_imid == imid) @@ -182,6 +168,32 @@ impl Service for IntermeshClient { } } +/// Resolve a mesh name to an IMID, deterministically selecting the minimum +/// IMID (per `Ord`) when the name maps to multiple candidates. +fn resolve_name_to_imid( + name: &Name, + name_to_imid: &BTreeMap>, +) -> Result { + name_to_imid + .get(name) + .and_then(|s| s.first()) + .cloned() + .ok_or_else(|| anyhow!("failed to resolve name to IMID: {name}")) +} + +/// Resolve an IMID to an IP address, deterministically selecting the minimum +/// IP (per `Ord`) when the IMID maps to multiple candidates. +fn resolve_imid_to_ip( + imid: &Imid, + imid_to_ip: &BTreeMap>, +) -> Result { + imid_to_ip + .get(imid) + .and_then(|s| s.first()) + .copied() + .ok_or_else(|| anyhow!("failed to find IP for IMID: {imid}")) +} + /// Creates a stream of TLS-wrapped TCP connections for serving gRPC requests. /// /// This function creates a `TlsStream` that uses the `IntermeshVerifier` to @@ -300,7 +312,7 @@ mod tests { use crate::proto::intermesh::GossipUpdate; use crate::test_utils::TestFixture; use futures::TryStreamExt; - use std::collections::BTreeMap; + use std::collections::{BTreeMap, BTreeSet}; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::sync::mpsc; use tokio::time::{timeout, Duration}; @@ -480,123 +492,65 @@ mod tests { server_task.await.assert(); } - #[tokio::test] - async fn test_connect_name_multi_imid() { - // Use deterministic keypairs so "server" sorts before "zzz", - // guaranteeing server_imid < other_imid in BTreeSet order. - let server_keypair = ImidKeypair::test_keypair("server"); - let server_imid = server_keypair.to_imid(); - - let client_keypair = ImidKeypair::test_keypair("client"); - let client_imid = client_keypair.to_imid(); - - let listener = TcpListener::bind("127.0.0.1:0").await.assert(); - let server_addr = listener.local_addr().assert(); - let (shutdown_tx, mut shutdown_rx) = mpsc::channel::<()>(1); - - let server_task = tokio::spawn(async move { - let server_verifier = Arc::new(IntermeshVerifier::new_permissive()); - let server_stream = intermesh_server_stream(&server_keypair, server_verifier, listener); - tokio::pin!(server_stream); - loop { - tokio::select! { - result = server_stream.next() => { - let mut tls_stream = result.assert().assert(); - let mut buf = [0u8; 1024]; - let n = tls_stream.read(&mut buf).await.assert(); - tls_stream.write_all(&buf[..n]).await.assert(); - } - _ = shutdown_rx.recv() => break, - } - } - }); - - let trust_engine = Arc::new(TrustEngine::new(client_imid)); - let other_imid = ImidKeypair::test_keypair("zzz").to_imid(); - trust_engine.set_ip_map(BTreeMap::from([ - (server_imid.clone(), vec![server_addr.ip()]), - (other_imid.clone(), vec![server_addr.ip()]), - ])); - trust_engine.set_name_map(BTreeMap::from([( - "multi-imid.test".parse().assert(), - vec![server_imid.clone(), other_imid], - )])); - - let client_verifier = Arc::new(IntermeshVerifier::new_permissive()); - let mut client = IntermeshClient::new(&client_keypair, client_verifier, trust_engine); + #[test] + fn test_resolve_name_to_imid_first() { + let a = ImidKeypair::test_keypair("a").to_imid(); + let b = ImidKeypair::test_keypair("b").to_imid(); + let (min, max) = if a < b { + (a.clone(), b.clone()) + } else { + (b.clone(), a.clone()) + }; - for _ in 0..10 { - let mut stream = client - .connect("multi-imid.test", server_addr.port()) - .await - .assert(); - stream.write_all(b"multi-imid").await.assert(); - let mut buf = [0u8; 1024]; - let n = stream.read(&mut buf).await.assert(); - assert_eq!(&buf[..n], b"multi-imid"); - } + let map = BTreeMap::from([( + "test.local".parse().assert(), + BTreeSet::from([min.clone(), max]), + )]); - shutdown_tx.send(()).await.assert(); - server_task.await.assert(); + let result = resolve_name_to_imid(&"test.local".parse().assert(), &map).assert(); + assert_eq!( + result, min, + "should select the minimum IMID per Ord (BTreeSet::first)" + ); } - #[tokio::test] - async fn test_connect_imid_multi_ip() { - let mut fix = TestFixture::new(); - - let server_keypair = fix.keypair("server").clone(); - let server_imid = server_keypair.to_imid(); - - let client_keypair = fix.keypair("client").clone(); - let client_imid = client_keypair.to_imid(); + #[test] + fn test_resolve_name_to_imid_not_found() { + let a = ImidKeypair::test_keypair("a").to_imid(); + let map = BTreeMap::from([("test.local".parse().assert(), BTreeSet::from([a]))]); - let listener = TcpListener::bind("127.0.0.1:0").await.assert(); - let server_addr = listener.local_addr().assert(); - let (shutdown_tx, mut shutdown_rx) = mpsc::channel::<()>(1); + let err = resolve_name_to_imid(&"other.local".parse().assert(), &map).unwrap_err(); + assert!( + err.to_string().contains("failed to resolve name to IMID"), + "unexpected error: {err}" + ); + } - let server_task = tokio::spawn(async move { - let server_verifier = Arc::new(IntermeshVerifier::new_permissive()); - let server_stream = intermesh_server_stream(&server_keypair, server_verifier, listener); - tokio::pin!(server_stream); - loop { - tokio::select! { - result = server_stream.next() => { - let mut tls_stream = result.assert().assert(); - let mut buf = [0u8; 1024]; - let n = tls_stream.read(&mut buf).await.assert(); - tls_stream.write_all(&buf[..n]).await.assert(); - } - _ = shutdown_rx.recv() => break, - } - } - }); + #[test] + fn test_resolve_imid_to_ip_first() { + let imid = ImidKeypair::test_keypair("test").to_imid(); + let ip: IpAddr = "127.0.0.1".parse().assert(); - let trust_engine = Arc::new(TrustEngine::new(client_imid)); - // Include a second IP that sorts after 127.0.0.1 in BTreeSet order - // to exercise the multi-IP selection path. BTreeSet::first() - // deterministically returns the minimum (127.0.0.1) which matches - // the server. - trust_engine.set_ip_map(BTreeMap::from([( - server_imid.clone(), - vec![server_addr.ip(), "127.0.0.2".parse().assert()], - )])); + let map = BTreeMap::from([(imid.clone(), BTreeSet::from([ip]))]); - let client_verifier = Arc::new(IntermeshVerifier::new_permissive()); - let mut client = IntermeshClient::new(&client_keypair, client_verifier, trust_engine); + let result = resolve_imid_to_ip(&imid, &map).assert(); + assert_eq!(result, ip); + } - for _ in 0..10 { - let mut stream = client - .connect_imid(&server_imid, server_addr.port()) - .await - .assert(); - stream.write_all(b"multi-ip").await.assert(); - let mut buf = [0u8; 1024]; - let n = stream.read(&mut buf).await.assert(); - assert_eq!(&buf[..n], b"multi-ip"); - } + #[test] + fn test_resolve_imid_to_ip_not_found() { + let imid = ImidKeypair::test_keypair("test").to_imid(); + let other = ImidKeypair::test_keypair("other").to_imid(); + let map = BTreeMap::from([( + imid, + BTreeSet::from(["127.0.0.1".parse::().assert()]), + )]); - shutdown_tx.send(()).await.assert(); - server_task.await.assert(); + let err = resolve_imid_to_ip(&other, &map).unwrap_err(); + assert!( + err.to_string().contains("failed to find IP for IMID"), + "unexpected error: {err}" + ); } #[tokio::test(flavor = "multi_thread")]