diff --git a/base/cvd/cuttlefish/host/commands/vhost_user_media/v4l2_stream_proxy/src/device.rs b/base/cvd/cuttlefish/host/commands/vhost_user_media/v4l2_stream_proxy/src/device.rs index e7a72332b80..2dd0303a5c9 100644 --- a/base/cvd/cuttlefish/host/commands/vhost_user_media/v4l2_stream_proxy/src/device.rs +++ b/base/cvd/cuttlefish/host/commands/vhost_user_media/v4l2_stream_proxy/src/device.rs @@ -12,15 +12,14 @@ // See the License for the specific language governing permissions and // limitations under the License. +use nix::sys::eventfd::{EfdFlags, EventFd}; use std::collections::VecDeque; use std::io::Result as IoResult; use std::os::fd::AsFd; use std::os::fd::BorrowedFd; -use std::sync::{Arc, Mutex}; use std::sync::mpsc::channel; +use std::sync::{Arc, Mutex}; use std::time::Instant; -use nix::sys::eventfd::{EventFd, EfdFlags}; - use v4l2r::PixelFormat; use v4l2r::QueueType; @@ -31,6 +30,9 @@ use v4l2r::bindings::v4l2_requestbuffers; use v4l2r::ioctl::BufferCapabilities; use v4l2r::ioctl::BufferField; use v4l2r::ioctl::BufferFlags; +use v4l2r::ioctl::CtrlId; +use v4l2r::ioctl::CtrlWhich; +use v4l2r::ioctl::QueryCtrlFlags; use v4l2r::ioctl::V4l2Buffer; use v4l2r::ioctl::V4l2PlanesWithBackingMut; use v4l2r::memory::MemoryType; @@ -94,11 +96,7 @@ impl Format { let w = width as usize; let h = height as usize; match self { - Self::Yuv420M => vec![ - w * h, - w * h / 4, - w * h / 4, - ], + Self::Yuv420M => vec![w * h, w * h / 4, w * h / 4], } } @@ -179,9 +177,11 @@ impl Buffer { tv_usec: (ts.tv_nsec / 1000) as bindings::__time_t, }); Self::unset_flag(&mut flags, BufferFlags::QUEUED); - + // Set bytesused to plane length for all planes - if let V4l2PlanesWithBackingMut::Mmap(planes) = self.v4l2_buffer.planes_with_backing_iter_mut() { + if let V4l2PlanesWithBackingMut::Mmap(planes) = + self.v4l2_buffer.planes_with_backing_iter_mut() + { for mut plane in planes { let len = *plane.length; *plane.bytesused = len; @@ -194,7 +194,9 @@ impl Buffer { } fn clear_bytesused(&mut self) { - if let V4l2PlanesWithBackingMut::Mmap(planes) = self.v4l2_buffer.planes_with_backing_iter_mut() { + if let V4l2PlanesWithBackingMut::Mmap(planes) = + self.v4l2_buffer.planes_with_backing_iter_mut() + { for mut plane in planes { *plane.bytesused = 0; } @@ -270,13 +272,16 @@ where } fn default_fmt(&self, queue: QueueType) -> v4l2_format { - let plane_sizes = self.config.format.plane_sizes(self.config.input_width, self.config.input_height); + let plane_sizes = self + .config + .format + .plane_sizes(self.config.input_width, self.config.input_height); let mut plane_fmt: [bindings::v4l2_plane_pix_format; 8] = Default::default(); for (i, &size) in plane_sizes.iter().enumerate() { plane_fmt[i].sizeimage = size as u32; plane_fmt[i].bytesperline = self.config.format.bytesperline(self.config.input_width, i); } - + let pix_mp = bindings::v4l2_pix_format_mplane { width: self.config.input_width, height: self.config.input_height, @@ -421,11 +426,7 @@ where Ok(self.default_fmtdesc(queue)) } - fn g_fmt( - &mut self, - _session: &Self::Session, - queue: QueueType, - ) -> IoctlResult { + fn g_fmt(&mut self, _session: &Self::Session, queue: QueueType) -> IoctlResult { if queue != self.queue_type() { return Err(libc::EINVAL); } @@ -549,19 +550,21 @@ where } } - let plane_sizes = self.config.format.plane_sizes(self.config.input_width, self.config.input_height); + let plane_sizes = self + .config + .format + .plane_sizes(self.config.input_width, self.config.input_height); let num_planes = plane_sizes.len(); state.buffers = (0..count) .map(|i| -> Result { let mut planes = Vec::new(); - + for &size in &plane_sizes { - let fd = MemFdBuffer::new(size as u64) - .map_err(|e| { - log::error!("failed to allocate MMAP buffer: {:#}", e); - libc::ENOMEM - })?; + let fd = MemFdBuffer::new(size as u64).map_err(|e| { + log::error!("failed to allocate MMAP buffer: {:#}", e); + libc::ENOMEM + })?; let offset = self .mmap_manager .register_buffer(None, size as u32) @@ -570,7 +573,7 @@ where } let mut v4l2_buffer = V4l2Buffer::new(expected_queue, i, MemoryType::Mmap); - + if num_planes > 1 { unsafe { (*v4l2_buffer.as_mut_ptr()).length = num_planes as u32; @@ -596,7 +599,7 @@ where panic!() } } - + v4l2_buffer.set_field(BufferField::None); v4l2_buffer.set_flags(BufferFlags::TIMESTAMP_MONOTONIC); @@ -638,7 +641,7 @@ where ) -> IoctlResult { let mut state = session.state.lock().unwrap(); let buf_id = qbuf.index() as usize; - + let buf_v4l2 = { let buffer = state.buffers.get_mut(buf_id).ok_or(libc::EINVAL)?; if buffer.state == BufferState::Incoming { @@ -647,9 +650,9 @@ where buffer.set_state(BufferState::Incoming); buffer.v4l2_buffer.clone() }; - + state.queued_buffers.push_back(buf_id); - + if session.streaming { if let Some(ref worker) = session.worker { let _ = worker.tx.send(WorkerCmd::BufferQueued); @@ -691,10 +694,20 @@ where let config_clone = self.config.clone(); let join_handle = std::thread::spawn(move || { - worker_thread_loop(config_clone, evt_queue_clone, state_clone, rx, event_fd_clone); + worker_thread_loop( + config_clone, + evt_queue_clone, + state_clone, + rx, + event_fd_clone, + ); }); - session.worker = Some(WorkerHandle { tx, event_fd, join_handle }); + session.worker = Some(WorkerHandle { + tx, + event_fd, + join_handle, + }); Ok(()) } @@ -799,4 +812,135 @@ where ..Default::default() }) } + + fn query_ext_ctrl( + &mut self, + _session: &Self::Session, + id: CtrlId, + flags: QueryCtrlFlags, + ) -> IoctlResult { + let requested_id: u32 = unsafe { std::mem::transmute(id) }; + + if flags.contains(QueryCtrlFlags::NEXT) { + if requested_id < CID_LENS_FACING { + return Ok(self.lens_facing_query_ext_ctrl()); + } + } else if requested_id == CID_LENS_FACING { + return Ok(self.lens_facing_query_ext_ctrl()); + } + + Err(libc::EINVAL) + } + + fn g_ctrl(&mut self, _session: &Self::Session, id: u32) -> IoctlResult { + if id == CID_LENS_FACING { + return Ok(bindings::v4l2_control { + id, + value: self.config.lens_facing as i32, + }); + } + Err(libc::EINVAL) + } + + fn g_ext_ctrls( + &mut self, + _session: &Self::Session, + _which: CtrlWhich, + _ctrls: &mut bindings::v4l2_ext_controls, + ctrl_array: &mut Vec, + _user_regions: Vec>, + ) -> IoctlResult<()> { + for ctrl in ctrl_array.iter_mut() { + if ctrl.id == CID_LENS_FACING { + ctrl.__bindgen_anon_1.value64 = self.config.lens_facing as i64; + } else { + return Err(libc::EINVAL); + } + } + Ok(()) + } +} + +/// https://developer.android.com/reference/android/hardware/camera2/CameraMetadata#LENS_FACING_FRONT +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum LensFacing { + Front = 0, + Back = 1, + External = 2, +} + +impl std::str::FromStr for LensFacing { + type Err = String; + + fn from_str(s: &str) -> Result { + match s { + "FRONT" => Ok(LensFacing::Front), + "BACK" => Ok(LensFacing::Back), + "EXTERNAL" => Ok(LensFacing::External), + _ => Err(format!( + "Invalid lens facing: {}. Expected FRONT, BACK, or EXTERNAL", + s + )), + } + } +} + +const CID_OFFSET: u32 = bindings::V4L2_CID_CAMERA_CLASS_BASE + 0x100; +const CID_LENS_FACING: u32 = CID_OFFSET + 1; + +fn ctrl_name(name: &str) -> [i8; 32] { + let mut array = [0i8; 32]; + let bytes = name.as_bytes(); + let len = std::cmp::min(bytes.len(), 31); + for i in 0..len { + array[i] = bytes[i] as i8; + } + array +} + +impl V4l2Stream { + fn lens_facing_query_ext_ctrl(&self) -> bindings::v4l2_query_ext_ctrl { + bindings::v4l2_query_ext_ctrl { + id: CID_LENS_FACING, + type_: bindings::v4l2_ctrl_type_V4L2_CTRL_TYPE_INTEGER, + name: ctrl_name("LENS_FACING"), + minimum: 0, + maximum: 2, + step: 1, + default_value: self.config.lens_facing as i64, + flags: bindings::V4L2_CTRL_FLAG_READ_ONLY, + elems: 1, + elem_size: std::mem::size_of::() as u32, + ..Default::default() + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_lens_facing_from_str() { + assert_eq!("FRONT".parse::().unwrap(), LensFacing::Front); + assert_eq!("BACK".parse::().unwrap(), LensFacing::Back); + assert_eq!( + "EXTERNAL".parse::().unwrap(), + LensFacing::External + ); + assert!("INVALID".parse::().is_err()); + } + + #[test] + fn test_ctrl_name_is_nul_padded() { + let name = ctrl_name("LENS_FACING"); + assert_eq!(&name[..11], b"LENS_FACING".map(|b| b as i8)); + assert!(name[11..].iter().all(|byte| *byte == 0)); + } + + #[test] + fn test_ctrl_name_truncates_and_stays_nul_terminated() { + let name = ctrl_name("This control name is definitely far too long to fit"); + assert_eq!(name[31], 0); + } } diff --git a/base/cvd/cuttlefish/host/commands/vhost_user_media/v4l2_stream_proxy/src/main.rs b/base/cvd/cuttlefish/host/commands/vhost_user_media/v4l2_stream_proxy/src/main.rs index fe5628c5c1b..1041bcf163e 100644 --- a/base/cvd/cuttlefish/host/commands/vhost_user_media/v4l2_stream_proxy/src/main.rs +++ b/base/cvd/cuttlefish/host/commands/vhost_user_media/v4l2_stream_proxy/src/main.rs @@ -49,6 +49,9 @@ struct CmdLineArgs { /// Frames per second (e.g. 30 or 30000/1001). #[clap(long = "input_fps", value_name = "INPUT_FPS")] input_fps: String, + /// Lens facing configuration: FRONT, BACK, or EXTERNAL. + #[clap(long, value_name = "LENS_FACING", default_value = "EXTERNAL")] + lens_facing: String, } fn parse_fps_to_interval(fps_str: &str) -> Option<(u32, u32)> { @@ -72,14 +75,20 @@ pub struct Config { pub input_height: u32, pub fps_interval: (u32, u32), pub format: device::Format, + pub lens_facing: device::LensFacing, } impl TryFrom for Config { type Error = Error; fn try_from(args: CmdLineArgs) -> Result { - let fps_interval = parse_fps_to_interval(&args.input_fps) - .ok_or_else(|| Error::InvalidArgument(format!("Invalid FPS format: {}", args.input_fps)))?; + let fps_interval = parse_fps_to_interval(&args.input_fps).ok_or_else(|| { + Error::InvalidArgument(format!("Invalid FPS format: {}", args.input_fps)) + })?; + let lens_facing = args + .lens_facing + .parse::() + .map_err(Error::InvalidArgument)?; Ok(Config { socket_path: args.socket_path, input_path: args.input_path, @@ -87,6 +96,7 @@ impl TryFrom for Config { input_height: args.input_height, fps_interval, format: device::Format::Yuv420M, + lens_facing, }) } } @@ -106,7 +116,7 @@ fn start_backend(config: Config) -> Result<()> { let mut card = [0u8; 32]; let card_name = "v4l2_stream_proxy"; card[0..card_name.len()].copy_from_slice(card_name.as_bytes()); - + loop { let caps = Capabilities::VIDEO_CAPTURE_MPLANE | Capabilities::STREAMING; let device_config = VirtioMediaDeviceConfig { @@ -133,6 +143,55 @@ fn start_backend(config: Config) -> Result<()> { } } +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_cmdline_args_default_lens_facing() { + let args = CmdLineArgs::try_parse_from(&[ + "v4l2_stream_proxy", + "--socket-path", + "/tmp/sock", + "--input_path", + "/tmp/fifo", + "--input_width", + "640", + "--input_height", + "480", + "--input_fps", + "30", + ]) + .unwrap(); + assert_eq!(args.lens_facing, "EXTERNAL"); + let config = Config::try_from(args).unwrap(); + assert_eq!(config.lens_facing, device::LensFacing::External); + } + + #[test] + fn test_cmdline_args_custom_lens_facing() { + let args = CmdLineArgs::try_parse_from(&[ + "v4l2_stream_proxy", + "--socket-path", + "/tmp/sock", + "--input_path", + "/tmp/fifo", + "--input_width", + "640", + "--input_height", + "480", + "--input_fps", + "30", + "--lens-facing", + "BACK", + ]) + .unwrap(); + assert_eq!(args.lens_facing, "BACK"); + let config = Config::try_from(args).unwrap(); + assert_eq!(config.lens_facing, device::LensFacing::Back); + } +} + fn main() -> Result<()> { let args = CmdLineArgs::parse(); diff --git a/base/cvd/cuttlefish/host/commands/vhost_user_media/v4l2_stream_proxy/src/worker.rs b/base/cvd/cuttlefish/host/commands/vhost_user_media/v4l2_stream_proxy/src/worker.rs index 3f55fd189f1..b19265917fe 100644 --- a/base/cvd/cuttlefish/host/commands/vhost_user_media/v4l2_stream_proxy/src/worker.rs +++ b/base/cvd/cuttlefish/host/commands/vhost_user_media/v4l2_stream_proxy/src/worker.rs @@ -13,16 +13,16 @@ // limitations under the License. use std::fs::File; -use std::io::{Read, Write, Seek, SeekFrom}; +use std::io::{Read, Seek, SeekFrom, Write}; use std::os::fd::AsFd; -use std::sync::{Arc, Mutex}; +use std::os::unix::fs::OpenOptionsExt; use std::sync::mpsc::{Receiver, Sender}; +use std::sync::{Arc, Mutex}; use std::thread::JoinHandle; use std::time::Instant; -use std::os::unix::fs::OpenOptionsExt; +use nix::poll::{PollFd, PollFlags, PollTimeout, poll}; use nix::sys::eventfd::EventFd; -use nix::poll::{poll, PollFd, PollFlags, PollTimeout}; use virtio_media::VirtioMediaEventQueue; use virtio_media::protocol::{DequeueBufferEvent, V4l2Event}; @@ -46,9 +46,7 @@ enum WorkerState { /// FIFO is closed. Trying to open it. Unopened, /// FIFO is open, but waiting for buffers. - Idle { - fifo_file: File, - }, + Idle { fifo_file: File }, /// FIFO is open and streaming into a buffer. Streaming { fifo_file: File, @@ -131,7 +129,11 @@ fn handle_idle( } } - if poll_fds[0].revents().unwrap_or(PollFlags::empty()).contains(PollFlags::POLLIN) { + if poll_fds[0] + .revents() + .unwrap_or(PollFlags::empty()) + .contains(PollFlags::POLLIN) + { if should_stop(rx, event_fd) { return WorkerState::Stopped; } @@ -170,21 +172,43 @@ fn handle_streaming( match poll(&mut poll_fds, PollTimeout::NONE) { Ok(_) => {} Err(e) if e == nix::Error::EINTR => { - return WorkerState::Streaming { fifo_file, buffer_idx, plane_idx, plane_offset }; + return WorkerState::Streaming { + fifo_file, + buffer_idx, + plane_idx, + plane_offset, + }; } Err(e) => { log::error!("Poll error in Streaming: {:?}", e); - return WorkerState::Streaming { fifo_file, buffer_idx, plane_idx, plane_offset }; + return WorkerState::Streaming { + fifo_file, + buffer_idx, + plane_idx, + plane_offset, + }; } } - if poll_fds[0].revents().unwrap_or(PollFlags::empty()).contains(PollFlags::POLLIN) { + if poll_fds[0] + .revents() + .unwrap_or(PollFlags::empty()) + .contains(PollFlags::POLLIN) + { if should_stop(rx, event_fd) { return WorkerState::Stopped; } } - if poll_fds[1].revents().unwrap_or(PollFlags::empty()).contains(PollFlags::POLLIN) { + // When the FIFO writer disconnects and the buffer is drained, Linux poll() + // reports only POLLHUP (without POLLIN). Handling POLLHUP ensures we call + // read() to consume remaining bytes and receive Ok(0) (EOF) to cleanly + // transition to WorkerState::Unopened rather than busy-looping on poll(). + if poll_fds[1] + .revents() + .unwrap_or(PollFlags::empty()) + .intersects(PollFlags::POLLIN | PollFlags::POLLHUP) + { let plane_size = plane_sizes[plane_idx]; let remaining = plane_size - plane_offset; let read_chunk = std::cmp::min(local_buf.len(), remaining); @@ -199,12 +223,16 @@ fn handle_streaming( let session_id = s_state.id; let buffer = &mut s_state.buffers[buffer_idx]; let plane = &mut buffer.planes[plane_idx]; - - if let Err(e) = plane.fd.as_file().seek(SeekFrom::Start(plane_offset as u64)) { + + if let Err(e) = plane + .fd + .as_file() + .seek(SeekFrom::Start(plane_offset as u64)) + { log::error!("Seek error: {:?}", e); return WorkerState::Stopped; } - + if let Err(e) = plane.fd.as_file().write_all(&local_buf[..bytes_read]) { log::error!("Write error: {:?}", e); return WorkerState::Stopped; @@ -214,7 +242,7 @@ fn handle_streaming( if plane_offset == plane_size { plane_idx += 1; plane_offset = 0; - + if plane_idx == plane_sizes.len() { let sequence = s_state.sequence; s_state.sequence += 1; @@ -222,15 +250,23 @@ fn handle_streaming( let delta = now.duration_since(s_state.last_frame_time); s_state.last_frame_time = now; - log::info!("Frame completed: session {}, seq {}, buf_idx {}, delta {:?}", session_id, sequence, buffer_idx, delta); + log::info!( + "Frame completed: session {}, seq {}, buf_idx {}, delta {:?}", + session_id, + sequence, + buffer_idx, + delta + ); s_state.buffers[buffer_idx].set_state(BufferState::Outgoing { sequence }); let v4l2_buf = s_state.buffers[buffer_idx].v4l2_buffer.clone(); - evt_queue.lock().unwrap().send_event(V4l2Event::DequeueBuffer(DequeueBufferEvent::new( - session_id, - v4l2_buf, - ))); + evt_queue + .lock() + .unwrap() + .send_event(V4l2Event::DequeueBuffer(DequeueBufferEvent::new( + session_id, v4l2_buf, + ))); if let Some(next_buf_idx) = s_state.queued_buffers.pop_front() { WorkerState::Streaming { @@ -243,22 +279,40 @@ fn handle_streaming( WorkerState::Idle { fifo_file } } } else { - WorkerState::Streaming { fifo_file, buffer_idx, plane_idx, plane_offset } + WorkerState::Streaming { + fifo_file, + buffer_idx, + plane_idx, + plane_offset, + } } } else { - WorkerState::Streaming { fifo_file, buffer_idx, plane_idx, plane_offset } + WorkerState::Streaming { + fifo_file, + buffer_idx, + plane_idx, + plane_offset, + } } } - Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => { - WorkerState::Streaming { fifo_file, buffer_idx, plane_idx, plane_offset } - } + Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => WorkerState::Streaming { + fifo_file, + buffer_idx, + plane_idx, + plane_offset, + }, Err(e) => { log::error!("Read error: {:?}", e); WorkerState::Unopened } } } else { - WorkerState::Streaming { fifo_file, buffer_idx, plane_idx, plane_offset } + WorkerState::Streaming { + fifo_file, + buffer_idx, + plane_idx, + plane_offset, + } } } @@ -271,16 +325,16 @@ pub(crate) fn worker_thread_loop( ) { log::info!("Worker thread started for FIFO: {:?}", config.input_path); - let plane_sizes = config.format.plane_sizes(config.input_width, config.input_height); + let plane_sizes = config + .format + .plane_sizes(config.input_width, config.input_height); let mut local_buf = [0u8; 4096]; - + let mut state = WorkerState::Unopened; while !matches!(state, WorkerState::Stopped) { state = match state { - WorkerState::Unopened => { - handle_unopened(&config, &session_state, &rx, &event_fd) - } + WorkerState::Unopened => handle_unopened(&config, &session_state, &rx, &event_fd), WorkerState::Idle { fifo_file } => { handle_idle(fifo_file, &session_state, &rx, &event_fd) } @@ -289,23 +343,21 @@ pub(crate) fn worker_thread_loop( buffer_idx, plane_idx, plane_offset, - } => { - handle_streaming( - fifo_file, - buffer_idx, - plane_idx, - plane_offset, - &session_state, - &rx, - &event_fd, - &plane_sizes, - &mut local_buf, - &*evt_queue, - ) - } + } => handle_streaming( + fifo_file, + buffer_idx, + plane_idx, + plane_offset, + &session_state, + &rx, + &event_fd, + &plane_sizes, + &mut local_buf, + &*evt_queue, + ), WorkerState::Stopped => WorkerState::Stopped, }; } - + log::info!("Worker thread stopped"); } diff --git a/base/cvd/cuttlefish/host/commands/vhost_user_media/vhu_media/src/lib.rs b/base/cvd/cuttlefish/host/commands/vhost_user_media/vhu_media/src/lib.rs index d476c4321ac..45ecb36fe3d 100644 --- a/base/cvd/cuttlefish/host/commands/vhost_user_media/vhu_media/src/lib.rs +++ b/base/cvd/cuttlefish/host/commands/vhost_user_media/vhu_media/src/lib.rs @@ -157,7 +157,7 @@ impl VirtioMediaHostMemoryMapper for HostMemoryMapper { padding: [0, 0, 0, 0, 0, 0, 0], fd_offset: 0, shm_offset: shm_offset, - len: 1, + len: end - shm_offset + 1, flags: 0, };