diff --git a/native/opennow-streamer/crates/opennow-streamer-core/src/lib.rs b/native/opennow-streamer/crates/opennow-streamer-core/src/lib.rs index 311681aaa..7446273e2 100644 --- a/native/opennow-streamer/crates/opennow-streamer-core/src/lib.rs +++ b/native/opennow-streamer/crates/opennow-streamer-core/src/lib.rs @@ -2246,6 +2246,39 @@ fn forward_nvst_media_feedback( }), )); } + MediaFeedback::AudioDecoderError { + message, + consecutive, + } => { + let _ = output.send(event( + "log", + json!({ + "event": "decoder-error", + "codec": "opus", + "consecutive": consecutive, + "level": "warn", + "message": format!("Opus decoder error: {message}") + }), + )); + } + MediaFeedback::AudioUnavailable { + backend, + reason, + rejected, + } => { + let message = format!("Audio is unavailable for the rest of this session: {reason}"); + opennow_streamer_protocol::log::log_line("WARN", "media-audio", &message); + let _ = output.send(event( + "log", + json!({ + "event": "audio-unavailable", + "backend": backend, + "rejectedPackets": rejected, + "level": "warn", + "message": message + }), + )); + } MediaFeedback::OutputError { message } => { let _ = output.send(event( "error", @@ -3019,6 +3052,92 @@ mod tests { ); } + #[test] + fn audio_decode_feedback_stays_non_fatal_and_audio_scoped() { + let (sender, receiver) = std::sync::mpsc::channel(); + let sender = EventSender::unbounded(sender); + let lifecycle = connected_lifecycle(); + let resources = TestNvstResources::default(); + let mut state = NvstMediaFeedbackState::new(true); + + forward_nvst_media_feedback( + &sender, + &lifecycle, + 7, + &resources, + MediaFeedback::AudioDecoderError { + message: "corrupted stream".to_owned(), + consecutive: 3, + }, + &mut state, + ); + + let message = receiver.recv().expect("audio decode log"); + assert_eq!(message["type"], "log"); + assert_eq!(message["event"], "decoder-error"); + assert_eq!(message["codec"], "opus"); + assert_eq!(message["level"], "warn"); + assert_eq!(message["consecutive"], 3); + assert!( + message["message"] + .as_str() + .is_some_and(|message| message.contains("corrupted stream")) + ); + + forward_nvst_media_feedback( + &sender, + &lifecycle, + 7, + &resources, + MediaFeedback::AudioUnavailable { + backend: "ALSA", + reason: "audio output was lost".to_owned(), + rejected: 0, + }, + &mut state, + ); + + let message = receiver.recv().expect("audio unavailable log"); + assert_eq!(message["type"], "log"); + assert_eq!(message["event"], "audio-unavailable"); + assert_eq!(message["backend"], "ALSA"); + assert_eq!(message["rejectedPackets"], 0); + assert_eq!(message["level"], "warn"); + assert_eq!(resources.stops.load(Ordering::Relaxed), 0); + assert_eq!(resources.keyframe_requests.load(Ordering::Relaxed), 0); + } + + #[test] + fn audio_output_device_loss_feedback_stays_non_fatal() { + let (sender, receiver) = std::sync::mpsc::channel(); + let sender = EventSender::unbounded(sender); + let lifecycle = connected_lifecycle(); + let resources = TestNvstResources::default(); + let mut state = NvstMediaFeedbackState::new(true); + + forward_nvst_media_feedback( + &sender, + &lifecycle, + 7, + &resources, + MediaFeedback::DeviceLost { + subsystem: "ALSA", + recovered: false, + message: Some("Alsa device was lost: Broken pipe (os error 32)".to_owned()), + }, + &mut state, + ); + + let message = receiver.recv().expect("device state log"); + assert_eq!(message["type"], "log"); + assert_eq!(message["event"], "device-state"); + assert_eq!(message["subsystem"], "ALSA"); + assert_eq!(message["recovered"], false); + assert_eq!(message["level"], "warn"); + assert_eq!(resources.stops.load(Ordering::Relaxed), 0); + assert_eq!(resources.keyframe_requests.load(Ordering::Relaxed), 0); + } + #[test] fn color_format_feedback_reports_actual_output_without_restarting_media() { let (sender, receiver) = std::sync::mpsc::channel(); diff --git a/native/opennow-streamer/crates/opennow-streamer-platform-linux/README.md b/native/opennow-streamer/crates/opennow-streamer-platform-linux/README.md index bbdc0b5f6..8d9e3da67 100644 --- a/native/opennow-streamer/crates/opennow-streamer-platform-linux/README.md +++ b/native/opennow-streamer/crates/opennow-streamer-platform-linux/README.md @@ -8,7 +8,7 @@ The standalone development host retains its existing decoder order: with the `ff Runtime probing opens the real device or service. A compiled feature or a shared library by itself is never reported as an available decoder, presenter, or audio sink. -`LinuxSession` owns the bounded decode and audio workers. Submit complete encoded access units and Opus packets, drain decoded frames and typed events, and call `reconfigure` or `stop` explicitly. Queue overflow drops the oldest media, flushes video reference state, and emits `NeedKeyframe` instead of allowing corrupted prediction chains to continue. If a decoder fails, the session advances through the configured fallback order and reuses the current keyframe when possible. +`LinuxSession` owns the bounded decode and audio workers. Submit complete encoded access units and Opus packets, drain decoded frames and typed events, and call `reconfigure` or `stop` explicitly. Queue overflow drops the oldest media, flushes video reference state, and emits `NeedKeyframe` instead of allowing corrupted prediction chains to continue. If a decoder fails, the session advances through the configured fallback order and reuses the current keyframe when possible. Audio failures stay in the audio domain: the worker skips Opus packets libopus rejects, rebuilds its decoder after eight consecutive rejections, and after two rebuilds without recovery disables audio for the rest of the session and reports `AudioUnavailable` instead of failing the shared session. A sink that cannot be recovered, including a fixed output device that forbids fallback, ends the same way, so video continues without audio. Standalone presentation stays on the caller's window-system thread. Construct `NativeSurface::borrow_x11` or `NativeSurface::borrow_wayland` from handles that the caller owns, then create `VulkanPresenter` and feed it frames returned by the session. The standalone development runtime obtains either pair from SDL's raw window handles: X11 is reparented into the host-owned surface, while Wayland remains a compositor-managed top-level surface. The presenter owns its Vulkan instance, device, swapchain, and `VkSurfaceKHR`; it never destroys the borrowed X11 window, Xlib display, `wl_surface`, or `wl_display`. Its lifetime is tied to `NativeSurface`, and it is intentionally not `Send`. @@ -47,7 +47,7 @@ cargo check -p opennow-streamer-platform-linux --target aarch64-unknown-linux-gn The H.264 V4L2 fallback supports both single-planar and multi-planar stateful decoder nodes used on x86_64, aarch64, and Raspberry Pi 4. Bundled Linux aarch64 builds also include pinned Raspberry Pi FFmpeg support for HEVC V4L2 Request decoding on Pi 4 and Pi 5. That HEVC path is embedded-only, supports 8-bit 4:2:0 SDR, and requires Qt's explicit DMA-BUF buffer-import capability. SAND128 single-object NC12 and split-object Nc12 frames are copied into planar textures entirely on the GPU; packed P030/10-bit is rejected. See [Raspberry Pi requirements and acceptance](../../../../docs/raspberry-pi.md). -Actual decode tests require `/dev/video*` and, for HEVC Request, `/dev/media*`; VA-API tests require `/dev/dri/renderD*`; Vulkan presentation requires a live X11 or Wayland surface; audio requires an SDL-supported system audio service. The ignored `raspberry_pi_hevc_request_drm_only` test exercises Pi HEVC Request output, and `local_hardware_decodes_all_required_codecs` exercises Vulkan Video and CUDA/NVDEC, when the corresponding hardware and FFmpeg command-line encoder are installed. Ordinary unit tests do not claim those devices exist. +Actual decode tests require `/dev/video*` and, for HEVC Request, `/dev/media*`; VA-API tests require `/dev/dri/renderD*`; Vulkan presentation requires a live X11 or Wayland surface; audio requires an SDL-supported system audio service. The audio failure-boundary tests are ordinary unit tests: they point `ALSA_CONFIG_PATH` at `tests/fixtures/audio-failure.conf`, which defines an in-process ALSA `null` output and an ALSA `file` output writing to `/dev/full`, so they need `libasound.so.2`, `libopus.so.0`, and the standard `/dev/full` device node rather than audio hardware. The ignored `raspberry_pi_hevc_request_drm_only` test exercises Pi HEVC Request output, and `local_hardware_decodes_all_required_codecs` exercises Vulkan Video and CUDA/NVDEC, when the corresponding hardware and FFmpeg command-line encoder are installed. Ordinary unit tests do not claim those devices exist. The ignored shared-device HEVC tests exercise the embedded GPU-only snapshot path separately for NV12, P010, NV24, and P410. Run them on a Vulkan Video device supporting the corresponding HEVC profiles; 4:4:4 requires Range Extensions support. Compilation, ordinary unit tests, and synthetic conversion tests on Mesa software Vulkan do not establish hardware decode support: diff --git a/native/opennow-streamer/crates/opennow-streamer-platform-linux/src/audio.rs b/native/opennow-streamer/crates/opennow-streamer-platform-linux/src/audio.rs index 90cf0111d..f546597bb 100644 --- a/native/opennow-streamer/crates/opennow-streamer-platform-linux/src/audio.rs +++ b/native/opennow-streamer/crates/opennow-streamer-platform-linux/src/audio.rs @@ -250,6 +250,17 @@ impl OpusDecoder { Ok(&self.pcm[..samples_per_channel * self.channels]) } + pub fn record_dropped_packet(&mut self, packet: &AudioPacket) { + let frame_samples = self.last_frame_samples_per_channel; + if frame_samples == 0 { + return; + } + let frame_ticks = + (frame_samples as u64) * u64::from(packet.clock_rate_hz) / u64::from(self.sample_rate); + self.last_ssrc = Some(packet.ssrc); + self.last_timestamp = Some(packet.rtp_timestamp.wrapping_sub(frame_ticks as u32)); + } + pub fn conceal_before<'a>(&'a mut self, packet: &AudioPacket) -> Result<&'a [f32]> { if self.source_changed(packet.ssrc) { self.reset_decoder_state()?; @@ -865,6 +876,33 @@ mod tests { assert_eq!(decoder.decode(&one_lost).expect("decode").len(), 1_920); } + #[test] + fn dropped_packet_is_concealed_once_by_the_next_gap() { + let (mut decoder, mut encoder) = decoder_and_encoder(); + let first = audio_packet(encoded_frame(&mut encoder, 960, 440.0), 0); + assert_eq!(decoder.decode(&first).expect("decode").len(), 1_920); + + let malformed = AudioPacket::new(Arc::<[u8]>::from(vec![0xff; 40]), 1_920, 48_000, 7) + .expect("malformed packet"); + let first_concealment = decoder.conceal_before(&malformed).expect("concealment"); + assert_eq!(first_concealment.len(), 1_920); + assert!(decoder.decode(&malformed).is_err()); + decoder.record_dropped_packet(&malformed); + + let next = audio_packet(encoded_frame(&mut encoder, 960, 440.0), 2_880); + let second_concealment = decoder.conceal_before(&next).expect("concealment"); + assert_eq!(second_concealment.len(), 1_920); + assert_eq!(decoder.decode(&next).expect("decode").len(), 1_920); + + let contiguous = audio_packet(encoded_frame(&mut encoder, 960, 440.0), 3_840); + assert!( + decoder + .conceal_before(&contiguous) + .expect("contiguous") + .is_empty() + ); + } + #[test] fn burst_losses_conceal_every_missing_frame_duration() { let (mut decoder, mut encoder) = decoder_and_encoder(); diff --git a/native/opennow-streamer/crates/opennow-streamer-platform-linux/src/session.rs b/native/opennow-streamer/crates/opennow-streamer-platform-linux/src/session.rs index 1c717d9b4..a3d72e836 100644 --- a/native/opennow-streamer/crates/opennow-streamer-platform-linux/src/session.rs +++ b/native/opennow-streamer/crates/opennow-streamer-platform-linux/src/session.rs @@ -129,6 +129,7 @@ pub enum PushOutcome { Queued, DroppedOldest, Paused, + AudioDisabled, } #[derive(Debug)] @@ -141,6 +142,23 @@ pub enum BackendEvent { reason: String, }, AudioSelected(AudioBackend), + AudioDecodeError { + message: String, + consecutive: u32, + }, + AudioOutputError { + backend: AudioBackend, + message: String, + }, + AudioOutputRecovered { + from: AudioBackend, + to: AudioBackend, + }, + AudioUnavailable { + backend: AudioBackend, + reason: String, + rejected: u64, + }, FormatChanged(StreamFormat), NeedKeyframe, QueueOverflow { @@ -173,6 +191,7 @@ pub struct LinuxSession { video_commands: Arc>, decoded_frames: Arc>, audio_packets: Option>>, + audio_unavailable: Arc, events: EventQueue, video_generation: AtomicU64, video_needs_keyframe: AtomicBool, @@ -225,17 +244,18 @@ impl LinuxSession { } } + let audio_unavailable = Arc::new(AtomicBool::new(false)); let (audio_packets, audio_worker) = if let Some(audio_config) = config.audio.clone() { let queue = Arc::new(BoundedQueue::new(audio_config.queue_depth)); let (audio_start_tx, audio_start_rx) = mpsc::sync_channel(1); let worker = { let queue = Arc::clone(&queue); let events = Arc::clone(&events); - let state = Arc::clone(&state); + let unavailable = Arc::clone(&audio_unavailable); match thread::Builder::new() .name("opennow-linux-audio".to_owned()) .spawn(move || { - run_audio_worker(audio_config, queue, state, events, audio_start_tx) + run_audio_worker(audio_config, queue, events, unavailable, audio_start_tx) }) { Ok(worker) => worker, Err(error) => { @@ -275,6 +295,7 @@ impl LinuxSession { video_commands, decoded_frames, audio_packets, + audio_unavailable, events, video_generation: AtomicU64::new(0), video_needs_keyframe: AtomicBool::new(false), @@ -318,21 +339,16 @@ impl LinuxSession { if self.state() != LifecycleState::Running { return Err(Error::NotRunning); } + if self.audio_unavailable.load(Ordering::Acquire) { + return Ok(PushOutcome::AudioDisabled); + } if self.paused.load(Ordering::Acquire) { return Ok(PushOutcome::Paused); } let queue = self.audio_packets.as_ref().ok_or_else(|| { Error::unavailable(Subsystem::Session, "audio is disabled for this session") })?; - match queue.push_latest(packet) { - QueuePush::Added => Ok(PushOutcome::Queued), - QueuePush::DroppedOldest => { - emit(&self.events, BackendEvent::QueueOverflow { media: "audio" }); - Ok(PushOutcome::DroppedOldest) - } - QueuePush::Full => unreachable!(), - QueuePush::Closed => Err(Error::QueueClosed), - } + submit_audio_packet(queue, &self.audio_unavailable, &self.events, packet) } pub fn reconfigure(&self, format: StreamFormat) -> Result<()> { @@ -733,11 +749,96 @@ fn run_video_worker( finish_decoder_drain(flush_decoder(&mut decoder), &decoded, &events, "shutdown"); } +const AUDIO_DECODE_STRIKES_BEFORE_RESET: u32 = 8; +const AUDIO_DECODER_RESET_LIMIT: u32 = 2; + +struct AudioRecovery { + consecutive_rejections: u32, + rejections: u64, + resets: u32, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum AudioRecoveryAction { + Skip, + ResetDecoder, + GiveUp, +} + +impl AudioRecovery { + fn record_decoded(&mut self) { + self.consecutive_rejections = 0; + } + + fn record_rejection(&mut self) -> AudioRecoveryAction { + self.consecutive_rejections += 1; + self.rejections += 1; + if self.consecutive_rejections < AUDIO_DECODE_STRIKES_BEFORE_RESET { + return AudioRecoveryAction::Skip; + } + if self.resets >= AUDIO_DECODER_RESET_LIMIT { + return AudioRecoveryAction::GiveUp; + } + AudioRecoveryAction::ResetDecoder + } + + fn record_reset(&mut self) { + self.consecutive_rejections = 0; + self.resets += 1; + } +} + +fn audio_unavailable( + packets: &BoundedQueue, + events: &EventQueue, + unavailable: &AtomicBool, + backend: AudioBackend, + recovery: &AudioRecovery, + reason: String, +) { + unavailable.store(true, Ordering::Release); + packets.close(); + emit( + events, + BackendEvent::AudioUnavailable { + backend, + reason, + rejected: recovery.rejections, + }, + ); +} + +fn submit_audio_packet( + queue: &BoundedQueue, + unavailable: &AtomicBool, + events: &EventQueue, + packet: AudioPacket, +) -> Result { + if unavailable.load(Ordering::Acquire) { + return Ok(PushOutcome::AudioDisabled); + } + match queue.push_latest(packet) { + QueuePush::Added => Ok(PushOutcome::Queued), + QueuePush::DroppedOldest => { + emit(events, BackendEvent::QueueOverflow { media: "audio" }); + Ok(PushOutcome::DroppedOldest) + } + QueuePush::Full => unreachable!(), + QueuePush::Closed => { + if unavailable.load(Ordering::Acquire) { + Ok(PushOutcome::AudioDisabled) + } else { + Err(Error::QueueClosed) + } + } + } +} + fn run_audio_worker( config: AudioConfig, packets: Arc>, - state: Arc>, events: EventQueue, + unavailable: Arc, startup: mpsc::SyncSender>, ) { let mut opus = match OpusDecoder::open(&config) { @@ -757,7 +858,33 @@ fn run_audio_worker( let mut backend = sink.backend(); let _ = startup.send(Ok(backend)); emit(&events, BackendEvent::AudioSelected(backend)); + let mut recovery = AudioRecovery { + consecutive_rejections: 0, + rejections: 0, + resets: 0, + }; + let mut rebuild_decoder = false; loop { + if rebuild_decoder { + rebuild_decoder = false; + match OpusDecoder::open(&config) { + Ok(decoder) => { + opus = decoder; + recovery.record_reset(); + } + Err(error) => { + audio_unavailable( + &packets, + &events, + &unavailable, + backend, + &recovery, + format!("the Opus decoder could not be rebuilt: {error}"), + ); + return; + } + } + } let packet = match packets.wait_pop(Duration::from_millis(100)) { QueuePop::Item(packet) => packet, QueuePop::TimedOut => continue, @@ -767,9 +894,23 @@ fn run_audio_worker( let concealed = match opus.conceal_before(&packet) { Ok(pcm) => pcm, Err(error) => { - packets.close(); - report_worker_error(&state, &events, error); - return; + opus.record_dropped_packet(&packet); + match record_audio_decode_failure(&events, &mut recovery, error.to_string()) { + AudioRecoveryAction::Skip => {} + AudioRecoveryAction::ResetDecoder => rebuild_decoder = true, + AudioRecoveryAction::GiveUp => { + audio_unavailable( + &packets, + &events, + &unavailable, + backend, + &recovery, + format!("the audio decoder rejected {} packets", recovery.rejections), + ); + return; + } + } + continue; } }; if !concealed.is_empty() { @@ -783,36 +924,84 @@ fn run_audio_worker( ) { Ok(()) => {} Err(AudioWriteError::Closed) => return, - Err(AudioWriteError::Failed(error)) => { - packets.close(); - report_worker_error(&state, &events, error); + Err(AudioWriteError::Failed { backend, reason }) => { + audio_unavailable( + &packets, + &events, + &unavailable, + backend, + &recovery, + format!("audio output was lost and could not be replaced: {reason}"), + ); return; } } } let pcm = match opus.decode(&packet) { - Ok(pcm) => pcm, + Ok(pcm) => { + recovery.record_decoded(); + pcm + } Err(error) => { - packets.close(); - report_worker_error(&state, &events, error); - return; + opus.record_dropped_packet(&packet); + match record_audio_decode_failure(&events, &mut recovery, error.to_string()) { + AudioRecoveryAction::Skip => {} + AudioRecoveryAction::ResetDecoder => rebuild_decoder = true, + AudioRecoveryAction::GiveUp => { + audio_unavailable( + &packets, + &events, + &unavailable, + backend, + &recovery, + format!("the audio decoder rejected {} packets", recovery.rejections), + ); + return; + } + } + continue; } }; match write_audio(&mut sink, &mut backend, &config, &events, pcm, &cancelled) { Ok(()) => {} Err(AudioWriteError::Closed) => return, - Err(AudioWriteError::Failed(error)) => { - packets.close(); - report_worker_error(&state, &events, error); + Err(AudioWriteError::Failed { backend, reason }) => { + audio_unavailable( + &packets, + &events, + &unavailable, + backend, + &recovery, + format!("audio output was lost and could not be replaced: {reason}"), + ); return; } } } } +fn record_audio_decode_failure( + events: &EventQueue, + recovery: &mut AudioRecovery, + message: String, +) -> AudioRecoveryAction { + let action = recovery.record_rejection(); + emit( + events, + BackendEvent::AudioDecodeError { + message, + consecutive: recovery.consecutive_rejections, + }, + ); + action +} + enum AudioWriteError { Closed, - Failed(Error), + Failed { + backend: AudioBackend, + reason: String, + }, } fn write_audio( @@ -827,27 +1016,60 @@ fn write_audio( if cancelled() { return Err(AudioWriteError::Closed); } + let message = error.to_string(); + emit( + events, + BackendEvent::AudioOutputError { + backend: *backend, + message: message.clone(), + }, + ); match open_audio_fallback(config, *backend) { Ok(mut fallback) => { + let fallback_backend = fallback.backend(); if let Err(fallback_error) = fallback.write(pcm, cancelled) { if cancelled() { return Err(AudioWriteError::Closed); } - return Err(AudioWriteError::Failed(fallback_error)); + emit( + events, + BackendEvent::AudioOutputError { + backend: fallback_backend, + message: fallback_error.to_string(), + }, + ); + return Err(AudioWriteError::Failed { + backend: fallback_backend, + reason: fallback_error.to_string(), + }); } - *backend = fallback.backend(); + let lost_backend = *backend; + *backend = fallback_backend; *sink = fallback; emit(events, BackendEvent::AudioSelected(*backend)); + emit( + events, + BackendEvent::AudioOutputRecovered { + from: lost_backend, + to: fallback_backend, + }, + ); return Ok(()); } Err(fallback_error) => { emit( events, - BackendEvent::Error(format!("audio fallback failed: {fallback_error}")), + BackendEvent::AudioOutputError { + backend: *backend, + message: fallback_error.to_string(), + }, ); } } - return Err(AudioWriteError::Failed(error)); + return Err(AudioWriteError::Failed { + backend: *backend, + reason: message, + }); } Ok(()) } @@ -1477,4 +1699,616 @@ mod tests { assert!(!reset); assert_eq!(inter_generation, keyframe_generation); } + + const VALID_OPUS_PACKET: [u8; 1] = [0x08]; + const MALFORMED_OPUS_PACKET: [u8; 40] = [0xff; 40]; + + fn audio_worker_config( + alsa_device: &str, + preference: crate::AudioBackendPreference, + ) -> super::AudioConfig { + super::AudioConfig { + alsa_device: alsa_device.to_owned(), + preference, + ..super::AudioConfig::default() + } + } + + fn use_audio_failure_fixture() { + unsafe { + std::env::set_var( + "ALSA_CONFIG_PATH", + concat!( + env!("CARGO_MANIFEST_DIR"), + "/tests/fixtures/audio-failure.conf" + ), + ) + }; + } + + struct RejectingAudioSink { + backend: crate::AudioBackend, + } + + impl super::AudioSink for RejectingAudioSink { + fn backend(&self) -> crate::AudioBackend { + self.backend + } + + fn write(&mut self, _: &[f32], _: &dyn Fn() -> bool) -> crate::Result<()> { + Err(crate::Error::backend( + crate::Subsystem::Alsa, + "test sink rejected PCM", + )) + } + } + + fn rejecting_output_sink(backend: crate::AudioBackend) -> Box { + Box::new(RejectingAudioSink { backend }) + } + + fn audio_packet(data: &[u8], rtp_timestamp: u32) -> super::AudioPacket { + super::AudioPacket::new(Arc::<[u8]>::from(data.to_vec()), rtp_timestamp, 48_000, 7) + .expect("valid audio packet input") + } + + fn spawn_audio_worker( + config: super::AudioConfig, + ) -> ( + Arc>, + Arc, + super::EventQueue, + thread::JoinHandle<()>, + ) { + use_audio_failure_fixture(); + let packets = Arc::new(super::BoundedQueue::new(64)); + let unavailable = Arc::new(AtomicBool::new(false)); + let events: super::EventQueue = Arc::new(super::BoundedQueue::new(512)); + let (startup_tx, startup_rx) = mpsc::sync_channel(1); + let worker = { + let packets = Arc::clone(&packets); + let unavailable = Arc::clone(&unavailable); + let events = Arc::clone(&events); + thread::Builder::new() + .name("opennow-test-audio-policy".to_owned()) + .spawn(move || { + super::run_audio_worker(config, packets, events, unavailable, startup_tx) + }) + .expect("audio worker") + }; + let startup = startup_rx.recv_timeout(Duration::from_secs(3)); + assert!( + matches!(startup, Ok(Ok(_))), + "audio worker startup: {startup:?}" + ); + (packets, unavailable, events, worker) + } + + fn collect_events(events: &super::EventQueue) -> Vec { + let mut collected = Vec::new(); + while let Some(event) = events.try_pop() { + collected.push(event); + } + collected + } + + fn wait_for_event( + events: &super::EventQueue, + timeout: Duration, + predicate: impl Fn(&super::BackendEvent) -> bool, + ) -> Vec { + let deadline = std::time::Instant::now() + timeout; + let mut collected = Vec::new(); + loop { + match events.wait_pop(Duration::from_millis(5)) { + super::QueuePop::Item(event) => { + let matched = predicate(&event); + collected.push(event); + if matched { + break; + } + } + super::QueuePop::TimedOut => { + if std::time::Instant::now() >= deadline { + break; + } + } + super::QueuePop::Closed => break, + } + } + collected.extend(collect_events(events)); + collected + } + + #[test] + fn terminal_audio_state_is_observed_across_queue_closure() { + let events: super::EventQueue = Arc::new(super::BoundedQueue::new(8)); + let queue = super::BoundedQueue::new(4); + let unavailable = AtomicBool::new(false); + queue.close(); + unavailable.store(true, Ordering::Release); + assert_eq!( + super::submit_audio_packet( + &queue, + &unavailable, + &events, + audio_packet(&VALID_OPUS_PACKET, 0), + ) + .expect("terminal audio state"), + super::PushOutcome::AudioDisabled + ); + assert!(collect_events(&events).is_empty()); + + let queue = super::BoundedQueue::new(4); + let unavailable = AtomicBool::new(false); + queue.close(); + assert!( + super::submit_audio_packet( + &queue, + &unavailable, + &events, + audio_packet(&VALID_OPUS_PACKET, 0), + ) + .is_err(), + "a closed queue without a terminal audio state stays an error" + ); + + let queue = super::BoundedQueue::new(4); + let unavailable = AtomicBool::new(true); + assert_eq!( + super::submit_audio_packet( + &queue, + &unavailable, + &events, + audio_packet(&VALID_OPUS_PACKET, 0), + ) + .expect("terminal audio state"), + super::PushOutcome::AudioDisabled + ); + } + + #[test] + fn concurrent_audio_submission_never_reports_a_queue_error() { + let queue = Arc::new(super::BoundedQueue::new(8)); + let unavailable = Arc::new(AtomicBool::new(false)); + let events: super::EventQueue = Arc::new(super::BoundedQueue::new(64)); + let publisher = { + let queue = Arc::clone(&queue); + let unavailable = Arc::clone(&unavailable); + thread::spawn(move || { + unavailable.store(true, Ordering::Release); + queue.close(); + }) + }; + let mut outcomes = Vec::new(); + let mut attempts = 0_u32; + while !queue.is_closed() && attempts < 100_000 { + outcomes.push(super::submit_audio_packet( + &queue, + &unavailable, + &events, + audio_packet(&VALID_OPUS_PACKET, 0), + )); + attempts += 1; + } + publisher.join().expect("publisher thread"); + for _ in 0..8 { + outcomes.push(super::submit_audio_packet( + &queue, + &unavailable, + &events, + audio_packet(&VALID_OPUS_PACKET, 0), + )); + } + assert!( + outcomes.iter().all(|outcome| outcome.is_ok()), + "a submitter racing the terminal audio state must not see a queue error: {outcomes:?}" + ); + assert!( + outcomes + .iter() + .rev() + .take(8) + .all(|outcome| *outcome.as_ref().expect("terminal outcome") + == super::PushOutcome::AudioDisabled), + "after the terminal audio state the submitter must report disabled: {outcomes:?}" + ); + } + + #[test] + fn audio_recovery_policy_bounds_decoder_rebuilds_and_gives_up() { + let mut recovery = super::AudioRecovery { + consecutive_rejections: 0, + rejections: 0, + resets: 0, + }; + for expected_consecutive in 1..super::AUDIO_DECODE_STRIKES_BEFORE_RESET { + assert_eq!( + recovery.record_rejection(), + super::AudioRecoveryAction::Skip + ); + assert_eq!(recovery.consecutive_rejections, expected_consecutive); + } + assert_eq!( + recovery.record_rejection(), + super::AudioRecoveryAction::ResetDecoder + ); + assert_eq!(recovery.rejections, 8); + recovery.record_reset(); + assert_eq!(recovery.consecutive_rejections, 0); + assert_eq!(recovery.resets, 1); + + recovery.record_decoded(); + assert_eq!(recovery.consecutive_rejections, 0); + for _ in 1..super::AUDIO_DECODE_STRIKES_BEFORE_RESET { + assert_eq!( + recovery.record_rejection(), + super::AudioRecoveryAction::Skip + ); + } + assert_eq!( + recovery.record_rejection(), + super::AudioRecoveryAction::ResetDecoder + ); + recovery.record_reset(); + for _ in 1..super::AUDIO_DECODE_STRIKES_BEFORE_RESET { + assert_eq!( + recovery.record_rejection(), + super::AudioRecoveryAction::Skip + ); + } + assert_eq!( + recovery.record_rejection(), + super::AudioRecoveryAction::GiveUp + ); + assert_eq!( + recovery.rejections, + 3 * u64::from(super::AUDIO_DECODE_STRIKES_BEFORE_RESET) + ); + assert_eq!(recovery.resets, super::AUDIO_DECODER_RESET_LIMIT); + } + + #[test] + fn malformed_opus_does_not_fail_the_shared_session() { + let (packets, unavailable, events, worker) = spawn_audio_worker(audio_worker_config( + "opennow_test_output", + crate::AudioBackendPreference::AlsaOnly, + )); + packets.push(audio_packet(&VALID_OPUS_PACKET, 0)); + packets.push(audio_packet(&[0x0c], 960)); + thread::sleep(Duration::from_millis(50)); + let startup = collect_events(&events); + assert!( + startup + .iter() + .all(|event| matches!(event, super::BackendEvent::AudioSelected(_))), + "valid Opus packets must not raise errors, saw {startup:?}" + ); + + packets.push(audio_packet(&MALFORMED_OPUS_PACKET, 1920)); + let observed = wait_for_event(&events, Duration::from_secs(3), |event| { + matches!(event, super::BackendEvent::AudioDecodeError { .. }) + }); + assert!( + observed.iter().any(|event| matches!( + event, + super::BackendEvent::AudioDecodeError { message, consecutive } + if message.contains("corrupted stream") && *consecutive == 1 + )), + "expected an Opus-scoped decode error, saw {observed:?}" + ); + assert!( + !observed + .iter() + .any(|event| matches!(event, super::BackendEvent::StateChanged(_))), + "audio decode errors must not change the shared session state, saw {observed:?}" + ); + assert!(!packets.is_closed(), "audio must keep running"); + assert!(!unavailable.load(Ordering::Acquire)); + + packets.push(audio_packet(&VALID_OPUS_PACKET, 2880)); + thread::sleep(Duration::from_millis(50)); + assert!( + collect_events(&events).is_empty(), + "audio must decode normally after a rejected packet" + ); + assert!(!packets.is_closed()); + + packets.close(); + worker.join().expect("audio worker thread"); + } + + #[test] + fn continuous_opus_rejections_disable_audio_only() { + let (packets, unavailable, events, worker) = spawn_audio_worker(audio_worker_config( + "opennow_test_output", + crate::AudioBackendPreference::AlsaOnly, + )); + let rejections = 3 * super::AUDIO_DECODE_STRIKES_BEFORE_RESET; + for index in 0..rejections { + packets.push(audio_packet(&MALFORMED_OPUS_PACKET, index * 960)); + } + let observed = wait_for_event(&events, Duration::from_secs(5), |event| { + matches!(event, super::BackendEvent::AudioUnavailable { .. }) + }); + let decode_errors = observed + .iter() + .filter(|event| matches!(event, super::BackendEvent::AudioDecodeError { .. })) + .count(); + assert_eq!( + decode_errors, rejections as usize, + "every rejected packet is reported, saw {observed:?}" + ); + assert!( + observed.iter().any(|event| matches!( + event, + super::BackendEvent::AudioUnavailable { backend, rejected, .. } + if *backend == crate::AudioBackend::Alsa + && *rejected == u64::from(rejections) + )), + "expected an audio-only escalation, saw {observed:?}" + ); + assert!( + !observed + .iter() + .any(|event| matches!(event, super::BackendEvent::StateChanged(_))), + "audio escalation must not fail the shared session, saw {observed:?}" + ); + assert!(packets.is_closed()); + assert!(unavailable.load(Ordering::Acquire)); + worker.join().expect("audio worker thread"); + } + + fn drive_until_audio_unavailable( + packets: &Arc>, + events: &super::EventQueue, + ) -> Vec { + let deadline = std::time::Instant::now() + Duration::from_secs(5); + let mut observed = Vec::new(); + let mut pushed = 0_u64; + while std::time::Instant::now() < deadline { + observed.extend(collect_events(events)); + if observed + .iter() + .any(|event| matches!(event, super::BackendEvent::AudioUnavailable { .. })) + { + break; + } + if pushed < 64 { + packets.push(audio_packet(&VALID_OPUS_PACKET, (pushed as u32) * 960)); + pushed += 1; + } + thread::sleep(Duration::from_millis(2)); + } + observed.extend(collect_events(events)); + observed + } + + #[test] + fn failing_audio_output_disables_audio_only() { + let (packets, unavailable, events, worker) = spawn_audio_worker(audio_worker_config( + "opennow_test_failing_output", + crate::AudioBackendPreference::AlsaOnly, + )); + let observed = drive_until_audio_unavailable(&packets, &events); + assert!( + observed.iter().any(|event| matches!( + event, + super::BackendEvent::AudioOutputError { backend, message } + if *backend == crate::AudioBackend::Alsa + && message.contains("Input/output error") + )), + "expected an audio output error, saw {observed:?}" + ); + assert!( + observed.iter().any(|event| matches!( + event, + super::BackendEvent::AudioUnavailable { backend, .. } + if *backend == crate::AudioBackend::Alsa + )), + "expected audio to be disabled, saw {observed:?}" + ); + assert!( + !observed + .iter() + .any(|event| matches!(event, super::BackendEvent::StateChanged(_))), + "an audio output failure must not fail the shared session, saw {observed:?}" + ); + assert!( + observed + .iter() + .all(|event| !matches!(event, super::BackendEvent::AudioOutputRecovered { .. })), + "a terminal output failure must not report recovery, saw {observed:?}" + ); + assert!(packets.is_closed()); + assert!(unavailable.load(Ordering::Acquire)); + worker.join().expect("audio worker thread"); + } + + #[test] + fn fixed_audio_output_never_falls_back_and_stays_isolated() { + let config = super::AudioConfig { + output_device: "alsa:opennow_test_failing_output".to_owned(), + ..audio_worker_config("unused", crate::AudioBackendPreference::AlsaOnly) + }; + let (packets, unavailable, events, worker) = spawn_audio_worker(config); + let observed = drive_until_audio_unavailable(&packets, &events); + assert!( + observed.iter().any(|event| matches!( + event, + super::BackendEvent::AudioOutputError { message, .. } + if message.contains("forbids fallback") + )), + "a fixed route must refuse fallback, saw {observed:?}" + ); + assert_eq!( + observed + .iter() + .filter(|event| matches!(event, super::BackendEvent::AudioSelected(_))) + .count(), + 1, + "a fixed route must never select another sink, saw {observed:?}" + ); + assert!( + !observed + .iter() + .any(|event| matches!(event, super::BackendEvent::StateChanged(_))), + "a fixed output failure must not fail the shared session, saw {observed:?}" + ); + assert!(packets.is_closed()); + assert!(unavailable.load(Ordering::Acquire)); + worker.join().expect("audio worker thread"); + } + + #[test] + fn audio_output_recovery_is_reported_once_the_replacement_sink_accepts_output() { + use_audio_failure_fixture(); + let events: super::EventQueue = Arc::new(super::BoundedQueue::new(32)); + let config = super::AudioConfig { + alsa_device: "opennow_test_output".to_owned(), + preference: crate::AudioBackendPreference::PipeWireThenAlsa, + ..super::AudioConfig::default() + }; + let pcm = [0.0f32; 960 * 2]; + let mut sink = rejecting_output_sink(crate::AudioBackend::PipeWire); + let mut backend = crate::AudioBackend::PipeWire; + assert!( + super::write_audio(&mut sink, &mut backend, &config, &events, &pcm, &|| false).is_ok(), + "the ALSA fallback must accept output" + ); + assert_eq!(backend, crate::AudioBackend::Alsa); + let observed = collect_events(&events); + assert_eq!( + observed + .iter() + .filter(|event| matches!(event, super::BackendEvent::AudioOutputRecovered { .. })) + .count(), + 1, + "one accepted replacement is one recovery observation, saw {observed:?}" + ); + assert!( + observed.iter().any(|event| matches!( + event, + super::BackendEvent::AudioOutputError { backend, message } + if *backend == crate::AudioBackend::PipeWire + && message.contains("test sink rejected PCM") + )), + "the lost sink is reported before recovery, saw {observed:?}" + ); + assert!( + observed.iter().any(|event| matches!( + event, + super::BackendEvent::AudioOutputRecovered { from, to } + if *from == crate::AudioBackend::PipeWire && *to == crate::AudioBackend::Alsa + )), + "the recovery pairs the lost sink with the accepting one, saw {observed:?}" + ); + assert!( + observed.iter().any(|event| matches!( + event, + super::BackendEvent::AudioSelected(backend) + if *backend == crate::AudioBackend::Alsa + )), + "the accepting sink becomes the selected sink, saw {observed:?}" + ); + } + + #[test] + fn repeated_audio_output_recoveries_report_each_accepted_replacement_once() { + use_audio_failure_fixture(); + let events: super::EventQueue = Arc::new(super::BoundedQueue::new(64)); + let config = super::AudioConfig { + alsa_device: "opennow_test_output".to_owned(), + preference: crate::AudioBackendPreference::PipeWireThenAlsa, + ..super::AudioConfig::default() + }; + let pcm = [0.0f32; 960 * 2]; + let mut sink = rejecting_output_sink(crate::AudioBackend::PipeWire); + let mut backend = crate::AudioBackend::PipeWire; + for _ in 0..2 { + assert!( + super::write_audio(&mut sink, &mut backend, &config, &events, &pcm, &|| false) + .is_ok(), + "the ALSA fallback must accept output" + ); + assert_eq!(backend, crate::AudioBackend::Alsa); + assert!( + super::write_audio(&mut sink, &mut backend, &config, &events, &pcm, &|| false) + .is_ok(), + "the selected sink must keep accepting output" + ); + let accepted = collect_events(&events); + assert_eq!( + accepted + .iter() + .filter(|event| matches!( + event, + super::BackendEvent::AudioOutputRecovered { .. } + )) + .count(), + 1, + "an accepted replacement reports recovery once, saw {accepted:?}" + ); + sink = rejecting_output_sink(crate::AudioBackend::PipeWire); + backend = crate::AudioBackend::PipeWire; + } + let observed = collect_events(&events); + assert!( + observed + .iter() + .all(|event| !matches!(event, super::BackendEvent::StateChanged(_))), + "audio output recovery must not fail the shared session, saw {observed:?}" + ); + } + + #[test] + fn audio_output_loss_without_an_accepting_replacement_reports_no_recovery() { + use_audio_failure_fixture(); + let pcm = [0.0f32; 960 * 2]; + for (scenario, config) in [ + ( + "a fixed output device refuses fallback", + super::AudioConfig { + output_device: "alsa:opennow_test_output".to_owned(), + ..audio_worker_config( + "opennow_test_output", + crate::AudioBackendPreference::PipeWireThenAlsa, + ) + }, + ), + ( + "the replacement sink does not open", + super::AudioConfig { + alsa_device: "opennow_test_missing_output".to_owned(), + ..audio_worker_config("", crate::AudioBackendPreference::PipeWireThenAlsa) + }, + ), + ] { + let events: super::EventQueue = Arc::new(super::BoundedQueue::new(32)); + let mut sink = rejecting_output_sink(crate::AudioBackend::PipeWire); + let mut backend = crate::AudioBackend::PipeWire; + assert!( + super::write_audio(&mut sink, &mut backend, &config, &events, &pcm, &|| false) + .is_err(), + "{scenario} must not accept output" + ); + assert_eq!(backend, crate::AudioBackend::PipeWire, "{scenario}"); + let observed = collect_events(&events); + assert!( + observed.iter().all(|event| !matches!( + event, + super::BackendEvent::AudioOutputRecovered { .. } + )), + "{scenario} must not report recovery, saw {observed:?}" + ); + assert!( + observed.iter().any(|event| matches!( + event, + super::BackendEvent::AudioOutputError { backend, .. } + if *backend == crate::AudioBackend::PipeWire + )), + "{scenario} must still report the loss, saw {observed:?}" + ); + } + } } diff --git a/native/opennow-streamer/crates/opennow-streamer-platform-linux/tests/fixtures/audio-failure.conf b/native/opennow-streamer/crates/opennow-streamer-platform-linux/tests/fixtures/audio-failure.conf new file mode 100644 index 000000000..33dff13b8 --- /dev/null +++ b/native/opennow-streamer/crates/opennow-streamer-platform-linux/tests/fixtures/audio-failure.conf @@ -0,0 +1,24 @@ +pcm.opennow_test_output { + type null + hint { + show on + description "OpenNOW isolated test output" + ioid "Output" + } +} + +pcm.opennow_test_failing_output { + type file + slave.pcm "null" + file "/dev/full" + format "raw" + hint { + show on + description "OpenNOW test output that fails once its buffer fills" + ioid "Output" + } +} + +pcm.null { + type null +} diff --git a/native/opennow-streamer/crates/opennow-streamer-platform/src/media.rs b/native/opennow-streamer/crates/opennow-streamer-platform/src/media.rs index 7fbbe508a..d718235cd 100644 --- a/native/opennow-streamer/crates/opennow-streamer-platform/src/media.rs +++ b/native/opennow-streamer/crates/opennow-streamer-platform/src/media.rs @@ -697,6 +697,15 @@ pub enum MediaFeedback { codec: &'static str, message: String, }, + AudioDecoderError { + message: String, + consecutive: u32, + }, + AudioUnavailable { + backend: &'static str, + reason: String, + rejected: u64, + }, QueueDropped { media: &'static str, count: usize, @@ -2895,6 +2904,9 @@ fn run_linux_video(shared: Arc, host_commands: Sender {} + Ok(opennow_streamer_platform_linux::PushOutcome::AudioDisabled) => { + unreachable!("audio-disabled outcomes are not produced by video submission") + } Err(reason) => trigger_linux_fallback( &shared, &host_commands, @@ -2953,6 +2965,9 @@ fn run_embedded_linux_video(shared: Arc) { request_linux_keyframe(&shared, "embedded Linux decoder queue overflow"); } Ok(opennow_streamer_platform_linux::PushOutcome::Paused) => {} + Ok(opennow_streamer_platform_linux::PushOutcome::AudioDisabled) => { + unreachable!("audio-disabled outcomes are not produced by video submission") + } Err(message) => { let _ = shared.feedback.send(MediaFeedback::DecoderError { codec: shared.linux_codec.label(), @@ -3023,7 +3038,8 @@ fn run_embedded_linux_audio(shared: Arc) { }); } Ok(opennow_streamer_platform_linux::PushOutcome::Queued) - | Ok(opennow_streamer_platform_linux::PushOutcome::Paused) => {} + | Ok(opennow_streamer_platform_linux::PushOutcome::Paused) + | Ok(opennow_streamer_platform_linux::PushOutcome::AudioDisabled) => {} Err(message) => { let _ = shared.feedback.send(MediaFeedback::DecoderError { codec: "opus", @@ -3142,6 +3158,34 @@ fn run_embedded_linux_monitor( stop_linux_session(&shared); return; } + opennow_streamer_platform_linux::BackendEvent::AudioDecodeError { + message, + consecutive, + } => { + let _ = shared.feedback.send(MediaFeedback::AudioDecoderError { + message, + consecutive, + }); + } + opennow_streamer_platform_linux::BackendEvent::AudioOutputError { + backend, + message, + } => forward_linux_audio_output_loss(&shared, backend, message), + opennow_streamer_platform_linux::BackendEvent::AudioOutputRecovered { + from, + to, + } => forward_linux_audio_output_recovery(&shared, from, to), + opennow_streamer_platform_linux::BackendEvent::AudioUnavailable { + backend, + reason, + rejected, + } => { + let _ = shared.feedback.send(MediaFeedback::AudioUnavailable { + backend: linux_audio_backend_name(backend), + reason, + rejected, + }); + } opennow_streamer_platform_linux::BackendEvent::FormatChanged(format) => { report_linux_color_format_change(&shared, &mut reported_color, format); } @@ -3237,6 +3281,34 @@ fn run_linux_monitor(shared: Arc, host_commands: Sender { + let _ = shared.feedback.send(MediaFeedback::AudioDecoderError { + message, + consecutive, + }); + } + opennow_streamer_platform_linux::BackendEvent::AudioOutputError { + backend, + message, + } => forward_linux_audio_output_loss(&shared, backend, message), + opennow_streamer_platform_linux::BackendEvent::AudioOutputRecovered { + from, + to, + } => forward_linux_audio_output_recovery(&shared, from, to), + opennow_streamer_platform_linux::BackendEvent::AudioUnavailable { + backend, + reason, + rejected, + } => { + let _ = shared.feedback.send(MediaFeedback::AudioUnavailable { + backend: linux_audio_backend_name(backend), + reason, + rejected, + }); + } opennow_streamer_platform_linux::BackendEvent::FormatChanged(format) => { report_linux_color_format_change(&shared, &mut reported_color, format); } @@ -3335,6 +3407,46 @@ const fn linux_decoder_name( } } +#[cfg(target_os = "linux")] +const fn linux_audio_backend_name( + backend: opennow_streamer_platform_linux::AudioBackend, +) -> &'static str { + match backend { + opennow_streamer_platform_linux::AudioBackend::PipeWire => "PipeWire", + opennow_streamer_platform_linux::AudioBackend::Alsa => "ALSA", + } +} + +#[cfg(target_os = "linux")] +fn forward_linux_audio_output_loss( + shared: &SharedPipeline, + backend: opennow_streamer_platform_linux::AudioBackend, + message: String, +) { + let _ = shared.feedback.send(MediaFeedback::DeviceLost { + subsystem: linux_audio_backend_name(backend), + recovered: false, + message: Some(message), + }); +} + +#[cfg(target_os = "linux")] +fn forward_linux_audio_output_recovery( + shared: &SharedPipeline, + from: opennow_streamer_platform_linux::AudioBackend, + to: opennow_streamer_platform_linux::AudioBackend, +) { + let _ = shared.feedback.send(MediaFeedback::DeviceLost { + subsystem: linux_audio_backend_name(from), + recovered: true, + message: Some(format!( + "{} accepted audio output after the {} sink failed", + linux_audio_backend_name(to), + linux_audio_backend_name(from) + )), + }); +} + #[cfg(target_os = "linux")] const fn linux_subsystem_name( subsystem: opennow_streamer_platform_linux::Subsystem, @@ -4235,6 +4347,83 @@ mod tests { } } + #[cfg(target_os = "linux")] + fn device_lost_feedback( + receiver: &Receiver, + ) -> (&'static str, bool, Option) { + match receiver + .try_recv() + .expect("a device state feedback for every audio output event") + { + MediaFeedback::DeviceLost { + subsystem, + recovered, + message, + } => (subsystem, recovered, message), + other => panic!("expected a device state feedback, saw {other:?}"), + } + } + + #[cfg(target_os = "linux")] + #[test] + fn linux_audio_output_feedback_pairs_loss_and_recovery_without_touching_video() { + use opennow_streamer_platform_linux::AudioBackend; + + let (shared, receiver) = software_test_pipeline(); + forward_linux_audio_output_loss(&shared, AudioBackend::PipeWire, "Broken pipe".to_owned()); + forward_linux_audio_output_recovery(&shared, AudioBackend::PipeWire, AudioBackend::Alsa); + + let loss = device_lost_feedback(&receiver); + assert_eq!(loss.0, "PipeWire"); + assert!(!loss.1); + assert_eq!(loss.2.as_deref(), Some("Broken pipe")); + + let recovery = device_lost_feedback(&receiver); + assert_eq!( + recovery.0, "PipeWire", + "recovery must clear the audio subsystem that lost output" + ); + assert!(recovery.1); + assert!( + recovery + .2 + .as_deref() + .is_some_and(|message| message.contains("ALSA") && message.contains("PipeWire")), + "recovery must name the sink that accepted output, saw {:?}", + recovery.2 + ); + + assert!( + receiver.try_recv().is_err(), + "one loss and one recovery only" + ); + assert!(!shared.stopped.load(Ordering::Acquire)); + assert!(!shared.keyframe_requested.load(Ordering::Acquire)); + assert!(shared.video_desynced.load(Ordering::Acquire)); + } + + #[cfg(target_os = "linux")] + #[test] + fn linux_audio_output_recovery_pairs_a_backend_that_recovered_itself() { + use opennow_streamer_platform_linux::AudioBackend; + + let (shared, receiver) = software_test_pipeline(); + forward_linux_audio_output_loss(&shared, AudioBackend::Alsa, "Device lost".to_owned()); + forward_linux_audio_output_recovery(&shared, AudioBackend::Alsa, AudioBackend::Alsa); + + let loss = device_lost_feedback(&receiver); + let recovery = device_lost_feedback(&receiver); + assert_eq!(loss, ("ALSA", false, Some("Device lost".to_owned()))); + assert_eq!( + recovery.0, loss.0, + "a backend that accepted output again clears its own loss" + ); + assert!(recovery.1); + assert!(receiver.try_recv().is_err()); + assert!(!shared.stopped.load(Ordering::Acquire)); + assert!(!shared.keyframe_requested.load(Ordering::Acquire)); + } + #[cfg(target_os = "windows")] #[test] fn embedded_playback_starts_only_once_after_a_successful_record() {