From 131c4637b887ab5f264a797962aca433f43c5398 Mon Sep 17 00:00:00 2001 From: Changyeon Jo Date: Tue, 1 Sep 2026 15:37:13 +0000 Subject: [PATCH 1/4] vhost_user_media: Run cargo fmt on baseline v4l2_stream_proxy --- .../v4l2_stream_proxy/src/device.rs | 74 +++++---- .../v4l2_stream_proxy/src/main.rs | 7 +- .../v4l2_stream_proxy/src/worker.rs | 142 ++++++++++++------ 3 files changed, 141 insertions(+), 82 deletions(-) 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..b0a0cfd3d47 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; @@ -94,11 +93,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 +174,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 +191,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 +269,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 +423,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 +547,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 +570,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 +596,7 @@ where panic!() } } - + v4l2_buffer.set_field(BufferField::None); v4l2_buffer.set_flags(BufferFlags::TIMESTAMP_MONOTONIC); @@ -638,7 +638,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 +647,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 +691,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(()) } 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..51d7c487080 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 @@ -78,8 +78,9 @@ 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)) + })?; Ok(Config { socket_path: args.socket_path, input_path: args.input_path, @@ -106,7 +107,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 { 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..a9620cec3e3 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,39 @@ 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) { + if poll_fds[1] + .revents() + .unwrap_or(PollFlags::empty()) + .contains(PollFlags::POLLIN) + { 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 +219,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 +238,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 +246,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 +275,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 +321,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 +339,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"); } From 0f0f3f7fb189f6ebb4722cf5ab88730a4aaf5efc Mon Sep 17 00:00:00 2001 From: Changyeon Jo Date: Tue, 1 Sep 2026 17:13:45 +0000 Subject: [PATCH 2/4] v4l2_stream_proxy: Implement V4L2 CID_LENS_FACING control Add LensFacing enum and --lens-facing CLI option to v4l2_stream_proxy, aligning it with emulated_camera_mplane and emulated_camera_splane. Implement V4L2 CID_LENS_FACING control in query_ext_ctrl, g_ctrl, and g_ext_ctrls returning the configured lens facing. Add unit tests for CLI argument parsing, lens facing string conversion, and control name padding. --- .../v4l2_stream_proxy/src/device.rs | 134 ++++++++++++++++++ .../v4l2_stream_proxy/src/main.rs | 58 ++++++++ 2 files changed, 192 insertions(+) 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 b0a0cfd3d47..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 @@ -30,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; @@ -809,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 51d7c487080..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,6 +75,7 @@ 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 { @@ -81,6 +85,10 @@ impl TryFrom for Config { 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, @@ -88,6 +96,7 @@ impl TryFrom for Config { input_height: args.input_height, fps_interval, format: device::Format::Yuv420M, + lens_facing, }) } } @@ -134,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(); From 1c18176614f801485d8121105e098808a04b3f8d Mon Sep 17 00:00:00 2001 From: Changyeon Jo Date: Tue, 1 Sep 2026 17:13:47 +0000 Subject: [PATCH 3/4] vhu_media: Correct shmem_unmap message length calculation The previous calculation passed len: 1 to shmem_unmap. Update it to unmap the entire allocated range (end - shm_offset + 1). --- .../host/commands/vhost_user_media/vhu_media/src/lib.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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, }; From 1da09336951bb4dae27a55dcc5f855d64fe301bb Mon Sep 17 00:00:00 2001 From: Changyeon Jo Date: Tue, 1 Sep 2026 17:13:49 +0000 Subject: [PATCH 4/4] v4l2_stream_proxy: Handle POLLHUP alongside POLLIN in worker loop When a FIFO writer disconnects and the buffer is drained, Linux poll() reports only POLLHUP (without POLLIN). Handling POLLHUP ensures read() is invoked, which consumes any remaining bytes and returns Ok(0) (EOF). This cleanly transitions the worker to WorkerState::Unopened rather than spinning in a 100% CPU poll busy-loop. --- .../vhost_user_media/v4l2_stream_proxy/src/worker.rs | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) 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 a9620cec3e3..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 @@ -200,10 +200,14 @@ fn handle_streaming( } } + // 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()) - .contains(PollFlags::POLLIN) + .intersects(PollFlags::POLLIN | PollFlags::POLLHUP) { let plane_size = plane_sizes[plane_idx]; let remaining = plane_size - plane_offset;