use super::color::uniform;
use super::source::{self, Layout, Source};
use crate::{Color, Error, Frame, Size};
const ZERO_COPY_STRIKES: u32 = 3;
#[cfg(all(target_os = "linux", feature = "dmabuf"))]
fn dma_buf_import_timed_out(error: &Error) -> bool {
let Error::Render(error) = error else {
return false;
};
error.chain().any(|cause| {
cause
.downcast_ref::<std::io::Error>()
.is_some_and(|error| error.kind() == std::io::ErrorKind::TimedOut)
})
}
#[derive(Clone, Debug)]
#[non_exhaustive]
pub struct Config {
pub size: Option<Size>,
pub format: wgpu::TextureFormat,
pub usage: wgpu::TextureUsages,
pub color: Option<Color>,
pub zero_copy: bool,
}
impl Default for Config {
fn default() -> Self {
Self {
size: None,
format: wgpu::TextureFormat::Rgba8Unorm,
usage: wgpu::TextureUsages::empty(),
color: None,
zero_copy: true,
}
}
}
impl Config {
pub fn new() -> Self {
Self::default()
}
}
pub struct Renderer {
device: wgpu::Device,
queue: wgpu::Queue,
#[cfg(all(target_os = "linux", feature = "dmabuf"))]
completion: Completion,
config: Config,
shader: Pipelines,
uniform: wgpu::Buffer,
color: Option<Color>,
source: source::Cache,
output: Option<wgpu::Texture>,
strikes: u32,
retired: bool,
imported: bool,
}
struct Pipelines {
layout: wgpu::BindGroupLayout,
sampler: wgpu::Sampler,
#[cfg(any(target_os = "macos", all(target_os = "linux", feature = "dmabuf")))]
nv12: wgpu::RenderPipeline,
#[cfg(all(target_os = "linux", feature = "dmabuf"))]
rgba: wgpu::RenderPipeline,
i420: wgpu::RenderPipeline,
}
#[cfg(all(target_os = "linux", feature = "dmabuf"))]
struct Completion {
device: wgpu::Device,
tx: Option<std::sync::mpsc::Sender<(wgpu::SubmissionIndex, Box<dyn Send + Sync>)>>,
thread: Option<std::thread::JoinHandle<()>>,
}
#[cfg(all(target_os = "linux", feature = "dmabuf"))]
impl Completion {
fn new(device: &wgpu::Device) -> Result<Self, Error> {
let (tx, rx) = std::sync::mpsc::channel();
let worker_device = device.clone();
let thread = std::thread::Builder::new()
.name("moq-video-gpu-completion".into())
.spawn(move || {
while let Ok((submission, keepalive)) = rx.recv() {
if let Err(err) = worker_device.poll(wgpu::PollType::Wait {
submission_index: Some(submission),
timeout: None,
}) {
tracing::warn!(%err, "waiting for imported GPU surface failed");
}
drop(keepalive);
}
})
.map_err(|err| Error::Render(anyhow::anyhow!("start GPU completion worker: {err}")))?;
Ok(Self {
device: device.clone(),
tx: Some(tx),
thread: Some(thread),
})
}
fn submit(&self, submission: wgpu::SubmissionIndex, keepalive: Box<dyn Send + Sync>) {
let tx = self.tx.as_ref().expect("completion sender lives until drop");
if let Err(err) = tx.send((submission, keepalive)) {
let (submission, keepalive) = err.0;
let _ = self.device.poll(wgpu::PollType::Wait {
submission_index: Some(submission),
timeout: None,
});
drop(keepalive);
}
}
}
#[cfg(all(target_os = "linux", feature = "dmabuf"))]
impl Drop for Completion {
fn drop(&mut self) {
drop(self.tx.take());
if let Some(thread) = self.thread.take() {
let _ = thread.join();
}
}
}
impl Renderer {
pub fn new(device: &wgpu::Device, queue: &wgpu::Queue, config: Config) -> Result<Self, Error> {
if let Some(size) = config.size {
size.validate_nonzero("render output")?;
}
let shader = Pipelines::new(device, config.format)?;
let uniform = device.create_buffer(&wgpu::BufferDescriptor {
label: Some("moq-video color conversion"),
size: std::mem::size_of::<[f32; 16]>() as u64,
usage: wgpu::BufferUsages::UNIFORM | wgpu::BufferUsages::COPY_DST,
mapped_at_creation: false,
});
Ok(Self {
device: device.clone(),
queue: queue.clone(),
#[cfg(all(target_os = "linux", feature = "dmabuf"))]
completion: Completion::new(device)?,
config,
shader,
uniform,
color: None,
source: source::Cache::default(),
output: None,
strikes: 0,
retired: false,
imported: false,
})
}
pub fn render(&mut self, frame: &Frame) -> Result<wgpu::Texture, Error> {
let mut source = self.source(frame)?;
if let Some(source_color) = source.color {
let color = self.config.color.unwrap_or(source_color);
if self.color != Some(color) {
self.queue
.write_buffer(&self.uniform, 0, bytemuck::cast_slice(&uniform(color)));
self.color = Some(color);
}
}
let output = self.output(frame.size())?;
let view = output.create_view(&wgpu::TextureViewDescriptor::default());
let bind = self.device.create_bind_group(&wgpu::BindGroupDescriptor {
label: Some("moq-video planes"),
layout: &self.shader.layout,
entries: &[
wgpu::BindGroupEntry {
binding: 0,
resource: self.uniform.as_entire_binding(),
},
wgpu::BindGroupEntry {
binding: 1,
resource: wgpu::BindingResource::Sampler(&self.shader.sampler),
},
wgpu::BindGroupEntry {
binding: 2,
resource: wgpu::BindingResource::TextureView(&source.plane0),
},
wgpu::BindGroupEntry {
binding: 3,
resource: wgpu::BindingResource::TextureView(&source.plane1),
},
wgpu::BindGroupEntry {
binding: 4,
resource: wgpu::BindingResource::TextureView(&source.plane2),
},
],
});
let pipeline = match source.layout {
#[cfg(all(target_os = "linux", feature = "dmabuf"))]
Layout::Rgba => &self.shader.rgba,
#[cfg(any(target_os = "macos", all(target_os = "linux", feature = "dmabuf")))]
Layout::Nv12 => &self.shader.nv12,
Layout::I420 => &self.shader.i420,
};
let mut encoder = self.device.create_command_encoder(&wgpu::CommandEncoderDescriptor {
label: Some("moq-video render"),
});
{
let mut pass = encoder.begin_render_pass(&wgpu::RenderPassDescriptor {
label: Some("moq-video yuv to rgb"),
color_attachments: &[Some(wgpu::RenderPassColorAttachment {
view: &view,
resolve_target: None,
depth_slice: None,
ops: wgpu::Operations {
load: wgpu::LoadOp::Load,
store: wgpu::StoreOp::Store,
},
})],
depth_stencil_attachment: None,
timestamp_writes: None,
occlusion_query_set: None,
multiview_mask: None,
});
pass.set_pipeline(pipeline);
pass.set_bind_group(0, &bind, &[]);
pass.draw(0..3, 0..1);
}
let submission = self.queue.submit([encoder.finish()]);
if let Some(keepalive) = source.keepalive.take() {
#[cfg(all(target_os = "linux", feature = "dmabuf"))]
self.completion.submit(submission, keepalive);
#[cfg(not(all(target_os = "linux", feature = "dmabuf")))]
drop((submission, keepalive));
}
Ok(output)
}
fn source(&mut self, frame: &Frame) -> Result<Source, Error> {
if self.config.zero_copy && !self.retired {
match self.source.import(&self.device, &frame.surface) {
Ok(Some(source)) => {
self.strikes = 0;
if !self.imported {
self.imported = true;
tracing::debug!(
layout = ?source.layout,
"drawing frames zero-copy; the picture reaches the GPU without a download"
);
}
return Ok(source);
}
Ok(None) => {}
#[cfg(all(target_os = "linux", feature = "dmabuf"))]
Err(err) if dma_buf_import_timed_out(&err) => return Err(err),
Err(err) => {
self.strikes += 1;
self.retired = self.strikes >= ZERO_COPY_STRIKES;
match self.retired {
true => tracing::warn!(%err, "zero-copy import failed repeatedly; using the CPU path"),
false => tracing::debug!(%err, "zero-copy import failed; falling back to the CPU"),
}
}
}
}
self.source.upload(&self.device, &self.queue, frame)
}
fn output(&mut self, frame: Size) -> Result<wgpu::Texture, Error> {
let size = self.config.size.unwrap_or(frame);
if let Some(output) = &self.output
&& output.width() == size.width
&& output.height() == size.height
{
return Ok(output.clone());
}
size.validate_nonzero("render output")?;
let format = self.config.format;
let sibling = match format == format.add_srgb_suffix() {
true => format.remove_srgb_suffix(),
false => format.add_srgb_suffix(),
};
let view_formats: &[wgpu::TextureFormat] = match sibling == format {
true => &[],
false => &[sibling],
};
let texture = self.device.create_texture(&wgpu::TextureDescriptor {
label: Some("moq-video output"),
size: wgpu::Extent3d {
width: size.width,
height: size.height,
depth_or_array_layers: 1,
},
mip_level_count: 1,
sample_count: 1,
dimension: wgpu::TextureDimension::D2,
format,
usage: wgpu::TextureUsages::RENDER_ATTACHMENT | wgpu::TextureUsages::TEXTURE_BINDING | self.config.usage,
view_formats,
});
self.output = Some(texture.clone());
Ok(texture)
}
}
impl Pipelines {
fn new(device: &wgpu::Device, format: wgpu::TextureFormat) -> Result<Self, Error> {
let module = device.create_shader_module(wgpu::ShaderModuleDescriptor {
label: Some("moq-video yuv"),
source: wgpu::ShaderSource::Wgsl(include_str!("shader.wgsl").into()),
});
let plane = |binding: u32| wgpu::BindGroupLayoutEntry {
binding,
visibility: wgpu::ShaderStages::FRAGMENT,
ty: wgpu::BindingType::Texture {
sample_type: wgpu::TextureSampleType::Float { filterable: true },
view_dimension: wgpu::TextureViewDimension::D2,
multisampled: false,
},
count: None,
};
let layout = device.create_bind_group_layout(&wgpu::BindGroupLayoutDescriptor {
label: Some("moq-video planes"),
entries: &[
wgpu::BindGroupLayoutEntry {
binding: 0,
visibility: wgpu::ShaderStages::FRAGMENT,
ty: wgpu::BindingType::Buffer {
ty: wgpu::BufferBindingType::Uniform,
has_dynamic_offset: false,
min_binding_size: None,
},
count: None,
},
wgpu::BindGroupLayoutEntry {
binding: 1,
visibility: wgpu::ShaderStages::FRAGMENT,
ty: wgpu::BindingType::Sampler(wgpu::SamplerBindingType::Filtering),
count: None,
},
plane(2),
plane(3),
plane(4),
],
});
let pipeline_layout = device.create_pipeline_layout(&wgpu::PipelineLayoutDescriptor {
label: Some("moq-video yuv"),
bind_group_layouts: &[Some(&layout)],
immediate_size: 0,
});
let pipeline = |entry: &str| {
device.create_render_pipeline(&wgpu::RenderPipelineDescriptor {
label: Some("moq-video yuv to rgb"),
layout: Some(&pipeline_layout),
vertex: wgpu::VertexState {
module: &module,
entry_point: Some("vertex"),
compilation_options: Default::default(),
buffers: &[],
},
primitive: wgpu::PrimitiveState::default(),
depth_stencil: None,
multisample: wgpu::MultisampleState::default(),
fragment: Some(wgpu::FragmentState {
module: &module,
entry_point: Some(entry),
compilation_options: Default::default(),
targets: &[Some(wgpu::ColorTargetState {
format,
blend: None,
write_mask: wgpu::ColorWrites::ALL,
})],
}),
multiview_mask: None,
cache: None,
})
};
let sampler = device.create_sampler(&wgpu::SamplerDescriptor {
label: Some("moq-video planes"),
address_mode_u: wgpu::AddressMode::ClampToEdge,
address_mode_v: wgpu::AddressMode::ClampToEdge,
address_mode_w: wgpu::AddressMode::ClampToEdge,
mag_filter: wgpu::FilterMode::Linear,
min_filter: wgpu::FilterMode::Linear,
mipmap_filter: wgpu::MipmapFilterMode::Nearest,
..Default::default()
});
Ok(Self {
#[cfg(all(target_os = "linux", feature = "dmabuf"))]
rgba: pipeline("rgba"),
#[cfg(any(target_os = "macos", all(target_os = "linux", feature = "dmabuf")))]
nv12: pipeline("nv12"),
i420: pipeline("i420"),
layout,
sampler,
})
}
}
#[cfg(test)]
mod tests {
use moq_net::Timestamp;
use super::*;
use crate::Surface;
#[cfg(all(target_os = "linux", feature = "dmabuf"))]
#[test]
fn dma_buf_fence_timeout_is_terminal_for_the_frame() {
let timed_out = Error::Render(anyhow::Error::new(std::io::Error::from(std::io::ErrorKind::TimedOut)));
let other = Error::Render(anyhow::Error::new(std::io::Error::from(
std::io::ErrorKind::PermissionDenied,
)));
assert!(dma_buf_import_timed_out(&timed_out));
assert!(!dma_buf_import_timed_out(&other));
}
#[cfg(all(target_os = "linux", feature = "pipewire"))]
#[tokio::test]
#[ignore = "needs a PipeWire desktop and Vulkan GPU"]
async fn packed_dmabuf_renders_through_vulkan() {
let instance = wgpu::Instance::default();
let adapter = instance
.request_adapter(&wgpu::RequestAdapterOptions::default())
.await
.expect("a GPU adapter");
assert_eq!(
adapter.get_info().backend,
wgpu::Backend::Vulkan,
"DMA-BUF import needs Vulkan"
);
let external_memory = wgpu::Features::VULKAN_EXTERNAL_MEMORY_DMA_BUF;
assert!(
adapter.features().contains(external_memory),
"Vulkan adapter does not support DMA-BUF external memory"
);
let (device, queue) = adapter
.request_device(&wgpu::DeviceDescriptor {
required_features: external_memory,
..Default::default()
})
.await
.expect("a DMA-BUF-capable GPU device");
let mut renderer = Renderer::new(&device, &queue, Config::new()).expect("a renderer");
let capture = crate::capture::Config {
source: crate::capture::Source::Display(None),
..Default::default()
};
let mut stream = crate::capture::open(&capture).await.expect("portal screen capture");
for index in 0..16 {
let surface = tokio::time::timeout(std::time::Duration::from_secs(5), stream.read())
.await
.unwrap_or_else(|_| panic!("timed out waiting for frame {index}"))
.unwrap_or_else(|error| panic!("capture failed before frame {index}: {error}"))
.unwrap_or_else(|| panic!("capture ended before frame {index}"));
let Surface::DmaBuf(buffer) = &surface else {
panic!("frame {index} used shared memory instead of DMA-BUF");
};
assert!(
matches!(
buffer.format(),
crate::DrmFormat::XRGB8888
| crate::DrmFormat::ARGB8888
| crate::DrmFormat::XBGR8888
| crate::DrmFormat::ABGR8888
),
"frame {index} negotiated unsupported DMA-BUF format {:#x}",
buffer.format().as_raw()
);
let frame = Frame::new(surface, Timestamp::ZERO);
let imported = renderer
.source
.import(&device, &frame.surface)
.expect("Vulkan DMA-BUF import")
.expect("a DMA-BUF import path");
assert_eq!(imported.layout, Layout::Rgba);
drop(imported);
let texture = renderer.render(&frame).expect("a zero-copy rendered frame");
assert_eq!(renderer.strikes, 0, "frame {index} fell back to the CPU");
assert!(!renderer.retired, "frame {index} retired the zero-copy path");
assert_eq!((texture.width(), texture.height()), (stream.width(), stream.height()));
}
drop(renderer);
device
.poll(wgpu::PollType::wait_indefinitely())
.expect("all imported frame reads completed");
}
async fn gpu() -> (wgpu::Device, wgpu::Queue) {
let instance = wgpu::Instance::default();
let adapter = instance
.request_adapter(&wgpu::RequestAdapterOptions::default())
.await
.expect("a GPU adapter");
adapter
.request_device(&wgpu::DeviceDescriptor::default())
.await
.expect("a GPU device")
}
fn solid(size: Size, rgba: [u8; 4]) -> Frame {
let pixels: Vec<u8> = rgba.iter().copied().cycle().take(size.pixels() as usize * 4).collect();
let surface = Surface::rgba(&pixels, size).expect("a valid RGBA frame");
Frame::new(surface, Timestamp::ZERO)
}
async fn readback(device: &wgpu::Device, queue: &wgpu::Queue, texture: &wgpu::Texture) -> Vec<[u8; 4]> {
let (width, height) = (texture.width(), texture.height());
let row = (width * 4).next_multiple_of(wgpu::COPY_BYTES_PER_ROW_ALIGNMENT);
let buffer = device.create_buffer(&wgpu::BufferDescriptor {
label: Some("readback"),
size: (row * height) as u64,
usage: wgpu::BufferUsages::COPY_DST | wgpu::BufferUsages::MAP_READ,
mapped_at_creation: false,
});
let mut encoder = device.create_command_encoder(&Default::default());
encoder.copy_texture_to_buffer(
wgpu::TexelCopyTextureInfo {
texture,
mip_level: 0,
origin: wgpu::Origin3d::ZERO,
aspect: wgpu::TextureAspect::All,
},
wgpu::TexelCopyBufferInfo {
buffer: &buffer,
layout: wgpu::TexelCopyBufferLayout {
offset: 0,
bytes_per_row: Some(row),
rows_per_image: Some(height),
},
},
wgpu::Extent3d {
width,
height,
depth_or_array_layers: 1,
},
);
queue.submit([encoder.finish()]);
let (send, recv) = tokio::sync::oneshot::channel();
buffer.map_async(wgpu::MapMode::Read, .., |result| {
let _ = send.send(result);
});
device
.poll(wgpu::PollType::wait_indefinitely())
.expect("the copy to complete");
recv.await.expect("a mapping result").expect("a mapped buffer");
let view = buffer.slice(..).get_mapped_range().expect("a mapped range");
let mut pixels = Vec::with_capacity((width * height) as usize);
for y in 0..height as usize {
let start = y * row as usize;
for x in 0..width as usize {
let px = &view[start + x * 4..start + x * 4 + 4];
pixels.push([px[0], px[1], px[2], px[3]]);
}
}
pixels
}
fn assert_close(actual: [u8; 4], expected: [u8; 4]) {
for channel in 0..3 {
let (a, e) = (actual[channel] as i32, expected[channel] as i32);
assert!((a - e).abs() <= 6, "got {actual:?}, expected about {expected:?}");
}
assert_eq!(actual[3], 255, "alpha should be opaque");
}
#[tokio::test]
#[ignore]
async fn the_declared_color_space_beats_the_size_guess() {
let (device, queue) = gpu().await;
let size = Size::new(1280, 720);
let config = Config {
usage: wgpu::TextureUsages::COPY_SRC,
..Config::new()
};
let mut renderer = Renderer::new(&device, &queue, config).expect("a renderer");
for rgba in [[255, 0, 0, 255], [0, 255, 0, 255], [0, 0, 255, 255]] {
let frame = solid(size, rgba);
assert_eq!(
frame.surface.color(),
Some(crate::Color::infer(size)),
"the RGB conversion reports the space it converted into"
);
let texture = renderer.render(&frame).expect("a rendered frame");
let pixels = readback(&device, &queue, &texture).await;
let center = (size.height as usize / 2) * size.width as usize + size.width as usize / 2;
assert_close(pixels[center], rgba);
let sd = solid(Size::new(640, 480), rgba);
assert_eq!(sd.surface.color(), Some(crate::Color::Bt601Limited));
let scaled = sd.resize(size).expect("scale past 576 lines");
assert_eq!(
scaled.surface.color(),
Some(crate::Color::Bt601Limited),
"resize carries the space across rather than re-guessing"
);
assert_ne!(scaled.surface.color(), Some(crate::Color::infer(size)));
let texture = renderer.render(&scaled).expect("a rendered frame");
let pixels = readback(&device, &queue, &texture).await;
assert_close(pixels[center], rgba);
}
}
#[tokio::test]
#[ignore]
async fn cpu_frames_survive_the_round_trip() {
let (device, queue) = gpu().await;
let size = Size::new(64, 64);
let config = Config {
usage: wgpu::TextureUsages::COPY_SRC,
..Config::new()
};
let mut renderer = Renderer::new(&device, &queue, config).expect("a renderer");
for rgba in [
[255, 0, 0, 255],
[0, 255, 0, 255],
[0, 0, 255, 255],
[255, 255, 255, 255],
[0, 0, 0, 255],
[77, 153, 230, 255],
] {
let texture = renderer.render(&solid(size, rgba)).expect("a rendered frame");
assert_eq!((texture.width(), texture.height()), (size.width, size.height));
let pixels = readback(&device, &queue, &texture).await;
assert_close(
pixels[(size.height as usize / 2) * size.width as usize + size.width as usize / 2],
rgba,
);
}
}
#[tokio::test]
#[ignore]
async fn config_size_overrides_the_frame_size() {
let (device, queue) = gpu().await;
let config = Config {
size: Some(Size::new(32, 16)),
usage: wgpu::TextureUsages::COPY_SRC,
..Config::new()
};
let mut renderer = Renderer::new(&device, &queue, config).expect("a renderer");
let texture = renderer
.render(&solid(Size::new(64, 64), [255, 0, 0, 255]))
.expect("a rendered frame");
assert_eq!((texture.width(), texture.height()), (32, 16));
let pixels = readback(&device, &queue, &texture).await;
assert_close(pixels[8 * 32 + 16], [255, 0, 0, 255]);
}
#[cfg(all(target_os = "linux", feature = "vaapi"))]
const PALETTE: [[u8; 3]; 8] = [
[255, 0, 0],
[0, 255, 0],
[0, 0, 255],
[255, 255, 0],
[0, 255, 255],
[255, 0, 255],
[255, 128, 0],
[255, 255, 255],
];
#[cfg(all(target_os = "linux", feature = "vaapi"))]
fn block(x: u32, y: u32) -> [u8; 3] {
PALETTE[((y / 32 * 3 + x / 32) % PALETTE.len() as u32) as usize]
}
#[cfg(all(target_os = "linux", feature = "vaapi"))]
fn to_yuv(color: Color, rgb: [u8; 3]) -> [u8; 3] {
let (kr, kb) = match color {
Color::Bt601Limited | Color::Bt601Full => (0.299f32, 0.114f32),
Color::Bt709Limited | Color::Bt709Full => (0.2126, 0.0722),
};
let kg = 1.0 - kr - kb;
let [r, g, b] = rgb.map(|c| c as f32 / 255.0);
let luma = kr * r + kg * g + kb * b;
let (y, cb, cr) = (
luma * 219.0 + 16.0,
(b - luma) / (2.0 * (1.0 - kb)) * 224.0 + 128.0,
(r - luma) / (2.0 * (1.0 - kr)) * 224.0 + 128.0,
);
[y, cb, cr].map(|c| c.round().clamp(0.0, 255.0) as u8)
}
#[cfg(all(target_os = "linux", feature = "vaapi"))]
fn pattern(format: crate::DrmFormat, size: Size, color: Color) -> Vec<u8> {
let (width, height) = (size.width, size.height);
let mut pixels = Vec::with_capacity((width * height * 3 / 2) as usize);
for y in 0..height {
for x in 0..width {
pixels.push(to_yuv(color, block(x, y))[0]);
}
}
let chroma = |channel: usize, pixels: &mut Vec<u8>| {
for y in (0..height).step_by(2) {
for x in (0..width).step_by(2) {
pixels.push(to_yuv(color, block(x, y))[channel]);
}
}
};
match format {
crate::DrmFormat::NV12 => {
for y in (0..height).step_by(2) {
for x in (0..width).step_by(2) {
let [_, cb, cr] = to_yuv(color, block(x, y));
pixels.push(cb);
pixels.push(cr);
}
}
}
crate::DrmFormat::YUV420 => {
chroma(1, &mut pixels);
chroma(2, &mut pixels);
}
format => panic!("no pattern for DMA-BUF format {:#x}", format.as_raw()),
}
pixels
}
#[cfg(all(target_os = "linux", feature = "vaapi"))]
async fn dmabuf_gpu() -> Option<(wgpu::Device, wgpu::Queue)> {
let instance = wgpu::Instance::default();
let adapter = instance
.request_adapter(&wgpu::RequestAdapterOptions::default())
.await
.expect("a GPU adapter");
let external_memory = wgpu::Features::VULKAN_EXTERNAL_MEMORY_DMA_BUF;
if adapter.get_info().backend != wgpu::Backend::Vulkan || !adapter.features().contains(external_memory) {
return None;
}
Some(
adapter
.request_device(&wgpu::DeviceDescriptor {
required_features: external_memory,
..Default::default()
})
.await
.expect("a DMA-BUF-capable GPU device"),
)
}
#[cfg(all(target_os = "linux", feature = "vaapi"))]
async fn imported_matches_uploaded(
device: &wgpu::Device,
queue: &wgpu::Queue,
format: crate::DrmFormat,
size: Size,
expected: Layout,
) -> bool {
let color = Color::infer(size);
assert_eq!(color, Color::Bt601Limited);
let Some(buffer) = crate::render::dmabuf::fixture::surface(format, size, &pattern(format, size, color)) else {
return false;
};
eprintln!(
"{size} {:?}: modifier {:#x}, planes {:?}",
expected,
buffer.modifier(),
buffer.planes()
);
let config = Config {
usage: wgpu::TextureUsages::COPY_SRC,
..Config::new()
};
let uploaded = {
let nv12 = pattern(crate::DrmFormat::NV12, size, color);
let i420 = crate::frame::I420::from_nv12(&nv12, size.width, size.height).expect("deinterleave NV12");
let frame = Frame::new(Surface::I420(i420), Timestamp::ZERO);
let mut renderer = Renderer::new(device, queue, config.clone()).expect("a renderer");
let texture = renderer.render(&frame).expect("a rendered frame");
readback(device, queue, &texture).await
};
let frame = Frame::new(Surface::DmaBuf(buffer), Timestamp::ZERO);
let mut renderer = Renderer::new(device, queue, config).expect("a renderer");
let imported = renderer
.source
.import(device, &frame.surface)
.expect("import the DMA-BUF")
.expect("a DMA-BUF import path");
assert_eq!(imported.layout, expected, "the frame imported as the wrong layout");
assert_eq!(imported.color, Some(color));
drop(imported);
let texture = renderer.render(&frame).expect("a rendered frame");
assert_eq!(renderer.strikes, 0, "the zero-copy import should not have failed");
assert!(!renderer.retired);
let zero_copy = readback(device, queue, &texture).await;
assert_eq!(zero_copy.len(), uploaded.len());
let mut worst = 0u8;
for (index, (&imported, &reference)) in zero_copy.iter().zip(&uploaded).enumerate() {
let drift = (0..4).map(|c| imported[c].abs_diff(reference[c])).max().unwrap_or(0);
worst = worst.max(drift);
let (x, y) = (index % size.width as usize, index / size.width as usize);
assert!(drift <= 2, "({x}, {y}): imported {imported:?}, uploaded {reference:?}");
}
eprintln!(
" imported vs uploaded: worst drift {worst} of 255 over {} pixels",
zero_copy.len(),
);
for y in (16..size.height).step_by(32) {
for x in (16..size.width).step_by(32) {
let rgb = block(x, y);
assert_close(zero_copy[(y * size.width + x) as usize], [rgb[0], rgb[1], rgb[2], 255]);
}
}
true
}
#[cfg(all(target_os = "linux", feature = "vaapi"))]
#[tokio::test]
#[ignore = "needs a Vulkan GPU and a VA-API device"]
async fn the_nv12_dmabuf_import_matches_the_cpu_path() {
let Some((device, queue)) = dmabuf_gpu().await else {
eprintln!("skipping: no Vulkan adapter with DMA-BUF external memory");
return;
};
let mut ran = false;
for size in [Size::new(256, 192), Size::new(200, 120), Size::new(62, 34)] {
ran |= imported_matches_uploaded(&device, &queue, crate::DrmFormat::NV12, size, Layout::Nv12).await;
}
assert!(ran, "no NV12 DMA-BUF could be allocated to import");
}
#[cfg(all(target_os = "linux", feature = "vaapi"))]
#[tokio::test]
#[ignore = "needs a Vulkan GPU and a VA-API device"]
async fn the_i420_dmabuf_import_matches_the_cpu_path() {
let Some((device, queue)) = dmabuf_gpu().await else {
eprintln!("skipping: no Vulkan adapter with DMA-BUF external memory");
return;
};
for size in [Size::new(256, 192), Size::new(200, 120)] {
if !imported_matches_uploaded(&device, &queue, crate::DrmFormat::YUV420, size, Layout::I420).await {
eprintln!("skipping: this driver does not do YU12 DMA-BUFs");
return;
}
}
}
#[cfg(all(target_os = "linux", feature = "vaapi"))]
#[tokio::test]
#[ignore = "needs a Vulkan GPU and a VA-API device"]
async fn decoded_frames_reach_the_gpu_without_a_download() {
use crate::decode::backend::{self, Codec, vaapi};
let Some((device, queue)) = dmabuf_gpu().await else {
eprintln!("skipping: no Vulkan adapter with DMA-BUF external memory");
return;
};
let decode = |gpu_frames| {
backend::open(
Codec::H264,
&crate::decode::Config {
kind: crate::decode::Kind::Named(vaapi::NAME.into()),
gpu_frames,
..crate::decode::Config::new()
},
)
};
let Ok(mut exporting) = decode(true) else {
eprintln!("skipping: no VA-API H.264 decoder");
return;
};
let mut downloading = decode(false).expect("a second decoder");
let size = Size::new(320, 240);
let (width, height) = (size.width, size.height);
let mut rgba = vec![0u8; (width * height * 4) as usize];
for y in 0..height {
for x in 0..width {
let index = ((y * width + x) * 4) as usize;
rgba[index] = (x * 255 / width) as u8;
rgba[index + 1] = (y * 255 / height) as u8;
rgba[index + 2] = ((x + y) * 255 / (width + height)) as u8;
rgba[index + 3] = 255;
}
}
let mut encoder = crate::encode::Encoder::new(&crate::encode::Config {
kind: crate::encode::Kind::Software,
..crate::encode::Config::new(width, height, 30)
})
.expect("a software H.264 encoder");
let mut exported = Vec::new();
let mut downloaded = Vec::new();
for index in 0..8u64 {
if index == 0 {
encoder.keyframe();
}
let surface = Surface::rgba(&rgba, size).expect("a valid RGBA frame");
let frame = Frame::new(surface, Timestamp::from_micros(index * 33_333).unwrap());
for unit in encoder.encode(&frame).expect("encode a picture") {
let timestamp = unit.timestamp;
exported.extend(
exporting
.decode(unit.payload.clone(), timestamp, index == 0)
.expect("decode to the GPU"),
);
downloaded.extend(
downloading
.decode(unit.payload, timestamp, index == 0)
.expect("decode to the CPU"),
);
}
}
assert!(!exported.is_empty(), "the decoder produced no pictures");
assert_eq!(exported.len(), downloaded.len(), "the two decoders disagreed");
let Surface::DmaBuf(first) = &exported[0].surface else {
panic!("gpu_frames did not produce a DMA-BUF surface");
};
eprintln!(
"decoded {} pictures, exported at modifier {:#x}",
exported.len(),
first.modifier()
);
let config = Config {
usage: wgpu::TextureUsages::COPY_SRC,
..Config::new()
};
let mut importing = Renderer::new(&device, &queue, config.clone()).expect("a renderer");
let mut uploading = Renderer::new(&device, &queue, config).expect("a renderer");
for (index, (gpu, cpu)) in exported.iter().zip(&downloaded).enumerate() {
let Surface::DmaBuf(buffer) = &gpu.surface else {
panic!("picture {index} did not come back GPU-resident");
};
if index == 0 {
eprintln!(" planes {:?}", buffer.planes());
}
assert!(
matches!(cpu.surface, Surface::I420(_)),
"picture {index} was not downloaded without gpu_frames"
);
let source = importing
.source
.import(&device, &gpu.surface)
.expect("import the decoded surface")
.expect("a DMA-BUF import path");
assert_eq!(source.layout, Layout::Nv12, "picture {index} did not import as NV12");
drop(source);
let texture = importing.render(gpu).expect("draw the imported picture");
assert_eq!(importing.strikes, 0, "picture {index} fell back to the CPU");
let zero_copy = readback(&device, &queue, &texture).await;
let texture = uploading.render(cpu).expect("draw the downloaded picture");
let cpu = readback(&device, &queue, &texture).await;
assert_eq!(
zero_copy.len(),
cpu.len(),
"picture {index} read back at a different size"
);
let mut worst = 0u8;
for (pixel, (&imported, &reference)) in zero_copy.iter().zip(&cpu).enumerate() {
let drift = (0..4).map(|c| imported[c].abs_diff(reference[c])).max().unwrap_or(0);
worst = worst.max(drift);
let (x, y) = (pixel % width as usize, pixel / width as usize);
assert!(
drift <= 2,
"picture {index} at ({x}, {y}): imported {imported:?}, downloaded {reference:?}"
);
}
if index == 0 {
eprintln!(
" imported vs downloaded: worst drift {worst} of 255 over {} pixels",
cpu.len()
);
}
}
}
#[cfg(target_os = "macos")]
fn pooled(size: Size, rgba: [u8; 4]) -> crate::Surface {
let uploaded = solid(size, rgba).surface.into_pixel_buffer().expect("a pixel buffer");
let planar =
crate::Surface::PixelBuffer(crate::frame::macos::PixelBuffer::new(uploaded, size.width, size.height));
planar.resize(size).expect("a transfer into the NV12 pool")
}
#[cfg(target_os = "macos")]
#[tokio::test]
#[ignore]
async fn imports_survive_decoder_pool_recycling() {
let (device, queue) = gpu().await;
let size = Size::new(256, 256);
let config = Config {
usage: wgpu::TextureUsages::COPY_SRC,
..Config::new()
};
let mut renderer = Renderer::new(&device, &queue, config).expect("a renderer");
let colors = [[255, 0, 0, 255], [0, 255, 0, 255], [0, 0, 255, 255], [255, 255, 0, 255]];
for round in 0..24 {
let rgba = colors[round % colors.len()];
let texture = renderer
.render(&Frame::new(pooled(size, rgba), Timestamp::ZERO))
.expect("a rendered frame");
assert_eq!(renderer.strikes, 0, "round {round} fell back to the CPU");
let pixels = readback(&device, &queue, &texture).await;
assert_close(pixels[128 * 256 + 128], rgba);
}
}
#[cfg(target_os = "macos")]
#[tokio::test]
#[ignore]
async fn the_metal_import_matches_the_cpu_path() {
let (device, queue) = gpu().await;
let size = Size::new(64, 64);
let rgba = [77, 153, 230, 255];
let config = Config {
usage: wgpu::TextureUsages::COPY_SRC,
..Config::new()
};
let uploaded = {
let mut renderer = Renderer::new(&device, &queue, config.clone()).expect("a renderer");
let texture = renderer.render(&solid(size, rgba)).expect("a rendered frame");
readback(&device, &queue, &texture).await
};
let surface = pooled(size, rgba);
let mut renderer = Renderer::new(&device, &queue, config).expect("a renderer");
let texture = renderer
.render(&Frame::new(surface, Timestamp::ZERO))
.expect("a rendered frame");
assert_eq!(renderer.strikes, 0, "the zero-copy import should not have failed");
assert!(!renderer.retired);
let imported = readback(&device, &queue, &texture).await;
assert_eq!(uploaded.len(), imported.len());
for (upload, import) in uploaded.iter().zip(imported.iter()) {
assert_close(*import, *upload);
}
}
}