use std::collections::VecDeque;
use std::time::Duration;
use neuromorphic_drivers as nd;
use neuromorphic_drivers::types::Polarity;
use neuromorphic_drivers::Adapter;
use crate::{EventStream, EventStreamBuilder};
const TIMESTAMP_SCALE_MS: f64 = 0.001;
const POLL_TIMEOUT: Duration = Duration::from_millis(5);
#[derive(Clone, Debug)]
pub struct CameraInfo {
pub kind: String,
pub name: String,
pub serial: Option<String>,
pub bus_number: u8,
pub address: u8,
pub speed: String,
}
pub fn list_cameras() -> Result<Vec<CameraInfo>, String> {
let listed = nd::list_devices().map_err(|error| error.to_string())?;
Ok(listed
.into_iter()
.map(|device| CameraInfo {
kind: device.device_type.to_string(),
name: device.device_type.name().to_owned(),
serial: device.serial.ok(),
bus_number: device.bus_number,
address: device.address,
speed: speed_name(device.speed).to_owned(),
})
.collect())
}
fn open_error(serial: Option<&str>, error: impl std::fmt::Display) -> String {
let visible = list_cameras().unwrap_or_default();
let matched = match serial {
Some(serial) => visible.iter().find(|camera| camera.serial.as_deref() == Some(serial)),
None => visible.first(),
};
match matched {
Some(camera) => format!(
"found {} but couldn't open it ({error}). It may already be open in this or another \
process — an earlier eventcv.stream()/EventCamera that is still alive holds the device \
until it is closed (use `with eventcv.stream() as cam:` or call `cam.close()`); or, on \
Linux, the USB udev rules may be missing. Only one handle can open a camera at a time.",
camera.name,
),
None => match serial {
Some(serial) => format!("no camera with serial {serial} found ({error})"),
None => format!("no event camera found ({error})"),
},
}
}
fn speed_name(speed: nd::usb::Speed) -> &'static str {
match speed {
nd::usb::Speed::Unknown => "unknown",
nd::usb::Speed::Low => "low",
nd::usb::Speed::Full => "full",
nd::usb::Speed::High => "high",
nd::usb::Speed::Super => "super",
nd::usb::Speed::SuperPlus => "super+",
}
}
#[derive(Clone, Copy, Debug, PartialEq)]
pub enum Window {
Duration(f64),
Count(usize),
}
#[derive(Clone, Debug)]
pub struct CaptureWindow {
pub stream: EventStream,
pub first_after_overflow: bool,
}
struct Windower {
width: usize,
height: usize,
window: Window,
dt_us: i64,
builder: EventStreamBuilder,
next_boundary_t: Option<i64>,
pending: VecDeque<CaptureWindow>,
overflow_pending: bool,
}
impl Windower {
fn new(width: usize, height: usize, window: Window) -> Self {
let dt_us = match window {
Window::Duration(dt_ms) => ((dt_ms * 1000.0).max(1.0)) as i64,
Window::Count(_) => 0,
};
Self {
width,
height,
window,
dt_us,
builder: EventStreamBuilder::new(width, height, TIMESTAMP_SCALE_MS),
next_boundary_t: None,
pending: VecDeque::new(),
overflow_pending: false,
}
}
fn mark_overflow(&mut self) {
self.overflow_pending = true;
}
fn push(&mut self, x: u16, y: u16, t: i64, polarity: bool) {
match self.window {
Window::Duration(_) => {
match self.next_boundary_t {
None => self.next_boundary_t = Some(t + self.dt_us),
Some(_) => self.advance_to(t),
}
self.builder.push(x, y, t, polarity);
}
Window::Count(n) => {
self.builder.push(x, y, t, polarity);
if self.builder.len() >= n {
self.seal();
}
}
}
}
fn advance_to(&mut self, t: i64) {
let Some(mut boundary) = self.next_boundary_t else {
return;
};
while t >= boundary {
if self.builder.is_empty() {
let steps = (t - boundary) / self.dt_us + 1;
boundary += steps * self.dt_us;
} else {
self.seal();
boundary += self.dt_us;
}
}
self.next_boundary_t = Some(boundary);
}
fn seal_until(&mut self, t_now: i64) {
self.advance_to(t_now);
}
fn seal(&mut self) {
let fresh = EventStreamBuilder::new(self.width, self.height, TIMESTAMP_SCALE_MS);
let builder = std::mem::replace(&mut self.builder, fresh);
self.pending.push_back(CaptureWindow {
stream: builder.build(),
first_after_overflow: self.overflow_pending,
});
self.overflow_pending = false;
}
fn flush(&mut self) {
if !self.builder.is_empty() {
self.seal();
}
}
}
pub struct Capture {
device: nd::Device,
adapter: Adapter,
_event_loop: std::sync::Arc<nd::UsbEventLoop>,
flag: nd::Flag<nd::Error, nd::UsbOverflow>,
windower: Windower,
width: usize,
height: usize,
name: String,
serial: String,
}
impl Capture {
pub fn open(serial: Option<&str>, window: Window) -> Result<Self, String> {
let (flag, event_loop) = nd::flag_and_event_loop().map_err(|error| error.to_string())?;
let selector = match serial {
Some(serial) => nd::SerialOrBusNumberAndAddress::Serial(serial),
None => nd::SerialOrBusNumberAndAddress::None,
};
let device = nd::open(selector, None, None, event_loop.clone(), flag.clone())
.map_err(|error| open_error(serial, error))?;
let (width, height) = device_dimensions(&device);
let name = device.name().to_owned();
let serial = device.serial();
let adapter = device.create_adapter();
Ok(Self {
device,
adapter,
_event_loop: event_loop,
flag,
windower: Windower::new(width, height, window),
width,
height,
name,
serial,
})
}
pub fn width(&self) -> usize {
self.width
}
pub fn height(&self) -> usize {
self.height
}
pub fn name(&self) -> &str {
&self.name
}
pub fn serial(&self) -> &str {
&self.serial
}
pub fn backlog(&self) -> usize {
self.device.backlog()
}
pub fn poll(&mut self) -> Result<Option<CaptureWindow>, String> {
if let Some(window) = self.windower.pending.pop_front() {
return Ok(Some(window));
}
self.flag.load_error().map_err(|error| error.to_string())?;
if let Some(view) = self.device.next_with_timeout(&POLL_TIMEOUT) {
if view.first_after_overflow {
self.windower.mark_overflow();
}
let slice = view.slice;
let windower = &mut self.windower;
let on_polarity = |event: nd::types::PolarityEvent<u64, u16, u16>| {
windower.push(
event.x,
event.y,
event.t as i64,
matches!(event.polarity, Polarity::On),
);
};
match &mut self.adapter {
Adapter::Evt3(adapter) => adapter.convert(slice, on_polarity, |_trigger| {}),
Adapter::Dvxplorer(adapter) => {
adapter.convert(slice, on_polarity, |_imu| {}, |_trigger| {})
}
Adapter::Davis346(adapter) => {
adapter.convert(slice, on_polarity, |_imu| {}, |_trigger| {}, |_frame| {})
}
}
let current_t = self.adapter.current_t() as i64;
self.windower.seal_until(current_t);
}
Ok(self.windower.pending.pop_front())
}
pub fn finish(&mut self) -> Option<CaptureWindow> {
self.windower.flush();
self.windower.pending.pop_front()
}
pub fn drain_events<F>(&mut self, mut on_event: F) -> Result<bool, String>
where
F: FnMut(u16, u16, i64, bool),
{
self.flag.load_error().map_err(|error| error.to_string())?;
let mut overflow = false;
while let Some(view) = self.device.next_with_timeout(&Duration::ZERO) {
if view.first_after_overflow {
overflow = true;
}
let slice = view.slice;
let handle = |event: nd::types::PolarityEvent<u64, u16, u16>| {
on_event(
event.x,
event.y,
event.t as i64,
matches!(event.polarity, Polarity::On),
);
};
match &mut self.adapter {
Adapter::Evt3(adapter) => adapter.convert(slice, handle, |_trigger| {}),
Adapter::Dvxplorer(adapter) => {
adapter.convert(slice, handle, |_imu| {}, |_trigger| {})
}
Adapter::Davis346(adapter) => {
adapter.convert(slice, handle, |_imu| {}, |_trigger| {}, |_frame| {})
}
}
}
Ok(overflow)
}
pub fn drain_events_budgeted<F>(
&mut self,
budget: Duration,
mut on_event: F,
) -> Result<bool, String>
where
F: FnMut(u16, u16, i64, bool),
{
self.flag.load_error().map_err(|error| error.to_string())?;
let start = std::time::Instant::now();
let mut overflow = false;
let mut decoding = true;
while let Some(view) = self.device.next_with_timeout(&Duration::ZERO) {
if view.first_after_overflow {
overflow = true;
}
if !decoding {
continue; }
let slice = view.slice;
let handle = |event: nd::types::PolarityEvent<u64, u16, u16>| {
on_event(
event.x,
event.y,
event.t as i64,
matches!(event.polarity, Polarity::On),
);
};
match &mut self.adapter {
Adapter::Evt3(adapter) => adapter.convert(slice, handle, |_trigger| {}),
Adapter::Dvxplorer(adapter) => {
adapter.convert(slice, handle, |_imu| {}, |_trigger| {})
}
Adapter::Davis346(adapter) => {
adapter.convert(slice, handle, |_imu| {}, |_trigger| {}, |_frame| {})
}
}
if start.elapsed() >= budget {
decoding = false;
}
}
Ok(overflow)
}
}
fn device_dimensions(device: &nd::Device) -> (usize, usize) {
use nd::Properties;
let (width, height) = match device.properties() {
Properties::InivationDavis346(properties) => (properties.width, properties.height),
Properties::InivationDvxplorer(properties) => (properties.width, properties.height),
Properties::PropheseeEvk3Hd(properties) => (properties.width, properties.height),
Properties::PropheseeEvk4(properties) => (properties.width, properties.height),
Properties::CenturyarksVga(properties) => (properties.width, properties.height),
};
(width as usize, height as usize)
}
#[cfg(test)]
mod tests {
use super::{Window, Windower};
#[test]
fn count_windows_split_by_event_count() {
let mut windower = Windower::new(10, 10, Window::Count(2));
for t in 0..5 {
windower.push(1, 1, t, true);
}
assert_eq!(windower.pending.len(), 2);
assert_eq!(windower.pending[0].stream.len(), 2);
assert_eq!(windower.pending[1].stream.len(), 2);
windower.flush();
assert_eq!(windower.pending.len(), 3);
assert_eq!(windower.pending[2].stream.len(), 1);
}
#[test]
fn duration_windows_split_by_event_time() {
let mut windower = Windower::new(10, 10, Window::Duration(1.0));
windower.push(1, 1, 0, true); windower.push(1, 1, 500, true); windower.push(1, 1, 1500, true); windower.push(1, 1, 2500, true);
assert_eq!(windower.pending.len(), 2);
assert_eq!(windower.pending[0].stream.len(), 2);
assert_eq!(windower.pending[1].stream.len(), 1);
windower.flush(); assert_eq!(windower.pending.len(), 3);
assert_eq!(windower.pending[2].stream.len(), 1);
}
#[test]
fn idle_gap_emits_one_window_not_many_empties() {
let mut windower = Windower::new(10, 10, Window::Duration(1.0));
windower.push(1, 1, 0, true); windower.push(1, 1, 10_000, true);
assert_eq!(windower.pending.len(), 1);
assert_eq!(windower.pending[0].stream.len(), 1);
}
#[test]
fn seal_until_flushes_a_stalled_window() {
let mut windower = Windower::new(10, 10, Window::Duration(1.0));
windower.push(1, 1, 100, true); assert_eq!(windower.pending.len(), 0);
windower.seal_until(2000); assert_eq!(windower.pending.len(), 1);
assert_eq!(windower.pending[0].stream.len(), 1);
}
#[test]
fn overflow_flag_rides_on_the_next_window() {
let mut windower = Windower::new(10, 10, Window::Count(1));
windower.mark_overflow();
windower.push(1, 1, 0, true); windower.push(2, 2, 1, true);
assert!(windower.pending[0].first_after_overflow);
assert!(!windower.pending[1].first_after_overflow);
}
}