Skip to content

Commit 35bdd91

Browse files
mariusaemeta-codesync[bot]
authored andcommitted
route switched I/O through completion-owned UDP (#4814)
Summary: Pull Request resolved: #4814 - add `RoutedUdpPacketIo` so quiche, CID routing, pacing, GSO/GRO, and io_uring waits share one endpoint owner - route non-local CID aggregates directly between completion-owned UDP carriers while retaining a Unix datagram fallback for task-local leaves - remove the threaded `CompletionUdpSocket` adapter and preserve route admission through actual submission - use the roofline-validated eight-segment switched GSO frontier and wait for response FIN before benchmark teardown Design document: https://mdoc.internalmeta.com/doc/e6918c35-b3d1-4118-ac39-6308b0cb55ba Walkthrough document: https://mdoc.internalmeta.com/doc/ef7df54e-b047-4e9a-96b5-6055b2a411da Reviewed By: shayne-fletcher Differential Revision: D117242784
1 parent 6ee8fb3 commit 35bdd91

6 files changed

Lines changed: 1017 additions & 66 deletions

File tree

‎chrysalis-scale/src/benchmark.rs‎

Lines changed: 15 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -84,6 +84,7 @@ const META_NETWORK_MAX_UDP_PAYLOAD: u16 = 1_450;
8484
const EXPERIMENT_POLL_INTERVAL: Duration = Duration::from_millis(250);
8585
const EXPERIMENT_LEASE_GRACE: Duration = Duration::from_secs(60);
8686
const WORKER_POLL_INTERVAL: Duration = Duration::from_millis(25);
87+
const STREAM_FIN_TIMEOUT: Duration = Duration::from_secs(5);
8788
const PERSISTENCE_SCHEMA_VERSION: u32 = 8;
8889
pub(crate) const DEFAULT_NODES_PER_TASK: usize = 100;
8990
pub(crate) const DEFAULT_IDENTITY_CONCURRENCY: usize = 8;
@@ -640,22 +641,18 @@ async fn create_local_node(
640641
) -> Result<LocalNode> {
641642
let udp_address = socket.address();
642643
let identity = issue_scale_identity(rank).await?;
644+
let quic_config = scale_quic_config()?;
643645
let transport = match unix_path {
644646
Some(path) => {
645-
let primary: Arc<dyn DatagramSocket> = Arc::new(socket);
646-
let alternatives: Vec<Arc<dyn DatagramSocket>> = vec![Arc::new(
647+
let fallback: Arc<dyn DatagramSocket> = Arc::new(
647648
UnixDatagramSocket::bind(path)
648649
.with_context(|| format!("bind task-head Unix carrier {}", path.display()))?,
649-
)];
650-
let sockets = Arc::new(
651-
DatagramSocketSet::new(primary, alternatives)
652-
.context("create task-head socket set")?,
653650
);
654-
TransportConfig::new(sockets, identity)
651+
TransportConfig::routed_udp(socket.into_std()?, Some(fallback), identity)?
655652
}
656653
None => TransportConfig::direct_udp(socket.into_std()?, identity)?,
657654
}
658-
.with_quic_config(scale_quic_config()?);
655+
.with_quic_config(quic_config);
659656
let mut config = NodeConfig::new(transport);
660657
let role = if is_root {
661658
NodeRole::Root
@@ -1753,7 +1750,16 @@ fn spawn_operation(
17531750
read_delivery_receipt(&mut recv, size).await?;
17541751
}
17551752
}
1756-
Ok(started.elapsed())
1753+
let elapsed = started.elapsed();
1754+
let mut trailing = [0];
1755+
anyhow::ensure!(
1756+
tokio::time::timeout(STREAM_FIN_TIMEOUT, recv.read(&mut trailing))
1757+
.await
1758+
.context("timed out waiting for benchmark response FIN")??
1759+
== 0,
1760+
"benchmark response contains trailing data"
1761+
);
1762+
Ok(elapsed)
17571763
});
17581764
}
17591765

‎chrysalis-scale/src/network_baseline.rs‎

Lines changed: 24 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@ use anyhow::Context;
1717
use anyhow::Result;
1818
use bytes::Bytes;
1919
use chrysalis::DatagramSocket;
20-
use chrysalis::DatagramSwitch;
2120
use chrysalis::Pid;
2221
use chrysalis::QuicConnectionStats;
2322
use chrysalis::QuicIoStats;
@@ -1127,24 +1126,23 @@ async fn receive_switched_quic(
11271126
timeout: Duration,
11281127
profile: NetworkBaselineProfile,
11291128
) -> Result<()> {
1130-
let physical = Arc::new(
1131-
UdpSocket::bind(local)
1132-
.await
1133-
.with_context(|| format!("bind switched QUIC receiver {local}"))?,
1134-
);
1129+
let quic_config = scale_quic_config()?;
11351130
let router = Arc::new(Router::new());
11361131
router.insert(
11371132
peer,
11381133
Route::permanent(UdpSocket::datagram_addr(peer_address)),
11391134
);
1140-
let datagram_switch = DatagramSwitch::spawn(physical, router);
1141-
let binding = Arc::new(
1142-
datagram_switch
1143-
.bind_routed(identity.pid())
1144-
.context("bind switched QUIC receiver PID")?,
1145-
);
1146-
let transport = QuicTransport::spawn_with_config(binding, identity, scale_quic_config()?)
1147-
.context("start switched QUIC receiver")?;
1135+
let transport = QuicTransport::spawn_routed_udp_with_config(
1136+
UdpSocket::bind(local)
1137+
.await
1138+
.with_context(|| format!("bind switched QUIC receiver {local}"))?
1139+
.into_std()?,
1140+
None,
1141+
router,
1142+
identity,
1143+
quic_config,
1144+
)
1145+
.context("start switched QUIC receiver")?;
11481146
control.write_all(&[PHASE_READY]).await?;
11491147

11501148
let (warmup_source, connection_before) = receive_quic_stream(
@@ -1181,9 +1179,6 @@ async fn receive_switched_quic(
11811179
signal_phase_complete(control).await?;
11821180
transport.shutdown();
11831181
transport.join().await;
1184-
drop(transport);
1185-
datagram_switch.shutdown();
1186-
datagram_switch.join().await;
11871182
Ok(())
11881183
}
11891184

@@ -1200,24 +1195,23 @@ async fn send_switched_quic(
12001195
let mut ready = [0];
12011196
control.read_exact(&mut ready).await?;
12021197
anyhow::ensure!(ready == [PHASE_READY], "invalid switched QUIC start marker");
1203-
let physical = Arc::new(
1204-
UdpSocket::bind(local)
1205-
.await
1206-
.with_context(|| format!("bind switched QUIC sender {local}"))?,
1207-
);
1198+
let quic_config = scale_quic_config()?;
12081199
let router = Arc::new(Router::new());
12091200
router.insert(
12101201
target,
12111202
Route::permanent(UdpSocket::datagram_addr(peer_address)),
12121203
);
1213-
let datagram_switch = DatagramSwitch::spawn(physical, router);
1214-
let binding = Arc::new(
1215-
datagram_switch
1216-
.bind_routed(identity.pid())
1217-
.context("bind switched QUIC sender PID")?,
1218-
);
1219-
let transport = QuicTransport::spawn_with_config(binding, identity, scale_quic_config()?)
1220-
.context("start switched QUIC sender")?;
1204+
let transport = QuicTransport::spawn_routed_udp_with_config(
1205+
UdpSocket::bind(local)
1206+
.await
1207+
.with_context(|| format!("bind switched QUIC sender {local}"))?
1208+
.into_std()?,
1209+
None,
1210+
router,
1211+
identity,
1212+
quic_config,
1213+
)
1214+
.context("start switched QUIC sender")?;
12211215
let destination = UdpSocket::datagram_addr(peer_address);
12221216

12231217
let connection_before = send_quic_stream(
@@ -1261,9 +1255,6 @@ async fn send_switched_quic(
12611255
wait_for_phase_complete(control).await?;
12621256
transport.shutdown();
12631257
transport.join().await;
1264-
drop(transport);
1265-
datagram_switch.shutdown();
1266-
datagram_switch.join().await;
12671258
Ok(())
12681259
}
12691260

0 commit comments

Comments
 (0)