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| {
format!(
"could not enumerate USB event cameras ({error}). On Linux this is usually missing \
udev rules for the device (see the neuromorphic-drivers README, then re-plug it); on \
Windows, a missing WinUSB driver for it"
)
})?;
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, Eq)]
pub enum RoiPlacement {
Hardware,
Host,
}
#[derive(Default)]
struct LimitsPlan {
configuration: Option<nd::Configuration>,
host_roi: Option<(usize, usize, usize, usize)>,
}
#[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>,
mask: Option<Vec<bool>>,
roi: Option<((usize, usize, usize, usize), RoiPlacement)>,
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, plan) = if limits.is_empty() {
let selector = match serial {
Some(serial) => nd::SerialOrBusNumberAndAddress::Serial(serial),
None => nd::SerialOrBusNumberAndAddress::None,
};
(selector, LimitsPlan::default())
} else {
let listed = select_device(serial)?;
let plan = plan_limits(&listed, limits)?;
let address = (listed.bus_number, listed.address);
(
nd::SerialOrBusNumberAndAddress::BusNumberAndAddress(address),
plan,
)
};
let device = nd::open(
selector,
plan.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()?;
let mut capture = Self {
device,
adapter,
_event_loop: event_loop,
flag,
windower: Windower::new(width, height, window),
bias,
mask: None,
roi: limits.roi.map(|rect| {
let placement = match plan.host_roi {
Some(_) => RoiPlacement::Host,
None => RoiPlacement::Hardware,
};
(rect, placement)
}),
width,
height,
name,
serial,
};
if let Some((x0, y0, w, h)) = plan.host_roi {
let rect = crate::mask::rect(width, height, x0 as f64, y0 as f64, w as f64, h as f64);
capture.set_mask(Some(rect))?;
}
Ok(capture)
}
pub fn roi(&self) -> Option<((usize, usize, usize, usize), RoiPlacement)> {
self.roi
}
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 mask(&self) -> Option<&[bool]> {
self.mask.as_deref()
}
pub fn set_mask(&mut self, mask: Option<Vec<bool>>) -> Result<(), String> {
if let Some(mask) = &mask {
if mask.len() != self.width * self.height {
return Err(format!(
"mask has {} pixels, expected {} for this {}x{} sensor",
mask.len(),
self.width * self.height,
self.width,
self.height
));
}
}
self.mask = mask;
Ok(())
}
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,
self.mask.as_deref(),
self.width,
|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,
self.mask.as_deref(),
self.width,
|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,
self.mask.as_deref(),
self.width,
|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,
self.mask.as_deref(),
self.width,
|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 select_device(serial: Option<&str>) -> Result<nd::devices::ListedDevice, String> {
let listed = nd::list_devices().map_err(|error| error.to_string())?;
match serial {
Some(serial) => listed
.into_iter()
.find(|device| device.serial.as_deref().ok() == Some(serial))
.ok_or_else(|| format!("no camera with serial {serial} found")),
None => listed
.into_iter()
.next()
.ok_or_else(|| "no event camera found".to_owned()),
}
}
fn plan_limits(device: &nd::devices::ListedDevice, limits: Limits) -> Result<LimitsPlan, String> {
macro_rules! regions {
($module:ident, $variant:ident) => {{
use nd::devices::$module::{RateLimiter, DEFAULT_CONFIGURATION, PROPERTIES};
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;
}
LimitsPlan {
configuration: Some(nd::Configuration::$variant(configuration)),
host_roi: None,
}
}};
}
Ok(match device.device_type {
nd::Type::PropheseeEvk4 => regions!(prophesee_evk4, PropheseeEvk4),
nd::Type::PropheseeEvk3Hd => regions!(prophesee_evk3_hd, PropheseeEvk3Hd),
nd::Type::CenturyarksVga => regions!(centuryarks_vga, CenturyarksVga),
other => {
if limits.max_event_rate.is_some() {
return Err(format!(
"{} has no on-chip event-rate controller — max_event_rate is a sensor feature, \
and capping the rate on the host would not save any of the work it exists to \
avoid. It is supported on the Prophesee EVK4 and EVK3 HD and the CenturyArks \
VGA",
other.name(),
));
}
LimitsPlan {
configuration: None,
host_roi: limits.roi,
}
}
})
}
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<const NX: usize, const NY: usize>(
roi: (usize, usize, usize, usize),
sensor_width: usize,
sensor_height: usize,
) -> Result<([u64; NX], [u64; NY]), 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; NX];
let mut y_mask = [0_u64; NY];
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],
mask: Option<&[bool]>,
width: usize,
mut on_event: impl FnMut(u16, u16, i64, bool),
) {
let handle = |event: nd::types::PolarityEvent<u64, u16, u16>| {
if let Some(mask) = mask {
if mask.get(event.y as usize * width + event.x as usize) != Some(&true) {
return;
}
}
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::<20, 12>((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::<20, 12>((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_fit_the_centuryarks_bitmaps() {
let (x_mask, y_mask) = region_masks::<10, 8>((320, 240, 320, 240), 640, 480).unwrap();
assert_eq!((x_mask.len(), y_mask.len()), (10, 8));
let bit = |mask: &[u64], index: usize| mask[index / 64] & (1 << (index % 64)) != 0;
assert!(bit(&x_mask, 319) && !bit(&x_mask, 320)); assert!(bit(&y_mask, 480 - 1 - 239) && !bit(&y_mask, 480 - 1 - 240)); assert!(!bit(&y_mask, 480));
}
#[test]
fn region_masks_reject_a_rectangle_off_the_sensor() {
assert!(region_masks::<20, 12>((0, 0, 1280, 720), 1280, 720).is_ok());
assert!(region_masks::<20, 12>((1, 0, 1280, 720), 1280, 720).is_err());
assert!(region_masks::<20, 12>((0, 0, 0, 100), 1280, 720).is_err());
assert!(region_masks::<10, 8>((0, 0, 1280, 720), 640, 480).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);
}
}