1use std::convert::TryInto;
2use std::time::Duration;
3use std::{io, mem, sync::Arc};
4
5use crate::buffer::{Metadata, Type};
6use crate::device::{Device, Handle};
7use crate::io::mmap::arena::Arena;
8use crate::io::traits::{CaptureStream, OutputStream, Stream as StreamTrait};
9use crate::memory::Memory;
10use crate::v4l2;
11use crate::v4l_sys::*;
12
13pub struct Stream<'a> {
17 handle: Arc<Handle>,
18 arena: Arena<'a>,
19 arena_index: usize,
20 buf_type: Type,
21 buf_meta: Vec<Metadata>,
22 timeout: Option<i32>,
23
24 active: bool,
25}
26
27impl<'a> Stream<'a> {
28 pub fn new(dev: &Device, buf_type: Type) -> io::Result<Self> {
48 Stream::with_buffers(dev, buf_type, 4)
49 }
50
51 pub fn with_buffers(dev: &Device, buf_type: Type, buf_count: u32) -> io::Result<Self> {
52 let mut arena = Arena::new(dev.handle(), buf_type);
53 let count = arena.allocate(buf_count)?;
54 let mut buf_meta = Vec::new();
55 buf_meta.resize(count as usize, Metadata::default());
56
57 Ok(Stream {
58 handle: dev.handle(),
59 arena,
60 arena_index: 0,
61 buf_type,
62 buf_meta,
63 active: false,
64 timeout: None,
65 })
66 }
67
68 pub fn handle(&self) -> Arc<Handle> {
70 self.handle.clone()
71 }
72
73 pub fn set_timeout(&mut self, duration: Duration) {
75 self.timeout = Some(duration.as_millis().try_into().unwrap());
76 }
77
78 pub fn clear_timeout(&mut self) {
80 self.timeout = None;
81 }
82
83 fn buffer_desc(&self) -> v4l2_buffer {
84 v4l2_buffer {
85 type_: self.buf_type as u32,
86 memory: Memory::Mmap as u32,
87 ..unsafe { mem::zeroed() }
88 }
89 }
90}
91
92impl<'a> Drop for Stream<'a> {
93 fn drop(&mut self) {
94 if let Err(e) = self.stop() {
95 if let Some(code) = e.raw_os_error() {
96 if code == 19 {
100 return;
102 }
103 }
104
105 panic!("{:?}", e)
106 }
107 }
108}
109
110impl<'a> StreamTrait for Stream<'a> {
111 type Item = [u8];
112
113 fn start(&mut self) -> io::Result<()> {
114 unsafe {
115 let mut typ = self.buf_type as u32;
116 v4l2::ioctl(
117 self.handle.fd(),
118 v4l2::vidioc::VIDIOC_STREAMON,
119 &mut typ as *mut _ as *mut std::os::raw::c_void,
120 )?;
121 }
122
123 self.active = true;
124 Ok(())
125 }
126
127 fn stop(&mut self) -> io::Result<()> {
128 unsafe {
129 let mut typ = self.buf_type as u32;
130 v4l2::ioctl(
131 self.handle.fd(),
132 v4l2::vidioc::VIDIOC_STREAMOFF,
133 &mut typ as *mut _ as *mut std::os::raw::c_void,
134 )?;
135 }
136
137 self.active = false;
138 Ok(())
139 }
140}
141
142impl<'a, 'b> CaptureStream<'b> for Stream<'a> {
143 fn queue(&mut self, index: usize) -> io::Result<()> {
144 let mut v4l2_buf = v4l2_buffer {
145 index: index as u32,
146 ..self.buffer_desc()
147 };
148
149 unsafe {
150 v4l2::ioctl(
151 self.handle.fd(),
152 v4l2::vidioc::VIDIOC_QBUF,
153 &mut v4l2_buf as *mut _ as *mut std::os::raw::c_void,
154 )?;
155 }
156
157 Ok(())
158 }
159
160 fn dequeue(&mut self) -> io::Result<usize> {
161 let mut v4l2_buf = self.buffer_desc();
162
163 if self.handle.poll(libc::POLLIN, self.timeout.unwrap_or(-1))? == 0 {
164 return Err(io::Error::new(io::ErrorKind::TimedOut, "VIDIOC_DQBUF"));
168 }
169
170 unsafe {
171 v4l2::ioctl(
172 self.handle.fd(),
173 v4l2::vidioc::VIDIOC_DQBUF,
174 &mut v4l2_buf as *mut _ as *mut std::os::raw::c_void,
175 )?;
176 }
177 self.arena_index = v4l2_buf.index as usize;
178
179 self.buf_meta[self.arena_index] = Metadata {
180 bytesused: v4l2_buf.bytesused,
181 flags: v4l2_buf.flags.into(),
182 field: v4l2_buf.field,
183 timestamp: v4l2_buf.timestamp.into(),
184 sequence: v4l2_buf.sequence,
185 };
186
187 Ok(self.arena_index)
188 }
189
190 fn next(&'b mut self) -> io::Result<(&Self::Item, &Metadata)> {
191 if !self.active {
192 for index in 0..self.arena.bufs.len() {
194 CaptureStream::queue(self, index)?;
195 }
196
197 self.start()?;
198 } else {
199 CaptureStream::queue(self, self.arena_index)?;
200 }
201
202 self.arena_index = CaptureStream::dequeue(self)?;
203
204 let bytes = &self.arena.bufs[self.arena_index];
207 let meta = &self.buf_meta[self.arena_index];
208 Ok((bytes, meta))
209 }
210}
211
212impl<'a, 'b> OutputStream<'b> for Stream<'a> {
213 fn queue(&mut self, index: usize) -> io::Result<()> {
214 let mut v4l2_buf = v4l2_buffer {
215 index: index as u32,
216 ..self.buffer_desc()
217 };
218 unsafe {
219 v4l2_buf.bytesused = self.buf_meta[index].bytesused;
225 v4l2_buf.field = self.buf_meta[index].field;
226
227 if self
228 .handle
229 .poll(libc::POLLOUT, self.timeout.unwrap_or(-1))?
230 == 0
231 {
232 return Err(io::Error::new(io::ErrorKind::TimedOut, "VIDIOC_QBUF"));
236 }
237
238 v4l2::ioctl(
239 self.handle.fd(),
240 v4l2::vidioc::VIDIOC_QBUF,
241 &mut v4l2_buf as *mut _ as *mut std::os::raw::c_void,
242 )
243 }
244 }
245
246 fn dequeue(&mut self) -> io::Result<usize> {
247 let mut v4l2_buf = self.buffer_desc();
248
249 unsafe {
250 v4l2::ioctl(
251 self.handle.fd(),
252 v4l2::vidioc::VIDIOC_DQBUF,
253 &mut v4l2_buf as *mut _ as *mut std::os::raw::c_void,
254 )?;
255 }
256 self.arena_index = v4l2_buf.index as usize;
257
258 self.buf_meta[self.arena_index] = Metadata {
259 bytesused: v4l2_buf.bytesused,
260 flags: v4l2_buf.flags.into(),
261 field: v4l2_buf.field,
262 timestamp: v4l2_buf.timestamp.into(),
263 sequence: v4l2_buf.sequence,
264 };
265
266 Ok(self.arena_index)
267 }
268
269 fn next(&'b mut self) -> io::Result<(&mut Self::Item, &mut Metadata)> {
270 let init = !self.active;
271 if !self.active {
272 self.start()?;
273 }
274
275 if !init {
279 OutputStream::queue(self, self.arena_index)?;
280 self.arena_index = OutputStream::dequeue(self)?;
281 }
282
283 let bytes = &mut self.arena.bufs[self.arena_index];
286 let meta = &mut self.buf_meta[self.arena_index];
287 Ok((bytes, meta))
288 }
289}