Skip to content
Merged
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
24 changes: 21 additions & 3 deletions electron/native/pipewire-capture/build.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,26 @@ use std::path::{Path, PathBuf};
fn main() {
let root = PathBuf::from(std::env::var("CARGO_MANIFEST_DIR").expect("CARGO_MANIFEST_DIR"));
build_pipewire_shim(&root);
link_ffmpeg(&root);
let ffmpeg = ffmpeg_dir(&root);
build_ffmpeg_accessors(&root, &ffmpeg);
link_ffmpeg(&ffmpeg);
}

fn build_ffmpeg_accessors(root: &Path, ffmpeg: &Path) {
let include = ffmpeg.join("include");
let source = root.join("csrc/ffmpeg_accessors.c");
assert!(
include.join("libavformat/avformat.h").is_file(),
"FFmpeg headers are missing at {}",
include.display()
);

cc::Build::new()
.file(&source)
.include(include)
.warnings(true)
.compile("openscreen_ffmpeg_accessors");
println!("cargo:rerun-if-changed={}", source.display());
}

fn build_pipewire_shim(root: &Path) {
Expand Down Expand Up @@ -115,8 +134,7 @@ fn ffmpeg_dir(root: &Path) -> PathBuf {
root.join("../../../crates/thirdparty/ffmpeg-linux64-lgpl-shared")
}

fn link_ffmpeg(root: &Path) {
let dir = ffmpeg_dir(root);
fn link_ffmpeg(dir: &Path) {
let include = dir.join("include");
let lib = dir.join("lib");

Expand Down
21 changes: 21 additions & 0 deletions electron/native/pipewire-capture/csrc/ffmpeg_accessors.c
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
#include <libavformat/avformat.h>

/* FFmpeg 8 exposes AVFormatContext as opaque to bindgen on some clang builds.
* Keep these field reads in C, compiled against the exact headers being linked. */
AVIOContext **osc_avformat_pb(AVFormatContext *context)
{
return context != NULL ? &context->pb : NULL;
}

AVStream *osc_avformat_stream(AVFormatContext *context, unsigned int index)
{
if (context == NULL || index >= context->nb_streams) {
return NULL;
}
return context->streams[index];
}

unsigned int osc_avformat_nb_streams(AVFormatContext *context)
{
return context != NULL ? context->nb_streams : 0;
}
21 changes: 15 additions & 6 deletions electron/native/pipewire-capture/src/encoder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1159,7 +1159,11 @@ impl Muxer {

let muxer = Self { fmt, tracks: Vec::new(), header_written: false };

let opened = ff::avio_open(&mut (*fmt).pb, path_c.as_ptr(), ff::AVIO_FLAG_WRITE as i32);
let pb = ff::osc_avformat_pb(fmt);
if pb.is_null() {
return Err("FFmpeg did not expose the output IO context".to_owned());
}
let opened = ff::avio_open(pb, path_c.as_ptr(), ff::AVIO_FLAG_WRITE as i32);
if opened < 0 {
return Err(format!(
"avio_open({}): {}",
Expand Down Expand Up @@ -1236,7 +1240,10 @@ impl Muxer {
// the value it CHOSE, not the value we asked for — using ours
// produces a file whose duration is wrong by the ratio between them.
for track in &mut self.tracks {
let stream = *(*self.fmt).streams.add(track.index as usize);
let stream = ff::osc_avformat_stream(self.fmt, track.index as u32);
if stream.is_null() {
return Err(format!("FFmpeg output stream {} disappeared", track.index));
}
track.stream_time_base = (*stream).time_base;
}
}
Expand Down Expand Up @@ -1292,8 +1299,9 @@ impl Drop for Muxer {
if self.header_written {
ff::av_write_trailer(self.fmt);
}
if !(*self.fmt).pb.is_null() {
ff::avio_closep(&mut (*self.fmt).pb);
let pb = ff::osc_avformat_pb(self.fmt);
if !pb.is_null() && !(*pb).is_null() {
ff::avio_closep(pb);
}
ff::avformat_free_context(self.fmt);
self.fmt = ptr::null_mut();
Expand Down Expand Up @@ -1567,7 +1575,8 @@ mod tests {
let mut video_frames = 0;
let mut second_frame_at = 0.0;
while ff::av_read_frame(input, packet) >= 0 {
let stream = *(*input).streams.add((*packet).stream_index as usize);
let stream = ff::osc_avformat_stream(input, (*packet).stream_index as u32);
assert!(!stream.is_null(), "packet stream index must be valid");
if (*(*stream).codecpar).codec_type == ff::AVMEDIA_TYPE_VIDEO {
if video_frames == 1 {
let base = (*stream).time_base;
Expand All @@ -1580,7 +1589,7 @@ mod tests {
}
let mut packet = packet;
ff::av_packet_free(&mut packet);
let streams = (*input).nb_streams;
let streams = ff::osc_avformat_nb_streams(input);
ff::avformat_close_input(&mut input);
(streams, video_frames, second_frame_at)
};
Expand Down
6 changes: 6 additions & 0 deletions electron/native/pipewire-capture/src/ffmpeg.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,12 @@

include!(concat!(env!("OUT_DIR"), "/ffmpeg_sys.rs"));

extern "C" {
pub fn osc_avformat_pb(context: *mut AVFormatContext) -> *mut *mut AVIOContext;
pub fn osc_avformat_stream(context: *mut AVFormatContext, index: u32) -> *mut AVStream;
pub fn osc_avformat_nb_streams(context: *mut AVFormatContext) -> u32;
}

/// `MKTAG(a,b,c,d)` — a FourCC packed little-endian, as libavutil defines it.
const fn mktag(a: u8, b: u8, c: u8, d: u8) -> i32 {
(a as i32) | ((b as i32) << 8) | ((c as i32) << 16) | ((d as i32) << 24)
Expand Down
180 changes: 177 additions & 3 deletions electron/native/pipewire-capture/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ mod shim;

use std::collections::HashSet;
use std::io::{BufRead, Write};
use std::os::fd::OwnedFd;
use std::path::PathBuf;
use std::sync::mpsc::{self, RecvTimeoutError, Sender};
use std::sync::Arc;
Expand Down Expand Up @@ -544,10 +545,20 @@ fn begin_stream<W: Write>(
granted_kind: &mut Option<portal::SourceKind>,
stream: portal::PortalStream,
prefer_dmabuf: bool,
shared_memory_fallback: &mut Option<SharedMemoryFallback>,
) -> Result<(), ()> {
// The fd is consumed by libpipewire; the rest is kept for the
// `stream-started` event, emitted once the format is negotiated.
let portal::PortalStream { fd, node_id, position, source_kind, .. } = stream;
let portal::PortalStream { fd, fallback_fd, node_id, position, size, source_kind } = stream;
*shared_memory_fallback = if prefer_dmabuf {
fallback_fd.map(|fd| SharedMemoryFallback {
fd,
node_id,
position,
size,
source_kind,
})
} else {
None
};
let forward = sender.clone();
match shim::Session::start(
fd,
Expand Down Expand Up @@ -582,6 +593,40 @@ struct StreamInfo {
source_kind: Option<portal::SourceKind>,
}

/// Separate portal remote kept only until a DMA-BUF-first stream starts.
struct SharedMemoryFallback {
fd: OwnedFd,
node_id: u32,
position: Option<(i32, i32)>,
size: Option<(i32, i32)>,
source_kind: Option<portal::SourceKind>,
}

fn should_retry_shared_memory(
state: &str,
error: Option<&str>,
capture_started: bool,
fallback_available: bool,
) -> bool {
state == "error"
&& !capture_started
&& fallback_available
&& error.is_some_and(|message| message.contains("alloc buffers"))
}

/// A stream error ends the helper only before capture has started, or the app
/// waits on a recording that never starts. "Started" is the first staged frame
/// in video mode, and `streaming` in cursor-only mode, which has no `Capture`.
fn stream_error_is_fatal(
state: &str,
video: bool,
recording_started: bool,
stream_was_live: bool,
) -> bool {
let started = if video { recording_started } else { stream_was_live };
state == "error" && !started
}

fn run<W: Write>(
emitter: &mut Emitter<W>,
receiver: mpsc::Receiver<Message>,
Expand All @@ -602,6 +647,9 @@ fn run<W: Write>(
let mut crop_change_reported = false;
// A `deferStart` grant held while the caller runs its countdown.
let mut pending_portal: Option<portal::PortalStream> = None;
// Retain the portal grant while the preferred DMA-BUF format is starting,
// so a buffer-allocation failure can retry shared-memory-first once.
let mut shared_memory_fallback: Option<SharedMemoryFallback> = None;
// `record` has been received. Latched, because it can arrive before the
// picker has been answered.
let mut armed = false;
Expand Down Expand Up @@ -969,6 +1017,7 @@ fn run<W: Write>(
&mut granted_kind,
stream,
prefer_dmabuf,
&mut shared_memory_fallback,
) {
exit_code = 1;
break;
Expand Down Expand Up @@ -1029,6 +1078,7 @@ fn run<W: Write>(
&mut granted_kind,
stream,
prefer_dmabuf,
&mut shared_memory_fallback,
) {
exit_code = 1;
break;
Expand Down Expand Up @@ -1163,10 +1213,72 @@ fn run<W: Write>(
if streaming_since.is_none() {
streaming_since = Some(timestamp_ms());
}
// No retry is appropriate after the compositor accepted
// the stream; release the extra portal fd immediately.
shared_memory_fallback = None;
} else if state == "unconnected" {
streaming_since = None;
}
let recording_started = capture.as_ref().is_some_and(Capture::started);
let stream_was_live = streaming_since.is_some();
if should_retry_shared_memory(
&state,
error.as_deref(),
recording_started || stream_was_live,
shared_memory_fallback.is_some(),
) {
let fallback = shared_memory_fallback.take().expect("checked above");
// Drop and join the failed stream before reconnecting with
// the portal's independent remote. The second attempt
// prefers shared memory; it will not arm another retry.
drop(session.take());
portal_stream = None;
let _ = emitter.emit(&Event::Warning {
code: "capture-retry-shared-memory".to_owned(),
message: "PipeWire buffer allocation failed with DMA-BUF; retrying screen capture with shared-memory buffers."
.to_owned(),
});
let retry = portal::PortalStream {
fd: fallback.fd,
fallback_fd: None,
node_id: fallback.node_id,
position: fallback.position,
size: fallback.size,
source_kind: fallback.source_kind,
};
if let Err(()) = begin_stream(
emitter,
&sender,
&frames,
&mut session,
&mut portal_stream,
&mut granted_kind,
retry,
false,
&mut shared_memory_fallback,
) {
exit_code = 1;
break;
}
// `begin_stream` already replaced the session and info;
// any queued messages from the old stream precede events
// from this new PipeWire thread.
continue;
}
if let Some(error) = error {
if stream_error_is_fatal(
&state,
frames.is_some(),
recording_started,
stream_was_live,
) {
let _ = emitter.emit(&Event::Error {
code: "pipewire-capture-failed".to_owned(),
message: format!("OpenScreen could not start screen capture: {error}"),
});
exit_code = 1;
break;
}
let _ = emitter.emit(&Event::Warning {
code: "stream-error".to_owned(),
message: format!("PipeWire stream reported an error in state {state}: {error}"),
Expand Down Expand Up @@ -1720,3 +1832,65 @@ mod microphone_resolution_tests {
);
}
}

#[cfg(test)]
mod capture_fallback_tests {
use super::{should_retry_shared_memory, stream_error_is_fatal};

#[test]
fn a_stream_error_is_fatal_only_before_capture_starts() {
// Video: started means a staged frame.
assert!(stream_error_is_fatal("error", true, false, true));
assert!(!stream_error_is_fatal("error", true, true, true));
// Cursor-only: started means the stream reached `streaming`.
assert!(stream_error_is_fatal("error", false, false, false));
assert!(!stream_error_is_fatal("error", false, false, true));
assert!(!stream_error_is_fatal("paused", true, false, false));
}

#[test]
fn retries_only_for_pre_recording_pipewire_buffer_allocation_errors() {
assert!(should_retry_shared_memory(
"error",
Some("error alloc buffers: Invalid argument"),
false,
true,
));
}

#[test]
fn does_not_retry_after_capture_has_started() {
assert!(!should_retry_shared_memory(
"error",
Some("error alloc buffers: Invalid argument"),
true,
true,
));
}

#[test]
fn does_not_retry_without_a_retained_portal_grant() {
assert!(!should_retry_shared_memory(
"error",
Some("error alloc buffers: Invalid argument"),
false,
false,
));
}

#[test]
fn does_not_retry_unrelated_pipewire_errors_or_non_error_states() {
assert!(!should_retry_shared_memory(
"error",
Some("target not found"),
false,
true,
));
assert!(!should_retry_shared_memory(
"paused",
Some("error alloc buffers: Invalid argument"),
false,
true,
));
}
}
8 changes: 8 additions & 0 deletions electron/native/pipewire-capture/src/portal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,8 @@ impl SourceKind {
/// Everything the PipeWire half needs, plus what the helper reports upward.
pub struct PortalStream {
pub fd: OwnedFd,
/// An independently opened remote for a one-time shared-memory retry.
pub fallback_fd: Option<OwnedFd>,
pub node_id: u32,
pub position: Option<(i32, i32)>,
pub size: Option<(i32, i32)>,
Expand Down Expand Up @@ -260,8 +262,14 @@ pub async fn negotiate(cursor_mode: CursorMode) -> Result<PortalStream, PortalEr
.await
.map_err(|error| failed("OpenPipeWireRemote", error))?;

// A duplicated Unix fd would share one PipeWire protocol connection and
// cannot be used by a second context. Ask the portal for a fresh remote
// while this screen-cast session is still alive, for the allocation retry.
let fallback_fd = proxy.open_pipe_wire_remote(&session).await.ok();

Ok(PortalStream {
fd,
fallback_fd,
node_id: stream.pipe_wire_node_id(),
position: stream.position(),
size: stream.size(),
Expand Down
Loading