//! A long run with faults, holding the daemon to zero unrecovered disconnections (SC-007). //! //! Ignored by default because it runs for as long as it is told to: //! //! ```text //! HARBOR_SOAK_SECS=86400 cargo test --release -p midi-harbor-daemon --test endurance -- --ignored --nocapture //! ``` //! //! Two daemons in one process are joined by a loopback session, with a keyboard on one routed //! through the session to a synth on the other, and MIDI flowing throughout. Faults are injected //! in turn: the far session disappears and returns, the keyboard is unplugged and replugged, and //! the synth port is switched off and on. After each one, a message has to make it all the way //! across within SC-003's ten seconds, or the run fails. #![allow( clippy::expect_used, clippy::indexing_slicing, clippy::panic, clippy::unwrap_used )] mod common; use midi_harbor_core::endpoint::{Direction, InvitationPolicy}; use midi_harbor_core::fingerprint::DeviceFingerprint; use midi_harbor_core::midi::{Channel, MidiMessage}; use midi_harbor_daemon::Daemon; use midi_harbor_platform::fake::{FakeMidiPlatform, Outgoing}; use midi_harbor_platform::midi::{DiscoveredDevice, MidiPlatform}; use std::net::SocketAddr; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::time::{Duration, Instant}; /// SC-003: how long a recovery may take. const RECOVERY_BUDGET: Duration = Duration::from_secs(10); /// How long a fault lasts before it is undone. const OUTAGE: Duration = Duration::from_secs(2); /// How long the run lasts when `HARBOR_SOAK_SECS` does not say. const DEFAULT_SOAK: Duration = Duration::from_secs(120); /// Starts a daemon over a scratch directory, with sessions kept off the network. async fn machine(label: &str) -> (Arc, Arc) { let root = common::scratch("midi-harbor-endurance").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, ) .await .expect("the daemon starts over a scratch directory"); (daemon, platform) } /// Returns the hardware keyboard that is unplugged and plugged back in. fn keyboard() -> DiscoveredDevice { DiscoveredDevice { fingerprint: DeviceFingerprint { name: "Soak Keys".to_owned(), unique_id: Some(0x50AC), ..DeviceFingerprint::default() }, direction: Direction::Bidirectional, claimed_by: None, software: false, } } /// The two machines and what joins them. struct Rig { near: Arc, near_platform: Arc, far: Arc, far_platform: Arc, far_session: midi_harbor_core::ids::EndpointId, near_session: midi_harbor_core::ids::EndpointId, synth: midi_harbor_core::ids::EndpointId, /// Set while a probe is out, so the background traffic does not land after it and hide it. probing: Arc, } impl Rig { /// Describes both machines as they stand, so a fault that was not recovered can be explained /// from the run that hit it: the soak runs with logging off, and a day is too long to /// reproduce on demand. async fn diagnose(&self) -> String { let mut report = String::new(); for (side, daemon, session) in [ ("near", &self.near, self.near_session), ("far", &self.far, self.far_session), ] { let status = daemon.session_status(session).await.map(|status| { format!( "{:?} attempt {} peer {:?}", status.state.phase(), status.state.attempt(), status.peer_address ) }); report.push_str(&format!("{side} session: {status:?}\n")); for route in daemon.router().await.routes() { report.push_str(&format!( "{side} route {} -> {}: {:?}, enabled {}\n", route.from, route.to, route.validity, route.enabled )); } let history = daemon.events(None, usize::MAX).await; for event in history.iter().rev().take(30).rev() { report.push_str(&format!( "{side} event {:?} {:?}: {}\n", event.at, event.kind, event.detail )); } } report.push_str(&format!( "keyboard handle: {:?}\n", self.near_platform.device_handle("Soak Keys") )); report } /// Sends a probe from the keyboard and waits for it at the synth, returning how long that /// took, or nothing if it never arrived within the budget. async fn probe(&self, value: u8) -> Option { let message = MidiMessage::ControlChange { channel: Channel::new(15).expect("channel 15 is a valid MIDI channel"), controller: 119, value, }; self.probing.store(true, Ordering::Relaxed); let arrived = self.await_probe(message).await; self.probing.store(false, Ordering::Relaxed); arrived } /// Feeds `message` until it arrives at the synth or `RECOVERY_BUDGET` runs out, returning how /// long it took. async fn await_probe(&self, message: MidiMessage) -> Option { let started = Instant::now(); while started.elapsed() < RECOVERY_BUDGET { // Resent until it arrives, because a message sent mid-outage is lost by design. if let Some(keys) = self.near_platform.device_handle("Soak Keys") { let _ = self.near_platform.feed(keys, &[message]); } for _ in 0..10 { tokio::time::sleep(Duration::from_millis(20)).await; let arrived = self .far_platform .port_handle("Soak Synth") .and_then(|synth| self.far_platform.last_sent(synth)); if arrived == Some(Outgoing::Message(message)) { return Some(started.elapsed()); } } } None } } /// Proves SC-007: across a long run with faults, every disconnection is recovered, each within /// SC-003's ten seconds. /// /// Three faults take turns, each lasting two seconds before it is undone: the far session is /// switched off and on (the near session has to notice and reconnect by itself), the keyboard is /// unplugged and replugged, and the synth port is switched off and on. After each, a probe must /// cross the whole path from keyboard to synth within the budget while background notes play /// throughout, or the run fails with a description of both machines. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[ignore = "runs for as long as HARBOR_SOAK_SECS says"] async fn a_long_run_with_faults_never_stays_down() { let soak = std::env::var("HARBOR_SOAK_SECS") .ok() .and_then(|secs| secs.parse().ok()) .map_or(DEFAULT_SOAK, Duration::from_secs); // Set up the far side: a session accepting the near one, feeding a synth. let (far, far_platform) = machine("far").await; let synth = far .create_virtual_port("Soak Synth", 1, 1) .await .expect("the Soak Synth port is created"); let far_session = far .create_network_session("Soak In", 0, InvitationPolicy::AcceptAll) .await .expect("the far session is created"); far.create_route("Soak In", "Soak Synth") .await .expect("the route from Soak In to Soak Synth is created"); let port = far .session_status(far_session.id) .await .expect("the far session reports its status") .control_port; // Set up the near side: a keyboard routed into a session connected to the far one. let (near, near_platform) = machine("near").await; let near_session = near .create_network_session("Soak Out", 0, InvitationPolicy::Prompt) .await .expect("the near session is created"); near_platform.attach(keyboard()); near.refresh_devices().await; near.create_route("Soak Keys", "Soak Out") .await .expect("the route from Soak Keys to Soak Out is created"); near.connect_peer(near_session.id, SocketAddr::from(([127, 0, 0, 1], port))) .await .expect("the near session is pointed at the far one"); let rig = Rig { near: Arc::clone(&near), near_platform: Arc::clone(&near_platform), far: Arc::clone(&far), far_platform, far_session: far_session.id, near_session: near_session.id, synth: synth.id, probing: Arc::new(AtomicBool::new(false)), }; assert!( rig.probe(0).await.is_some(), "the rig never carried MIDI before any fault was injected" ); // Play background traffic, as a player would, for the whole run. let playing = Arc::new(AtomicBool::new(true)); let played = Arc::new(AtomicU64::new(0)); let traffic = { let playing = Arc::clone(&playing); let played = Arc::clone(&played); let probing = Arc::clone(&rig.probing); let platform = Arc::clone(&near_platform); tokio::spawn(async move { let channel = Channel::new(0).expect("channel 0 is a valid MIDI channel"); let mut note = 36u8; while playing.load(Ordering::Relaxed) { if !probing.load(Ordering::Relaxed) && let Some(keys) = platform.device_handle("Soak Keys") { let on = MidiMessage::NoteOn { channel, note, velocity: 80, }; let off = MidiMessage::NoteOff { channel, note, velocity: 0, }; if platform.feed(keys, &[on, off]) { played.fetch_add(2, Ordering::Relaxed); } } note = if note >= 84 { 36 } else { note + 1 }; tokio::time::sleep(Duration::from_millis(10)).await; } }) }; // Inject faults, in turn, until the time is up. let started = Instant::now(); let (mut cycles, mut slowest) = (0u64, Duration::ZERO); let mut value = 1u8; while started.elapsed() < soak { let fault = cycles % 3; match fault { // The far machine's session disappears and returns; the near one has to notice and // reconnect by itself. 0 => { rig.far .set_enabled(rig.far_session, false) .await .expect("the far session is switched off"); tokio::time::sleep(OUTAGE).await; rig.far .set_enabled(rig.far_session, true) .await .expect("the far session is switched back on"); } // The keyboard is unplugged and plugged back in. 1 => { rig.near_platform.detach("Soak Keys"); rig.near.refresh_devices().await; tokio::time::sleep(OUTAGE).await; rig.near_platform.attach(keyboard()); rig.near.refresh_devices().await; } // The synth port is switched off and on. _ => { rig.far .set_enabled(rig.synth, false) .await .expect("the synth port is switched off"); tokio::time::sleep(OUTAGE).await; rig.far .set_enabled(rig.synth, true) .await .expect("the synth port is switched back on"); } } value = if value >= 127 { 1 } else { value + 1 }; let Some(took) = rig.probe(value).await else { panic!( "cycle {cycles} (fault {fault}) was not recovered within {RECOVERY_BUDGET:?}, \ {:?} into the run\n{}", started.elapsed(), rig.diagnose().await ); }; slowest = slowest.max(took); cycles += 1; if cycles % 50 == 0 { println!( "{:?}: {cycles} faults recovered, slowest {slowest:?}, {} messages played", started.elapsed(), played.load(Ordering::Relaxed) ); } tokio::time::sleep(Duration::from_secs(1)).await; } playing.store(false, Ordering::Relaxed); let _ = traffic.await; println!( "done: {cycles} faults over {:?}, every one recovered, slowest {slowest:?}, {} messages played", started.elapsed(), played.load(Ordering::Relaxed) ); }