use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
use libspa::buffer::Data;
use pipewire as pw;
use pipewire::spa::{
param::ParamType,
utils::{Direction, SpaTypes},
};
use crate::{error::Error, pod, Config, Format, Negotiated};
pub struct Plane<'f> {
pub stride: u32,
pub height: u32,
pub data: &'f mut [u8],
}
pub struct Frame<'f> {
pub width: u32,
pub height: u32,
pub format: Format,
pub fps_num: u32,
pub fps_denom: u32,
pub planes: Vec<Plane<'f>>,
}
impl<'f> Frame<'f> {
pub fn fill_black(&mut self) {
let planes = std::mem::take(&mut self.planes);
for (i, p) in planes.into_iter().enumerate() {
let Plane {
stride,
height,
data,
} = p;
let mut plane = Plane {
stride,
height,
data,
};
match (self.format, i) {
(Format::I420, 0) | (Format::Nv12, 0) | (Format::Nv21, 0) => plane.fill(0),
(Format::I420, 1) | (Format::I420, 2) => plane.fill(128),
(Format::Nv12, 1) | (Format::Nv21, 1) => plane.fill_interleaved(128, 128),
_ => plane.fill(0),
}
}
}
}
impl Plane<'_> {
fn fill(&mut self, value: u8) {
let n = (self.stride * self.height) as usize;
self.data[..n].fill(value);
}
fn fill_interleaved(&mut self, a: u8, b: u8) {
let n = (self.stride * self.height) as usize;
let mut i = 0;
while i + 1 < n {
self.data[i] = a;
self.data[i + 1] = b;
i += 2;
}
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum State {
Disconnected { error: Option<String> },
Paused { node_id: u32 },
Streaming { node_id: u32 },
}
struct Pacing {
streaming: AtomicBool,
period_ns: AtomicU64,
next_due: AtomicU64,
}
static START: std::sync::OnceLock<Instant> = std::sync::OnceLock::new();
fn now_ns() -> u64 {
let start = START.get_or_init(Instant::now);
start.elapsed().as_nanos() as u64
}
type NegotiateAcceptCb = Box<dyn FnMut(&Negotiated) -> Result<(), String>>;
type NegotiatedCb = Box<dyn FnMut(&Negotiated)>;
type FillCb = Box<dyn FnMut(&mut Frame, &Negotiated)>;
struct Inner {
pacing: Arc<Pacing>,
negotiated: Option<Negotiated>,
max_buffers: u32,
state_cb: Box<dyn FnMut(State)>,
neg_cb_accept: NegotiateAcceptCb,
neg_cb: NegotiatedCb,
fill: FillCb,
}
pub struct Camera {
mainloop: pw::main_loop::MainLoopRc,
stream: pw::stream::StreamRc,
pacing: Arc<Pacing>,
max_buffers: u32,
state_cb: Box<dyn FnMut(State)>,
neg_cb_accept: NegotiateAcceptCb,
neg_cb: NegotiatedCb,
}
#[derive(Clone)]
pub struct QuitHandle {
mainloop: pw::main_loop::MainLoopRc,
}
impl QuitHandle {
pub fn quit(&self) {
self.mainloop.quit();
}
}
fn validate(config: &Config) -> Result<(), Error> {
if config.name.is_empty() {
return Err(Error::InvalidConfig("name must not be empty".into()));
}
if config.modes.is_empty() {
return Err(Error::InvalidConfig("at least one mode is required".into()));
}
for mode in &config.modes {
if mode.width == 0 || mode.height == 0 {
return Err(Error::InvalidConfig("mode size must be nonzero".into()));
}
if mode.fps.is_empty() {
return Err(Error::InvalidConfig("mode needs at least one fps".into()));
}
if mode.fps.contains(&0) {
return Err(Error::InvalidConfig(
"mode fps values must be nonzero".into(),
));
}
if mode.formats.is_empty() {
return Err(Error::InvalidConfig(
"each mode needs at least one format".into(),
));
}
}
Ok(())
}
fn pw_stream_update_params_ptr(stream: *mut pw::sys::pw_stream, blobs: &[Vec<u8>]) -> i32 {
let mut ptrs: Vec<*const pw::spa::sys::spa_pod> =
blobs.iter().map(|b| b.as_ptr() as *const _).collect();
unsafe { pw::sys::pw_stream_update_params(stream, ptrs.as_mut_ptr().cast(), ptrs.len() as u32) }
}
impl Camera {
pub fn new(config: Config) -> Result<Self, Error> {
validate(&config)?;
let mainloop =
pw::main_loop::MainLoopRc::new(None).map_err(|e| Error::Connect(e.to_string()))?;
let context = pw::context::ContextRc::new(&mainloop, None)
.map_err(|e| Error::Connect(format!("{e:?}")))?;
let core = context
.connect_rc(None)
.map_err(|e| Error::Connect(format!("{e:?}")))?;
let properties = pw::properties::properties! {
*pw::keys::NODE_NAME => config.name.as_str(),
*pw::keys::NODE_DESCRIPTION => config.media_name.as_str(),
*pw::keys::NODE_NICK => config.media_name.as_str(),
*pw::keys::MEDIA_NAME => config.media_name.as_str(),
*pw::keys::MEDIA_CLASS => "Video/Source",
*pw::keys::MEDIA_TYPE => "Video",
*pw::keys::MEDIA_CATEGORY => "Capture",
*pw::keys::MEDIA_ROLE => "Camera",
};
let stream = pw::stream::StreamRc::new(core, &config.name, properties)
.map_err(|e| Error::Connect(e.to_string()))?;
let param_blobs = pod::advertised_param_blobs(&config);
let mut pod_ptrs: Vec<*const pw::spa::sys::spa_pod> =
param_blobs.iter().map(|b| b.as_ptr() as *const _).collect();
let r = unsafe {
pw::sys::pw_stream_connect(
stream.as_raw_ptr(),
Direction::Output.as_raw(),
pw::constants::ID_ANY,
pw::stream::StreamFlags::DRIVER.bits(),
pod_ptrs.as_mut_ptr(),
pod_ptrs.len() as u32,
)
};
if r < 0 {
return Err(Error::Connect(format!("pw_stream_connect: {r}")));
}
Ok(Camera {
max_buffers: config.max_buffers,
mainloop,
stream,
pacing: Arc::new(Pacing {
streaming: AtomicBool::new(false),
period_ns: AtomicU64::new(0),
next_due: AtomicU64::new(0),
}),
state_cb: Box::new(|_| {}),
neg_cb_accept: Box::new(|_| Ok(())),
neg_cb: Box::new(|_| {}),
})
}
pub fn on_state(mut self, cb: impl FnMut(State) + 'static) -> Self {
self.state_cb = Box::new(cb);
self
}
pub fn quit(&self) {
self.mainloop.quit();
}
pub fn quit_handle(&self) -> QuitHandle {
QuitHandle {
mainloop: self.mainloop.clone(),
}
}
pub fn on_negotiated(mut self, cb: impl FnMut(&Negotiated) + 'static) -> Self {
self.neg_cb = Box::new(cb);
self
}
pub fn on_negotiate_accept(
mut self,
cb: impl FnMut(&Negotiated) -> Result<(), String> + 'static,
) -> Self {
self.neg_cb_accept = Box::new(cb);
self
}
pub fn run(self, fill: impl FnMut(&mut Frame, &Negotiated) + 'static) -> Result<(), Error> {
let inner = Inner {
pacing: self.pacing.clone(),
negotiated: None,
max_buffers: self.max_buffers,
state_cb: self.state_cb,
neg_cb_accept: self.neg_cb_accept,
neg_cb: self.neg_cb,
fill: Box::new(fill),
};
let _timer = driver_timer(&self.mainloop, &self.stream, self.pacing.clone());
let _quit_signals = quit_signals(&self.mainloop);
let _mainloop = self.mainloop.clone();
let _listener = self
.stream
.add_local_listener_with_user_data(inner)
.state_changed(on_state_changed)
.param_changed(on_param_changed)
.process(on_process)
.register();
self.mainloop.run();
Ok(())
}
}
fn driver_timer<'l>(
mainloop: &'l pw::main_loop::MainLoopRc,
stream: &pw::stream::StreamRc,
pacing: Arc<Pacing>,
) -> pw::loop_::TimerSource<'l> {
let stream = stream.clone();
let timer = mainloop.loop_().add_timer(move |_exps| {
if !pacing.streaming.load(Ordering::SeqCst) {
return;
}
let period = pacing.period_ns.load(Ordering::SeqCst);
if period == 0 {
return;
}
let now = now_ns();
let next = pacing.next_due.load(Ordering::SeqCst);
if now < next {
return;
}
let missed = (now - next) / period;
pacing
.next_due
.store(next + (missed + 1) * period, Ordering::SeqCst);
let _ = stream.trigger_process();
});
let _ = timer.update_timer(
Some(Duration::from_millis(1)),
Some(Duration::from_millis(1)),
);
timer
}
fn quit_signals(
mainloop: &pw::main_loop::MainLoopRc,
) -> (pw::loop_::SignalSource<'_>, pw::loop_::SignalSource<'_>) {
let loop_ = mainloop.loop_();
let ml = mainloop.clone();
let sig_int = loop_.add_signal_local(pw::loop_::Signal::INT, move || ml.quit());
let ml = mainloop.clone();
let sig_term = loop_.add_signal_local(pw::loop_::Signal::TERM, move || ml.quit());
(sig_int, sig_term)
}
fn on_state_changed(
stream: &pw::stream::Stream,
inner: &mut Inner,
_old: pw::stream::StreamState,
new: pw::stream::StreamState,
) {
let node_id = stream.node_id();
match new {
pw::stream::StreamState::Unconnected | pw::stream::StreamState::Connecting => {
inner.pacing.streaming.store(false, Ordering::SeqCst);
(inner.state_cb)(State::Disconnected { error: None });
}
pw::stream::StreamState::Paused => {
inner.pacing.streaming.store(false, Ordering::SeqCst);
(inner.state_cb)(State::Paused { node_id });
}
pw::stream::StreamState::Streaming => {
let driving = stream.is_driving();
let lazy = unsafe { pw::sys::pw_stream_is_lazy(stream.as_raw_ptr()) };
inner.pacing.next_due.store(now_ns(), Ordering::SeqCst);
inner
.pacing
.streaming
.store(driving && !lazy, Ordering::SeqCst);
(inner.state_cb)(State::Streaming { node_id });
}
pw::stream::StreamState::Error(msg) => {
inner.pacing.streaming.store(false, Ordering::SeqCst);
println!("stream state: \"error\" {msg}");
(inner.state_cb)(State::Disconnected { error: Some(msg) });
}
}
}
fn on_param_changed(
stream: &pw::stream::Stream,
inner: &mut Inner,
id: u32,
param: Option<&pw::spa::pod::Pod>,
) {
if id != ParamType::Format.as_raw() {
return;
}
let Some(param) = param else { return };
if param.type_().as_raw() != SpaTypes::Object.as_raw() {
return;
}
let mut info = libspa::param::video::VideoInfoRaw::default();
if info.parse(param).is_err() {
eprintln!("pipewire-vircam: failed to parse negotiated Format");
return;
}
let Some(format) = Format::from_spa_id(info.format().0) else {
eprintln!(
"pipewire-vircam: negotiated format {} is not supported",
info.format().0
);
return;
};
let (stride, _h) = format.planes(info.size().width, info.size().height)[0];
let neg = Negotiated {
format,
width: info.size().width,
height: info.size().height,
fps_num: info.framerate().num,
fps_denom: info.framerate().denom,
stride,
node_id: stream.node_id(),
};
inner.pacing.period_ns.store(
neg.fps_denom as u64 * 1_000_000_000 / neg.fps_num.max(1) as u64,
Ordering::SeqCst,
);
inner.negotiated = Some(neg);
match (inner.neg_cb_accept)(&neg) {
Ok(()) => {}
Err(msg) => {
eprintln!("pipewire-vircam: rejecting negotiated format: {msg}");
return;
}
}
(inner.neg_cb)(&neg);
reply_buffers(stream, neg, inner.max_buffers);
}
fn reply_buffers(stream: &pw::stream::Stream, neg: Negotiated, max_buffers: u32) {
let num_planes = neg.format.planes(neg.width, neg.height).len() as u32;
let blobs: Vec<Vec<u8>> = vec![
pod::meta_pod(),
pod::buffers_pod(neg.stride, neg.height, num_planes, max_buffers),
];
let _ = pw_stream_update_params_ptr(stream.as_raw_ptr(), &blobs);
}
fn on_process(stream: &pw::stream::Stream, inner: &mut Inner) {
let Some(mut buf) = stream.dequeue_buffer() else {
return;
};
let Some(neg) = inner.negotiated else {
return; };
let layout = neg.format.planes(neg.width, neg.height);
let Some(planes) = build_planes(buf.datas_mut(), &layout) else {
return;
};
let mut frame = Frame {
width: neg.width,
height: neg.height,
format: neg.format,
fps_num: neg.fps_num,
fps_denom: neg.fps_denom,
planes,
};
(inner.fill)(&mut frame, &neg);
}
fn build_planes<'b>(datas: &'b mut [Data], layout: &[(u32, u32)]) -> Option<Vec<Plane<'b>>> {
let n = datas.len().min(layout.len());
let datas_ptr: *mut Data = datas.as_mut_ptr();
let mut planes: Vec<Plane> = Vec::with_capacity(n);
for (i, &(stride, height)) in layout.iter().take(n).enumerate() {
let d: &mut Data = unsafe { &mut *datas_ptr.add(i) };
{
let c = d.chunk_mut();
*c.offset_mut() = 0;
*c.size_mut() = stride * height;
*c.stride_mut() = stride as i32;
}
let data = d.data()?;
planes.push(Plane {
stride,
height,
data,
});
}
Some(planes)
}