From 8179d9fa88f6ca5201f738559af000544d521f41 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Tue, 6 Oct 2026 18:23:01 +0500 Subject: [PATCH 1/3] fix(net): retire paths stalled without acknowledgement progress --- Cargo.lock | 2 + crates/rds-net/Cargo.toml | 2 + crates/rds-net/src/ack_progress.rs | 347 ++++++++++++++++++ crates/rds-net/src/backends/iroh.rs | 13 +- crates/rds-net/src/backends/iroh/latency.rs | 42 ++- .../src/backends/iroh/latency_tests.rs | 180 +++++++++ crates/rds-net/src/backends/noq/mod.rs | 16 +- .../src/backends/noq/path_blackhole_tests.rs | 229 ++++++++++++ crates/rds-net/src/backends/noq/policy.rs | 50 ++- crates/rds-net/src/lib.rs | 1 + docs/architecture.md | 8 + .../reports/rds-path-ack-progress-20261006.md | 56 +++ vendor/iroh/RDS-PATCH.md | 8 + .../src/socket/remote_map/remote_state.rs | 33 ++ 14 files changed, 981 insertions(+), 6 deletions(-) create mode 100644 crates/rds-net/src/ack_progress.rs create mode 100644 crates/rds-net/src/backends/iroh/latency_tests.rs create mode 100644 crates/rds-net/src/backends/noq/path_blackhole_tests.rs create mode 100644 docs/reports/rds-path-ack-progress-20261006.md diff --git a/Cargo.lock b/Cargo.lock index 364b4a6..bc21d58 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4324,7 +4324,9 @@ dependencies = [ "data-encoding", "ed25519-dalek", "iroh", + "iroh-base", "iroh-relay", + "n0-watcher", "noq", "noq-proto", "postcard", diff --git a/crates/rds-net/Cargo.toml b/crates/rds-net/Cargo.toml index 2c31ad0..5de9faf 100644 --- a/crates/rds-net/Cargo.toml +++ b/crates/rds-net/Cargo.toml @@ -33,6 +33,8 @@ tokio-util = { workspace = true, optional = true } tracing.workspace = true [dev-dependencies] +n0-watcher = "1.0.0" +iroh-base = "1.3.0" iroh-relay = { workspace = true, features = ["server"] } proptest.workspace = true turmoil = "0.7.2" diff --git a/crates/rds-net/src/ack_progress.rs b/crates/rds-net/src/ack_progress.rs new file mode 100644 index 0000000..73fdc31 --- /dev/null +++ b/crates/rds-net/src/ack_progress.rs @@ -0,0 +1,347 @@ +//! Passive per-controller acknowledgement progress, for opt-in path policy. +//! Every congestion callback and metric delegates unchanged to the selected CCA. +use noq_proto::RttEstimator; +use noq_proto::congestion::{Controller, ControllerFactory, ControllerMetrics}; +use std::any::Any; +use std::sync::Arc; +use std::time::{Duration, Instant}; + +pub(crate) type Clock = Arc Instant + Send + Sync>; + +#[derive(Clone, Copy, Debug)] +pub(crate) struct Observation { + progress: Progress, + now: Instant, +} +impl Observation { + pub(crate) fn confirmed(self, rtt: Duration) -> bool { + self.progress.confirmed(self.now, rtt) + } + pub(crate) fn stalled(self, rtt: Duration) -> bool { + self.progress.stalled(self.now, rtt) + } + pub(crate) fn needs_probe(self) -> bool { + self.progress.needs_probe(self.now) + } + pub(crate) fn pending_age(self) -> Option { + self.progress + .pending_since + .map(|at| self.now.saturating_duration_since(at)) + } + pub(crate) fn confirmation_age(self) -> Option { + self.progress + .confirmed_at + .map(|at| self.now.saturating_duration_since(at)) + } +} + +#[derive(Clone, Copy, Debug, Default)] +pub(crate) struct Progress { + pending_since: Option, + confirmed_at: Option, +} +impl Progress { + fn sent(&mut self, now: Instant, bytes: u64) { + if bytes > 0 && self.pending_since.is_none() { + self.pending_since = Some(now); + } + } + fn ack(&mut self, now: Instant, bytes: u64) { + if bytes > 0 { + self.confirmed_at = Some(now); + if self.pending_since.is_some() { + self.pending_since = Some(now); + } + } + } + fn end_acks(&mut self, now: Instant, in_flight: u64) { + if in_flight == 0 { + self.pending_since = None; + } else if self.pending_since.is_none() { + self.pending_since = Some(now); + } + } + pub(crate) fn confirmed(self, now: Instant, rtt: Duration) -> bool { + let budget = rtt.saturating_mul(8).max(Duration::from_secs(2)); + self.confirmed_at + .is_some_and(|at| now.saturating_duration_since(at) < budget) + } + pub(crate) fn needs_probe(self, now: Instant) -> bool { + self.pending_since.is_none() + && self + .confirmed_at + .is_none_or(|at| now.saturating_duration_since(at) >= Duration::from_secs(1)) + } + pub(crate) fn stalled(self, now: Instant, rtt: Duration) -> bool { + let budget = rtt.saturating_mul(4).max(Duration::from_millis(500)); + self.pending_since + .is_some_and(|since| now.saturating_duration_since(since) >= budget) + } +} + +struct Factory(Arc, Clock); +impl std::fmt::Debug for Factory { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("AckProgressFactory").finish_non_exhaustive() + } +} +impl ControllerFactory for Factory { + fn build(self: Arc, now: Instant, mtu: u16) -> Box { + Box::new(Tracked { + inner: self.0.clone().build(now, mtu), + progress: Progress::default(), + reset_on_mutation: false, + clock: self.1.clone(), + }) + } +} +pub(crate) fn factory( + inner: Arc, + clock: Clock, +) -> Arc { + Arc::new(Factory(inner, clock)) +} +pub(crate) fn snapshot(controller: Box) -> Option { + controller + .into_any() + .downcast::() + .ok() + .map(|tracked| Observation { + progress: tracked.progress, + now: (tracked.clock)(), + }) +} + +struct Tracked { + inner: Box, + progress: Progress, + // Noq uses clone_box both for snapshots and live migration. Reading the + // clone preserves snapshot evidence; its first live mutation starts a + // new proof domain without guessing across QUIC packet-number spaces. + reset_on_mutation: bool, + clock: Clock, +} +impl std::fmt::Debug for Tracked { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("AckProgressController") + .field("inner", &self.inner) + .field("progress", &self.progress) + .finish_non_exhaustive() + } +} +impl Tracked { + fn prepare(&mut self) { + if self.reset_on_mutation { + self.progress = Progress::default(); + self.reset_on_mutation = false; + } + } +} +impl Controller for Tracked { + fn on_sent(&mut self, now: Instant, bytes: u64, packet: u64) { + self.prepare(); + self.inner.on_sent(now, bytes, packet); + self.progress.sent(now, bytes); + } + fn on_packet_sent(&mut self, now: Instant, bytes: u16, packet: u64) { + self.prepare(); + self.inner.on_packet_sent(now, bytes, packet); + } + fn on_cwnd_limited(&mut self) { + self.prepare(); + self.inner.on_cwnd_limited(); + } + fn on_ack( + &mut self, + now: Instant, + sent: Instant, + bytes: u64, + packet: u64, + app_limited: bool, + rtt: &RttEstimator, + ) { + self.prepare(); + self.inner + .on_ack(now, sent, bytes, packet, app_limited, rtt); + self.progress.ack(now, bytes); + } + fn on_end_acks( + &mut self, + now: Instant, + in_flight: u64, + app_limited: bool, + largest: Option, + ) { + self.prepare(); + self.inner.on_end_acks(now, in_flight, app_limited, largest); + self.progress.end_acks(now, in_flight); + } + fn on_congestion_event( + &mut self, + now: Instant, + sent: Instant, + persistent: bool, + ecn: bool, + bytes: u64, + packet: u64, + ) { + self.prepare(); + self.inner + .on_congestion_event(now, sent, persistent, ecn, bytes, packet); + } + fn on_packet_lost(&mut self, bytes: u16, packet: u64, now: Instant) { + self.prepare(); + self.inner.on_packet_lost(bytes, packet, now); + } + fn on_spurious_congestion_event(&mut self) { + self.prepare(); + self.inner.on_spurious_congestion_event(); + } + fn on_mtu_update(&mut self, mtu: u16) { + self.prepare(); + self.inner.on_mtu_update(mtu); + } + fn on_ack_frequency_update(&mut self, threshold: u64, delay: Duration) { + self.prepare(); + self.inner.on_ack_frequency_update(threshold, delay); + } + fn window(&self) -> u64 { + self.inner.window() + } + fn metrics(&self) -> ControllerMetrics { + self.inner.metrics() + } + fn clone_box(&self) -> Box { + Box::new(Self { + inner: self.inner.clone_box(), + progress: self.progress, + reset_on_mutation: true, + clock: self.clock.clone(), + }) + } + fn initial_window(&self) -> u64 { + self.inner.initial_window() + } + fn into_any(self: Box) -> Box { + self + } +} + +#[cfg(test)] +mod tests { + use super::*; + #[test] + fn a_high_rtt_path_is_not_failed_before_its_own_round_trip_budget() { + let now = Instant::now(); + let rtt = Duration::from_secs(3); + let mut progress = Progress::default(); + progress.sent(now, 100); + assert!(!progress.stalled(now + Duration::from_secs(11), rtt)); + assert!(progress.stalled(now + Duration::from_secs(12), rtt)); + progress.ack(now, 100); + assert!(progress.confirmed(now + Duration::from_secs(23), rtt)); + assert!(!progress.confirmed(now + Duration::from_secs(24), rtt)); + } + #[test] + fn observations_use_the_controller_runtime_clock_not_the_callers_wall_clock() { + use std::sync::atomic::{AtomicU64, Ordering}; + let epoch = Instant::now() + Duration::from_secs(3600); + let elapsed = Arc::new(AtomicU64::new(0)); + let ticks = elapsed.clone(); + let clock: Clock = + Arc::new(move || epoch + Duration::from_millis(ticks.load(Ordering::Relaxed))); + let mut controller = + factory(crate::CongestionControl::Cubic.factory(), clock).build(epoch, 1200); + controller.on_sent(epoch, 100, 0); + elapsed.store(499, Ordering::Relaxed); + assert!( + !snapshot(controller.clone_box()) + .unwrap() + .stalled(Duration::from_millis(70)) + ); + elapsed.store(500, Ordering::Relaxed); + assert!( + snapshot(controller.clone_box()) + .unwrap() + .stalled(Duration::from_millis(70)) + ); + } + #[test] + fn snapshots_keep_evidence_but_a_live_clone_requires_its_own_confirmation() { + let now = Instant::now(); + let inner = crate::CongestionControl::Cubic.factory().build(now, 1200); + let mut tracked = Tracked { + inner, + progress: Progress::default(), + reset_on_mutation: false, + clock: Arc::new(move || now), + }; + tracked.progress.ack(now, 100); + let expected_window = tracked.window(); + let mut cloned = tracked + .clone_box() + .into_any() + .downcast::() + .unwrap(); + assert!(cloned.progress.confirmed(now, Duration::from_millis(70))); + assert_eq!(cloned.window(), expected_window); + cloned.on_mtu_update(1200); + assert!(!cloned.progress.confirmed(now, Duration::from_millis(70))); + assert!(tracked.progress.confirmed(now, Duration::from_millis(70))); + assert_eq!(cloned.window(), tracked.window()); + } + #[test] + fn stale_idle_confirmation_requires_a_probe_and_cannot_retire_a_sibling() { + let now = Instant::now(); + let mut progress = Progress::default(); + progress.ack(now, 100); + assert!(!progress.needs_probe(now)); + assert!(progress.needs_probe(now + Duration::from_secs(1))); + assert!(!progress.confirmed(now + Duration::from_secs(2), Duration::from_millis(70))); + assert!(!progress.stalled(now + Duration::from_secs(30), Duration::from_millis(70))); + } + #[test] + fn only_outstanding_ack_eliciting_work_can_stall_and_idle_resumption_gets_its_own_budget() { + let start = Instant::now(); + let mut progress = Progress::default(); + progress.sent(start, 0); + assert!(!progress.stalled(start + Duration::from_secs(10), Duration::ZERO)); + progress.sent(start, 100); + assert!(!progress.stalled( + start + Duration::from_millis(499), + Duration::from_millis(70) + )); + assert!(progress.stalled( + start + Duration::from_millis(500), + Duration::from_millis(70) + )); + progress.ack(start + Duration::from_millis(600), 100); + progress.end_acks(start + Duration::from_millis(600), 0); + assert!(progress.confirmed( + start + Duration::from_millis(600), + Duration::from_millis(70) + )); + progress.sent(start + Duration::from_secs(60), 100); + assert!(!progress.stalled(start + Duration::from_secs(60), Duration::from_millis(70))); + } + #[test] + fn fresh_ack_progress_and_migration_reset_are_explicit() { + let start = Instant::now(); + let mut progress = Progress::default(); + progress.sent(start, 100); + progress.ack(start + Duration::from_millis(450), 100); + progress.end_acks(start + Duration::from_millis(450), 100); + assert!(!progress.stalled( + start + Duration::from_millis(900), + Duration::from_millis(70) + )); + assert!(progress.stalled( + start + Duration::from_millis(950), + Duration::from_millis(70) + )); + progress = Progress::default(); + progress.sent(start + Duration::from_secs(2), 100); + assert!(!progress.confirmed(start + Duration::from_secs(2), Duration::from_millis(70))); + assert!(!progress.stalled(start + Duration::from_secs(2), Duration::from_millis(70))); + } +} diff --git a/crates/rds-net/src/backends/iroh.rs b/crates/rds-net/src/backends/iroh.rs index f20cbb6..0309c9b 100644 --- a/crates/rds-net/src/backends/iroh.rs +++ b/crates/rds-net/src/backends/iroh.rs @@ -11,6 +11,8 @@ use iroh::{Endpoint, RelayMap, RelayMode}; use crate::{EndpointAddr, EndpointConfig, EndpointId, RelayUrl, TransportAddr}; mod latency; +#[cfg(all(test, feature = "transport-noq"))] +mod latency_tests; /// Adapter conversions between the owned shared types and iroh-base. /// @@ -126,8 +128,17 @@ pub async fn bind_endpoint(config: EndpointConfig) -> anyhow::Result { // 100 ms; a larger per-stream window keeps a big keyframe or sync // chunk stream from stalling on high-BDP links. // - 32 MiB connection send window keeps several bulk streams busy. + let controller = config.congestion_control.factory(); + let controller = if config.path_preference == crate::PathPreference::Latency { + crate::ack_progress::factory( + controller, + std::sync::Arc::new(|| tokio::time::Instant::now().into_std()), + ) + } else { + controller + }; let mut transport = iroh::endpoint::QuicTransportConfig::builder() - .congestion_controller_factory(config.congestion_control.factory()) + .congestion_controller_factory(controller) .stream_receive_window(noq_proto::VarInt::from_u32(4 * 1024 * 1024)) .send_window(32 * 1024 * 1024); if config.packetization == crate::Packetization::Conservative { diff --git a/crates/rds-net/src/backends/iroh/latency.rs b/crates/rds-net/src/backends/iroh/latency.rs index bb0e172..0632c3c 100644 --- a/crates/rds-net/src/backends/iroh/latency.rs +++ b/crates/rds-net/src/backends/iroh/latency.rs @@ -14,8 +14,46 @@ impl PathSelector for LatencySelector { } fn select(&self, ctx: &PathSelectionContext<'_>) -> PathSelection { - let choice = choose(ctx.paths().filter_map(|path| { - let rtt = path.stats()?.rtt; + let paths: Vec<_> = ctx + .paths() + .filter_map(|path| { + let rtt = path.stats()?.rtt; + let progress = path + .congestion_state() + .and_then(crate::ack_progress::snapshot); + if progress.is_some_and(|state| state.needs_probe()) { + path.ping(); + } + Some((path, rtt, progress)) + }) + .collect(); + for (failed, rtt, progress) in &paths { + if !progress.is_some_and(|state| state.stalled(*rtt)) { + continue; + } + for (fallback, other_rtt, state) in &paths { + if state + .is_some_and(|state| state.confirmed(*other_rtt) && !state.stalled(*other_rtt)) + && failed.abandon_with_fallback(fallback) + { + tracing::warn!(target:"rds_net::path_policy", + pending_ack_age_ms=?progress.and_then(|state| state.pending_age()).map(|age| age.as_millis()), + fallback_ack_age_ms=?state.and_then(|state| state.confirmation_age()).map(|age| age.as_millis()), + path_rtt_ms=rtt.as_millis(), fallback_rtt_ms=other_rtt.as_millis(), + "unresponsive path retired with a confirmed sibling; reliable streams retained"); + break; + } + } + } + let has_confirmed = paths.iter().any(|(_, rtt, state)| { + state.is_some_and(|state| state.confirmed(*rtt) && !state.stalled(*rtt)) + }); + let choice = choose(paths.into_iter().filter_map(|(path, rtt, state)| { + if state.is_some_and(|state| state.stalled(rtt)) + || (has_confirmed && state.is_some_and(|state| !state.confirmed(rtt))) + { + return None; + } let current = Some(path.network_path()) == ctx.current(); Some((path, rtt, current)) })); diff --git a/crates/rds-net/src/backends/iroh/latency_tests.rs b/crates/rds-net/src/backends/iroh/latency_tests.rs new file mode 100644 index 0000000..d12a03c --- /dev/null +++ b/crates/rds-net/src/backends/iroh/latency_tests.rs @@ -0,0 +1,180 @@ +//! Real Iroh actor over a bounded controllable custom link and UDP fallback. +use iroh::TransportAddr; +use iroh::endpoint::transports::{ + CustomEndpoint, CustomSender, CustomTransport, RecvInfo, Transmit, +}; +use iroh_base::CustomAddr; +use std::collections::HashMap; +use std::io::{self, IoSliceMut}; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::{Arc, Mutex}; +use std::task::{Context, Poll}; +use std::time::{Duration, Instant}; +use tokio::sync::mpsc; + +#[derive(Debug)] +struct Packet { + from: CustomAddr, + data: Vec, +} +type Fabric = Arc>>>; +#[derive(Debug)] +struct Link { + addr: CustomAddr, + fabric: Fabric, + drop_packets: Arc, + dropped: Arc, +} +impl CustomTransport for Link { + fn bind(&self) -> io::Result> { + let (tx, rx) = mpsc::channel(64); + self.fabric.lock().unwrap().insert(self.addr.clone(), tx); + Ok(Box::new(LinkEndpoint { + sender: Arc::new(LinkSender { + addr: self.addr.clone(), + fabric: self.fabric.clone(), + drop_packets: self.drop_packets.clone(), + dropped: self.dropped.clone(), + }), + rx, + addresses: n0_watcher::Watchable::new(vec![self.addr.clone()]), + })) + } +} +#[derive(Debug)] +struct LinkEndpoint { + sender: Arc, + rx: mpsc::Receiver, + addresses: n0_watcher::Watchable>, +} +impl Drop for LinkEndpoint { + fn drop(&mut self) { + self.sender.fabric.lock().unwrap().remove(&self.sender.addr); + } +} +impl CustomEndpoint for LinkEndpoint { + fn watch_local_addrs(&self) -> n0_watcher::Direct> { + self.addresses.watch() + } + fn create_sender(&self) -> Arc { + self.sender.clone() + } + fn poll_recv( + &mut self, + cx: &mut Context<'_>, + bufs: &mut [IoSliceMut<'_>], + metas: &mut [noq::udp::RecvMeta], + infos: &mut [RecvInfo], + ) -> Poll> { + match self.rx.poll_recv(cx) { + Poll::Ready(Some(packet)) => { + assert!(packet.data.len() <= bufs[0].len()); + bufs[0][..packet.data.len()].copy_from_slice(&packet.data); + let mut meta = noq::udp::RecvMeta::default(); + meta.len = packet.data.len(); + meta.stride = packet.data.len(); + metas[0] = meta; + infos[0] = RecvInfo::new(packet.from, Some(self.sender.addr.clone())); + Poll::Ready(Ok(1)) + } + Poll::Ready(None) => Poll::Ready(Err(io::ErrorKind::BrokenPipe.into())), + Poll::Pending => Poll::Pending, + } + } +} +#[derive(Debug)] +struct LinkSender { + addr: CustomAddr, + fabric: Fabric, + drop_packets: Arc, + dropped: Arc, +} +impl CustomSender for LinkSender { + fn is_valid_send_addr(&self, addr: &CustomAddr) -> bool { + addr.id() == self.addr.id() + } + fn poll_send( + &self, + _cx: &mut Context<'_>, + dst: &CustomAddr, + src: Option<&CustomAddr>, + transmit: &Transmit<'_>, + ) -> Poll> { + if self.drop_packets.load(Ordering::Acquire) { + self.dropped.fetch_add(1, Ordering::Relaxed); + return Poll::Ready(Ok(())); + } + let sender = self + .fabric + .lock() + .unwrap() + .get(dst) + .cloned() + .ok_or(io::ErrorKind::NotConnected); + let result = sender.and_then(|sender| { + sender + .try_send(Packet { + from: src.unwrap_or(&self.addr).clone(), + data: transmit.contents.to_vec(), + }) + .map_err(|_| io::ErrorKind::Other) + }); + Poll::Ready(result.map_err(Into::into)) + } +} + +async fn endpoint(link: Link) -> iroh::Endpoint { + let mut transport = iroh::endpoint::QuicTransportConfig::builder() + .congestion_controller_factory(crate::ack_progress::factory( + crate::CongestionControl::Cubic.factory(), + Arc::new(|| tokio::time::Instant::now().into_std()), + )); + transport = transport + .initial_mtu(1200) + .min_mtu(1200) + .mtu_discovery_config(None) + .enable_segmentation_offload(false); + iroh::Endpoint::builder(iroh::endpoint::presets::Minimal) + .clear_ip_transports() + .bind_addr("127.0.0.1:0".parse::().unwrap()) + .unwrap() + .portmapper_config(iroh::endpoint::PortmapperConfig::Disabled) + .add_custom_transport(Arc::new(link)) + .path_selector(Arc::new(super::latency::LatencySelector)) + .transport_config(transport.build()) + .alpns(vec![rds_core::ALPN.to_vec()]) + .bind() + .await + .unwrap() +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn actual_iroh_actor_retires_blackholed_preferred_link_and_delivers_pending_stream_bytes() { + tokio::time::timeout(Duration::from_secs(15),async { + let fabric=Fabric::default(); + let dropping=Arc::new(AtomicBool::new(false));let dropped=Arc::new(AtomicU64::new(0)); + let a_addr=CustomAddr::from_parts(0x72647374657374,b"a");let b_addr=CustomAddr::from_parts(0x72647374657374,b"b"); + let a=endpoint(Link {addr:a_addr,fabric:fabric.clone(),drop_packets:Arc::new(AtomicBool::new(false)),dropped:Arc::new(AtomicU64::new(0))}).await; + let b=endpoint(Link {addr:b_addr.clone(),fabric,drop_packets:dropping.clone(),dropped:dropped.clone()}).await; + let target=iroh::EndpointAddr::from_parts(b.id(),[TransportAddr::Custom(b_addr)]); + let (client,server)=tokio::join!(a.connect(target,rds_core::ALPN),async {b.accept().await.unwrap().await}); + let (client,server)=(client.unwrap(),server.unwrap()); + let (mut request,mut reply)=client.open_bi().await.unwrap();request.write_all(b"warm").await.unwrap(); + let (mut send,mut recv)=server.accept_bi().await.unwrap();let mut warm=[0;4];recv.read_exact(&mut warm).await.unwrap(); + send.write_all(b"ok").await.unwrap();let mut ack=[0;2];reply.read_exact(&mut ack).await.unwrap(); + tokio::time::timeout(Duration::from_secs(5),async { + loop { + let paths=server.paths(); + if paths.iter().any(|p|p.is_ip()) && paths.iter().any(|p|p.is_selected() && matches!(p.remote_addr(),TransportAddr::Custom(_))) { break } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }).await.expect("fixture did not establish its UDP standby while the custom link was preferred"); + dropping.store(true,Ordering::Release);let started=Instant::now();send.write_all(b"pending reliable bytes").await.unwrap(); + let mut body=[0;22];let result=tokio::time::timeout(Duration::from_secs(3),reply.read_exact(&mut body)).await; + let elapsed=started.elapsed().as_millis();let injected=dropped.load(Ordering::Relaxed); + tokio::join!(a.close(),b.close()); + assert!(injected>0,"blackhole fixture did not drop a packet"); + result.expect("Iroh left the existing reliable stream stalled after its preferred path blackholed").unwrap(); + assert_eq!(&body,b"pending reliable bytes");assert!(elapsed<3000); + }).await.expect("isolated Iroh blackhole fixture did not terminate"); +} diff --git a/crates/rds-net/src/backends/noq/mod.rs b/crates/rds-net/src/backends/noq/mod.rs index 3c67551..b231a59 100644 --- a/crates/rds-net/src/backends/noq/mod.rs +++ b/crates/rds-net/src/backends/noq/mod.rs @@ -22,6 +22,8 @@ mod candidates; mod dial; mod drivers; mod hmac; +#[cfg(test)] +mod path_blackhole_tests; pub mod policy; pub mod relay; pub mod socket; @@ -62,6 +64,8 @@ fn transport_config( max_multipath_paths: Option, packetization: crate::Packetization, congestion_control: crate::CongestionControl, + path_preference: crate::PathPreference, + clock: crate::ack_progress::Clock, ) -> Arc { let mut cfg = noq::TransportConfig::default(); cfg.keep_alive_interval(Some(HEARTBEAT_INTERVAL)); @@ -75,7 +79,13 @@ fn transport_config( // Same controller selection as the iroh backend (default BBRv3), // and windows above the 100Mbps x 100ms // defaults so bulk streams do not stall on high-BDP links. - cfg.congestion_controller_factory(congestion_control.factory()); + let controller = congestion_control.factory(); + let controller = if path_preference == crate::PathPreference::Latency { + crate::ack_progress::factory(controller, clock) + } else { + controller + }; + cfg.congestion_controller_factory(controller); cfg.stream_receive_window(noq_proto::VarInt::from_u32(4 * 1024 * 1024)); cfg.send_window(32 * 1024 * 1024); if packetization == crate::Packetization::Conservative { @@ -219,10 +229,14 @@ async fn bind_socket( let endpoint_config = noq::EndpointConfig::new(Arc::new(hmac::Blake3HmacKey::new(&mut rand::rng()))); + let progress_runtime = runtime.clone(); + let progress_clock: crate::ack_progress::Clock = Arc::new(move || progress_runtime.now()); let transport = transport_config( config.max_multipath_paths, config.packetization, config.congestion_control, + config.path_preference, + progress_clock, ); let mut server_config = noq::ServerConfig::with_crypto(Arc::new(server_crypto)); server_config.transport = transport.clone(); diff --git a/crates/rds-net/src/backends/noq/path_blackhole_tests.rs b/crates/rds-net/src/backends/noq/path_blackhole_tests.rs new file mode 100644 index 0000000..317b895 --- /dev/null +++ b/crates/rds-net/src/backends/noq/path_blackhole_tests.rs @@ -0,0 +1,229 @@ +//! Isolated engine boundary: pending reliable bytes on a blackholed path. +use super::{Endpoint, bind_with_socket, socket, tls}; +use crate::{Backend, CongestionControl, EndpointConfig, Packetization, SecretKey}; +use noq::{AsyncUdpSocket, PathStatus, Runtime, UdpSender}; +use std::io::{self, IoSliceMut}; +use std::net::SocketAddr; +use std::pin::Pin; +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::task::{Context, Poll}; +use std::time::{Duration, Instant}; + +#[derive(Debug)] +struct LossSocket { + inner: Box, + enabled: Arc, + dropped: Arc, +} +impl AsyncUdpSocket for LossSocket { + fn create_sender(&self) -> Pin> { + Box::pin(LossSender { + inner: self.inner.create_sender(), + enabled: self.enabled.clone(), + dropped: self.dropped.clone(), + }) + } + fn poll_recv( + &mut self, + cx: &mut Context<'_>, + bufs: &mut [IoSliceMut<'_>], + meta: &mut [noq::udp::RecvMeta], + ) -> Poll> { + self.inner.poll_recv(cx, bufs, meta) + } + fn local_addr(&self) -> io::Result { + self.inner.local_addr() + } + fn may_fragment(&self) -> bool { + self.inner.may_fragment() + } +} +#[derive(Debug)] +struct LossSender { + inner: Pin>, + enabled: Arc, + dropped: Arc, +} +impl UdpSender for LossSender { + fn poll_send( + self: Pin<&mut Self>, + transmit: &noq::udp::Transmit<'_>, + cx: &mut Context<'_>, + ) -> Poll> { + let this = self.get_mut(); + if this.enabled.load(Ordering::Acquire) { + this.dropped.fetch_add(1, Ordering::Relaxed); + Poll::Ready(Ok(())) + } else { + this.inner.as_mut().poll_send(transmit, cx) + } + } +} + +async fn endpoints() -> ( + Endpoint, + Endpoint, + SocketAddr, + Arc, + Arc, +) { + let runtime = Arc::new(noq::TokioRuntime); + let bind = |addr| { + runtime + .wrap_udp_socket(std::net::UdpSocket::bind(addr).unwrap()) + .unwrap() + }; + let enabled = Arc::new(AtomicBool::new(false)); + let dropped = Arc::new(AtomicU64::new(0)); + let primary = LossSocket { + inner: bind("127.0.0.1:0"), + enabled: enabled.clone(), + dropped: dropped.clone(), + }; + let secondary = bind("[::1]:0"); + let secondary_addr = secondary.local_addr().unwrap(); + let primary_addr = primary.local_addr().unwrap(); + let server = bind_with_socket( + EndpointConfig { + backend: Backend::Noq, + secret_key: Some(SecretKey::from_bytes(&[112; 32])), + discovery: false, + observed_address_reports: false, + congestion_control: CongestionControl::Cubic, + packetization: Packetization::Conservative, + path_preference: crate::PathPreference::Latency, + ..Default::default() + }, + Box::new(socket::Mux::new(vec![Box::new(primary), secondary]).unwrap()), + vec![primary_addr], + runtime.clone(), + Vec::new(), + ) + .await + .unwrap(); + let client_mux = socket::Mux::new(vec![bind("127.0.0.1:0"), bind("[::1]:0")]).unwrap(); + let client_addrs = client_mux.local_addrs(); + let client = bind_with_socket( + EndpointConfig { + backend: Backend::Noq, + secret_key: Some(SecretKey::from_bytes(&[113; 32])), + discovery: false, + observed_address_reports: false, + congestion_control: CongestionControl::Cubic, + packetization: Packetization::Conservative, + path_preference: crate::PathPreference::Latency, + ..Default::default() + }, + Box::new(client_mux), + client_addrs, + runtime, + Vec::new(), + ) + .await + .unwrap(); + (client, server, secondary_addr, enabled, dropped) +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn retiring_a_blackholed_path_preserves_already_sent_reliable_bytes() { + exercise(false).await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn acknowledgement_starvation_policy_recovers_pending_bytes_without_reconnecting() { + exercise(true).await; +} + +async fn exercise(automatic: bool) { + tokio::time::timeout(Duration::from_secs(12), async { + let (client, server, secondary_addr, enabled, dropped) = endpoints().await; + // Raw connections deliberately bypass QNT advertisement and the policy + // driver, so no extra route can hide the injected loss or manual switch. + let cfg = client.client_configs.get(rds_core::ALPN).unwrap().clone(); + let name = tls::name::encode(server.id()); + let (a, b) = tokio::join!( + client + .inner + .connect_with(cfg, server.local_addr(), &name) + .unwrap(), + async { server.inner.accept().await.unwrap().await } + ); + let (a, b) = (a.unwrap(), b.unwrap()); + let telemetry = if automatic { + Some(( + client.wire_connection(&a, Vec::new(), Vec::new()).unwrap(), + server.wire_connection(&b, Vec::new(), Vec::new()).unwrap(), + )) + } else { + None + }; + let (mut request, mut reply) = a.open_bi().await.unwrap(); + request.write_all(b"warm").await.unwrap(); + let (mut send, mut recv) = b.accept_bi().await.unwrap(); + let mut warm = [0; 4]; + recv.read_exact(&mut warm).await.unwrap(); + assert_eq!(&warm, b"warm"); + send.write_all(b"ok").await.unwrap(); + let mut ack = [0; 2]; + reply.read_exact(&mut ack).await.unwrap(); + assert_eq!(&ack, b"ok"); + // Handshake completion may precede the additional CID credit needed + // to open a path. Await that protocol precondition, not a fixed sleep. + let secondary = tokio::time::timeout(Duration::from_secs(2), async { + loop { + match a.open_path_ensure(secondary_addr, PathStatus::Backup).await { + Ok(path) => break path, + Err(noq::PathError::RemoteCidsExhausted) => { + tokio::time::sleep(Duration::from_millis(10)).await; + } + Err(error) => panic!("secondary path validation failed: {error}"), + } + } + }) + .await + .expect("peer did not provide secondary path CID credit"); + let remote_secondary = b.path(secondary.id()).unwrap(); + remote_secondary.set_status(PathStatus::Backup).unwrap(); + b.path(noq::PathId::ZERO) + .unwrap() + .set_status(PathStatus::Available) + .unwrap(); + let before = b.path_stats(noq::PathId::ZERO).unwrap().frame_tx.stream; + enabled.store(true, Ordering::Release); + send.write_all(b"pending reliable bytes").await.unwrap(); + tokio::time::timeout(Duration::from_secs(2), async { + while dropped.load(Ordering::Relaxed) == 0 + || b.path_stats(noq::PathId::ZERO).unwrap().frame_tx.stream <= before + { + tokio::time::sleep(Duration::from_millis(5)).await; + } + }) + .await + .expect("fixture never emitted and dropped the pending stream frame"); + let mut body = [0; 22]; + assert!( + tokio::time::timeout(Duration::from_millis(100), reply.read(&mut body)) + .await + .is_err() + ); + let started = Instant::now(); + if !automatic { + remote_secondary.set_status(PathStatus::Available).unwrap(); + b.path(noq::PathId::ZERO).unwrap().close().unwrap(); + } + let result = + tokio::time::timeout(Duration::from_secs(2), reply.read_exact(&mut body)).await; + let delivered_ms = started.elapsed().as_millis(); + // Clean up the isolated endpoints even when the engine boundary fails. + tokio::join!(client.close(), server.close()); + result + .expect("retired path left reliable bytes stranded") + .unwrap(); + assert_eq!(&body, b"pending reliable bytes"); + assert!(delivered_ms < 2000); + drop(telemetry); + }) + .await + .expect("isolated blackhole fixture did not terminate"); +} diff --git a/crates/rds-net/src/backends/noq/policy.rs b/crates/rds-net/src/backends/noq/policy.rs index ec0dae5..f2597b3 100644 --- a/crates/rds-net/src/backends/noq/policy.rs +++ b/crates/rds-net/src/backends/noq/policy.rs @@ -461,6 +461,9 @@ fn reselect( if !conn.is_alive() { return; } + let Some(connection) = conn.upgrade() else { + return; + }; // A WeakPathHandle can upgrade even after its path closes: it retains // final statistics until dropped. Upgrade alone is not a liveness check. paths.retain(|_, weak| weak.upgrade().is_some_and(|path| path.status().is_ok())); @@ -487,7 +490,7 @@ fn reselect( return; } - let rtts: Vec<(noq::PathId, Duration)> = paths + let candidates: Vec<_> = paths .iter() .filter_map(|(id, weak)| { let path = weak.upgrade()?; @@ -509,7 +512,50 @@ fn reselect( } else { path.stats().rtt }; - Some((*id, rtt)) + let progress = connection + .congestion_state(*id) + .and_then(crate::ack_progress::snapshot); + if progress.is_some_and(|state| state.needs_probe()) { + let _ = path.ping(); + } + Some((*id, rtt, progress)) + }) + .collect(); + for (id, rtt, state) in &candidates { + if !state.is_some_and(|state| state.stalled(*rtt)) { + continue; + } + let fallback = candidates.iter().find(|(other, other_rtt, state)| { + other != id + && state + .is_some_and(|state| state.confirmed(*other_rtt) && !state.stalled(*other_rtt)) + }); + if let Some((other, other_rtt, fallback_state)) = fallback + && let Some(path) = connection.path(*other) + && path.set_status(noq::PathStatus::Available).is_ok() + && let Some(failed) = connection.path(*id) + && failed.close().is_ok() + { + tracing::warn!(target:"rds_net::path_policy", path_id=%id, fallback_path_id=%other, + pending_ack_age_ms=?state.and_then(|state| state.pending_age()).map(|age| age.as_millis()), + fallback_ack_age_ms=?fallback_state.and_then(|state| state.confirmation_age()).map(|age| age.as_millis()), + path_rtt_ms=rtt.as_millis(), fallback_rtt_ms=other_rtt.as_millis(), + "unresponsive path retired with a confirmed sibling; reliable streams retained"); + } + } + let has_confirmed = candidates.iter().any(|(_, rtt, state)| { + state.is_some_and(|state| state.confirmed(*rtt) && !state.stalled(*rtt)) + }); + let rtts: Vec<_> = candidates + .into_iter() + .filter_map(|(id, rtt, state)| { + if state.is_some_and(|state| state.stalled(rtt)) + || (has_confirmed && state.is_some_and(|state| !state.confirmed(rtt))) + { + None + } else { + Some((id, rtt)) + } }) .collect(); let Some((mut choice, best_rtt)) = rtts.iter().min_by_key(|(_, rtt)| *rtt).copied() else { diff --git a/crates/rds-net/src/lib.rs b/crates/rds-net/src/lib.rs index aadb0d3..9d7401e 100644 --- a/crates/rds-net/src/lib.rs +++ b/crates/rds-net/src/lib.rs @@ -42,6 +42,7 @@ pub mod backends { #[cfg(feature = "transport-noq")] pub mod noq; } +mod ack_progress; pub mod deadline; mod identity; pub mod metrics; diff --git a/docs/architecture.md b/docs/architecture.md index 5fea4ca..b33c7dc 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -10,6 +10,14 @@ between any two enrolled devices, assisted by a GDS-operated server. Design goals, in order: **minimum interactive latency**, **maximum connection stability**, defense in depth, no inbound firewall changes. +Explicit latency path preference now observes acknowledgement progress as well +as RTT. A passive adapter delegates the chosen congestion controller unchanged; +each established path owns its pending-work proof and runtime clock. A path +with outstanding ack-eliciting work and no progress for `max(4 RTT, 500ms)` may +retire only with a recently confirmed established sibling. The existing reliable +stream transfers to that sibling through Noq's ordinary path-abandon machinery; +the last open path is never retired. See [the recovery evidence](reports/rds-path-ack-progress-20261006.md). + Status: v0.1 foundation. This document records the protocol and stack research and the decisions that fall out of it. The deeper second-pass research — iroh 1.2/noq internals (multipath, path selectors, hooks), diff --git a/docs/reports/rds-path-ack-progress-20261006.md b/docs/reports/rds-path-ack-progress-20261006.md new file mode 100644 index 0000000..6421f58 --- /dev/null +++ b/docs/reports/rds-path-ack-progress-20261006.md @@ -0,0 +1,56 @@ +# Acknowledgement progress during path selection — 2026-10-06 + +This W6.8/W10 increment handles a validated path whose old RTT remains attractive +while new reliable data no longer receives acknowledgements. RTT alone is a +measurement of prior delivery; received bytes are also insufficient because +multipath ACKs may arrive on a different path. + +## Policy and ownership + +Only explicit `PathPreference::Latency` installs the passive controller adapter +on both backends. Every controller callback, congestion window, pacing metric, +MTU and ACK-frequency update delegates unchanged to the configured BBRv3/Cubic +controller. Outstanding ack-eliciting work begins a proof budget; fresh positive +ACK progress advances it, and an empty flight ends it. Idle time does not consume +a new send's budget. Snapshots preserve evidence, while a live controller clone +starts a new proof domain on its first mutation. Runtime clock injection avoids +comparing callbacks against an unrelated wall/virtual clock. + +No ACK progress for `max(4 RTT, 500ms)` makes a path ineligible while a healthy +alternative exists. A sibling's positive ACK must be younger than `max(8 RTT, +2s)`, and it must have no overdue work. Idle candidates receive standard probes +to establish/renew that evidence. The original5ms switching stickiness remains +for healthy paths. A high RTT cannot be truncated into premature failure by a +fixed timeout. Unknown instrumentation retains the prior policy. + +Retirement activates the confirmed established sibling in the same connection +before closing the failed path. Noq refuses closing the last open path; the +connection, endpoint, application stream, authorization and input sequence are +retained. This is not input replay or QUIC loss/crypto reimplementation. Iroh's +three small context methods remain in the existing single-file vendor patch. +Ordinary/default and pinned selectors preserve their behavior. + +## Native isolated evidence + +The Noq fixture uses real IPv4/IPv6 loopback sockets, authenticated connections, +a validated standby and a bounded switch that silently drops outgoing primary +packets after handshake. It proves a pending stream frame was emitted on that +path before switching. Manual retirement delivered the already sent22bytes in +the same stream; the automatic policy also recovered without reconnecting. + +The actual Iroh actor fixture starts over a bounded64-packet custom link, learns +its UDP standby, then drops the preferred custom link. Its existing stream +delivered the pending22bytes via the actor/policy in about1.1seconds. The fixture +uses Iroh's declared custom-transport types; `iroh-base`1.3.0 and `n0-watcher`1.0.0 +are test-only direct dependencies on already resolved transitive packages. + +The first expanded Mac unit run passed100tests, including both native recovery +fixtures and controller-clock/clone/idle boundaries. A high-RTT boundary was +added afterward. Final exact-source fmt/strict Clippy, complete regression, +Linux lane and installed two-device qualification are still pending. An earlier +full-workspace link failed on temporary-cache ENOSPC after Clippy had passed; +only owned inactive generated artifacts were removed, and subsequent checks use +a disk guard. The original failure is retained, not reported as passed. + +No deployment or full latency/stability acceptance is claimed by these fixtures. +They establish bounded same-connection recovery for this specific failure mode. diff --git a/vendor/iroh/RDS-PATCH.md b/vendor/iroh/RDS-PATCH.md index 31b3c5f..bf7bc3b 100644 --- a/vendor/iroh/RDS-PATCH.md +++ b/vendor/iroh/RDS-PATCH.md @@ -14,3 +14,11 @@ only weak connection references; refresh does not add connection ownership. No default behavior change for other selectors. This temporary source patch keeps the working substrate while the owned Noq backend retains its existing periodic lowest-RTT policy. Upstream qualification and convergence remain work. + +The RDS opt-in latency policy additionally reads the already-public Noq +congestion-controller snapshot, probes an established path, and can retire a +failed path after activating an established sibling in the same connection. +These small `PathSelectionData` methods remain in the same patched source file. +Noq's last-open-path guard and reliable retransmission remain unchanged. No +default selector, controller, handshake, TLS, transport eligibility or wire format +is replaced; the RDS controller adapter delegates the actual congestion algorithm. diff --git a/vendor/iroh/src/socket/remote_map/remote_state.rs b/vendor/iroh/src/socket/remote_map/remote_state.rs index adaac0a..3731805 100644 --- a/vendor/iroh/src/socket/remote_map/remote_state.rs +++ b/vendor/iroh/src/socket/remote_map/remote_state.rs @@ -1470,6 +1470,39 @@ impl<'a> PathSelectionData<'a> { StatsSource::Test(stats) => stats.as_deref().copied(), } } + + /// Snapshot the selected congestion controller without retaining I/O. + pub fn congestion_state(&self) -> Option> { + match &self.source { + StatsSource::Live { path_id, conn } => conn.congestion_state(*path_id), + #[cfg(test)] + StatsSource::Test(_) => None, + } + } + + /// Send an acknowledgement-eliciting probe on this established path. + pub fn ping(&self) -> bool { + match &self.source { + StatsSource::Live { path_id, conn } => conn.path(*path_id).is_some_and(|path| path.ping().is_ok()), + #[cfg(test)] + StatsSource::Test(_) => false, + } + } + + /// Retire only this path after activating an established alternative in the + /// same connection. Noq refuses closing the last open path. Outstanding + /// reliable bytes retain the engine's ordinary retransmission semantics. + pub fn abandon_with_fallback(&self, fallback: &Self) -> bool { + match (&self.source, &fallback.source) { + (StatsSource::Live { path_id, conn }, StatsSource::Live { path_id: other, conn: sibling }) + if path_id != other && conn.stable_id() == sibling.stable_id() => { + let Some(alternative) = conn.path(*other) else { return false }; + if alternative.set_status(PathStatus::Available).is_err() { return false } + conn.path(*path_id).is_some_and(|path| path.close().is_ok()) + } + _ => false, + } + } } /// Trait to configure path selection. From f0eebcf661bb74920820534b61e22d7a4beada36 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Tue, 6 Oct 2026 18:48:11 +0500 Subject: [PATCH 2/3] test(net): stage deterministic preferred-link blackhole --- .../src/backends/iroh/latency_tests.rs | 46 +++++++++++++++++-- 1 file changed, 41 insertions(+), 5 deletions(-) diff --git a/crates/rds-net/src/backends/iroh/latency_tests.rs b/crates/rds-net/src/backends/iroh/latency_tests.rs index d12a03c..35116eb 100644 --- a/crates/rds-net/src/backends/iroh/latency_tests.rs +++ b/crates/rds-net/src/backends/iroh/latency_tests.rs @@ -1,7 +1,8 @@ //! Real Iroh actor over a bounded controllable custom link and UDP fallback. use iroh::TransportAddr; use iroh::endpoint::transports::{ - CustomEndpoint, CustomSender, CustomTransport, RecvInfo, Transmit, + CustomEndpoint, CustomSender, CustomTransport, PathSelection, PathSelectionContext, + PathSelector, RecvInfo, Transmit, }; use iroh_base::CustomAddr; use std::collections::HashMap; @@ -123,7 +124,33 @@ impl CustomSender for LinkSender { } } -async fn endpoint(link: Link) -> iroh::Endpoint { +#[derive(Debug)] +struct SetupSelector(Arc); +impl PathSelector for SetupSelector { + fn refresh_interval(&self) -> Option { + Some(Duration::from_secs(1)) + } + fn select(&self, ctx: &PathSelectionContext<'_>) -> PathSelection { + if self.0.load(Ordering::Acquire) { + return super::latency::LatencySelector.select(ctx); + } + // Establish the failure precondition deterministically. On a busy CI + // runner the genuinely faster UDP path may otherwise win before loss + // is injected, leaving this a test of setup timing rather than recovery. + let mut choice = PathSelection::none(); + if let Some(path) = ctx.paths().find(|path| { + matches!( + path.network_path().remote(), + iroh::endpoint::transports::Addr::Custom(_) + ) + }) { + choice.set(&path); + } + choice + } +} + +async fn endpoint(link: Link, armed: Arc) -> iroh::Endpoint { let mut transport = iroh::endpoint::QuicTransportConfig::builder() .congestion_controller_factory(crate::ack_progress::factory( crate::CongestionControl::Cubic.factory(), @@ -140,7 +167,7 @@ async fn endpoint(link: Link) -> iroh::Endpoint { .unwrap() .portmapper_config(iroh::endpoint::PortmapperConfig::Disabled) .add_custom_transport(Arc::new(link)) - .path_selector(Arc::new(super::latency::LatencySelector)) + .path_selector(Arc::new(SetupSelector(armed))) .transport_config(transport.build()) .alpns(vec![rds_core::ALPN.to_vec()]) .bind() @@ -152,10 +179,11 @@ async fn endpoint(link: Link) -> iroh::Endpoint { async fn actual_iroh_actor_retires_blackholed_preferred_link_and_delivers_pending_stream_bytes() { tokio::time::timeout(Duration::from_secs(15),async { let fabric=Fabric::default(); + let armed=Arc::new(AtomicBool::new(false)); let dropping=Arc::new(AtomicBool::new(false));let dropped=Arc::new(AtomicU64::new(0)); let a_addr=CustomAddr::from_parts(0x72647374657374,b"a");let b_addr=CustomAddr::from_parts(0x72647374657374,b"b"); - let a=endpoint(Link {addr:a_addr,fabric:fabric.clone(),drop_packets:Arc::new(AtomicBool::new(false)),dropped:Arc::new(AtomicU64::new(0))}).await; - let b=endpoint(Link {addr:b_addr.clone(),fabric,drop_packets:dropping.clone(),dropped:dropped.clone()}).await; + let a=endpoint(Link {addr:a_addr,fabric:fabric.clone(),drop_packets:Arc::new(AtomicBool::new(false)),dropped:Arc::new(AtomicU64::new(0))},armed.clone()).await; + let b=endpoint(Link {addr:b_addr.clone(),fabric,drop_packets:dropping.clone(),dropped:dropped.clone()},armed.clone()).await; let target=iroh::EndpointAddr::from_parts(b.id(),[TransportAddr::Custom(b_addr)]); let (client,server)=tokio::join!(a.connect(target,rds_core::ALPN),async {b.accept().await.unwrap().await}); let (client,server)=(client.unwrap(),server.unwrap()); @@ -169,7 +197,15 @@ async fn actual_iroh_actor_retires_blackholed_preferred_link_and_delivers_pendin tokio::time::sleep(Duration::from_millis(10)).await; } }).await.expect("fixture did not establish its UDP standby while the custom link was preferred"); + let selected=server.paths().iter().find(|path|path.is_selected()).unwrap().id(); + let before=server.paths().get(selected).unwrap().stats().frame_tx.stream; dropping.store(true,Ordering::Release);let started=Instant::now();send.write_all(b"pending reliable bytes").await.unwrap(); + tokio::time::timeout(Duration::from_secs(1),async { + while dropped.load(Ordering::Relaxed)==0 || server.paths().get(selected).unwrap().stats().frame_tx.stream<=before { + tokio::time::sleep(Duration::from_millis(5)).await; + } + }).await.expect("fixture did not emit and drop the already-pending stream frame"); + armed.store(true,Ordering::Release); let mut body=[0;22];let result=tokio::time::timeout(Duration::from_secs(3),reply.read_exact(&mut body)).await; let elapsed=started.elapsed().as_millis();let injected=dropped.load(Ordering::Relaxed); tokio::join!(a.close(),b.close()); From 426daba39ed563417fe1e9bbaf22bc223d173e36 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Tue, 6 Oct 2026 19:31:24 +0500 Subject: [PATCH 3/3] test(net): establish standby ACK proof before path loss --- .../src/backends/noq/path_blackhole_tests.rs | 78 ++++++++++++++++--- .../reports/rds-path-ack-progress-20261006.md | 23 ++++-- 2 files changed, 85 insertions(+), 16 deletions(-) diff --git a/crates/rds-net/src/backends/noq/path_blackhole_tests.rs b/crates/rds-net/src/backends/noq/path_blackhole_tests.rs index 317b895..2345d0f 100644 --- a/crates/rds-net/src/backends/noq/path_blackhole_tests.rs +++ b/crates/rds-net/src/backends/noq/path_blackhole_tests.rs @@ -150,14 +150,11 @@ async fn exercise(automatic: bool) { async { server.inner.accept().await.unwrap().await } ); let (a, b) = (a.unwrap(), b.unwrap()); - let telemetry = if automatic { - Some(( - client.wire_connection(&a, Vec::new(), Vec::new()).unwrap(), - server.wire_connection(&b, Vec::new(), Vec::new()).unwrap(), - )) - } else { - None - }; + // Subscribe before establishing the standby. Do not start discovery + // during setup: QNT can legitimately add/promote routes before the + // intended failure, making the injected socket no longer primary. + let path_events = b.path_events(); + let qnt = b.nat_traversal_updates(); let (mut request, mut reply) = a.open_bi().await.unwrap(); request.write_all(b"warm").await.unwrap(); let (mut send, mut recv) = b.accept_bi().await.unwrap(); @@ -189,6 +186,23 @@ async fn exercise(automatic: bool) { .unwrap() .set_status(PathStatus::Available) .unwrap(); + // Validation completion and positive delivery acknowledgement are + // separate boundaries. Prove the standby can acknowledge before + // blackholing the path used to deliver its validation/ACK traffic. + remote_secondary.ping().unwrap(); + tokio::time::timeout(Duration::from_secs(2), async { + loop { + if b.congestion_state(remote_secondary.id()) + .and_then(crate::ack_progress::snapshot) + .is_some_and(|state| state.confirmed(remote_secondary.stats().rtt)) + { + break; + } + tokio::time::sleep(Duration::from_millis(5)).await; + } + }) + .await + .expect("standby validation did not yield acknowledgement proof"); let before = b.path_stats(noq::PathId::ZERO).unwrap().frame_tx.stream; enabled.store(true, Ordering::Release); send.write_all(b"pending reliable bytes").await.unwrap(); @@ -207,6 +221,31 @@ async fn exercise(automatic: bool) { .await .is_err() ); + let (_relay_sender, relay_dead) = tokio::sync::watch::channel(0u64); + let driver = if automatic { + let telemetry = super::telemetry::Telemetry::new(&b); + let mut validated = telemetry.paths(); + // The completed open_path future and the real remote path prove + // this fixture's already-established standby, not guessed credit. + validated.insert(remote_secondary.id(), remote_secondary.weak_handle()); + telemetry.publish(&validated, Some(noq::PathId::ZERO)); + Some(tokio::spawn(super::policy::connection_driver_observed( + b.weak_handle(), + qnt, + super::policy::Observer { + events: path_events, + telemetry, + transport: None, + }, + crate::metrics::Registry::default(), + Vec::new(), + Vec::new(), + relay_dead, + true, + ))) + } else { + None + }; let started = Instant::now(); if !automatic { remote_secondary.set_status(PathStatus::Available).unwrap(); @@ -215,6 +254,22 @@ async fn exercise(automatic: bool) { let result = tokio::time::timeout(Duration::from_secs(2), reply.read_exact(&mut body)).await; let delivered_ms = started.elapsed().as_millis(); + if result.is_err() { + for id in [noq::PathId::ZERO, remote_secondary.id()] { + let Some(path) = b.path(id) else { + eprintln!("recovery path {id} no longer open"); + continue; + }; + eprintln!( + "recovery path {} status={:?} stats={:?} progress={:?}", + path.id(), + path.status(), + path.stats(), + b.congestion_state(path.id()) + .and_then(crate::ack_progress::snapshot) + ); + } + } // Clean up the isolated endpoints even when the engine boundary fails. tokio::join!(client.close(), server.close()); result @@ -222,7 +277,12 @@ async fn exercise(automatic: bool) { .unwrap(); assert_eq!(&body, b"pending reliable bytes"); assert!(delivered_ms < 2000); - drop(telemetry); + if let Some(driver) = driver { + tokio::time::timeout(Duration::from_secs(1), driver) + .await + .expect("policy survived connection closure") + .unwrap(); + } }) .await .expect("isolated blackhole fixture did not terminate"); diff --git a/docs/reports/rds-path-ack-progress-20261006.md b/docs/reports/rds-path-ack-progress-20261006.md index 6421f58..691592b 100644 --- a/docs/reports/rds-path-ack-progress-20261006.md +++ b/docs/reports/rds-path-ack-progress-20261006.md @@ -44,13 +44,22 @@ delivered the pending22bytes via the actor/policy in about1.1seconds. The fixtur uses Iroh's declared custom-transport types; `iroh-base`1.3.0 and `n0-watcher`1.0.0 are test-only direct dependencies on already resolved transitive packages. -The first expanded Mac unit run passed100tests, including both native recovery -fixtures and controller-clock/clone/idle boundaries. A high-RTT boundary was -added afterward. Final exact-source fmt/strict Clippy, complete regression, -Linux lane and installed two-device qualification are still pending. An earlier -full-workspace link failed on temporary-cache ENOSPC after Clippy had passed; -only owned inactive generated artifacts were removed, and subsequent checks use -a disk guard. The original failure is retained, not reported as passed. +The final local Mac net run passed101tests, including both native recovery +fixtures and controller-clock/clone/idle/high-RTT boundaries; fmt and strict +all-target net Clippy also passed. A previous complete workspace regression +passed after raising the check process file-descriptor limit to4096. +Cross-platform CI exposed fixture setup races: Iroh must await an available +preferred link, and Noq must establish standby ACK evidence before injected +loss. The Noq fixture now starts the real production policy only after explicit +validation, positive acknowledgement and proof that the primary emitted the +lost stream frame. Discovery during setup cannot promote an unrelated route. +Linux and refreshed CI qualification remain pending for this fixture revision. + +An earlier full-workspace link failed on temporary-cache ENOSPC after Clippy had +passed; only owned inactive generated artifacts were removed, and subsequent +checks use a disk guard. A later workspace attempt hit the default256descriptor +limit. Original failures are retained, not reported as passed. No deployment or +installed two-device qualification has occurred. No deployment or full latency/stability acceptance is claimed by these fixtures. They establish bounded same-connection recovery for this specific failure mode.