Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 0 additions & 2 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

5 changes: 5 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
[workspace]
exclude = ["vendor/iroh"]
resolver = "3"
members = [
"crates/rds-agent",
Expand Down Expand Up @@ -104,3 +105,7 @@ rust_2018_idioms = "warn"
[profile.release]
lto = "thin"
codegen-units = 1

# Exact published Iroh1.3 source plus opt-in periodic selector refresh.
[patch.crates-io]
iroh = { path = "vendor/iroh" }
50 changes: 50 additions & 0 deletions crates/rds-desktop/src/control_observation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,18 @@ impl ControlObservation {
// diagnostic budget. Truncation must remain explicit.
snapshot.paths.sort_by_key(|path| !path.selected);
snapshot.paths.truncate(MAX_PATHS);
if selection_changed(&snapshot.paths, &self.previous) {
tracing::info!(target:"rds_desktop::control_timing", control_instance=self.instance,
heartbeat_seq=seq, coverage=?snapshot.coverage, observed_paths=observed,
paths_truncated=observed>MAX_PATHS,
"desktop selected transmit paths changed");
for path in snapshot.paths.iter().filter(|path| path.selected) {
tracing::info!(target:"rds_desktop::control_timing", control_instance=self.instance,
heartbeat_seq=seq, path_id=path.path_id, via_relay=path.via_relay,
path_rtt_ms=path.rtt.as_millis(),
"desktop selected transmit path");
}
}
if rtt >= SLOW_RESPONSE {
tracing::warn!(target:"rds_desktop::control_timing", control_instance=self.instance,
heartbeat_seq=seq, rtt_ms=rtt.as_millis(),
Expand Down Expand Up @@ -140,6 +152,22 @@ impl ControlObservation {
}
}

fn selection_changed(current: &[PathStats], previous: &[PathStats]) -> bool {
fn same(path: &PathStats, candidates: &[PathStats]) -> bool {
candidates.iter().any(|other| {
other.selected && path.path_id == other.path_id && path.via_relay == other.via_relay
})
}
current
.iter()
.filter(|path| path.selected)
.any(|path| !same(path, previous))
|| previous
.iter()
.filter(|path| path.selected)
.any(|path| !same(path, current))
}

struct Deltas {
sent: Option<u64>,
lost: Option<u64>,
Expand Down Expand Up @@ -177,6 +205,28 @@ mod tests {
}
}

#[test]
fn selection_observation_tracks_route_changes_without_rtt_chatter() {
let direct = path(5, 0);
let mut relay = path(10, 0);
relay.path_id = 2;
relay.via_relay = true;
let mut newer_direct = path(20, 1);
newer_direct.rtt = Duration::from_millis(150);
assert!(selection_changed(std::slice::from_ref(&direct), &[]));
assert!(!selection_changed(
&[newer_direct],
std::slice::from_ref(&direct)
));
assert!(selection_changed(
std::slice::from_ref(&relay),
std::slice::from_ref(&direct)
));
assert!(selection_changed(&[], std::slice::from_ref(&direct)));
assert!(!selection_changed(&[], &[]));
assert!(!selection_changed(&[direct, relay], &[relay, direct]));
}

#[test]
fn unobserved_or_reset_counters_are_unknown_not_zero_loss() {
let current = path(5, 2);
Expand Down
4 changes: 4 additions & 0 deletions crates/rds-net/src/backends/iroh.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ use std::str::FromStr;
use iroh::{Endpoint, RelayMap, RelayMode};

use crate::{EndpointAddr, EndpointConfig, EndpointId, RelayUrl, TransportAddr};
mod latency;

/// Adapter conversions between the owned shared types and iroh-base.
///
Expand Down Expand Up @@ -115,6 +116,9 @@ pub async fn bind_endpoint(config: EndpointConfig) -> anyhow::Result<Endpoint> {
crate::Transports::DirectOnly => builder.clear_relay_transports(),
crate::Transports::RelayOnly => builder.clear_ip_transports(),
};
if config.path_preference == crate::PathPreference::Latency {
builder = builder.path_selector(std::sync::Arc::new(latency::LatencySelector));
}
// Tuning on top of iroh's multipath-aware defaults:
// - BBRv3 remains our default; explicit Cubic selection enables
// same-path qualification without changing priorities or windows.
Expand Down
101 changes: 101 additions & 0 deletions crates/rds-net/src/backends/iroh/latency.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
//! Rank existing paths by RTT without making a relay permanently secondary.
use std::time::Duration;

use iroh::endpoint::transports::{PathSelection, PathSelectionContext, PathSelector};

const SWITCH_GAIN: Duration = Duration::from_millis(5);

#[derive(Debug)]
pub(super) struct LatencySelector;

impl PathSelector for LatencySelector {
fn refresh_interval(&self) -> Option<Duration> {
Some(Duration::from_secs(1))
}

fn select(&self, ctx: &PathSelectionContext<'_>) -> PathSelection {
let choice = choose(ctx.paths().filter_map(|path| {
let rtt = path.stats()?.rtt;
let current = Some(path.network_path()) == ctx.current();
Some((path, rtt, current))
}));
let mut selection = PathSelection::none();
if let Some(path) = choice {
selection.set(&path);
}
selection
}
}

fn choose<T>(paths: impl Iterator<Item = (T, Duration, bool)>) -> Option<T> {
let mut best: Option<(T, Duration)> = None;
let mut current: Option<Duration> = None;
for (path, rtt, selected) in paths {
if selected && current.is_none_or(|old| rtt < old) {
current = Some(rtt);
}
if best.as_ref().is_none_or(|(_, old)| rtt < *old) {
best = Some((path, rtt));
}
}
let (path, rtt) = best?;
if current.is_none_or(|old| old.saturating_sub(rtt) >= SWITCH_GAIN) {
Some(path)
} else {
// The Iroh selector contract interprets an empty selection as keep current.
None
}
}

#[cfg(test)]
mod tests {
use super::*;

#[derive(Debug, PartialEq, Clone, Copy)]
enum Route {
Direct,
Relay,
}

#[test]
fn usable_direct_path_does_not_hide_a_faster_relay() {
let paths = [
(Route::Direct, Duration::from_millis(170), true),
(Route::Relay, Duration::from_millis(65), false),
];
assert_eq!(choose(paths.into_iter()), Some(Route::Relay));
assert_eq!(choose(paths.into_iter().rev()), Some(Route::Relay));
}

#[test]
fn stickiness_ties_missing_current_and_extreme_values_are_explicit() {
let ms = Duration::from_millis;
assert_eq!(
choose([(Route::Direct, ms(69), true), (Route::Relay, ms(65), false)].into_iter()),
None
);
assert_eq!(
choose([(Route::Direct, ms(70), true), (Route::Relay, ms(65), false)].into_iter()),
Some(Route::Relay)
);
assert_eq!(
choose([(Route::Direct, ms(65), true), (Route::Relay, ms(65), false)].into_iter()),
None
);
assert_eq!(
choose([(Route::Relay, ms(65), false)].into_iter()),
Some(Route::Relay)
);
assert_eq!(
choose(
[
(Route::Direct, Duration::MAX, true),
(Route::Relay, Duration::MAX, false)
]
.into_iter()
),
None
);
assert_eq!(choose(std::iter::empty::<(Route, Duration, bool)>()), None);
}
}
7 changes: 7 additions & 0 deletions crates/rds-net/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,8 @@ pub struct EndpointSettings {
pub transports: crate::Transports,
#[serde(default, skip_serializing_if = "is_adaptive_packetization")]
pub packetization: crate::Packetization,
#[serde(default, skip_serializing_if = "is_backend_path_preference")]
pub path_preference: crate::PathPreference,
#[serde(default, skip_serializing_if = "is_bbr3")]
pub congestion_control: crate::CongestionControl,
}
Expand All @@ -143,6 +145,9 @@ fn is_adaptive_packetization(value: &crate::Packetization) -> bool {
fn is_bbr3(value: &crate::CongestionControl) -> bool {
*value == crate::CongestionControl::Bbr3
}
fn is_backend_path_preference(value: &crate::PathPreference) -> bool {
*value == crate::PathPreference::BackendDefault
}

impl Default for EndpointSettings {
fn default() -> Self {
Expand All @@ -154,6 +159,7 @@ impl Default for EndpointSettings {
max_multipath_paths: None,
transports: crate::Transports::default(),
packetization: crate::Packetization::default(),
path_preference: crate::PathPreference::default(),
congestion_control: crate::CongestionControl::default(),
}
}
Expand Down Expand Up @@ -256,6 +262,7 @@ impl EndpointSettings {
max_multipath_paths: self.max_multipath_paths,
transports: self.transports,
packetization: self.packetization,
path_preference: self.path_preference,
congestion_control: self.congestion_control,
discovery: self.backend == Backend::Iroh,
..Default::default()
Expand Down
31 changes: 31 additions & 0 deletions crates/rds-net/src/config/tests.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,36 @@
use super::*;

#[test]
fn latency_path_preference_is_explicit_and_cannot_expand_transport_scope() {
let legacy = EndpointSettings::from_json(br#"{"schema_version":1}"#).unwrap();
assert_eq!(
legacy.path_preference,
crate::PathPreference::BackendDefault
);
assert!(
!serde_json::to_string(&legacy)
.unwrap()
.contains("path_preference")
);
let latency = EndpointSettings::from_json(br#"{"schema_version":1,"path_preference":"latency","max_multipath_paths":1,"transports":"relay-only","relay":{"mode":"iroh","urls":["http://127.0.0.1:3340"]}}"#).unwrap();
assert_eq!(
EndpointSettings::from_json(&serde_json::to_vec(&latency).unwrap()).unwrap(),
latency
);
let lowered = latency
.apply(EndpointOverrides::default())
.unwrap()
.into_endpoint()
.unwrap();
assert_eq!(lowered.path_preference, crate::PathPreference::Latency);
assert_eq!(lowered.transports, crate::Transports::RelayOnly);
assert_eq!(lowered.max_multipath_paths, Some(1));
assert!(
EndpointSettings::from_json(br#"{"schema_version":1,"path_preference":"unknown"}"#)
.is_err()
);
}

#[test]
fn explicit_congestion_selection_preserves_legacy_defaults_and_lowers_exactly() {
let legacy = EndpointSettings::from_json(br#"{"schema_version":1}"#).unwrap();
Expand Down
15 changes: 15 additions & 0 deletions crates/rds-net/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,18 @@ pub enum Packetization {
Conservative,
}

/// Preference among already permitted, established transport paths.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum PathPreference {
/// Retain the backend's policy (Iroh prefers direct over relay).
#[default]
BackendDefault,
/// Prefer lower measured RTT regardless of direct/relay kind, with stickiness.
/// The owned Noq backend already uses this policy.
Latency,
}

/// Explicit congestion-controller selection for measured path qualification.
/// This does not change stream priority, path eligibility or packetization.
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
Expand Down Expand Up @@ -192,6 +204,8 @@ pub struct EndpointConfig {
pub transports: Transports,
/// Packet-size/offload policy. Default `Adaptive`.
pub packetization: Packetization,
/// Path ranking only; cannot expand permitted transport kinds or an explicit pin.
pub path_preference: PathPreference,
/// Local controller factory; absent file settings retain BBRv3.
pub congestion_control: CongestionControl,
/// Owned-relay attachments (`noq` backend only): each entry is a
Expand Down Expand Up @@ -223,6 +237,7 @@ impl Default for EndpointConfig {
observed_address_reports: true,
transports: Transports::default(),
packetization: Packetization::default(),
path_preference: PathPreference::default(),
congestion_control: CongestionControl::default(),
#[cfg(feature = "transport-noq")]
relay_endpoints: Vec::new(),
Expand Down
Loading
Loading