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
18 changes: 18 additions & 0 deletions crates/rds-net/src/backends/noq/policy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -388,6 +388,24 @@ pub(super) async fn connection_driver_observed(
}
// Lost/lagged path events cannot accumulate stale QNT history.
qnt_paths.retain(|id| owner.path(*id).is_some());
// Established/Abandoned delivery is a bounded broadcast: a lagged
// event must not hide a live path from selection nor keep a dead
// one selected. The negotiated PathId space is bounded, so
// reconcile by rescan rather than trusting event completeness.
// Adoption keeps the Established boundary: a path joins selection
// only after our challenge was answered, the same condition that
// emits the event — unvalidated candidates stay out.
for raw in 0..super::MAX_MULTIPATH_PATHS {
let id = noq::PathId::from(raw);
if paths.contains_key(&id) {
continue;
}
if let Some(path) = owner.path(id)
&& path.stats().frame_rx.path_response > 0
{
paths.insert(id, path.weak_handle());
}
}
drop(owner);
reselect(
&conn,
Expand Down
6 changes: 3 additions & 3 deletions crates/rds-net/tests/mux_isolation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -246,7 +246,7 @@ async fn exercise(receive: bool) {
// PathId. Both policies must observe a validated IPv6 sibling, not
// necessarily the particular path returned by our explicit open.
let facade = rds_net::Connection::from(b.clone());
tokio::time::timeout(Duration::from_secs(2), async {
tokio::time::timeout(Duration::from_secs(8), async {
while observed_path_to(&a, secondary_addr, false).is_none()
|| observed_path_to(&b, client_secondary_addr, false).is_none()
{
Expand All @@ -269,7 +269,7 @@ async fn exercise(receive: bool) {
let mut prefix = [0; 6];
request.read_exact(&mut prefix).await.unwrap();
assert_eq!(&prefix, b"before");
tokio::time::timeout(Duration::from_secs(2), async {
tokio::time::timeout(Duration::from_secs(5), async {
while !a
.inner()
.get_remote_nat_traversal_addresses()
Expand Down Expand Up @@ -337,7 +337,7 @@ async fn exercise(receive: bool) {
assert_eq!(health.snapshot()[0].failed, Some(io::ErrorKind::BrokenPipe));
assert!(health.snapshot()[1].failed.is_none());
assert!(observed_path_to(&b, client_secondary_addr, true).is_some());
tokio::time::timeout(Duration::from_secs(2), async {
tokio::time::timeout(Duration::from_secs(5), async {
while a
.inner()
.get_remote_nat_traversal_addresses()
Expand Down
Loading