From c339fac004a689e8e7f9f3dd8b4e6bb0a470ff42 Mon Sep 17 00:00:00 2001 From: zortos293 <65777760+zortos293@users.noreply.github.com> Date: Tue, 15 Sep 2026 18:33:42 +0000 Subject: [PATCH] Isolate embedded Linux audio failures from video --- .../crates/opennow-streamer-core/src/lib.rs | 119 +++ .../opennow-streamer-platform-linux/README.md | 4 +- .../src/audio.rs | 38 + .../src/session.rs | 699 +++++++++++++++++- .../tests/fixtures/audio-failure.conf | 24 + .../opennow-streamer-platform/src/media.rs | 88 ++- 6 files changed, 939 insertions(+), 33 deletions(-) create mode 100644 native/opennow-streamer/crates/opennow-streamer-platform-linux/tests/fixtures/audio-failure.conf 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 92f08f9c0..1a45ff443 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,19 @@ pub enum BackendEvent { reason: String, }, AudioSelected(AudioBackend), + AudioDecodeError { + message: String, + consecutive: u32, + }, + AudioOutputError { + backend: AudioBackend, + message: String, + }, + AudioUnavailable { + backend: AudioBackend, + reason: String, + rejected: u64, + }, FormatChanged(StreamFormat), NeedKeyframe, QueueOverflow { @@ -173,6 +187,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 +240,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 +291,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 +335,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<()> { @@ -724,11 +736,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) { @@ -748,7 +845,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, @@ -758,9 +881,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() { @@ -774,36 +911,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( @@ -818,15 +1003,34 @@ 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(); + *backend = fallback_backend; *sink = fallback; emit(events, BackendEvent::AudioSelected(*backend)); return Ok(()); @@ -834,11 +1038,17 @@ fn write_audio( 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(()) } @@ -1415,4 +1625,433 @@ 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 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<()>, + ) { + unsafe { + std::env::set_var( + "ALSA_CONFIG_PATH", + concat!( + env!("CARGO_MANIFEST_DIR"), + "/tests/fixtures/audio-failure.conf" + ), + ) + }; + 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!(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"); + } } 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 f8ce94186..ac00e5553 100644 --- a/native/opennow-streamer/crates/opennow-streamer-platform/src/media.rs +++ b/native/opennow-streamer/crates/opennow-streamer-platform/src/media.rs @@ -683,6 +683,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, @@ -2881,6 +2890,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, @@ -2939,6 +2951,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(), @@ -3009,7 +3024,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", @@ -3127,6 +3143,36 @@ 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, + } => { + let _ = shared.feedback.send(MediaFeedback::DeviceLost { + subsystem: linux_audio_backend_name(backend), + recovered: false, + message: Some(message), + }); + } + 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::StateChanged(_) | opennow_streamer_platform_linux::BackendEvent::DecoderSelected(_) | opennow_streamer_platform_linux::BackendEvent::AudioSelected(_) @@ -3219,6 +3265,36 @@ 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, + } => { + let _ = shared.feedback.send(MediaFeedback::DeviceLost { + subsystem: linux_audio_backend_name(backend), + recovered: false, + message: Some(message), + }); + } + 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::StateChanged(_) | opennow_streamer_platform_linux::BackendEvent::DecoderSelected(_) | opennow_streamer_platform_linux::BackendEvent::AudioSelected(_) @@ -3296,6 +3372,16 @@ 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")] const fn linux_subsystem_name( subsystem: opennow_streamer_platform_linux::Subsystem,