use std::collections::VecDeque;
use std::time::{Duration, Instant};
use neuromorphic_drivers as nd;
use neuromorphic_drivers::types::Polarity;
use neuromorphic_drivers::Adapter;
use crate::bias::{BiasConfig, BiasController, BiasOverrides, BiasState, BiasValues};
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+",
}
}
const RATE_LIMIT_PERIOD_US: u16 = 1000;
const RATE_LIMIT_MAX_PER_PERIOD: u64 = (1 << 22) - 1;
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct Limits {
pub max_event_rate: Option<u64>,
pub roi: Option<(usize, usize, usize, usize)>,
}
impl Limits {
fn is_empty(&self) -> bool {
*self == Self::default()
}
}
#[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::with_capacity(
self.width,
self.height,
TIMESTAMP_SCALE_MS,
self.builder.len(),
);
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();
}
}
}
struct AdaptiveBias {
controller: BiasController,
last_pass: Instant,
}
impl AdaptiveBias {
fn new(overrides: BiasOverrides, device: &nd::Device) -> Result<Self, String> {
let (base, start) = match device.current_configuration() {
nd::Configuration::InivationDavis346(configuration) => (
BiasConfig::davis346(),
davis346_biases(&configuration)?,
),
nd::Configuration::PropheseeEvk4(configuration) => (
BiasConfig::prophesee_evk4(),
evk4_biases(&configuration),
),
_ => {
return Err(format!(
"adaptive_bias is not supported on the {} — it is implemented for the \
iniVation DAVIS346 and the Prophesee EVK4 so far",
device.name(),
))
}
};
let config = overrides.apply(base);
config.validate()?;
Ok(Self {
controller: BiasController::new(config, start),
last_pass: Instant::now(),
})
}
}
pub struct Capture {
device: nd::Device,
adapter: Adapter,
_event_loop: std::sync::Arc<nd::UsbEventLoop>,
flag: nd::Flag<nd::Error, nd::UsbOverflow>,
windower: Windower,
bias: Option<AdaptiveBias>,
width: usize,
height: usize,
name: String,
serial: String,
}
impl Capture {
pub fn open(
serial: Option<&str>,
window: Window,
limits: Limits,
bias: Option<BiasOverrides>,
) -> 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 configuration = limits_configuration(serial, limits)?;
let device = nd::open(selector, configuration, 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();
let bias = bias
.map(|overrides| AdaptiveBias::new(overrides, &device))
.transpose()?;
Ok(Self {
device,
adapter,
_event_loop: event_loop,
flag,
windower: Windower::new(width, height, window),
bias,
width,
height,
name,
serial,
})
}
pub fn bias_state(&self) -> Option<BiasState> {
self.bias.as_ref().map(|bias| bias.controller.state())
}
fn tick_bias(&mut self, events: u64, undercounted: bool) {
let update = {
let Some(bias) = self.bias.as_mut() else {
return;
};
let now = Instant::now();
let elapsed = now.duration_since(bias.last_pass);
bias.last_pass = now;
if undercounted {
bias.controller.invalidate();
}
bias.controller.observe(events, elapsed)
};
let Some(values) = update else {
return;
};
let _ = self.apply_biases(values);
}
fn apply_biases(&mut self, values: BiasValues) -> Result<(), String> {
let updated = match self.device.current_configuration() {
nd::Configuration::InivationDavis346(configuration) => {
nd::Configuration::InivationDavis346(davis346_with_biases(&configuration, values)?)
}
nd::Configuration::PropheseeEvk4(configuration) => {
nd::Configuration::PropheseeEvk4(evk4_with_biases(&configuration, values))
}
_ => return Err("adaptive biasing is not implemented for this camera".to_owned()),
};
self.device
.update_configuration(updated)
.map_err(|error| error.to_string())
}
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())?;
let mut events = 0;
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;
convert_buffer(&mut self.adapter, slice, |x, y, t, polarity| {
events += 1;
windower.push(x, y, t, polarity)
});
let current_t = self.adapter.current_t() as i64;
self.windower.seal_until(current_t);
}
self.tick_bias(events, false);
Ok(self.windower.pending.pop_front())
}
pub fn decode_next(&mut self, timeout: Duration) -> Result<bool, String> {
self.flag.load_error().map_err(|error| error.to_string())?;
let mut events = 0;
let decoded = if let Some(view) = self.device.next_with_timeout(&timeout) {
if view.first_after_overflow {
self.windower.mark_overflow();
}
let slice = view.slice;
let windower = &mut self.windower;
convert_buffer(&mut self.adapter, slice, |x, y, t, polarity| {
events += 1;
windower.push(x, y, t, polarity)
});
true
} else {
false
};
let current_t = self.adapter.current_t() as i64;
self.windower.seal_until(current_t);
self.tick_bias(events, false);
Ok(decoded)
}
pub fn take_pending(&mut self) -> Option<CaptureWindow> {
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;
let mut events = 0;
while let Some(view) = self.device.next_with_timeout(&Duration::ZERO) {
if view.first_after_overflow {
overflow = true;
}
convert_buffer(&mut self.adapter, view.slice, |x, y, t, polarity| {
events += 1;
on_event(x, y, t, polarity)
});
}
self.tick_bias(events, false);
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 = Instant::now();
let mut overflow = false;
let mut decoding = true;
let mut events = 0;
let mut dropped = false;
while let Some(view) = self.device.next_with_timeout(&Duration::ZERO) {
if view.first_after_overflow {
overflow = true;
}
if !decoding {
dropped = true;
continue;
}
convert_buffer(&mut self.adapter, view.slice, |x, y, t, polarity| {
events += 1;
on_event(x, y, t, polarity)
});
if start.elapsed() >= budget {
decoding = false;
}
}
self.tick_bias(events, dropped);
Ok(overflow)
}
}
fn limits_configuration(
serial: Option<&str>,
limits: Limits,
) -> Result<Option<nd::Configuration>, String> {
if limits.is_empty() {
return Ok(None);
}
let listed = nd::list_devices().map_err(|error| error.to_string())?;
let device = match serial {
Some(serial) => listed
.iter()
.find(|device| device.serial.as_deref().ok() == Some(serial))
.ok_or_else(|| format!("no camera with serial {serial} found"))?,
None => listed
.first()
.ok_or_else(|| "no event camera found".to_owned())?,
};
macro_rules! prophesee {
($module:ident, $variant:ident) => {{
use nd::devices::$module::{DEFAULT_CONFIGURATION, PROPERTIES, RateLimiter};
let mut configuration = DEFAULT_CONFIGURATION;
configuration.rate_limiter = rate_limiter(limits.max_event_rate).map(
|(reference_period_us, maximum_events_per_period)| RateLimiter {
reference_period_us,
maximum_events_per_period,
},
);
if let Some(roi) = limits.roi {
let (x_mask, y_mask) =
region_masks(roi, PROPERTIES.width as usize, PROPERTIES.height as usize)?;
configuration.x_mask = x_mask;
configuration.y_mask = y_mask;
configuration.mask_intersection_only = false;
}
nd::Configuration::$variant(configuration)
}};
}
Ok(Some(match device.device_type {
nd::Type::PropheseeEvk4 => prophesee!(prophesee_evk4, PropheseeEvk4),
nd::Type::PropheseeEvk3Hd => prophesee!(prophesee_evk3_hd, PropheseeEvk3Hd),
other => {
return Err(format!(
"{} has no on-chip event-rate controller or region masks — max_event_rate and roi \
are supported on Prophesee sensors (EVK4, EVK3 HD) only",
other.name(),
))
}
}))
}
fn rate_limiter(max_event_rate: Option<u64>) -> Option<(u16, u32)> {
let rate = max_event_rate?;
let per_period = (rate * RATE_LIMIT_PERIOD_US as u64 / 1_000_000)
.clamp(1, RATE_LIMIT_MAX_PER_PERIOD) as u32;
Some((RATE_LIMIT_PERIOD_US, per_period))
}
fn region_masks(
roi: (usize, usize, usize, usize),
sensor_width: usize,
sensor_height: usize,
) -> Result<([u64; 20], [u64; 12]), String> {
let (x0, y0, width, height) = roi;
if width == 0 || height == 0 {
return Err("roi width and height must be at least 1".to_owned());
}
let (x1, y1) = (x0 + width, y0 + height);
if x1 > sensor_width || y1 > sensor_height {
return Err(format!(
"roi ({x0}, {y0}, {width}, {height}) reaches ({x1}, {y1}), outside the \
{sensor_width}x{sensor_height} sensor"
));
}
let mut x_mask = [0_u64; 20];
let mut y_mask = [0_u64; 12];
for column in (0..sensor_width).filter(|column| !(x0..x1).contains(column)) {
x_mask[column / 64] |= 1 << (column % 64);
}
for row in (0..sensor_height).filter(|row| !(y0..y1).contains(row)) {
let bit = sensor_height - 1 - row;
y_mask[bit / 64] |= 1 << (bit % 64);
}
Ok((x_mask, y_mask))
}
const DAVIS346_BIASES: [&str; 5] = ["refrbp", "prbp", "prsfbp", "onbn", "offbn"];
fn davis346_biases(
configuration: &nd::devices::inivation_davis346::Configuration,
) -> Result<BiasValues, String> {
let value = serde_json::to_value(&configuration.biases)
.map_err(|error| format!("could not read the camera's biases ({error})"))?;
let read = |field: &str| {
value
.get(field)
.and_then(serde_json::Value::as_u64)
.and_then(|value| u16::try_from(value).ok())
.ok_or_else(|| format!("the camera's biases have no readable {field}"))
};
Ok(BiasValues {
refractory: read(DAVIS346_BIASES[0])?,
photoreceptor: read(DAVIS346_BIASES[1])?,
follower: read(DAVIS346_BIASES[2])?,
on_threshold: read(DAVIS346_BIASES[3])?,
off_threshold: read(DAVIS346_BIASES[4])?,
})
}
fn davis346_with_biases(
configuration: &nd::devices::inivation_davis346::Configuration,
values: BiasValues,
) -> Result<nd::devices::inivation_davis346::Configuration, String> {
let mut biases = serde_json::to_value(&configuration.biases)
.map_err(|error| format!("could not read the camera's biases ({error})"))?;
let written = [
values.refractory,
values.photoreceptor,
values.follower,
values.on_threshold,
values.off_threshold,
];
for (field, value) in DAVIS346_BIASES.iter().zip(written) {
biases[field] = serde_json::Value::from(value);
}
let mut configuration = configuration.clone();
configuration.biases = serde_json::from_value(biases)
.map_err(|error| format!("could not apply the new biases ({error})"))?;
Ok(configuration)
}
fn evk4_biases(configuration: &nd::devices::prophesee_evk4::Configuration) -> BiasValues {
let biases = &configuration.biases;
BiasValues {
refractory: biases.refr.into(),
photoreceptor: biases.pr.into(),
follower: biases.fo.into(),
on_threshold: biases.diff_on.into(),
off_threshold: biases.diff_off.into(),
}
}
fn evk4_with_biases(
configuration: &nd::devices::prophesee_evk4::Configuration,
values: BiasValues,
) -> nd::devices::prophesee_evk4::Configuration {
let mut configuration = configuration.clone();
let byte = |value: u16| value.min(u16::from(u8::MAX)) as u8;
configuration.biases.refr = byte(values.refractory);
configuration.biases.pr = byte(values.photoreceptor);
configuration.biases.fo = byte(values.follower);
configuration.biases.diff_on = byte(values.on_threshold);
configuration.biases.diff_off = byte(values.off_threshold);
configuration
}
fn convert_buffer(
adapter: &mut Adapter,
slice: &[u8],
mut on_event: impl FnMut(u16, u16, i64, bool),
) {
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 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| {})
}
}
}
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::{rate_limiter, region_masks, Window, Windower};
#[test]
fn rate_limiter_converts_events_per_second_to_a_period_budget() {
assert_eq!(rate_limiter(Some(50_000_000)), Some((1000, 50_000)));
assert_eq!(rate_limiter(None), None);
assert_eq!(rate_limiter(Some(100)), Some((1000, 1)));
assert_eq!(rate_limiter(Some(u64::MAX / 1000)), Some((1000, (1 << 22) - 1)));
}
#[test]
fn region_masks_block_everything_outside_the_rectangle() {
let (x_mask, y_mask) = region_masks((100, 50, 200, 100), 1280, 720).unwrap();
let bit = |mask: &[u64], index: usize| mask[index / 64] & (1 << (index % 64)) != 0;
let column_blocked = |column: usize| bit(&x_mask, column);
let row_blocked = |row: usize| bit(&y_mask, 720 - 1 - row);
assert!(column_blocked(99) && column_blocked(300)); assert!(!column_blocked(100) && !column_blocked(299)); assert!(row_blocked(49) && row_blocked(150));
assert!(!row_blocked(50) && !row_blocked(149));
assert!(!bit(&y_mask, 720));
}
#[test]
fn region_masks_store_rows_bottom_up() {
let (_, y_mask) = region_masks((0, 0, 200, 200), 1280, 720).unwrap();
let bit = |index: usize| y_mask[index / 64] & (1 << (index % 64)) != 0;
assert!(bit(0) && bit(519), "rows 200..719 must be blocked, low bits set");
assert!(!bit(520) && !bit(719), "rows 0..199 must be kept, high bits clear");
}
#[test]
fn region_masks_reject_a_rectangle_off_the_sensor() {
assert!(region_masks((0, 0, 1280, 720), 1280, 720).is_ok());
assert!(region_masks((1, 0, 1280, 720), 1280, 720).is_err());
assert!(region_masks((0, 0, 0, 100), 1280, 720).is_err());
}
#[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);
}
}