From 637b817a33d3c25d5b639726fa30729bc0aff00b Mon Sep 17 00:00:00 2001 From: Lewis Date: Wed, 10 Jun 2026 09:28:24 +0300 Subject: [PATCH] ripple: cluster-key oomf authentication Lewis: May this revision serve well! --- Cargo.lock | 3 + Cargo.toml | 1 + crates/tranquil-ripple/Cargo.toml | 3 + crates/tranquil-ripple/src/engine.rs | 2 +- crates/tranquil-ripple/src/transport.rs | 431 ++++++++++++++++++++++-- 5 files changed, 405 insertions(+), 35 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index c842fc3..d4a02ec 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8000,6 +8000,7 @@ dependencies = [ "bincode 2.0.1", "foca", "futures", + "hkdf", "metrics", "parking_lot", "quinn", @@ -8007,6 +8008,7 @@ dependencies = [ "rcgen", "rustls 0.23.37", "serde", + "sha2", "thiserror 2.0.18", "tokio", "tokio-util", @@ -8015,6 +8017,7 @@ dependencies = [ "tranquil-config", "tranquil-infra", "uuid", + "x509-parser 0.18.1", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 978097a..fc77ee0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -141,6 +141,7 @@ urlencoding = "2.1" uuid = { version = "1.19", features = ["v4", "v5", "v7", "fast-rng", "serde"] } webauthn-rs = { version = "0.5", features = ["danger-allow-state-serialisation", "danger-user-presence-only-security-keys", "conditional-ui"] } webauthn-rs-proto = "0.5" +x509-parser = "0.18" zip = { version = "7.0", default-features = false, features = ["deflate"] } ciborium = "0.2" diff --git a/crates/tranquil-ripple/Cargo.toml b/crates/tranquil-ripple/Cargo.toml index 7181d20..a6feafa 100644 --- a/crates/tranquil-ripple/Cargo.toml +++ b/crates/tranquil-ripple/Cargo.toml @@ -12,6 +12,7 @@ async-trait = { workspace = true } backon = { workspace = true } bincode = { workspace = true } foca = { workspace = true } +hkdf = { workspace = true } metrics = { workspace = true } parking_lot = { workspace = true } quinn = { workspace = true } @@ -19,10 +20,12 @@ rand = "0.9" rcgen = { workspace = true } rustls = { workspace = true } serde = { workspace = true } +sha2 = { workspace = true } thiserror = { workspace = true } tokio = { workspace = true, features = ["net", "io-util", "sync", "time"] } tokio-util = { workspace = true } tracing = { workspace = true } +x509-parser = { workspace = true } [dev-dependencies] futures = { workspace = true } diff --git a/crates/tranquil-ripple/src/engine.rs b/crates/tranquil-ripple/src/engine.rs index fd404e5..334bfe8 100644 --- a/crates/tranquil-ripple/src/engine.rs +++ b/crates/tranquil-ripple/src/engine.rs @@ -22,7 +22,7 @@ impl RippleEngine { let store = Arc::new(ShardedCrdtStore::new(config.machine_id)); let (transport, incoming_rx) = - Transport::bind(config.bind_addr, shutdown.clone()) + Transport::bind(config.bind_addr, None, shutdown.clone()) .await .map_err(|e| RippleStartError::Bind(e.to_string()))?; diff --git a/crates/tranquil-ripple/src/transport.rs b/crates/tranquil-ripple/src/transport.rs index dcb78ba..60d9ed1 100644 --- a/crates/tranquil-ripple/src/transport.rs +++ b/crates/tranquil-ripple/src/transport.rs @@ -1,10 +1,12 @@ use backon::{ExponentialBuilder, Retryable}; +use hkdf::Hkdf; use quinn::crypto::rustls::{QuicClientConfig, QuicServerConfig}; use quinn::{ ClientConfig, Connection, Endpoint, IdleTimeout, RecvStream, ServerConfig, TransportConfig, VarInt, }; use rustls::pki_types::{CertificateDer, PrivateKeyDer, PrivatePkcs8KeyDer}; +use sha2::Sha256; use std::collections::HashMap; use std::net::{IpAddr, SocketAddr}; use std::sync::Arc; @@ -32,6 +34,20 @@ const WRITE_TIMEOUT: Duration = Duration::from_secs(10); const READ_TIMEOUT: Duration = Duration::from_secs(30); const RIPPLE_ALPN: &[u8] = b"ripple/1"; const RIPPLE_SERVER_NAME: &str = "ripple"; +const CLUSTER_HKDF_SALT: &[u8] = b"ripple-cluster-identity-v1"; +const CLUSTER_HKDF_INFO: &[u8] = b"ripple-node-ed25519"; +const ED25519_PKCS8_V1_PREFIX: [u8; 16] = [ + 0x30, 0x2e, 0x02, 0x01, 0x00, 0x30, 0x05, 0x06, 0x03, 0x2b, 0x65, 0x70, 0x04, 0x22, 0x04, 0x20, +]; + +#[derive(Clone)] +pub struct ClusterKey(Arc); + +impl ClusterKey { + pub fn new(secret: impl Into>) -> Self { + Self(secret.into()) + } +} struct NodeIdentity { cert: CertificateDer<'static>, @@ -98,11 +114,28 @@ pub struct Transport { impl Transport { pub async fn bind( addr: SocketAddr, + cluster_key: Option, shutdown: CancellationToken, ) -> Result<(Self, mpsc::Receiver), std::io::Error> { - let server_config = build_server_config() + let identity = match &cluster_key { + Some(key) => Some( + derive_identity(key) + .map_err(|e| std::io::Error::other(format!("ripple cluster identity: {e}")))?, + ), + None => { + if !addr.ip().is_loopback() { + tracing::warn!( + addr = %addr, + "ripple peer authentication disabled, set RIPPLE_CLUSTER_KEY to authenticate the cluster" + ); + } + None + } + }; + + let server_config = build_server_config(identity.as_ref()) .map_err(|e| std::io::Error::other(format!("ripple server config: {e}")))?; - let client_config = build_client_config() + let client_config = build_client_config(identity.as_ref()) .map_err(|e| std::io::Error::other(format!("ripple client config: {e}")))?; let mut endpoint = Endpoint::server(server_config, addr)?; @@ -169,7 +202,11 @@ impl Transport { endpoint.close(0u32.into(), b"shutdown"); }); - tracing::info!(addr = %local_addr, "ripple quic transport bound"); + tracing::info!( + addr = %local_addr, + authenticated = identity.is_some(), + "ripple quic transport bound" + ); Ok((transport, incoming_rx)) } @@ -569,6 +606,30 @@ fn transport_config() -> TransportConfig { type BoxError = Box; +fn derive_identity(key: &ClusterKey) -> Result { + let hk = Hkdf::::new(Some(CLUSTER_HKDF_SALT), key.0.as_bytes()); + let mut seed = [0u8; 32]; + hk.expand(CLUSTER_HKDF_INFO, &mut seed) + .map_err(|_| "hkdf expand failed")?; + + let mut pkcs8 = Vec::with_capacity(ED25519_PKCS8_V1_PREFIX.len() + seed.len()); + pkcs8.extend_from_slice(&ED25519_PKCS8_V1_PREFIX); + pkcs8.extend_from_slice(&seed); + let pkcs8 = PrivatePkcs8KeyDer::from(pkcs8); + + let keypair = rcgen::KeyPair::from_pkcs8_der_and_sign_algo(&pkcs8, &rcgen::PKCS_ED25519)?; + let mut params = rcgen::CertificateParams::new(vec![RIPPLE_SERVER_NAME.to_string()])?; + params.serial_number = Some(rcgen::SerialNumber::from(1u64)); + params.not_before = rcgen::date_time_ymd(2024, 1, 1); + params.not_after = rcgen::date_time_ymd(2099, 1, 1); + let cert = params.self_signed(&keypair)?; + + Ok(NodeIdentity { + cert: cert.der().clone(), + key: PrivateKeyDer::Pkcs8(PrivatePkcs8KeyDer::from(keypair.serialize_der())), + }) +} + fn ephemeral_identity() -> Result { let cert = rcgen::generate_simple_self_signed(vec![RIPPLE_SERVER_NAME.to_string()])?; Ok(NodeIdentity { @@ -577,14 +638,24 @@ fn ephemeral_identity() -> Result { }) } -fn build_server_config() -> Result { +fn build_server_config(identity: Option<&NodeIdentity>) -> Result { let provider = Arc::new(rustls::crypto::ring::default_provider()); - let builder = rustls::ServerConfig::builder_with_provider(provider) + let builder = rustls::ServerConfig::builder_with_provider(provider.clone()) .with_protocol_versions(&[&rustls::version::TLS13])?; - let ephemeral = ephemeral_identity()?; - let mut crypto = builder - .with_no_client_auth() - .with_single_cert(vec![ephemeral.cert], ephemeral.key)?; + let mut crypto = match identity { + Some(id) => { + let verifier = Arc::new(PinnedPeer::new(provider, &id.cert)?); + builder + .with_client_cert_verifier(verifier) + .with_single_cert(vec![id.cert.clone()], id.key.clone_key())? + } + None => { + let ephemeral = ephemeral_identity()?; + builder + .with_no_client_auth() + .with_single_cert(vec![ephemeral.cert], ephemeral.key)? + } + }; crypto.alpn_protocols = vec![RIPPLE_ALPN.to_vec()]; let mut config = ServerConfig::with_crypto(Arc::new(QuicServerConfig::try_from(crypto)?)); @@ -592,14 +663,20 @@ fn build_server_config() -> Result { Ok(config) } -fn build_client_config() -> Result { +fn build_client_config(identity: Option<&NodeIdentity>) -> Result { let provider = Arc::new(rustls::crypto::ring::default_provider()); - let verifier = Arc::new(SkipServerVerification::new(provider.clone())); - let mut crypto = rustls::ClientConfig::builder_with_provider(provider) + let verifier: Arc = match identity { + Some(id) => Arc::new(PinnedPeer::new(provider.clone(), &id.cert)?), + None => Arc::new(SkipServerVerification::new(provider.clone())), + }; + let builder = rustls::ClientConfig::builder_with_provider(provider) .with_protocol_versions(&[&rustls::version::TLS13])? .dangerous() - .with_custom_certificate_verifier(verifier) - .with_no_client_auth(); + .with_custom_certificate_verifier(verifier); + let mut crypto = match identity { + Some(id) => builder.with_client_auth_cert(vec![id.cert.clone()], id.key.clone_key())?, + None => builder.with_no_client_auth(), + }; crypto.alpn_protocols = vec![RIPPLE_ALPN.to_vec()]; let mut config = ClientConfig::new(Arc::new(QuicClientConfig::try_from(crypto)?)); @@ -661,6 +738,142 @@ impl rustls::client::danger::ServerCertVerifier for SkipServerVerification { } } +fn spki_der(cert: &CertificateDer<'_>) -> Result, BoxError> { + let (_, parsed) = x509_parser::parse_x509_certificate(cert.as_ref()) + .map_err(|e| format!("parse certificate: {e}"))?; + Ok(parsed.public_key().raw.to_vec()) +} + +#[derive(Debug)] +struct PinnedPeer { + provider: Arc, + expected_spki: Vec, + root_hint_subjects: Vec, +} + +impl PinnedPeer { + fn new( + provider: Arc, + expected: &CertificateDer<'_>, + ) -> Result { + Ok(Self { + provider, + expected_spki: spki_der(expected)?, + root_hint_subjects: Vec::new(), + }) + } + + fn check(&self, presented: &CertificateDer<'_>) -> Result<(), rustls::Error> { + let presented_spki = spki_der(presented).map_err(|_| { + rustls::Error::InvalidCertificate(rustls::CertificateError::BadEncoding) + })?; + match presented_spki == self.expected_spki { + true => Ok(()), + false => Err(rustls::Error::InvalidCertificate( + rustls::CertificateError::ApplicationVerificationFailure, + )), + } + } +} + +impl rustls::client::danger::ServerCertVerifier for PinnedPeer { + fn verify_server_cert( + &self, + end_entity: &CertificateDer<'_>, + _intermediates: &[CertificateDer<'_>], + _server_name: &rustls::pki_types::ServerName<'_>, + _ocsp_response: &[u8], + _now: rustls::pki_types::UnixTime, + ) -> Result { + self.check(end_entity)?; + Ok(rustls::client::danger::ServerCertVerified::assertion()) + } + + fn verify_tls12_signature( + &self, + message: &[u8], + cert: &CertificateDer<'_>, + dss: &rustls::DigitallySignedStruct, + ) -> Result { + rustls::crypto::verify_tls12_signature( + message, + cert, + dss, + &self.provider.signature_verification_algorithms, + ) + } + + fn verify_tls13_signature( + &self, + message: &[u8], + cert: &CertificateDer<'_>, + dss: &rustls::DigitallySignedStruct, + ) -> Result { + rustls::crypto::verify_tls13_signature( + message, + cert, + dss, + &self.provider.signature_verification_algorithms, + ) + } + + fn supported_verify_schemes(&self) -> Vec { + self.provider + .signature_verification_algorithms + .supported_schemes() + } +} + +impl rustls::server::danger::ClientCertVerifier for PinnedPeer { + fn root_hint_subjects(&self) -> &[rustls::DistinguishedName] { + &self.root_hint_subjects + } + + fn verify_client_cert( + &self, + end_entity: &CertificateDer<'_>, + _intermediates: &[CertificateDer<'_>], + _now: rustls::pki_types::UnixTime, + ) -> Result { + self.check(end_entity)?; + Ok(rustls::server::danger::ClientCertVerified::assertion()) + } + + fn verify_tls12_signature( + &self, + message: &[u8], + cert: &CertificateDer<'_>, + dss: &rustls::DigitallySignedStruct, + ) -> Result { + rustls::crypto::verify_tls12_signature( + message, + cert, + dss, + &self.provider.signature_verification_algorithms, + ) + } + + fn verify_tls13_signature( + &self, + message: &[u8], + cert: &CertificateDer<'_>, + dss: &rustls::DigitallySignedStruct, + ) -> Result { + rustls::crypto::verify_tls13_signature( + message, + cert, + dss, + &self.provider.signature_verification_algorithms, + ) + } + + fn supported_verify_schemes(&self) -> Vec { + self.provider + .signature_verification_algorithms + .supported_schemes() + } +} + #[cfg(test)] mod tests { use super::*; @@ -669,11 +882,11 @@ mod tests { async fn quic_frame_roundtrip() { let shutdown = CancellationToken::new(); let (sender, _rx_sender) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind sender"); let (receiver, mut rx_receiver) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind receiver"); let target = receiver.local_addr(); @@ -696,11 +909,11 @@ mod tests { async fn distinct_channels_roundtrip() { let shutdown = CancellationToken::new(); let (sender, _rx_sender) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind sender"); let (receiver, mut rx_receiver) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind receiver"); let target = receiver.local_addr(); @@ -722,15 +935,165 @@ mod tests { shutdown.cancel(); } + #[test] + fn cluster_identity_is_deterministic() { + let a = derive_identity(&ClusterKey::new("periwinkle")).expect("derive a"); + let b = derive_identity(&ClusterKey::new("periwinkle")).expect("derive b"); + assert_eq!(a.cert.as_ref(), b.cert.as_ref()); + let c = derive_identity(&ClusterKey::new("limpet")).expect("derive c"); + assert_ne!(a.cert.as_ref(), c.cert.as_ref()); + } + + #[test] + fn spki_pin_accepts_same_key_rejects_others() { + let provider = Arc::new(rustls::crypto::ring::default_provider()); + let a = derive_identity(&ClusterKey::new("periwinkle")).expect("derive a"); + let a2 = derive_identity(&ClusterKey::new("periwinkle")).expect("derive a2"); + let c = derive_identity(&ClusterKey::new("limpet")).expect("derive c"); + let pin = PinnedPeer::new(provider, &a.cert).expect("build pin"); + assert!(pin.check(&a2.cert).is_ok()); + assert!(pin.check(&c.cert).is_err()); + } + + #[test] + fn spki_pin_survives_cert_reencoding() { + let provider = Arc::new(rustls::crypto::ring::default_provider()); + let canonical = derive_identity(&ClusterKey::new("periwinkle")).expect("derive canonical"); + let pin = PinnedPeer::new(provider, &canonical.cert).expect("build pin"); + + let hk = Hkdf::::new(Some(CLUSTER_HKDF_SALT), b"periwinkle"); + let mut seed = [0u8; 32]; + hk.expand(CLUSTER_HKDF_INFO, &mut seed).expect("expand"); + let mut pkcs8 = Vec::with_capacity(ED25519_PKCS8_V1_PREFIX.len() + seed.len()); + pkcs8.extend_from_slice(&ED25519_PKCS8_V1_PREFIX); + pkcs8.extend_from_slice(&seed); + let keypair = rcgen::KeyPair::from_pkcs8_der_and_sign_algo( + &PrivatePkcs8KeyDer::from(pkcs8), + &rcgen::PKCS_ED25519, + ) + .expect("keypair"); + let mut params = + rcgen::CertificateParams::new(vec![RIPPLE_SERVER_NAME.to_string()]).expect("params"); + params.serial_number = Some(rcgen::SerialNumber::from(99u64)); + params.not_before = rcgen::date_time_ymd(2030, 6, 1); + params.not_after = rcgen::date_time_ymd(2040, 6, 1); + let reencoded = params.self_signed(&keypair).expect("self signed"); + + assert_ne!( + canonical.cert.as_ref(), + reencoded.der().as_ref(), + "test must actually vary the cert encoding" + ); + assert!( + pin.check(reencoded.der()).is_ok(), + "same cluster key must authenticate across a cert re-encoding" + ); + } + + #[tokio::test] + async fn matching_cluster_key_roundtrip() { + let shutdown = CancellationToken::new(); + let key = || Some(ClusterKey::new("nautilus-shared-secret")); + let (sender, _rx_sender) = + Transport::bind("127.0.0.1:0".parse().unwrap(), key(), shutdown.clone()) + .await + .expect("bind sender"); + let (receiver, mut rx_receiver) = + Transport::bind("127.0.0.1:0".parse().unwrap(), key(), shutdown.clone()) + .await + .expect("bind receiver"); + let target = receiver.local_addr(); + + sender + .send(target, ChannelTag::Gossip, b"authenticated") + .await; + + let frame = tokio::time::timeout(Duration::from_secs(5), rx_receiver.recv()) + .await + .expect("frame arrives before timeout") + .expect("channel open"); + assert_eq!(frame.data, b"authenticated"); + + shutdown.cancel(); + } + + #[tokio::test] + async fn mismatched_cluster_key_rejects() { + let shutdown = CancellationToken::new(); + let (sender, _rx_sender) = Transport::bind( + "127.0.0.1:0".parse().unwrap(), + Some(ClusterKey::new("squid-key")), + shutdown.clone(), + ) + .await + .expect("bind sender"); + let (receiver, mut rx_receiver) = Transport::bind( + "127.0.0.1:0".parse().unwrap(), + Some(ClusterKey::new("whelk-key")), + shutdown.clone(), + ) + .await + .expect("bind receiver"); + let target = receiver.local_addr(); + + sender + .send(target, ChannelTag::Gossip, b"should be rejected") + .await; + + let result = tokio::time::timeout(Duration::from_secs(2), rx_receiver.recv()).await; + assert!( + result.is_err(), + "frame must not cross a cluster-key boundary" + ); + + shutdown.cancel(); + } + + #[tokio::test] + async fn keyless_client_cannot_inject_into_keyed_node() { + let shutdown = CancellationToken::new(); + let (receiver, mut rx_receiver) = Transport::bind( + "127.0.0.1:0".parse().unwrap(), + Some(ClusterKey::new("conch-secret")), + shutdown.clone(), + ) + .await + .expect("bind receiver"); + let target = receiver.local_addr(); + + let client_config = build_client_config(None).expect("client config"); + let mut endpoint = + Endpoint::client("127.0.0.1:0".parse().unwrap()).expect("client endpoint"); + endpoint.set_default_client_config(client_config); + + if let Ok(connecting) = endpoint.connect(target, RIPPLE_SERVER_NAME) + && let Ok(Ok(conn)) = tokio::time::timeout(Duration::from_secs(3), connecting).await + && let Ok(mut s) = conn.open_uni().await + { + let _ = s.write_all(&[ChannelTag::Gossip as u8]).await; + let _ = s.write_all(b"attacker payload").await; + let _ = s.finish(); + let _ = tokio::time::timeout(Duration::from_secs(2), conn.closed()).await; + } + + let res = tokio::time::timeout(Duration::from_secs(2), rx_receiver.recv()).await; + assert!( + res.is_err(), + "no frame may cross from an unauthenticated client into a keyed node" + ); + + shutdown.cancel(); + } + #[tokio::test] async fn max_size_frame_roundtrips_and_oversize_is_refused() { let shutdown = CancellationToken::new(); let (sender, _rx_sender) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind sender"); let (receiver, mut rx_receiver) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind receiver"); let target = receiver.local_addr(); @@ -755,12 +1118,12 @@ mod tests { async fn stalled_streams_capped_per_peer() { let shutdown = CancellationToken::new(); let (receiver, _rx_receiver) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind receiver"); let target = receiver.local_addr(); - let client_config = build_client_config().expect("client config"); + let client_config = build_client_config(None).expect("client config"); let mut endpoint = Endpoint::client("127.0.0.1:0".parse().unwrap()).expect("client endpoint"); endpoint.set_default_client_config(client_config); @@ -812,11 +1175,11 @@ mod tests { async fn inbound_byte_budget_held_until_consumed() { let shutdown = CancellationToken::new(); let (sender, _rx_sender) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind sender"); let (receiver, _rx_receiver) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind receiver"); let target = receiver.local_addr(); @@ -852,11 +1215,11 @@ mod tests { async fn incoming_frame_from_matches_peer_listen_addr() { let shutdown = CancellationToken::new(); let (sender, _rx_sender) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind sender"); let (receiver, mut rx_receiver) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind receiver"); @@ -883,11 +1246,11 @@ mod tests { let shutdown = CancellationToken::new(); let (sender, _rx_sender) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind sender"); let (receiver, rx_receiver) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind receiver"); let target = receiver.local_addr(); @@ -925,11 +1288,11 @@ mod tests { async fn concurrent_sends_to_fresh_peer_all_delivered() { let shutdown = CancellationToken::new(); let (sender, _rx_sender) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind sender"); let (receiver, rx_receiver) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind receiver"); let target = receiver.local_addr(); @@ -967,11 +1330,11 @@ mod tests { async fn timed_out_write_resets_stream_instead_of_truncating() { let shutdown = CancellationToken::new(); let (sender, _rx_sender) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind sender"); - let server_config = build_server_config().expect("server config"); + let server_config = build_server_config(None).expect("server config"); let stalled_server = Endpoint::server(server_config, "127.0.0.1:0".parse().unwrap()) .expect("bind stalled server"); let target = stalled_server.local_addr().expect("local addr"); @@ -1010,11 +1373,11 @@ mod tests { let shutdown = CancellationToken::new(); let (sender, _rx_sender) = - Transport::bind("127.0.0.1:0".parse().unwrap(), shutdown.clone()) + Transport::bind("127.0.0.1:0".parse().unwrap(), None, shutdown.clone()) .await .expect("bind sender"); - let server_config = build_server_config().expect("server config"); + let server_config = build_server_config(None).expect("server config"); let mute_server = Endpoint::server(server_config, "127.0.0.1:0".parse().unwrap()) .expect("bind mute server"); let target = mute_server.local_addr().expect("local addr");