midi-harbor/crates/daemon/tests/network_latency.rs
2026-09-28 13:59:10 -05:00

228 lines
8.6 KiB
Rust

//! How long a message takes to cross a network session (SC-009).
//!
//! Ignored by default because it is a measurement, not a check that belongs in every run:
//!
//! ```text
//! cargo test --release -p midi-harbor-daemon --test network_latency -- --ignored --nocapture
//! ```
//!
//! Two daemons in one process are joined by a real RTP-MIDI session over loopback. A port on one
//! is routed into the session, and the session on the other is routed to a port there. The
//! platform at each end is the in-memory one, so what is timed is the daemon and the session:
//! handing a message to the supervisor, packing it, the datagram, reading it, and routing it on.
//! The network between two machines adds its own time on top.
#![allow(
clippy::expect_used,
clippy::indexing_slicing,
clippy::panic,
clippy::unwrap_used
)]
mod common;
use midi_harbor_core::endpoint::InvitationPolicy;
use midi_harbor_core::midi::{Channel, MidiMessage};
use midi_harbor_core::state::ConnectionPhase;
use midi_harbor_daemon::Daemon;
use midi_harbor_platform::fake::{FakeMidiPlatform, Outgoing};
use midi_harbor_platform::midi::{MidiPlatform, PortHandle};
use std::net::SocketAddr;
use std::sync::Arc;
use std::time::{Duration, Instant};
/// Messages timed, after the warm-up.
const SAMPLES: usize = 2_000;
/// Messages sent and discarded first, while the clock exchange settles.
const WARM_UP: usize = 100;
/// SC-009: the 99th percentile a message may take across a network session.
const P99_BUDGET: Duration = Duration::from_millis(5);
/// Starts a daemon over a scratch directory, with sessions kept off the network.
async fn machine(label: &str) -> (Arc<Daemon>, Arc<FakeMidiPlatform>) {
let root = common::scratch("midi-harbor-network-latency")
.join(format!("{label}-{}", uuid::Uuid::new_v4()));
let platform = Arc::new(FakeMidiPlatform::new());
let daemon = Daemon::start(
common::quiet(root),
Arc::clone(&platform) as Arc<dyn MidiPlatform>,
)
.await
.expect("the daemon starts over a scratch directory");
(daemon, platform)
}
/// Returns how long each of `SAMPLES` messages took from `keys` on one platform to `synth` on
/// the other, after `WARM_UP` messages that are not timed.
///
/// Runs on a blocking thread so the wait for each delivery does not hold a runtime worker. The
/// gaps between messages are uneven (1 ms plus up to 4 ms from a xorshift sequence), so the
/// result is not tied to one point in either side's cycle.
async fn time_deliveries(
from: Arc<FakeMidiPlatform>,
keys: PortHandle,
to: Arc<FakeMidiPlatform>,
synth: PortHandle,
) -> Vec<Duration> {
tokio::task::spawn_blocking(move || {
let channel = Channel::new(0).expect("channel 0 is a valid MIDI channel");
let mut timings = Vec::with_capacity(SAMPLES);
let mut seed: u32 = 0x2545_F491;
for index in 0..WARM_UP + SAMPLES {
let message = MidiMessage::ControlChange {
channel,
controller: 20,
value: u8::try_from(index % 128).expect("a value modulo 128 fits a data byte"),
};
let sent = Instant::now();
assert!(
from.feed(keys, &[message]),
"message {index} was refused by the sending port"
);
loop {
if to.last_sent(synth) == Some(Outgoing::Message(message)) {
break;
}
assert!(
sent.elapsed() < Duration::from_secs(1),
"message {index} was never delivered"
);
std::thread::sleep(Duration::from_micros(20));
}
if index >= WARM_UP {
timings.push(sent.elapsed());
}
seed ^= seed << 13;
seed ^= seed >> 17;
seed ^= seed << 5;
std::thread::sleep(Duration::from_micros(1_000 + u64::from(seed % 4_000)));
}
timings
})
.await
.expect("the timing thread finishes without panicking")
}
/// Returns the timing at `fraction` of the way through `sorted`, clamped to the last one.
fn percentile(sorted: &[Duration], fraction: f64) -> Duration {
sorted[((sorted.len() as f64 * fraction) as usize).min(sorted.len() - 1)]
}
/// Proves SC-009: the 99th percentile of messages crosses a loopback RTP-MIDI session within
/// 5 ms.
///
/// Times 2,000 messages after 100 of warm-up while the clock exchange settles. Loopback removes
/// the network, so the figure is the daemon's and the session's share; `the_same_route_without_a_session_for_comparison`
/// separates the two.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[ignore = "a measurement; run it deliberately"]
async fn a_message_crosses_a_network_session_within_budget() {
// Set up the far side: a session accepting the near one, feeding a synth.
let (far, far_platform) = machine("far").await;
far.create_virtual_port("Synth", 1, 1)
.await
.expect("the Synth port is created");
let incoming = far
.create_network_session("Stage In", 0, InvitationPolicy::AcceptAll)
.await
.expect("the far session is created");
far.create_route("Stage In", "Synth")
.await
.expect("the route from Stage In to Synth is created");
let port = far
.session_status(incoming.id)
.await
.expect("the far session reports its status")
.control_port;
// Set up the near side: a keyboard port routed into a session connected to the far one.
let (near, near_platform) = machine("near").await;
near.create_virtual_port("Keys", 1, 1)
.await
.expect("the Keys port is created");
let outgoing = near
.create_network_session("Stage Out", 0, InvitationPolicy::Prompt)
.await
.expect("the near session is created");
near.create_route("Keys", "Stage Out")
.await
.expect("the route from Keys to Stage Out is created");
near.connect_peer(outgoing.id, SocketAddr::from(([127, 0, 0, 1], port)))
.await
.expect("the near session is pointed at the far one");
for _ in 0..100 {
let phase = near
.session_status(outgoing.id)
.await
.map(|status| status.state.phase());
if phase == Some(ConnectionPhase::Connected) {
break;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
// Time the messages.
let keys = near_platform
.port_handle("Keys")
.expect("the Keys port has a platform handle");
let synth = far_platform
.port_handle("Synth")
.expect("the Synth port has a platform handle");
let mut sorted = time_deliveries(near_platform, keys, far_platform, synth).await;
sorted.sort();
// Report and judge them.
let total: Duration = sorted.iter().sum();
let mean = total / u32::try_from(sorted.len()).expect("the sample count fits a u32");
let (p50, p99, max) = (
percentile(&sorted, 0.50),
percentile(&sorted, 0.99),
percentile(&sorted, 1.0),
);
println!(
"{} messages over a loopback session: mean {mean:?}, p50 {p50:?}, p99 {p99:?}, max {max:?}",
sorted.len()
);
assert!(
p99 < P99_BUDGET,
"the 99th percentile of {p99:?} is over SC-009's {P99_BUDGET:?} budget"
);
}
/// Times messages from one port to another on a single daemon, with no session in between.
///
/// The control for the measurement above: whatever this shows is the daemon's own share, and the
/// difference is the session's. It prints rather than asserts, since SC-009 budgets the session.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[ignore = "a measurement; run it deliberately"]
async fn the_same_route_without_a_session_for_comparison() {
let (daemon, platform) = machine("local").await;
for name in ["Keys", "Synth"] {
daemon
.create_virtual_port(name, 1, 1)
.await
.expect("the virtual port is created");
}
daemon
.create_route("Keys", "Synth")
.await
.expect("the route from Keys to Synth is created");
let keys = platform
.port_handle("Keys")
.expect("the Keys port has a platform handle");
let synth = platform
.port_handle("Synth")
.expect("the Synth port has a platform handle");
let mut sorted = time_deliveries(Arc::clone(&platform), keys, platform, synth).await;
sorted.sort();
println!(
"{} messages through one daemon: p50 {:?}, p99 {:?}, max {:?}",
sorted.len(),
percentile(&sorted, 0.50),
percentile(&sorted, 0.99),
percentile(&sorted, 1.0)
);
}