use std::time::Instant;
use ad_core_rs::ndarray::{NDArray, NDDataBuffer, NDDataType};
use ad_core_rs::ndarray_pool::NDArrayPool;
use ad_core_rs::plugin::runtime::{
NDPluginProcess, ParamChangeResult, ParamUpdate, PluginParamSnapshot, ProcessResult,
};
use asyn_rs::param::ParamType;
use asyn_rs::port::PortDriverBase;
const DEFAULT_NUM_TSPOINTS: usize = 2048;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum AcquireMode {
Fixed,
Circular,
}
struct Params {
ts_acquire: usize,
ts_read: usize,
ts_num_points: usize,
ts_current_point: usize,
ts_time_per_point: usize,
ts_averaging_time: usize,
ts_num_average: usize,
ts_elapsed_time: usize,
ts_acquire_mode: usize,
ts_time_axis: usize,
ts_timestamp: usize,
ts_time_series: usize,
}
impl Params {
const fn sentinel() -> Self {
Self {
ts_acquire: usize::MAX,
ts_read: usize::MAX,
ts_num_points: usize::MAX,
ts_current_point: usize::MAX,
ts_time_per_point: usize::MAX,
ts_averaging_time: usize::MAX,
ts_num_average: usize::MAX,
ts_elapsed_time: usize::MAX,
ts_acquire_mode: usize::MAX,
ts_time_axis: usize::MAX,
ts_timestamp: usize::MAX,
ts_time_series: usize::MAX,
}
}
}
#[inline]
fn sample_f64(data: &NDDataBuffer, idx: usize) -> f64 {
match data {
NDDataBuffer::I8(v) => v[idx] as f64,
NDDataBuffer::U8(v) => v[idx] as f64,
NDDataBuffer::I16(v) => v[idx] as f64,
NDDataBuffer::U16(v) => v[idx] as f64,
NDDataBuffer::I32(v) => v[idx] as f64,
NDDataBuffer::U32(v) => v[idx] as f64,
NDDataBuffer::I64(v) => v[idx] as f64,
NDDataBuffer::U64(v) => v[idx] as f64,
NDDataBuffer::F32(v) => v[idx] as f64,
NDDataBuffer::F64(v) => v[idx],
}
}
fn averaged_value(sum: f64, num_averaged: usize, dt: NDDataType) -> f64 {
let mean = sum / num_averaged.max(1) as f64;
match dt {
NDDataType::Float32 => mean as f32 as f64,
NDDataType::Float64 => mean,
_ => mean.trunc(),
}
}
fn coalesce_updates(updates: Vec<ParamUpdate>) -> Vec<ParamUpdate> {
fn key(u: &ParamUpdate) -> (usize, i32) {
match u {
ParamUpdate::Int32 { reason, addr, .. }
| ParamUpdate::Float64 { reason, addr, .. }
| ParamUpdate::Octet { reason, addr, .. }
| ParamUpdate::Float64Array { reason, addr, .. } => (*reason, *addr),
}
}
let mut seen = std::collections::HashSet::new();
let mut out = Vec::with_capacity(updates.len());
for u in updates.into_iter().rev() {
if seen.insert(key(&u)) {
out.push(u);
}
}
out.reverse();
out
}
pub struct TimeSeriesProcessor {
max_signals: usize,
num_signals_in: i64,
num_signals: usize,
data_type: NDDataType,
num_time_points: usize,
current_time_point: usize,
num_average: usize,
num_averaged: usize,
average_store: Vec<f64>,
time_per_point: f64,
averaging_time_requested: f64,
averaging_time_actual: f64,
acquire_mode: AcquireMode,
acquiring: bool,
circular: Vec<f64>,
time_stamp: Vec<f64>,
start_time: Instant,
p: Params,
}
impl TimeSeriesProcessor {
pub fn new(max_signals: usize) -> Self {
let max_signals = max_signals.max(1);
let num_time_points = DEFAULT_NUM_TSPOINTS;
Self {
max_signals,
num_signals_in: -1,
num_signals: max_signals,
data_type: NDDataType::Float64,
num_time_points,
current_time_point: 0,
num_average: 1,
num_averaged: 0,
average_store: vec![0.0; max_signals],
time_per_point: 0.0,
averaging_time_requested: 1.0,
averaging_time_actual: 1.0,
acquire_mode: AcquireMode::Fixed,
acquiring: false,
circular: vec![0.0; max_signals * num_time_points],
time_stamp: vec![0.0; num_time_points],
start_time: Instant::now(),
p: Params::sentinel(),
}
}
fn create_axis_array(&self) -> Vec<ParamUpdate> {
let axis: Vec<f64> = (0..self.num_time_points)
.map(|i| match self.acquire_mode {
AcquireMode::Fixed => i as f64 * self.averaging_time_actual,
AcquireMode::Circular => {
-(((self.num_time_points - 1) - i) as f64) * self.averaging_time_actual
}
})
.collect();
vec![ParamUpdate::float64_array(self.p.ts_time_axis, axis)]
}
fn acquire_reset(&mut self) -> Vec<ParamUpdate> {
self.circular.iter_mut().for_each(|v| *v = 0.0);
self.time_stamp.iter_mut().for_each(|v| *v = 0.0);
self.current_time_point = 0;
self.start_time = Instant::now();
vec![ParamUpdate::int32(self.p.ts_current_point, 0)]
}
fn allocate_arrays(&mut self) -> Vec<ParamUpdate> {
self.circular = vec![0.0; self.num_signals * self.num_time_points];
self.time_stamp = vec![0.0; self.num_time_points];
let mut updates = self.create_axis_array();
updates.extend(self.acquire_reset());
updates
}
fn compute_num_average(&mut self) -> Vec<ParamUpdate> {
if self.time_per_point == 0.0 {
self.num_average = 1;
self.averaging_time_actual = self.averaging_time_requested;
} else {
let n = (self.averaging_time_requested / self.time_per_point + 0.5) as i64;
self.num_average = if n < 1 { 1 } else { n as usize };
self.averaging_time_actual = self.time_per_point * self.num_average as f64;
}
self.num_averaged = 0;
let mut updates = vec![
ParamUpdate::float64(self.p.ts_averaging_time, self.averaging_time_actual),
ParamUpdate::int32(self.p.ts_num_average, self.num_average as i32),
];
updates.extend(self.create_axis_array());
updates
}
fn do_time_series_callbacks(&self) -> Vec<ParamUpdate> {
let ntp = self.num_time_points;
let mut updates = Vec::with_capacity(self.num_signals);
match self.acquire_mode {
AcquireMode::Fixed => {
for signal in 0..self.num_signals {
let start = signal * ntp;
let series = self.circular[start..start + self.current_time_point].to_vec();
updates.push(ParamUpdate::float64_array_addr(
self.p.ts_time_series,
signal as i32,
series,
));
}
}
AcquireMode::Circular => {
for signal in 0..self.num_signals {
let base = signal * ntp;
let mut series = Vec::with_capacity(ntp);
let mut time_in = self.current_time_point;
for _ in 0..ntp {
series.push(self.circular[base + time_in]);
time_in += 1;
if time_in >= ntp {
time_in = 0;
}
}
updates.push(ParamUpdate::float64_array_addr(
self.p.ts_time_series,
signal as i32,
series,
));
}
}
}
updates
}
fn add_to_time_series(&mut self, array: &NDArray) -> Vec<ParamUpdate> {
let mut updates = Vec::new();
let num_signals_in = self.num_signals_in.max(0) as usize;
if num_signals_in == 0 {
return updates;
}
let data = &array.data;
let mut num_times = if array.dims.len() == 2 {
array.dims[1].size
} else {
1
};
let max_times = data.len() / num_signals_in;
if num_times > max_times {
num_times = max_times;
}
let ntp = self.num_time_points;
for i in 0..num_times {
let base = i * num_signals_in;
for s in 0..self.num_signals {
self.average_store[s] += sample_f64(data, base + s);
}
self.num_averaged += 1;
if self.num_averaged < self.num_average {
continue;
}
for s in 0..self.num_signals {
let avg = averaged_value(self.average_store[s], self.num_averaged, self.data_type);
self.circular[s * ntp + self.current_time_point] = avg;
self.average_store[s] = 0.0;
}
self.num_averaged = 0;
self.time_stamp[self.current_time_point] = array.time_stamp;
self.current_time_point += 1;
if self.current_time_point >= ntp {
match self.acquire_mode {
AcquireMode::Fixed => {
self.acquiring = false;
updates.push(ParamUpdate::int32(self.p.ts_acquire, 0));
updates.extend(self.do_time_series_callbacks());
break;
}
AcquireMode::Circular => {
self.current_time_point = 0;
}
}
}
}
updates.push(ParamUpdate::int32(
self.p.ts_current_point,
self.current_time_point as i32,
));
let elapsed = self.start_time.elapsed().as_secs_f64();
updates.push(ParamUpdate::float64(self.p.ts_elapsed_time, elapsed));
updates
}
}
impl NDPluginProcess for TimeSeriesProcessor {
fn process_array(&mut self, array: &NDArray, _pool: &NDArrayPool) -> ProcessResult {
let ndims = array.dims.len();
if !(1..=2).contains(&ndims) {
return ProcessResult::empty();
}
let mut updates: Vec<ParamUpdate> = Vec::new();
let num_signals_in = array.dims[0].size;
let dtype = array.data.data_type();
if dtype != self.data_type || (num_signals_in as i64) != self.num_signals_in {
self.data_type = dtype;
self.num_signals_in = num_signals_in as i64;
self.num_signals = num_signals_in.min(self.max_signals);
updates.extend(self.allocate_arrays());
}
if self.acquiring {
updates.extend(self.add_to_time_series(array));
}
ProcessResult::sink(coalesce_updates(updates))
}
fn plugin_type(&self) -> &str {
"NDPluginTimeSeries"
}
fn register_params(&mut self, base: &mut PortDriverBase) -> asyn_rs::error::AsynResult<()> {
self.p.ts_acquire = base.create_param("TS_ACQUIRE", ParamType::Int32)?;
self.p.ts_read = base.create_param("TS_READ", ParamType::Int32)?;
self.p.ts_num_points = base.create_param("TS_NUM_POINTS", ParamType::Int32)?;
self.p.ts_current_point = base.create_param("TS_CURRENT_POINT", ParamType::Int32)?;
self.p.ts_time_per_point = base.create_param("TS_TIME_PER_POINT", ParamType::Float64)?;
self.p.ts_averaging_time = base.create_param("TS_AVERAGING_TIME", ParamType::Float64)?;
self.p.ts_num_average = base.create_param("TS_NUM_AVERAGE", ParamType::Int32)?;
self.p.ts_elapsed_time = base.create_param("TS_ELAPSED_TIME", ParamType::Float64)?;
self.p.ts_acquire_mode = base.create_param("TS_ACQUIRE_MODE", ParamType::Int32)?;
self.p.ts_time_axis = base.create_param("TS_TIME_AXIS", ParamType::Float64Array)?;
self.p.ts_timestamp = base.create_param("TS_TIMESTAMP", ParamType::Float64Array)?;
self.p.ts_time_series = base.create_param("TS_TIME_SERIES", ParamType::Float64Array)?;
base.set_int32_param(self.p.ts_num_points, 0, self.num_time_points as i32)?;
base.set_int32_param(self.p.ts_num_average, 0, self.num_average as i32)?;
base.set_int32_param(self.p.ts_acquire, 0, 0)?;
base.set_int32_param(self.p.ts_acquire_mode, 0, 0)?;
base.set_int32_param(self.p.ts_current_point, 0, 0)?;
base.set_float64_param(self.p.ts_averaging_time, 0, self.averaging_time_actual)?;
base.set_float64_param(self.p.ts_time_per_point, 0, self.time_per_point)?;
let axis: Vec<f64> = (0..self.num_time_points)
.map(|i| i as f64 * self.averaging_time_actual)
.collect();
base.params
.set_float64_array(self.p.ts_time_axis, 0, axis)?;
Ok(())
}
fn on_param_change(
&mut self,
reason: usize,
params: &PluginParamSnapshot,
) -> ParamChangeResult {
let mut updates = Vec::new();
if reason == self.p.ts_num_points {
self.num_time_points = params.value.as_i32().max(1) as usize;
updates.extend(self.allocate_arrays());
} else if reason == self.p.ts_acquire_mode {
self.acquire_mode = if params.value.as_i32() == 0 {
AcquireMode::Fixed
} else {
AcquireMode::Circular
};
updates.extend(self.acquire_reset());
updates.extend(self.create_axis_array());
} else if reason == self.p.ts_acquire {
if params.value.as_i32() != 0 {
self.acquiring = true;
updates.extend(self.acquire_reset());
} else {
self.acquiring = false;
updates.extend(self.do_time_series_callbacks());
}
} else if reason == self.p.ts_read {
updates.extend(self.do_time_series_callbacks());
} else if reason == self.p.ts_time_per_point {
self.time_per_point = params.value.as_f64();
updates.extend(self.compute_num_average());
} else if reason == self.p.ts_averaging_time {
self.averaging_time_requested = params.value.as_f64();
updates.extend(self.compute_num_average());
}
ParamChangeResult::updates(updates)
}
}
#[cfg(test)]
mod tests {
use super::*;
use ad_core_rs::ndarray::NDDimension;
use asyn_rs::port::{PortDriverBase, PortFlags};
#[test]
fn test_averaged_value_uint8_divides_before_narrowing() {
assert_eq!(averaged_value(600.0, 3, NDDataType::UInt8), 200.0); }
#[test]
fn test_averaged_value_int8_negative_divides_before_narrowing() {
assert_eq!(averaged_value(-600.0, 3, NDDataType::Int8), -200.0); }
#[test]
fn test_averaged_value_uint16_does_not_wrap() {
assert_eq!(averaged_value(70000.0, 2, NDDataType::UInt16), 35000.0); }
#[test]
fn test_averaged_value_narrows_the_mean_toward_zero() {
assert_eq!(averaged_value(7.0, 2, NDDataType::Int8), 3.0);
assert_eq!(averaged_value(-7.0, 2, NDDataType::Int8), -3.0);
assert_eq!(averaged_value(7.0, 2, NDDataType::UInt16), 3.0);
}
#[test]
fn test_averaged_value_int32_in_range_no_wrap() {
assert_eq!(averaged_value(600.0, 3, NDDataType::Int32), 200.0);
}
#[test]
fn test_averaged_value_float_types_exact() {
assert_eq!(averaged_value(600.0, 3, NDDataType::Float64), 200.0);
assert_eq!(averaged_value(600.0, 3, NDDataType::Float32), 200.0);
assert_eq!(averaged_value(7.0, 2, NDDataType::Float64), 3.5);
}
#[test]
fn test_averaged_value_numaverage_one_is_passthrough() {
assert_eq!(averaged_value(200.0, 1, NDDataType::UInt8), 200.0);
assert_eq!(averaged_value(-50.0, 1, NDDataType::Int8), -50.0);
}
fn make_proc(max_signals: usize, port: &str) -> TimeSeriesProcessor {
let mut proc = TimeSeriesProcessor::new(max_signals);
let mut base = PortDriverBase::new(port, max_signals + 1, PortFlags::default());
proc.register_params(&mut base).unwrap();
proc
}
fn find_array(res: &ProcessResult, reason: usize, addr: i32) -> Option<Vec<f64>> {
res.param_updates.iter().find_map(|u| match u {
ParamUpdate::Float64Array {
reason: r,
addr: a,
value,
} if *r == reason && *a == addr => Some(value.clone()),
_ => None,
})
}
fn find_int(res: &ProcessResult, reason: usize) -> Option<i32> {
res.param_updates.iter().find_map(|u| match u {
ParamUpdate::Int32 {
reason: r, value, ..
} if *r == reason => Some(*value),
_ => None,
})
}
#[test]
fn test_process_array_uint8_average_per_signal() {
let mut proc = make_proc(2, "TST_TS_U8");
proc.time_per_point = 1.0;
proc.averaging_time_requested = 3.0;
let _ = proc.compute_num_average();
assert_eq!(proc.num_average, 3);
proc.acquiring = true;
let pool = NDArrayPool::new(1_000_000);
let arr = NDArray::with_data(
vec![NDDimension::new(2), NDDimension::new(3)],
NDDataBuffer::U8(vec![200; 6]),
);
let res = proc.process_array(&arr, &pool);
assert_eq!(proc.current_time_point, 1);
assert_eq!(proc.circular[0 * proc.num_time_points], 200.0);
assert_eq!(proc.circular[proc.num_time_points], 200.0); assert_eq!(find_int(&res, proc.p.ts_current_point), Some(1));
assert!(find_array(&res, proc.p.ts_time_series, 0).is_none());
}
#[test]
fn test_fixed_mode_fills_stops_and_emits_waveforms() {
let mut proc = make_proc(1, "TST_TS_FIX");
proc.num_time_points = 2; proc.acquiring = true;
let pool = NDArrayPool::new(1_000_000);
let arr = NDArray::with_data(
vec![NDDimension::new(1), NDDimension::new(3)],
NDDataBuffer::F64(vec![10.0, 20.0, 30.0]),
);
let res = proc.process_array(&arr, &pool);
assert!(!proc.acquiring);
assert_eq!(proc.current_time_point, 2);
assert_eq!(find_int(&res, proc.p.ts_acquire), Some(0));
let wf = find_array(&res, proc.p.ts_time_series, 0).expect("waveform emitted");
assert_eq!(wf, vec![10.0, 20.0]);
}
#[test]
fn test_circular_mode_wraps_and_rotates_oldest_first() {
let mut proc = make_proc(1, "TST_TS_CIRC");
proc.num_time_points = 3;
proc.acquire_mode = AcquireMode::Circular;
proc.acquiring = true;
let pool = NDArrayPool::new(1_000_000);
let arr = NDArray::with_data(
vec![NDDimension::new(1), NDDimension::new(5)],
NDDataBuffer::F64(vec![1.0, 2.0, 3.0, 4.0, 5.0]),
);
proc.process_array(&arr, &pool);
assert!(proc.acquiring);
assert_eq!(proc.current_time_point, 2);
let updates = proc.do_time_series_callbacks();
let wf = updates
.iter()
.find_map(|u| match u {
ParamUpdate::Float64Array {
reason,
addr,
value,
} if *reason == proc.p.ts_time_series && *addr == 0 => Some(value.clone()),
_ => None,
})
.unwrap();
assert_eq!(wf, vec![3.0, 4.0, 5.0]);
}
#[test]
fn test_one_d_array_is_single_time_point_across_signals() {
let mut proc = make_proc(3, "TST_TS_1D");
proc.acquiring = true;
let pool = NDArrayPool::new(1_000_000);
let arr = NDArray::with_data(
vec![NDDimension::new(3)],
NDDataBuffer::F64(vec![11.0, 22.0, 33.0]),
);
proc.process_array(&arr, &pool);
assert_eq!(proc.num_signals, 3);
assert_eq!(proc.current_time_point, 1);
let ntp = proc.num_time_points;
assert_eq!(proc.circular[0], 11.0);
assert_eq!(proc.circular[ntp], 22.0);
assert_eq!(proc.circular[2 * ntp], 33.0);
}
#[test]
fn test_num_signals_capped_at_max_signals() {
let mut proc = make_proc(2, "TST_TS_CAP");
proc.acquiring = true;
let pool = NDArrayPool::new(1_000_000);
let arr = NDArray::with_data(
vec![NDDimension::new(4), NDDimension::new(1)],
NDDataBuffer::F64(vec![1.0, 2.0, 3.0, 4.0]),
);
proc.process_array(&arr, &pool);
assert_eq!(proc.num_signals, 2);
assert_eq!(proc.circular[0], 1.0);
assert_eq!(proc.circular[proc.num_time_points], 2.0);
}
#[test]
fn test_invalid_ndims_is_ignored() {
let mut proc = make_proc(1, "TST_TS_BAD");
proc.acquiring = true;
let pool = NDArrayPool::new(1_000_000);
let arr = NDArray::with_data(
vec![
NDDimension::new(2),
NDDimension::new(2),
NDDimension::new(2),
],
NDDataBuffer::F64(vec![0.0; 8]),
);
let res = proc.process_array(&arr, &pool);
assert!(res.param_updates.is_empty());
assert_eq!(proc.current_time_point, 0);
}
#[test]
fn test_acquire_mode_flips_time_axis() {
let mut proc = make_proc(1, "TST_TS_AXIS");
proc.num_time_points = 4;
let _ = proc.allocate_arrays();
let fixed = proc.create_axis_array();
let fixed_axis = match &fixed[0] {
ParamUpdate::Float64Array { value, .. } => value.clone(),
_ => panic!("expected axis"),
};
assert_eq!(fixed_axis, vec![0.0, 1.0, 2.0, 3.0]);
proc.acquire_mode = AcquireMode::Circular;
let circ = proc.create_axis_array();
let circ_axis = match &circ[0] {
ParamUpdate::Float64Array { value, .. } => value.clone(),
_ => panic!("expected axis"),
};
assert_eq!(circ_axis, vec![-3.0, -2.0, -1.0, 0.0]);
}
#[test]
fn test_compute_num_average_from_averaging_time() {
let mut proc = make_proc(1, "TST_TS_NAVG");
proc.time_per_point = 0.5;
proc.averaging_time_requested = 2.0;
proc.compute_num_average();
assert_eq!(proc.num_average, 4);
assert_eq!(proc.averaging_time_actual, 2.0);
proc.time_per_point = 0.0;
proc.averaging_time_requested = 7.0;
proc.compute_num_average();
assert_eq!(proc.num_average, 1);
assert_eq!(proc.averaging_time_actual, 7.0);
}
}