Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
119 changes: 119 additions & 0 deletions native/opennow-streamer/crates/opennow-streamer-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2246,6 +2246,39 @@ fn forward_nvst_media_feedback<R: NvstSessionResources>(
}),
));
}
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",
Expand Down Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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`.

Expand Down Expand Up @@ -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:

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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()?;
Expand Down Expand Up @@ -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();
Expand Down
Loading
Loading