use std::{
cell::{Cell, RefCell},
rc::Rc,
sync::{
Arc, Mutex,
atomic::{AtomicBool, AtomicU64, Ordering},
mpsc::{self, RecvTimeoutError, SyncSender},
},
thread::JoinHandle,
time::{Duration, Instant},
};
use crossbeam_channel::{Receiver, Sender};
use ffmpeg_next as ffmpeg;
use pipewire as pw;
use pw::{properties::properties, spa};
use spa::pod::Pod;
use thiserror::Error as ThisError;
use crate::pp_log::{PpLog, pp_info, pp_trace, pp_warn};
use crate::{
buffer::MediaBuffer,
control::ControlMsg,
element::{Element, ElementType, Sink, element_pp_log},
elements::AudioFormat,
error::Result,
platform::linux::pipewire::{
PipeWireAudioDevice, PipeWireAudioDeviceKind, PipeWireDeviceError,
},
playback_clock::{AudioMasterRegistration, PlaybackClock, PlaybackClockError},
};
const NEGOTIATION_TIMEOUT: Duration = Duration::from_secs(5);
const SEND_GRANULARITY: Duration = Duration::from_millis(100);
const DRAIN_SLACK: Duration = Duration::from_secs(1);
const DEVICE_DRAIN_TIMEOUT: Duration = Duration::from_secs(5);
const QUEUE_CAPACITY: usize = 8;
const PRIME_FRAMES: usize = QUEUE_CAPACITY / 2;
#[derive(Debug, ThisError)]
pub enum PipeWireAudioRendererError {
#[error("pipewire error: {0}")]
PipeWire(String),
#[error(transparent)]
Device(#[from] PipeWireDeviceError),
#[error("PipeWireAudioRenderer needs a Sink device, but {0:?} is a Source")]
NotAPlaybackDevice(String),
#[error("timed out waiting for the PipeWire audio stream to negotiate a format")]
NegotiationTimeout,
#[error("unsupported PipeWire audio sample format {0:?}")]
UnsupportedFormat(u32),
#[error("PipeWire negotiated an empty audio format ({rate}Hz, {channels} channel(s))")]
EmptyFormat { rate: u32, channels: u32 },
#[error(
"expected {expected:?} audio at {expected_rate}Hz x{expected_channels}, \
got {actual:?} at {actual_rate}Hz x{actual_channels}"
)]
FormatMismatch {
expected: ffmpeg::format::Sample,
expected_rate: u32,
expected_channels: u16,
actual: ffmpeg::format::Sample,
actual_rate: u32,
actual_channels: u16,
},
#[error("PipeWireAudioRenderer only accepts MediaBuffer::Audio, got {0}")]
UnexpectedBuffer(&'static str),
#[error("audio frames need a PTS when PipeWireAudioRenderer is the playback-clock master")]
MissingPts,
#[error("the PipeWire playback stream ended")]
StreamEnded,
#[error("audio frame data is too short: need {expected} byte(s), got {actual}")]
FrameDataTooShort { expected: usize, actual: usize },
#[error(transparent)]
PlaybackClock(#[from] PlaybackClockError),
#[error("this renderer is already bound to a playback clock")]
PlaybackClockAlreadyBound,
#[error("a playback clock must be bound before the first frame is rendered")]
PlaybackClockBoundAfterStart,
}
#[derive(Debug, Clone)]
pub struct PipeWireAudioRendererOptions {
pub device: PipeWireAudioDevice,
}
enum PlaybackClockBinding {
Unbound,
Deferred(Arc<PlaybackClock>),
Registered(AudioMasterRegistration),
}
impl PlaybackClockBinding {
fn is_bound(&self) -> bool {
!matches!(self, Self::Unbound)
}
fn registration(&self) -> Option<&AudioMasterRegistration> {
match self {
Self::Registered(master) => Some(master),
Self::Unbound | Self::Deferred(_) => None,
}
}
fn ensure_registered(&mut self) -> std::result::Result<(), PlaybackClockError> {
if let Self::Deferred(playback_clock) = self {
let registration = playback_clock.register_audio_master()?;
*self = Self::Registered(registration);
}
Ok(())
}
}
struct Playback {
played_frames: AtomicU64,
latency_frames: AtomicU64,
queued_frames: AtomicU64,
ended: AtomicBool,
}
enum Startup {
Ready(AudioFormat),
Failed(PipeWireAudioRendererError),
}
#[derive(Default)]
struct Pending {
bytes: Vec<u8>,
offset: usize,
}
impl Pending {
fn new(bytes: Vec<u8>) -> Self {
Self { bytes, offset: 0 }
}
fn is_empty(&self) -> bool {
self.offset >= self.bytes.len()
}
fn copy_into(&mut self, out: &mut [u8]) -> usize {
let take = out.len().min(self.bytes.len() - self.offset);
out[..take].copy_from_slice(&self.bytes[self.offset..self.offset + take]);
self.offset += take;
take
}
}
fn discard_queued(frames: &Receiver<Vec<u8>>, leftover: &Mutex<Pending>, playback: &Playback) {
while frames.try_recv().is_ok() {}
if let Ok(mut leftover) = leftover.lock() {
*leftover = Pending::default();
}
playback.queued_frames.store(0, Ordering::Release);
}
fn wait_for_queue(playback: &Playback, deadline: Instant, mut tick: impl FnMut()) -> bool {
loop {
if playback.queued_frames.load(Ordering::Acquire) == 0 {
return true;
}
if playback.ended.load(Ordering::Acquire) {
return false;
}
let now = Instant::now();
if now >= deadline {
return false;
}
tick();
std::thread::sleep(SEND_GRANULARITY.min(deadline.saturating_duration_since(now)));
}
}
type CommandResult = std::result::Result<(), String>;
type CommandReply = SyncSender<CommandResult>;
enum Command {
SetActive {
active: bool,
reply: CommandReply,
},
Flush(CommandReply),
Drain(CommandReply),
Terminate,
}
fn queue_command(
commands: &pw::channel::Sender<Command>,
build: impl FnOnce(CommandReply) -> Command,
) -> std::result::Result<mpsc::Receiver<CommandResult>, PipeWireAudioRendererError> {
let (reply_tx, reply_rx) = mpsc::sync_channel(1);
commands
.send(build(reply_tx))
.map_err(|_| PipeWireAudioRendererError::StreamEnded)?;
Ok(reply_rx)
}
fn wait_command(
reply: mpsc::Receiver<CommandResult>,
operation: &'static str,
) -> std::result::Result<(), PipeWireAudioRendererError> {
match reply.recv_timeout(NEGOTIATION_TIMEOUT) {
Ok(Ok(())) => Ok(()),
Ok(Err(error)) => Err(PipeWireAudioRendererError::PipeWire(error)),
Err(RecvTimeoutError::Timeout) => Err(PipeWireAudioRendererError::PipeWire(format!(
"timed out waiting for PipeWire to {operation}"
))),
Err(RecvTimeoutError::Disconnected) => Err(PipeWireAudioRendererError::StreamEnded),
}
}
fn complete_device_drain(
draining: &Cell<bool>,
playback: &Playback,
reply: &RefCell<Option<CommandReply>>,
) {
playback.latency_frames.store(0, Ordering::Release);
draining.set(false);
if let Some(reply) = reply.borrow_mut().take() {
let _ = reply.send(Ok(()));
}
}
pub struct PipeWireAudioRenderer {
name: Arc<str>,
pp_log: PpLog,
format: AudioFormat,
frames: Sender<Vec<u8>>,
playback: Arc<Playback>,
clock_binding: PlaybackClockBinding,
primed: bool,
timeline: Option<Timeline>,
commands: Option<pw::channel::Sender<Command>>,
worker: Option<JoinHandle<()>>,
}
struct Timeline {
media_origin_ns: i64,
submitted_until_ns: i64,
played_origin: u64,
}
impl PipeWireAudioRenderer {
pub fn list_devices()
-> std::result::Result<Vec<PipeWireAudioDevice>, PipeWireAudioRendererError> {
let mut devices = crate::platform::linux::pipewire::list_devices()?;
devices.retain(|device| device.kind == PipeWireAudioDeviceKind::Sink);
Ok(devices)
}
pub fn open(
name: impl Into<String>,
options: PipeWireAudioRendererOptions,
) -> std::result::Result<(Self, AudioFormat), PipeWireAudioRendererError> {
if options.device.kind != PipeWireAudioDeviceKind::Sink {
return Err(PipeWireAudioRendererError::NotAPlaybackDevice(
options.device.name.clone(),
));
}
let name = name.into();
let pp_log = element_pp_log(ElementType::PipeWireAudioRenderer, &name, None);
let (frame_tx, frame_rx) = crossbeam_channel::bounded(QUEUE_CAPACITY);
let (startup_tx, startup_rx) = mpsc::channel::<Startup>();
let (command_tx, command_rx) = pw::channel::channel::<Command>();
let playback = Arc::new(Playback {
played_frames: AtomicU64::new(0),
latency_frames: AtomicU64::new(0),
queued_frames: AtomicU64::new(0),
ended: AtomicBool::new(false),
});
let device = options.device.clone();
let worker = std::thread::Builder::new()
.name(format!("{name}-pipewire-render"))
.spawn({
let startup_tx = startup_tx.clone();
let playback = playback.clone();
move || {
if let Err(error) =
run_pipewire(device, frame_rx, playback.clone(), &startup_tx, command_rx)
{
let _ = startup_tx.send(Startup::Failed(error));
}
playback.ended.store(true, Ordering::Release);
}
})
.map_err(|e| PipeWireAudioRendererError::PipeWire(e.to_string()))?;
let startup = match startup_rx.recv_timeout(NEGOTIATION_TIMEOUT) {
Ok(startup) => startup,
Err(RecvTimeoutError::Timeout) => {
let _ = command_tx.send(Command::Terminate);
let _ = worker.join();
return Err(PipeWireAudioRendererError::NegotiationTimeout);
}
Err(RecvTimeoutError::Disconnected) => {
let _ = worker.join();
return Err(PipeWireAudioRendererError::PipeWire(
"the PipeWire thread exited before negotiating a format".into(),
));
}
};
let format = match startup {
Startup::Ready(format) => format,
Startup::Failed(error) => {
let _ = command_tx.send(Command::Terminate);
let _ = worker.join();
return Err(error);
}
};
if let Err(error) = queue_command(&command_tx, |reply| Command::SetActive {
active: false,
reply,
})
.and_then(|reply| wait_command(reply, "park the negotiated stream"))
{
let _ = command_tx.send(Command::Terminate);
let _ = worker.join();
return Err(error);
}
pp_info!(
pp_log: &pp_log,
"opened: device={:?} (node {}), {}Hz, {} channel(s), format={:?}",
options.device.description,
options.device.id,
format.sample_rate,
format.channels,
format.sample_format
);
Ok((
Self {
name: name.into(),
pp_log,
format,
frames: frame_tx,
playback,
clock_binding: PlaybackClockBinding::Unbound,
primed: false,
timeline: None,
commands: Some(command_tx),
worker: Some(worker),
},
format,
))
}
pub fn format(&self) -> AudioFormat {
self.format
}
pub fn bind_playback_clock(
&mut self,
playback_clock: Arc<PlaybackClock>,
) -> std::result::Result<(), PipeWireAudioRendererError> {
self.check_bindable()?;
let master = playback_clock.register_audio_master()?;
self.clock_binding = PlaybackClockBinding::Registered(master);
Ok(())
}
pub fn bind_playback_clock_deferred(
&mut self,
playback_clock: Arc<PlaybackClock>,
) -> std::result::Result<(), PipeWireAudioRendererError> {
self.check_bindable()?;
self.clock_binding = PlaybackClockBinding::Deferred(playback_clock);
Ok(())
}
fn check_bindable(&self) -> std::result::Result<(), PipeWireAudioRendererError> {
if self.clock_binding.is_bound() {
return Err(PipeWireAudioRendererError::PlaybackClockAlreadyBound);
}
if self.timeline.is_some() {
return Err(PipeWireAudioRendererError::PlaybackClockBoundAfterStart);
}
Ok(())
}
fn frames_duration(&self, frames: u64) -> Duration {
Duration::from_nanos(self.frames_ns(frames).max(0) as u64)
}
fn frames_ns(&self, frames: u64) -> i64 {
((u128::from(frames) * 1_000_000_000u128) / u128::from(self.format.sample_rate.max(1)))
.min(i64::MAX as u128) as i64
}
fn publish_position(&self, running: bool) -> Result<()> {
let (Some(master), Some(timeline)) = (self.clock_binding.registration(), &self.timeline)
else {
return Ok(());
};
let played = self
.playback
.played_frames
.load(Ordering::Acquire)
.saturating_sub(timeline.played_origin);
let latency = self.playback.latency_frames.load(Ordering::Acquire);
let audible = played.saturating_sub(latency);
let position_ns = timeline
.media_origin_ns
.saturating_add(self.frames_ns(audible));
master
.publish(position_ns, timeline.submitted_until_ns, running)
.map_err(PipeWireAudioRendererError::from)?;
Ok(())
}
fn audio_pts_ns(&self, frame: &ffmpeg::frame::Audio) -> Result<i64> {
let pts = frame.pts().ok_or(PipeWireAudioRendererError::MissingPts)?;
Ok(self.frames_ns(pts.max(0) as u64))
}
fn render(&mut self, frame: &ffmpeg::frame::Audio) -> Result<()> {
let actual_rate = frame.rate();
let actual_channels = frame.channel_layout().channels() as u16;
if frame.format() != self.format.sample_format
|| actual_rate != self.format.sample_rate
|| actual_channels != self.format.channels
{
return Err(PipeWireAudioRendererError::FormatMismatch {
expected: self.format.sample_format,
expected_rate: self.format.sample_rate,
expected_channels: self.format.channels,
actual: frame.format(),
actual_rate,
actual_channels,
}
.into());
}
if frame.samples() == 0 {
return Ok(());
}
self.clock_binding
.ensure_registered()
.map_err(PipeWireAudioRendererError::from)?;
let bytes_per_frame = self.format.channels as usize * self.format.sample_format.bytes();
let tight = frame.samples() * bytes_per_frame;
let plane = frame.data(0);
if plane.len() < tight {
return Err(PipeWireAudioRendererError::FrameDataTooShort {
expected: tight,
actual: plane.len(),
}
.into());
}
let payload = plane[..tight].to_vec();
let pts_ns = if self.clock_binding.registration().is_some() {
Some(self.audio_pts_ns(frame)?)
} else {
frame.pts().map(|pts| self.frames_ns(pts.max(0) as u64))
};
let played_origin = self.playback.played_frames.load(Ordering::Acquire);
self.playback
.queued_frames
.fetch_add(frame.samples() as u64, Ordering::AcqRel);
let mut pending = payload;
loop {
match self.frames.send_timeout(pending, SEND_GRANULARITY) {
Ok(()) => break,
Err(crossbeam_channel::SendTimeoutError::Timeout(returned)) => {
if self.playback.ended.load(Ordering::Acquire) {
self.rollback_queued(frame.samples() as u64);
return Err(PipeWireAudioRendererError::StreamEnded.into());
}
if let Err(error) = self.start_once_primed() {
self.rollback_queued(frame.samples() as u64);
return Err(error);
}
if let Err(error) = self.publish_position(true) {
self.rollback_queued(frame.samples() as u64);
return Err(error);
}
pending = returned;
}
Err(crossbeam_channel::SendTimeoutError::Disconnected(_)) => {
self.rollback_queued(frame.samples() as u64);
return Err(PipeWireAudioRendererError::StreamEnded.into());
}
}
}
if let Some(pts_ns) = pts_ns {
let end_ns = pts_ns.saturating_add(self.frames_ns(frame.samples() as u64));
match &mut self.timeline {
Some(timeline) => timeline.submitted_until_ns = end_ns,
None => {
self.timeline = Some(Timeline {
media_origin_ns: pts_ns,
submitted_until_ns: end_ns,
played_origin,
})
}
}
}
self.start_once_primed()?;
self.publish_position(true)?;
Ok(())
}
fn rollback_queued(&self, frames: u64) {
let _ = self.playback.queued_frames.fetch_update(
Ordering::AcqRel,
Ordering::Acquire,
|queued| Some(queued.saturating_sub(frames)),
);
}
fn start_once_primed(&mut self) -> Result<()> {
if !self.primed && self.frames.len() >= PRIME_FRAMES {
self.set_active(true)?;
self.primed = true;
}
Ok(())
}
fn drain(&mut self) -> Result<()> {
if !self.primed {
self.set_active(true)?;
self.primed = true;
}
let outstanding = self.playback.queued_frames.load(Ordering::Acquire);
let deadline = Instant::now() + self.frames_duration(outstanding) + DRAIN_SLACK;
let mut published = Ok(());
let drained = wait_for_queue(&self.playback, deadline, || {
if published.is_ok() {
published = self.publish_position(true);
}
});
published?;
if !drained && !self.playback.ended.load(Ordering::Acquire) {
pp_warn!(
self,
"the device stopped taking audio during drain: {} frame(s) never played",
self.playback.queued_frames.load(Ordering::Acquire)
);
self.flush()?;
}
if drained && !self.playback.ended.load(Ordering::Acquire) {
self.drain_device()?;
}
self.publish_position(false)?;
let final_position = self
.timeline
.as_ref()
.map(|timeline| timeline.submitted_until_ns);
if let (Some(master), Some(final_position)) =
(self.clock_binding.registration(), final_position)
{
master
.finish(final_position)
.map_err(PipeWireAudioRendererError::from)?;
}
Ok(())
}
fn request_command(
&self,
operation: &'static str,
build: impl FnOnce(CommandReply) -> Command,
) -> Result<()> {
let commands = self
.commands
.as_ref()
.ok_or(PipeWireAudioRendererError::StreamEnded)?;
let reply = queue_command(commands, build)?;
wait_command(reply, operation)?;
Ok(())
}
fn set_active(&self, active: bool) -> Result<()> {
self.request_command(
if active {
"activate the playback stream"
} else {
"deactivate the playback stream"
},
|reply| Command::SetActive { active, reply },
)
}
fn flush(&self) -> Result<()> {
self.request_command("flush the playback stream", Command::Flush)
}
fn drain_device(&self) -> Result<()> {
let commands = self
.commands
.as_ref()
.ok_or(PipeWireAudioRendererError::StreamEnded)?;
let reply = queue_command(commands, Command::Drain)?;
let deadline = Instant::now() + DEVICE_DRAIN_TIMEOUT;
loop {
let now = Instant::now();
if now >= deadline {
pp_warn!(self, "PipeWire did not report a completed device drain");
self.flush()?;
return Ok(());
}
match reply.recv_timeout(SEND_GRANULARITY.min(deadline.saturating_duration_since(now)))
{
Ok(Ok(())) => return Ok(()),
Ok(Err(error)) => {
return Err(PipeWireAudioRendererError::PipeWire(error).into());
}
Err(RecvTimeoutError::Timeout) => self.publish_position(true)?,
Err(RecvTimeoutError::Disconnected) => {
return Err(PipeWireAudioRendererError::StreamEnded.into());
}
}
}
}
}
impl Drop for PipeWireAudioRenderer {
fn drop(&mut self) {
if let Some(commands) = self.commands.take() {
let _ = commands.send(Command::Terminate);
}
if let Some(worker) = self.worker.take()
&& worker.join().is_err()
{
pp_warn!(self, "the PipeWire playback thread panicked");
}
}
}
impl Element for PipeWireAudioRenderer {
fn name(&self) -> Arc<str> {
self.name.clone()
}
fn element_type(&self) -> ElementType {
ElementType::PipeWireAudioRenderer
}
fn pp_log(&self) -> &PpLog {
&self.pp_log
}
fn pp_log_mut(&mut self) -> &mut PpLog {
&mut self.pp_log
}
}
impl Sink for PipeWireAudioRenderer {
fn consume(&mut self, buf: MediaBuffer) -> Result<()> {
match buf {
MediaBuffer::Audio(frame) => self.render(&frame),
MediaBuffer::Eos => {
pp_trace!(self, "event=eos phase=received");
let outcome = self.drain();
pp_trace!(
self,
"event=eos phase=drained outcome={}",
if outcome.is_ok() { "ok" } else { "error" }
);
outcome
}
other => Err(PipeWireAudioRendererError::UnexpectedBuffer(other.kind()).into()),
}
}
fn control(&mut self, msg: ControlMsg) -> Result<()> {
match msg {
ControlMsg::Pause => {
self.set_active(false)?;
self.publish_position(false)?;
}
ControlMsg::Resume => {
self.set_active(true)?;
self.publish_position(true)?;
}
ControlMsg::Stop => {
if let Err(error) = self.flush() {
pp_warn!(self, "failed to flush while stopping: {error}");
}
if let Err(error) = self.set_active(false) {
pp_warn!(self, "failed to deactivate while stopping: {error}");
}
let position_result = self.publish_position(false);
self.timeline = None;
self.primed = false;
position_result?;
}
ControlMsg::Seek(_) => {
self.flush()?;
self.timeline = None;
}
}
Ok(())
}
}
fn ffmpeg_sample_format(format: spa::param::audio::AudioFormat) -> Option<ffmpeg::format::Sample> {
use ffmpeg::format::sample::Type::Packed;
use spa::param::audio::AudioFormat as Spa;
Some(match format {
Spa::F32LE => ffmpeg::format::Sample::F32(Packed),
Spa::S16LE => ffmpeg::format::Sample::I16(Packed),
Spa::S32LE => ffmpeg::format::Sample::I32(Packed),
_ => return None,
})
}
fn stream_latency_frames(
queued_frames: u64,
buffered_frames: u64,
delay_ticks: i64,
rate_num: u32,
rate_denom: u32,
sample_rate: u32,
current_frames: u64,
) -> u64 {
let delay_frames = if delay_ticks <= 0 || rate_denom == 0 {
0
} else {
((delay_ticks as u128)
.saturating_mul(u128::from(rate_num))
.saturating_mul(u128::from(sample_rate))
/ u128::from(rate_denom))
.min(u128::from(u64::MAX)) as u64
};
queued_frames
.saturating_add(buffered_frames)
.saturating_add(delay_frames)
.saturating_add(current_frames)
}
fn run_pipewire(
device: PipeWireAudioDevice,
frames: Receiver<Vec<u8>>,
playback: Arc<Playback>,
startup: &mpsc::Sender<Startup>,
commands: pw::channel::Receiver<Command>,
) -> std::result::Result<(), PipeWireAudioRendererError> {
fn pw_err(error: impl std::fmt::Display) -> PipeWireAudioRendererError {
PipeWireAudioRendererError::PipeWire(error.to_string())
}
pw::init();
let mainloop = pw::main_loop::MainLoopRc::new(None).map_err(pw_err)?;
let context = pw::context::ContextRc::new(&mainloop, None).map_err(pw_err)?;
let core = context.connect_rc(None).map_err(pw_err)?;
let stream = pw::stream::StreamRc::new(
core.clone(),
"media-pp-audio-playback",
properties! {
*pw::keys::MEDIA_TYPE => "Audio",
*pw::keys::MEDIA_CATEGORY => "Playback",
*pw::keys::MEDIA_ROLE => "Production",
},
)
.map_err(pw_err)?;
let leftover = Arc::new(Mutex::new(Pending::default()));
let draining = Rc::new(Cell::new(false));
let drain_reply = Rc::new(RefCell::new(None::<CommandReply>));
{
let command_stream = stream.clone();
let command_leftover = leftover.clone();
let command_frames = frames.clone();
let command_playback = playback.clone();
let command_draining = draining.clone();
let command_drain_reply = drain_reply.clone();
let quit = mainloop.clone();
let _commands = commands.attach(mainloop.loop_(), move |command| match command {
Command::SetActive { active, reply } => {
let result = command_stream
.set_active(active)
.map_err(|error| error.to_string());
let _ = reply.send(result);
}
Command::Flush(reply) => {
command_draining.set(false);
if let Some(drain_reply) = command_drain_reply.borrow_mut().take() {
let _ = drain_reply.send(Err("the device drain was cancelled".into()));
}
discard_queued(&command_frames, &command_leftover, &command_playback);
let result = command_stream
.flush(false)
.map_err(|error| error.to_string());
if result.is_ok() {
command_playback.latency_frames.store(0, Ordering::Release);
}
let _ = reply.send(result);
}
Command::Drain(reply) => {
if let Some(previous) = command_drain_reply.borrow_mut().replace(reply) {
let _ = previous.send(Err("a newer device drain replaced this one".into()));
}
command_draining.set(true);
if let Err(error) = command_stream.flush(true) {
command_draining.set(false);
if let Some(reply) = command_drain_reply.borrow_mut().take() {
let _ = reply.send(Err(error.to_string()));
}
}
}
Command::Terminate => {
if let Some(reply) = command_drain_reply.borrow_mut().take() {
let _ = reply.send(Err("the playback stream terminated during drain".into()));
}
quit.quit();
}
});
let format = Arc::new(Mutex::new(None::<AudioFormat>));
let _listener = stream
.add_local_listener_with_user_data(())
.param_changed({
let format = format.clone();
let startup = startup.clone();
move |_, (), id, param| {
let Some(param) = param else { return };
if id != spa::param::ParamType::Format.as_raw() {
return;
}
let Ok((media_type, media_subtype)) =
spa::param::format_utils::parse_format(param)
else {
return;
};
if media_type != spa::param::format::MediaType::Audio
|| media_subtype != spa::param::format::MediaSubtype::Raw
{
return;
}
let mut info = spa::param::audio::AudioInfoRaw::new();
if info.parse(param).is_err() {
return;
}
let Some(sample_format) = ffmpeg_sample_format(info.format()) else {
let _ = startup.send(Startup::Failed(
PipeWireAudioRendererError::UnsupportedFormat(info.format().as_raw()),
));
return;
};
let (rate, channels) = (info.rate(), info.channels());
if rate == 0 || channels == 0 {
let _ = startup.send(Startup::Failed(
PipeWireAudioRendererError::EmptyFormat { rate, channels },
));
return;
}
let negotiated = AudioFormat::new(sample_format, rate, channels as u16);
if let Ok(mut format) = format.lock() {
*format = Some(negotiated);
}
let _ = startup.send(Startup::Ready(negotiated));
}
})
.drained({
let drain_reply = drain_reply.clone();
let draining = draining.clone();
let playback = playback.clone();
move |_, ()| {
complete_device_drain(&draining, &playback, &drain_reply);
}
})
.process({
let format = format.clone();
let leftover = leftover.clone();
let playback = playback.clone();
let frames = frames.clone();
let draining = draining.clone();
move |stream, ()| {
if draining.get() {
return;
}
let Some(mut buffer) = stream.dequeue_buffer() else {
return;
};
let Some(negotiated) = format.lock().ok().and_then(|f| *f) else {
return;
};
let bytes_per_frame =
negotiated.channels as usize * negotiated.sample_format.bytes();
if bytes_per_frame == 0 {
return;
}
let timing = stream.time().ok().map(|time| {
let rate = time.rate();
(
time.queued(),
time.buffered(),
time.delay(),
rate.num,
rate.denom,
)
});
let requested_frames = buffer.requested() as usize;
let written_bytes = {
let datas = buffer.datas_mut();
let Some(data) = datas.first_mut() else {
return;
};
let capacity_frames =
data.data().map(|d| d.len()).unwrap_or(0) / bytes_per_frame;
let want_frames = if requested_frames > 0 {
requested_frames.min(capacity_frames)
} else {
capacity_frames
};
let want = want_frames * bytes_per_frame;
if want == 0 {
return;
}
let mut written = 0usize;
{
let Ok(mut pending) = leftover.lock() else {
return;
};
let Some(out) = data.data() else { return };
while written < want {
if pending.is_empty() {
match frames.try_recv() {
Ok(next) => *pending = Pending::new(next),
Err(_) => break,
}
}
written += pending.copy_into(&mut out[written..want]);
}
if written < want {
out[written..want].fill(0);
}
}
let consumed = (written / bytes_per_frame) as u64;
playback.played_frames.fetch_add(consumed, Ordering::AcqRel);
if let Some((queued, buffered, delay, rate_num, rate_denom)) = timing {
playback.latency_frames.store(
stream_latency_frames(
queued,
buffered,
delay,
rate_num,
rate_denom,
negotiated.sample_rate,
consumed,
),
Ordering::Release,
);
}
let _ = playback.queued_frames.fetch_update(
Ordering::AcqRel,
Ordering::Acquire,
|queued| Some(queued.saturating_sub(consumed)),
);
want
};
let datas = buffer.datas_mut();
let data = &mut datas[0];
let chunk = data.chunk_mut();
*chunk.offset_mut() = 0;
*chunk.stride_mut() = bytes_per_frame as i32;
*chunk.size_mut() = written_bytes as u32;
}
})
.register()
.map_err(pw_err)?;
let mut audio_info = spa::param::audio::AudioInfoRaw::new();
audio_info.set_format(spa::param::audio::AudioFormat::F32LE);
let obj = spa::pod::Object {
type_: spa::utils::SpaTypes::ObjectParamFormat.as_raw(),
id: spa::param::ParamType::EnumFormat.as_raw(),
properties: audio_info.into(),
};
let values: Vec<u8> = spa::pod::serialize::PodSerializer::serialize(
std::io::Cursor::new(Vec::new()),
&spa::pod::Value::Object(obj),
)
.map_err(pw_err)?
.0
.into_inner();
let mut params = [Pod::from_bytes(&values).ok_or_else(|| {
PipeWireAudioRendererError::PipeWire("failed to build format pod".into())
})?];
stream
.connect(
spa::utils::Direction::Output,
Some(device.id),
pw::stream::StreamFlags::AUTOCONNECT | pw::stream::StreamFlags::MAP_BUFFERS,
&mut params,
)
.map_err(pw_err)?;
mainloop.run();
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
fn playback() -> Playback {
Playback {
played_frames: AtomicU64::new(0),
latency_frames: AtomicU64::new(0),
queued_frames: AtomicU64::new(0),
ended: AtomicBool::new(false),
}
}
#[test]
fn pipewire_timing_is_normalized_before_latency_is_combined() {
assert_eq!(
stream_latency_frames(
240, 32, 480, 1, 48_000, 48_000, 128, ),
240 + 32 + 480 + 128
);
assert_eq!(
stream_latency_frames(0, 0, 50, 1, 1_000, 48_000, 0),
2_400,
"graph-rate ticks are not necessarily audio frames"
);
assert_eq!(
stream_latency_frames(0, 0, -128, 1, 48_000, 48_000, 0),
0,
"negative graph delay is clamped as PipeWire recommends"
);
}
#[test]
fn a_completed_device_drain_leaves_the_stream_reusable() {
let playback = playback();
playback.latency_frames.store(512, Ordering::Release);
let draining = Cell::new(true);
let (reply_tx, reply_rx) = mpsc::sync_channel(1);
let reply = RefCell::new(Some(reply_tx));
complete_device_drain(&draining, &playback, &reply);
assert!(!draining.get(), "later process callbacks must be accepted");
assert_eq!(playback.latency_frames.load(Ordering::Acquire), 0);
assert_eq!(reply_rx.try_recv(), Ok(Ok(())));
}
#[test]
fn a_drain_waits_until_the_device_has_taken_everything() {
let playback = Arc::new(playback());
playback.queued_frames.store(1024, Ordering::Release);
let consumer = {
let playback = playback.clone();
std::thread::spawn(move || {
for _ in 0..4 {
std::thread::sleep(Duration::from_millis(20));
playback
.queued_frames
.fetch_sub(256, std::sync::atomic::Ordering::AcqRel);
}
})
};
let mut ticks = 0;
let drained = wait_for_queue(&playback, Instant::now() + Duration::from_secs(5), || {
ticks += 1
});
consumer.join().expect("the consumer thread finishes");
assert!(
drained,
"the drain must wait for the device to take the audio"
);
assert_eq!(playback.queued_frames.load(Ordering::Acquire), 0);
assert!(
ticks > 0,
"the caller keeps publishing its position while waiting"
);
}
#[test]
fn a_drain_gives_up_on_a_device_that_stopped_taking_audio() {
let playback = playback();
playback.queued_frames.store(1024, Ordering::Release);
let started = Instant::now();
let drained = wait_for_queue(&playback, started + Duration::from_millis(150), || {});
assert!(
!drained,
"a stalled device is reported, not waited on forever"
);
assert!(
started.elapsed() < Duration::from_secs(2),
"the wait ends at its deadline"
);
}
#[test]
fn a_flush_discards_the_queue_and_not_only_the_frame_in_flight() {
let (tx, rx) = crossbeam_channel::bounded::<Vec<u8>>(8);
for byte in 0..4u8 {
tx.send(vec![byte; 16]).expect("the queue has room");
}
let leftover = Mutex::new(Pending::new(vec![0xAB; 16]));
let playback = playback();
playback.queued_frames.store(320, Ordering::Release);
discard_queued(&rx, &leftover, &playback);
assert!(rx.is_empty(), "queued audio must not outlive the seek");
assert_eq!(
playback.queued_frames.load(Ordering::Acquire),
0,
"discarded audio must stop counting as outstanding, or the next \
drain waits for audio that no longer exists"
);
assert!(
leftover.lock().unwrap().is_empty(),
"the frame being copied must not outlive the seek either"
);
tx.send(vec![7; 16]).expect("the queue still accepts audio");
assert_eq!(rx.try_recv().expect("the new frame arrives"), vec![7; 16]);
}
#[test]
fn a_control_command_completes_only_after_the_pipewire_reply() {
let (reply_tx, reply_rx) = mpsc::sync_channel(1);
let worker = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(50));
reply_tx.send(Ok(())).expect("the caller is still waiting");
});
let started = Instant::now();
wait_command(reply_rx, "test the command").expect("the command succeeds");
worker.join().expect("the reply thread finishes");
assert!(
started.elapsed() >= Duration::from_millis(40),
"enqueueing alone must not acknowledge a synchronous control"
);
}
#[test]
fn a_pipewire_command_failure_reaches_its_caller() {
let (reply_tx, reply_rx) = mpsc::sync_channel(1);
reply_tx
.send(Err("set_active failed".into()))
.expect("the caller is waiting");
let error = wait_command(reply_rx, "activate the stream")
.expect_err("the PipeWire mutation failed");
assert!(matches!(
error,
PipeWireAudioRendererError::PipeWire(message)
if message == "set_active failed"
));
}
#[test]
fn only_interleaved_formats_are_accepted() {
use spa::param::audio::AudioFormat as Spa;
assert_eq!(
ffmpeg_sample_format(Spa::F32LE),
Some(ffmpeg::format::Sample::F32(
ffmpeg::format::sample::Type::Packed
))
);
assert!(ffmpeg_sample_format(Spa::F32P).is_none());
}
#[test]
fn a_capture_node_is_rejected_before_any_stream_is_opened() {
let device = PipeWireAudioDevice {
id: 7,
name: "some-mic".into(),
description: "Some Mic".into(),
kind: PipeWireAudioDeviceKind::Source,
is_default: true,
};
let Err(error) =
PipeWireAudioRenderer::open("out", PipeWireAudioRendererOptions { device })
else {
panic!("a Source device cannot be played to");
};
assert!(
matches!(error, PipeWireAudioRendererError::NotAPlaybackDevice(name) if name == "some-mic"),
"the mistake must surface as a typed construction error, \
not as a stream that silently never plays"
);
}
#[test]
fn seek_is_rejected_as_a_typed_error() {
let error: crate::Error = PipeWireAudioRendererError::StreamEnded.into();
assert!(matches!(
error,
crate::Error::PipeWireAudioRendererError(PipeWireAudioRendererError::StreamEnded)
));
}
}