1use std::collections::VecDeque;
2use std::sync::Arc;
3
4use ad_core_rs::ndarray::NDArray;
5use ad_core_rs::ndarray_pool::NDArrayPool;
6use ad_core_rs::plugin::runtime::{NDPluginProcess, ProcessResult};
7use epics_base_rs::calc;
8use parking_lot::{MappedMutexGuard, Mutex, MutexGuard};
9
10#[derive(Debug, Clone)]
19pub struct CalcExpression {
20 compiled: calc::CompiledExpr,
21}
22
23impl CalcExpression {
24 pub fn parse(expr: &str) -> Option<CalcExpression> {
28 calc::compile(expr)
29 .ok()
30 .map(|compiled| CalcExpression { compiled })
31 }
32
33 pub fn evaluate(&self, a: f64, b: f64) -> f64 {
36 let mut inputs = calc::NumericInputs::new();
37 inputs.vars[0] = a; inputs.vars[1] = b; calc::eval(&self.compiled, &mut inputs).unwrap_or(0.0)
40 }
41
42 pub fn evaluate_vars(&self, vars: &[f64; calc::CALC_NARGS]) -> f64 {
46 let mut inputs = calc::NumericInputs::with_vars(*vars);
47 calc::eval(&self.compiled, &mut inputs).unwrap_or(0.0)
48 }
49}
50
51#[derive(Debug, Clone)]
53pub enum TriggerCondition {
54 AttributeThreshold { name: String, threshold: f64 },
56 External,
58 Calc {
63 attr_a: String,
64 attr_b: String,
65 expression: CalcExpression,
66 },
67}
68
69#[derive(Debug, Clone, Copy, PartialEq, Eq)]
71pub enum BufferStatus {
72 Idle,
73 BufferFilling,
74 Flushing,
75 AcquisitionCompleted,
76}
77
78#[derive(Debug, Clone, Copy)]
81pub struct TriggerValues {
82 pub a: f64,
84 pub b: f64,
86 pub calc: f64,
88}
89
90#[derive(Debug, Default, Clone, PartialEq, Eq)]
108pub struct FrameParams {
109 pub triggered: Option<i32>,
112 pub current_image: Option<i32>,
115 pub post_count: Option<i32>,
118 pub actual_trigger_count: Option<i32>,
121 pub soft_trigger: Option<i32>,
124 pub control: Option<i32>,
127 pub status: Option<&'static str>,
130}
131
132#[derive(Debug, Default)]
136pub struct PushResult {
137 pub forward: Vec<Arc<NDArray>>,
139 pub sequence_done: bool,
141 pub trigger_values: Option<TriggerValues>,
145 pub params: FrameParams,
147}
148
149pub struct CircularBuffer {
151 control: bool,
164 pub(crate) pre_count: usize,
165 pub(crate) post_count: usize,
166 buffer: VecDeque<Arc<NDArray>>,
167 pub(crate) trigger_condition: TriggerCondition,
168 triggered: bool,
169 post_done: usize,
171 pre_flushed: bool,
173 preset_trigger_count: usize,
175 trigger_count: usize,
179 flush_on_soft_trigger: i32,
185 pub(crate) status: BufferStatus,
187}
188
189impl CircularBuffer {
190 pub fn new(pre_count: usize, post_count: usize, condition: TriggerCondition) -> Self {
191 Self {
192 control: false,
195 pre_count,
196 post_count,
197 buffer: VecDeque::with_capacity(pre_count + 1),
198 trigger_condition: condition,
199 triggered: false,
200 post_done: 0,
201 pre_flushed: false,
202 preset_trigger_count: 0,
203 trigger_count: 0,
204 flush_on_soft_trigger: 0,
205 status: BufferStatus::Idle,
206 }
207 }
208
209 pub fn set_preset_trigger_count(&mut self, count: usize) {
211 self.preset_trigger_count = count;
212 }
213
214 pub fn trigger_count(&self) -> usize {
218 self.trigger_count
219 }
220
221 pub fn status(&self) -> BufferStatus {
223 self.status
224 }
225
226 pub fn set_flush_on_soft_trigger(&mut self, flush_on: i32) {
229 self.flush_on_soft_trigger = flush_on;
230 }
231
232 pub fn flushes_on_soft_trigger(&self) -> bool {
237 self.flush_on_soft_trigger > 0
238 }
239
240 pub fn start(&mut self) {
244 self.reset();
245 self.control = true;
246 self.status = BufferStatus::BufferFilling;
247 }
248
249 pub fn stop(&mut self) {
253 self.control = false;
254 self.triggered = false;
255 self.status = BufferStatus::Idle;
256 }
257
258 pub fn is_running(&self) -> bool {
262 self.control
263 }
264
265 pub fn push(&mut self, array: Arc<NDArray>) -> PushResult {
273 let mut result = PushResult::default();
274
275 if !self.control {
283 return result;
284 }
285
286 if !self.triggered {
290 let fired = self.evaluate_trigger(&array, &mut result);
291 result.params.triggered = Some(i32::from(fired));
294 if fired {
295 self.trigger();
299 }
300 }
301
302 if !self.triggered {
303 self.buffer.push_back(array);
305 if self.buffer.len() > self.pre_count {
306 self.buffer.pop_front();
307 }
308 result.params.current_image = Some(self.buffer.len() as i32);
311 if self.buffer.len() == self.pre_count {
313 result.params.status = Some(if self.pre_count > 0 {
314 "Buffer Wrapping"
315 } else {
316 "Dropping frames"
317 });
318 }
319 } else {
320 result.params.status = Some("Flushing");
323 if !self.pre_flushed {
324 result.forward.extend(self.flush_pre_buffer());
325 }
326 result.forward.push(array);
327 self.post_done += 1;
328 result.params.post_count = Some(self.post_done as i32);
329 }
330
331 if self.post_done >= self.post_count {
339 self.complete_sequence(&mut result);
340 }
341
342 result
343 }
344
345 fn evaluate_trigger(&self, array: &NDArray, result: &mut PushResult) -> bool {
350 match &self.trigger_condition {
351 TriggerCondition::AttributeThreshold { name, threshold } => array
352 .attributes
353 .get(name)
354 .and_then(|a| a.value.as_f64())
355 .map(|v| v >= *threshold)
356 .unwrap_or(false),
357 TriggerCondition::External => false,
358 TriggerCondition::Calc {
359 attr_a,
360 attr_b,
361 expression,
362 } => {
363 let a = array
364 .attributes
365 .get(attr_a)
366 .and_then(|a| a.value.as_f64())
367 .unwrap_or(f64::NAN);
368 let b = array
369 .attributes
370 .get(attr_b)
371 .and_then(|a| a.value.as_f64())
372 .unwrap_or(f64::NAN);
373 let mut vars = [0.0f64; calc::CALC_NARGS];
376 vars[0] = a; vars[1] = b; vars[2] = self.pre_count as f64; vars[3] = self.post_count as f64; vars[4] = self.buffer.len() as f64; vars[5] = if self.triggered { 1.0 } else { 0.0 }; let calc = expression.evaluate_vars(&vars);
383 result.trigger_values = Some(TriggerValues { a, b, calc });
384 calc.is_finite() && calc != 0.0
390 }
391 }
392 }
393
394 fn flush_pre_buffer(&mut self) -> Vec<Arc<NDArray>> {
402 self.pre_flushed = true;
403 self.buffer.drain(..).collect()
404 }
405
406 fn complete_sequence(&mut self, result: &mut PushResult) {
410 self.triggered = false;
411 self.pre_flushed = false;
412 self.post_done = 0;
413 self.trigger_count += 1;
417 result.params.actual_trigger_count = Some(self.trigger_count as i32);
418 if self.preset_trigger_count > 0 && self.trigger_count >= self.preset_trigger_count {
419 self.control = false;
425 self.status = BufferStatus::AcquisitionCompleted;
426 result.params.triggered = Some(0);
427 result.params.control = Some(0);
428 result.params.status = Some("Acquisition Completed");
429 } else {
430 self.status = BufferStatus::BufferFilling;
433 result.params.control = Some(1);
434 result.params.soft_trigger = Some(0);
435 result.params.triggered = Some(0);
436 result.params.post_count = Some(0);
437 result.params.status = Some(if self.pre_count > 0 {
438 "Buffer filling"
439 } else {
440 "Dropping frames"
441 });
442 }
443 result.sequence_done = true;
444 }
445
446 pub fn trigger(&mut self) {
448 if !self.control {
453 return;
454 }
455
456 self.triggered = true;
457 self.post_done = 0;
458 self.pre_flushed = false;
461 self.status = BufferStatus::Flushing;
462 }
463
464 pub fn is_triggered(&self) -> bool {
465 self.triggered
466 }
467
468 pub fn pre_buffer_len(&self) -> usize {
469 self.buffer.len()
470 }
471
472 pub fn reset(&mut self) {
475 self.control = false;
476 self.buffer.clear();
477 self.triggered = false;
478 self.post_done = 0;
479 self.pre_flushed = false;
480 self.trigger_count = 0;
481 self.status = BufferStatus::Idle;
482 }
483}
484
485#[derive(Default)]
490struct CBParamIndices {
491 control: Option<usize>,
492 status: Option<usize>,
493 trigger_a: Option<usize>,
494 trigger_b: Option<usize>,
495 trigger_a_val: Option<usize>,
496 trigger_b_val: Option<usize>,
497 trigger_calc: Option<usize>,
498 trigger_calc_val: Option<usize>,
499 pre_trigger: Option<usize>,
500 post_trigger: Option<usize>,
501 current_image: Option<usize>,
502 post_count: Option<usize>,
503 soft_trigger: Option<usize>,
504 triggered: Option<usize>,
505 preset_trigger_count: Option<usize>,
506 actual_trigger_count: Option<usize>,
507 flush_on_soft_trigger: Option<usize>,
508}
509
510struct CircularBuffState {
513 buffer: CircularBuffer,
514 trigger_a_name: String,
516 trigger_b_name: String,
517 trigger_calc_expr: String,
518}
519
520pub struct CircularBuffProcessor {
521 state: Mutex<CircularBuffState>,
522 params: CBParamIndices,
523 max_buffers: usize,
527}
528
529impl CircularBuffProcessor {
530 pub fn new(
531 pre_count: usize,
532 post_count: usize,
533 condition: TriggerCondition,
534 max_buffers: usize,
535 ) -> Self {
536 Self {
537 state: Mutex::new(CircularBuffState {
538 buffer: CircularBuffer::new(pre_count, post_count, condition),
539 trigger_a_name: String::new(),
540 trigger_b_name: String::new(),
541 trigger_calc_expr: String::new(),
542 }),
543 params: CBParamIndices::default(),
544 max_buffers,
545 }
546 }
547
548 pub fn trigger(&self) {
549 self.state.lock().buffer.trigger();
550 }
551
552 pub fn start(&self) {
556 self.state.lock().buffer.start();
557 }
558
559 pub fn stop(&self) {
561 self.state.lock().buffer.stop();
562 }
563
564 pub fn buffer(&self) -> MappedMutexGuard<'_, CircularBuffer> {
565 MutexGuard::map(self.state.lock(), |s| &mut s.buffer)
566 }
567}
568
569impl CircularBuffState {
570 fn rebuild_trigger_condition(&mut self) {
572 if !self.trigger_calc_expr.is_empty() {
573 if let Some(expr) = CalcExpression::parse(&self.trigger_calc_expr) {
574 self.buffer.trigger_condition = TriggerCondition::Calc {
575 attr_a: self.trigger_a_name.clone(),
576 attr_b: self.trigger_b_name.clone(),
577 expression: expr,
578 };
579 return;
580 }
581 }
582 if !self.trigger_a_name.is_empty() {
583 self.buffer.trigger_condition = TriggerCondition::AttributeThreshold {
584 name: self.trigger_a_name.clone(),
585 threshold: 0.5,
586 };
587 } else {
588 self.buffer.trigger_condition = TriggerCondition::External;
589 }
590 }
591}
592
593impl NDPluginProcess for CircularBuffProcessor {
594 fn process_array(&self, array: &NDArray, _pool: &NDArrayPool) -> ProcessResult {
595 use ad_core_rs::plugin::runtime::ParamUpdate;
596
597 let mut state = self.state.lock();
598 let push_result = state.buffer.push(Arc::new(array.clone()));
599
600 let mut updates = Vec::new();
606 let p = &push_result.params;
607 if let (Some(idx), Some(s)) = (self.params.status, p.status) {
608 updates.push(ParamUpdate::octet(idx, s.to_string()));
610 }
611 for (index, value) in [
612 (self.params.triggered, p.triggered),
613 (self.params.current_image, p.current_image),
614 (self.params.post_count, p.post_count),
615 (self.params.actual_trigger_count, p.actual_trigger_count),
616 (self.params.soft_trigger, p.soft_trigger),
617 (self.params.control, p.control),
618 ] {
619 if let (Some(idx), Some(v)) = (index, value) {
620 updates.push(ParamUpdate::int32(idx, v));
621 }
622 }
623 if let Some(tv) = push_result.trigger_values {
626 if let Some(idx) = self.params.trigger_a_val {
627 updates.push(ParamUpdate::float64(idx, tv.a));
628 }
629 if let Some(idx) = self.params.trigger_b_val {
630 updates.push(ParamUpdate::float64(idx, tv.b));
631 }
632 if let Some(idx) = self.params.trigger_calc_val {
633 updates.push(ParamUpdate::float64(idx, tv.calc));
634 }
635 }
636
637 if push_result.forward.is_empty() {
641 ProcessResult::sink(updates)
642 } else {
643 let mut result = ProcessResult::arrays(push_result.forward);
644 result.param_updates = updates;
645 result
646 }
647 }
648
649 fn plugin_type(&self) -> &str {
650 "NDPluginCircularBuff"
651 }
652
653 fn register_params(
654 &mut self,
655 base: &mut asyn_rs::port::PortDriverBase,
656 ) -> asyn_rs::error::AsynResult<()> {
657 use asyn_rs::param::ParamType;
658 base.create_param("CIRC_BUFF_CONTROL", ParamType::Int32)?;
659 base.create_param("CIRC_BUFF_STATUS", ParamType::Octet)?;
662 base.create_param("CIRC_BUFF_TRIGGER_A", ParamType::Octet)?;
663 base.create_param("CIRC_BUFF_TRIGGER_B", ParamType::Octet)?;
664 base.create_param("CIRC_BUFF_TRIGGER_A_VAL", ParamType::Float64)?;
665 base.create_param("CIRC_BUFF_TRIGGER_B_VAL", ParamType::Float64)?;
666 base.create_param("CIRC_BUFF_TRIGGER_CALC", ParamType::Octet)?;
667 base.create_param("CIRC_BUFF_TRIGGER_CALC_VAL", ParamType::Float64)?;
668 base.create_param("CIRC_BUFF_PRE_TRIGGER", ParamType::Int32)?;
669 base.create_param("CIRC_BUFF_POST_TRIGGER", ParamType::Int32)?;
670 base.create_param("CIRC_BUFF_CURRENT_IMAGE", ParamType::Int32)?;
671 base.create_param("CIRC_BUFF_POST_COUNT", ParamType::Int32)?;
672 base.create_param("CIRC_BUFF_SOFT_TRIGGER", ParamType::Int32)?;
673 base.create_param("CIRC_BUFF_TRIGGERED", ParamType::Int32)?;
674 base.create_param("CIRC_BUFF_PRESET_TRIGGER_COUNT", ParamType::Int32)?;
675 base.create_param("CIRC_BUFF_ACTUAL_TRIGGER_COUNT", ParamType::Int32)?;
676 base.create_param("CIRC_BUFF_FLUSH_ON_SOFTTRIGGER", ParamType::Int32)?;
677
678 self.params.control = base.find_param("CIRC_BUFF_CONTROL");
679 self.params.status = base.find_param("CIRC_BUFF_STATUS");
680 self.params.trigger_a = base.find_param("CIRC_BUFF_TRIGGER_A");
681 self.params.trigger_b = base.find_param("CIRC_BUFF_TRIGGER_B");
682 self.params.trigger_a_val = base.find_param("CIRC_BUFF_TRIGGER_A_VAL");
683 self.params.trigger_b_val = base.find_param("CIRC_BUFF_TRIGGER_B_VAL");
684 self.params.trigger_calc = base.find_param("CIRC_BUFF_TRIGGER_CALC");
685 self.params.trigger_calc_val = base.find_param("CIRC_BUFF_TRIGGER_CALC_VAL");
686 self.params.pre_trigger = base.find_param("CIRC_BUFF_PRE_TRIGGER");
687 self.params.post_trigger = base.find_param("CIRC_BUFF_POST_TRIGGER");
688 self.params.current_image = base.find_param("CIRC_BUFF_CURRENT_IMAGE");
689 self.params.post_count = base.find_param("CIRC_BUFF_POST_COUNT");
690 self.params.soft_trigger = base.find_param("CIRC_BUFF_SOFT_TRIGGER");
691 self.params.triggered = base.find_param("CIRC_BUFF_TRIGGERED");
692 self.params.preset_trigger_count = base.find_param("CIRC_BUFF_PRESET_TRIGGER_COUNT");
693 self.params.actual_trigger_count = base.find_param("CIRC_BUFF_ACTUAL_TRIGGER_COUNT");
694 self.params.flush_on_soft_trigger = base.find_param("CIRC_BUFF_FLUSH_ON_SOFTTRIGGER");
695
696 if let Some(idx) = self.params.status {
699 base.set_string_param(idx, 0, "Idle")?;
700 }
701 Ok(())
702 }
703
704 fn on_param_change(
705 &self,
706 reason: usize,
707 params: &ad_core_rs::plugin::runtime::PluginParamSnapshot,
708 ) -> ad_core_rs::plugin::runtime::ParamChangeResult {
709 use ad_core_rs::plugin::runtime::{ParamChangeResult, ParamChangeValue, ParamUpdate};
710
711 let mut state = self.state.lock();
712 let mut updates = Vec::new();
713 if Some(reason) == self.params.control {
714 let v = params.value.as_i32();
715 if v == 1 {
716 state.buffer.start();
721 for (index, value) in [
722 (self.params.soft_trigger, 0),
723 (self.params.triggered, 0),
724 (self.params.post_count, 0),
725 (self.params.actual_trigger_count, 0),
726 ] {
727 if let Some(idx) = index {
728 updates.push(ParamUpdate::int32(idx, value));
729 }
730 }
731 if let Some(idx) = self.params.status {
734 let s = if state.buffer.pre_count > 0 {
735 "Buffer filling"
736 } else {
737 "Dropping frames"
738 };
739 updates.push(ParamUpdate::octet(idx, s.to_string()));
740 }
741 } else {
742 state.buffer.stop();
747 for (index, value) in [
748 (self.params.soft_trigger, 0),
749 (self.params.triggered, 0),
750 (self.params.current_image, 0),
751 ] {
752 if let Some(idx) = index {
753 updates.push(ParamUpdate::int32(idx, value));
754 }
755 }
756 if let Some(idx) = self.params.status {
759 updates.push(ParamUpdate::octet(idx, "Acquisition Stopped".to_string()));
760 }
761 }
762 } else if Some(reason) == self.params.pre_trigger {
763 let value = params.value.as_i32();
769 let reject_msg = if state.buffer.is_running() {
772 Some("Stop acquisition to set pre-count")
773 } else if value > self.max_buffers as i32 - 1 {
774 Some("Pre-count too high")
776 } else if value < 0 {
777 Some("Invalid pre-count value")
778 } else {
779 None
780 };
781 if let Some(msg) = reject_msg {
782 if let Some(idx) = self.params.status {
783 updates.push(ParamUpdate::octet(idx, msg.to_string()));
784 }
785 if let Some(idx) = self.params.pre_trigger {
788 updates.push(ParamUpdate::int32(idx, state.buffer.pre_count as i32));
789 }
790 } else {
791 state.buffer.pre_count = value as usize;
792 }
793 } else if Some(reason) == self.params.post_trigger {
794 state.buffer.post_count = params.value.as_i32().max(0) as usize;
795 } else if Some(reason) == self.params.preset_trigger_count {
796 state
797 .buffer
798 .set_preset_trigger_count(params.value.as_i32().max(0) as usize);
799 } else if Some(reason) == self.params.flush_on_soft_trigger {
800 state
801 .buffer
802 .set_flush_on_soft_trigger(params.value.as_i32());
803 } else if Some(reason) == self.params.soft_trigger {
804 if params.value.as_i32() != 0 {
834 state.buffer.trigger();
835 if let Some(idx) = self.params.triggered {
836 updates.push(ParamUpdate::int32(idx, 1));
837 }
838 if state.buffer.flushes_on_soft_trigger() {
842 let flushed = state.buffer.flush_pre_buffer();
843 if !flushed.is_empty() {
844 return ParamChangeResult::combined(flushed, updates);
845 }
846 }
847 }
848 } else if Some(reason) == self.params.trigger_a {
849 if let ParamChangeValue::Octet(s) = ¶ms.value {
850 state.trigger_a_name = s.clone();
851 state.rebuild_trigger_condition();
852 }
853 } else if Some(reason) == self.params.trigger_b {
854 if let ParamChangeValue::Octet(s) = ¶ms.value {
855 state.trigger_b_name = s.clone();
856 state.rebuild_trigger_condition();
857 }
858 } else if Some(reason) == self.params.trigger_calc {
859 if let ParamChangeValue::Octet(s) = ¶ms.value {
860 state.trigger_calc_expr = s.clone();
861 state.rebuild_trigger_condition();
862 }
863 }
864
865 ParamChangeResult::updates(updates)
866 }
867}
868
869#[cfg(test)]
870mod tests {
871 use super::*;
872 use ad_core_rs::attributes::{NDAttrSource, NDAttrValue, NDAttribute};
873 use ad_core_rs::ndarray::{NDDataType, NDDimension};
874
875 fn make_array(id: i32) -> Arc<NDArray> {
876 let mut arr = NDArray::new(vec![NDDimension::new(4)], NDDataType::UInt8);
877 arr.unique_id = id;
878 Arc::new(arr)
879 }
880
881 fn make_array_with_attr(id: i32, attr_val: f64) -> Arc<NDArray> {
882 let mut arr = NDArray::new(vec![NDDimension::new(4)], NDDataType::UInt8);
883 arr.unique_id = id;
884 arr.attributes.add(NDAttribute::new_static(
885 "trigger",
886 "",
887 NDAttrSource::Driver,
888 NDAttrValue::Float64(attr_val),
889 ));
890 Arc::new(arr)
891 }
892
893 fn make_array_with_attrs(id: i32, a_val: f64, b_val: f64) -> Arc<NDArray> {
894 let mut arr = NDArray::new(vec![NDDimension::new(4)], NDDataType::UInt8);
895 arr.unique_id = id;
896 arr.attributes.add(NDAttribute::new_static(
897 "attr_a",
898 "",
899 NDAttrSource::Driver,
900 NDAttrValue::Float64(a_val),
901 ));
902 arr.attributes.add(NDAttribute::new_static(
903 "attr_b",
904 "",
905 NDAttrSource::Driver,
906 NDAttrValue::Float64(b_val),
907 ));
908 Arc::new(arr)
909 }
910
911 #[test]
912 fn test_pre_trigger_buffering() {
913 let mut cb = CircularBuffer::new(3, 2, TriggerCondition::External);
914 cb.start(); for i in 0..5 {
917 cb.push(make_array(i));
918 }
919 assert_eq!(cb.pre_buffer_len(), 3);
921 }
922
923 #[test]
930 fn a_completed_sequence_retains_no_forwarded_frames() {
931 let mut cb = CircularBuffer::new(2, 2, TriggerCondition::External);
932 cb.start();
933
934 let pre = make_array(1);
935 cb.push(Arc::clone(&pre));
936 cb.trigger();
937
938 let post = make_array(2);
939 let r1 = cb.push(Arc::clone(&post));
940 let r2 = cb.push(make_array(3));
941 assert!(r2.sequence_done);
942
943 drop(r1);
945 drop(r2);
946
947 assert_eq!(
948 Arc::strong_count(&pre),
949 1,
950 "the flushed pre-trigger frame must not be retained after forwarding"
951 );
952 assert_eq!(
953 Arc::strong_count(&post),
954 1,
955 "the forwarded post-trigger frame must not be retained"
956 );
957 }
958
959 #[test]
960 fn test_external_trigger() {
961 let mut cb = CircularBuffer::new(2, 2, TriggerCondition::External);
962 cb.start(); cb.push(make_array(1));
965 cb.push(make_array(2));
966 cb.push(make_array(3));
967 cb.trigger();
970 assert!(cb.is_triggered());
971
972 let r1 = cb.push(make_array(4));
974 assert!(!r1.sequence_done);
975 let ids1: Vec<_> = r1.forward.iter().map(|a| a.unique_id).collect();
976 assert_eq!(ids1, vec![2, 3, 4]); let r2 = cb.push(make_array(5));
980 assert!(r2.sequence_done);
981 let ids2: Vec<_> = r2.forward.iter().map(|a| a.unique_id).collect();
982 assert_eq!(ids2, vec![5]);
983 }
984
985 #[test]
986 fn test_post_count_zero_no_underflow() {
987 let mut cb = CircularBuffer::new(2, 0, TriggerCondition::External);
990 cb.start(); cb.push(make_array(1));
992 cb.push(make_array(2));
993 cb.trigger();
994 assert!(cb.is_triggered());
995
996 let r = cb.push(make_array(3));
999 assert!(r.sequence_done);
1000 let ids: Vec<_> = r.forward.iter().map(|a| a.unique_id).collect();
1001 assert_eq!(ids, vec![1, 2, 3]);
1002 assert!(!cb.is_triggered());
1003 assert_eq!(cb.status(), BufferStatus::BufferFilling);
1004
1005 let r2 = cb.push(make_array(4));
1013 assert!(r2.sequence_done);
1014 assert!(r2.forward.is_empty());
1015 }
1016
1017 #[test]
1018 fn test_post_count_zero_completes_on_untriggered_frame() {
1019 let mut cb = CircularBuffer::new(2, 0, TriggerCondition::External);
1026 cb.start(); for (n, id) in (1..=3).enumerate() {
1029 let r = cb.push(make_array(id));
1030 assert!(r.forward.is_empty(), "frame {id} must not be forwarded");
1032 assert!(r.sequence_done, "frame {id} must complete a sequence");
1034 assert_eq!(r.params.actual_trigger_count, Some(n as i32 + 1));
1035 assert_eq!(cb.trigger_count(), n + 1);
1036 assert_eq!(r.params.control, Some(1));
1038 assert_eq!(r.params.soft_trigger, Some(0));
1039 assert_eq!(r.params.triggered, Some(0));
1040 assert_eq!(r.params.post_count, Some(0));
1041 assert_eq!(r.params.status, Some("Buffer filling"));
1042 }
1043 assert_eq!(cb.pre_buffer_len(), 2);
1046
1047 let mut cb = CircularBuffer::new(2, 1, TriggerCondition::External);
1050 cb.start(); let r = cb.push(make_array(1));
1052 assert!(!r.sequence_done);
1053 assert_eq!(r.params.actual_trigger_count, None);
1054 assert_eq!(cb.trigger_count(), 0);
1055 }
1056
1057 #[test]
1058 fn test_post_count_zero_untriggered_frames_reach_preset_trigger_count() {
1059 let mut cb = CircularBuffer::new(2, 0, TriggerCondition::External);
1064 cb.start(); cb.set_preset_trigger_count(2);
1066
1067 let r1 = cb.push(make_array(1));
1068 assert!(r1.sequence_done);
1069 assert_eq!(r1.params.actual_trigger_count, Some(1));
1070 assert_eq!(cb.status(), BufferStatus::BufferFilling);
1071
1072 let r2 = cb.push(make_array(2));
1073 assert_eq!(r2.params.actual_trigger_count, Some(2));
1074 assert_eq!(r2.params.control, Some(0));
1076 assert_eq!(r2.params.status, Some("Acquisition Completed"));
1077 assert_eq!(cb.status(), BufferStatus::AcquisitionCompleted);
1078 }
1079
1080 #[test]
1081 fn test_attribute_trigger_post_count_zero() {
1082 let mut cb = CircularBuffer::new(
1085 1,
1086 0,
1087 TriggerCondition::AttributeThreshold {
1088 name: "trigger".into(),
1089 threshold: 5.0,
1090 },
1091 );
1092 cb.start(); cb.push(make_array_with_attr(1, 1.0));
1094 let r = cb.push(make_array_with_attr(2, 9.0));
1095 assert!(r.sequence_done);
1096 let ids: Vec<_> = r.forward.iter().map(|a| a.unique_id).collect();
1097 assert_eq!(ids, vec![1, 2]); assert!(!cb.is_triggered());
1099 }
1100
1101 #[test]
1102 fn test_attribute_trigger() {
1103 let mut cb = CircularBuffer::new(
1104 1,
1105 2,
1106 TriggerCondition::AttributeThreshold {
1107 name: "trigger".into(),
1108 threshold: 5.0,
1109 },
1110 );
1111 cb.start(); cb.push(make_array_with_attr(1, 1.0));
1114 cb.push(make_array_with_attr(2, 2.0));
1115 assert!(!cb.is_triggered());
1116
1117 let r3 = cb.push(make_array_with_attr(3, 5.0));
1119 assert!(cb.is_triggered());
1120 let ids3: Vec<_> = r3.forward.iter().map(|a| a.unique_id).collect();
1122 assert_eq!(ids3, vec![2, 3]);
1123
1124 let r4 = cb.push(make_array(4));
1125 assert!(r4.sequence_done);
1126 let ids4: Vec<_> = r4.forward.iter().map(|a| a.unique_id).collect();
1127 assert_eq!(ids4, vec![4]);
1128 }
1129
1130 #[test]
1133 fn test_calc_trigger() {
1134 let expr = CalcExpression::parse("A>5").unwrap();
1136 let mut cb = CircularBuffer::new(
1137 1,
1138 2,
1139 TriggerCondition::Calc {
1140 attr_a: "attr_a".into(),
1141 attr_b: "attr_b".into(),
1142 expression: expr,
1143 },
1144 );
1145 cb.start(); let mut forwarded: Vec<i32> = Vec::new();
1148 let mut record = |r: PushResult| {
1149 forwarded.extend(r.forward.iter().map(|a| a.unique_id));
1150 r.sequence_done
1151 };
1152
1153 record(cb.push(make_array_with_attrs(1, 3.0, 0.0)));
1155 assert!(!cb.is_triggered());
1156
1157 record(cb.push(make_array_with_attrs(2, 6.0, 0.0)));
1159 assert!(cb.is_triggered());
1160
1161 assert!(record(cb.push(make_array(3))));
1162
1163 assert_eq!(forwarded, vec![1, 2, 3]);
1165 }
1166
1167 #[test]
1168 fn test_calc_trigger_values_surface() {
1169 let expr = CalcExpression::parse("A+B").unwrap();
1172 let mut cb = CircularBuffer::new(
1175 2,
1176 3,
1177 TriggerCondition::Calc {
1178 attr_a: "attr_a".into(),
1179 attr_b: "attr_b".into(),
1180 expression: expr,
1181 },
1182 );
1183 cb.start(); let r = cb.push(make_array_with_attrs(1, 3.0, 4.0));
1187 let tv = r.trigger_values.expect("calc path surfaces trigger values");
1188 assert_eq!(tv.a, 3.0);
1189 assert_eq!(tv.b, 4.0);
1190 assert_eq!(tv.calc, 7.0);
1191
1192 let r2 = cb.push(make_array(2));
1195 assert!(r2.trigger_values.is_none());
1196 }
1197
1198 #[test]
1199 fn test_calc_trigger_values_nan_when_attr_absent() {
1200 let expr = CalcExpression::parse("A").unwrap();
1203 let mut cb = CircularBuffer::new(
1204 2,
1205 1,
1206 TriggerCondition::Calc {
1207 attr_a: "missing_a".into(),
1208 attr_b: "missing_b".into(),
1209 expression: expr,
1210 },
1211 );
1212 cb.start(); let r = cb.push(make_array(1));
1214 let tv = r.trigger_values.expect("calc path surfaces trigger values");
1215 assert!(tv.a.is_nan());
1216 assert!(tv.b.is_nan());
1217 assert!(tv.calc.is_nan());
1218 }
1219
1220 #[test]
1221 fn test_calc_trigger_skips_nan_and_inf_results() {
1222 let push_calc = |val: f64| {
1229 let expr = CalcExpression::parse("A").unwrap();
1230 let mut cb = CircularBuffer::new(
1231 2,
1232 2,
1233 TriggerCondition::Calc {
1234 attr_a: "attr_a".into(),
1235 attr_b: "attr_b".into(),
1236 expression: expr,
1237 },
1238 );
1239 cb.start(); cb.push(make_array_with_attrs(1, val, 0.0));
1241 cb.is_triggered()
1242 };
1243 assert!(!push_calc(f64::NAN));
1245 assert!(!push_calc(f64::INFINITY));
1246 assert!(!push_calc(f64::NEG_INFINITY));
1247 assert!(push_calc(1.0));
1250 assert!(!push_calc(0.0));
1251 }
1252
1253 #[test]
1254 fn test_calc_expression_parse() {
1255 let expr = CalcExpression::parse("A>5").unwrap();
1257 assert_eq!(expr.evaluate(6.0, 0.0), 1.0);
1258 assert_eq!(expr.evaluate(4.0, 0.0), 0.0);
1259 assert_eq!(expr.evaluate(5.0, 0.0), 0.0); let expr = CalcExpression::parse("A>=5").unwrap();
1263 assert_eq!(expr.evaluate(5.0, 0.0), 1.0);
1264 assert_eq!(expr.evaluate(4.9, 0.0), 0.0);
1265
1266 let expr = CalcExpression::parse("A>3&&B<10").unwrap();
1268 assert_eq!(expr.evaluate(4.0, 5.0), 1.0);
1269 assert_eq!(expr.evaluate(2.0, 5.0), 0.0);
1270 assert_eq!(expr.evaluate(4.0, 15.0), 0.0);
1271
1272 let expr = CalcExpression::parse("(A>10)||(B>10)").unwrap();
1274 assert_eq!(expr.evaluate(11.0, 0.0), 1.0);
1275 assert_eq!(expr.evaluate(0.0, 11.0), 1.0);
1276 assert_eq!(expr.evaluate(0.0, 0.0), 0.0);
1277
1278 let expr = CalcExpression::parse("A!=0").unwrap();
1280 assert_eq!(expr.evaluate(1.0, 0.0), 1.0);
1281 assert_eq!(expr.evaluate(0.0, 0.0), 0.0);
1282
1283 let expr = CalcExpression::parse("A==B").unwrap();
1285 assert_eq!(expr.evaluate(5.0, 5.0), 1.0);
1286 assert_eq!(expr.evaluate(5.0, 6.0), 0.0);
1287
1288 let expr = CalcExpression::parse("!A").unwrap();
1290 assert_eq!(expr.evaluate(0.0, 0.0), 1.0);
1291 assert_eq!(expr.evaluate(1.0, 0.0), 0.0);
1292
1293 let expr = CalcExpression::parse("A=5").unwrap();
1296 assert_eq!(expr.evaluate(5.0, 0.0), 1.0);
1297 assert_eq!(expr.evaluate(4.0, 0.0), 0.0);
1298
1299 let expr = CalcExpression::parse("A&B").unwrap();
1300 assert_eq!(expr.evaluate(3.0, 1.0), 1.0);
1302
1303 let expr = CalcExpression::parse("ABS(A)").unwrap();
1305 assert_eq!(expr.evaluate(-5.0, 0.0), 5.0);
1306
1307 let expr = CalcExpression::parse("SQRT(A)").unwrap();
1308 assert!((expr.evaluate(9.0, 0.0) - 3.0).abs() < 1e-10);
1309
1310 let expr = CalcExpression::parse("A+B").unwrap();
1311 assert_eq!(expr.evaluate(3.0, 4.0), 7.0);
1312
1313 let expr = CalcExpression::parse("A-B").unwrap();
1314 assert_eq!(expr.evaluate(10.0, 3.0), 7.0);
1315
1316 let expr = CalcExpression::parse("A*B").unwrap();
1317 assert_eq!(expr.evaluate(3.0, 4.0), 12.0);
1318
1319 let expr = CalcExpression::parse("A/B").unwrap();
1320 assert_eq!(expr.evaluate(12.0, 4.0), 3.0);
1321
1322 let expr = CalcExpression::parse("A>5&&C>0").unwrap();
1324 let mut vars = [0.0f64; calc::CALC_NARGS];
1325 vars[0] = 6.0; vars[2] = 1.0; assert_eq!(expr.evaluate_vars(&vars), 1.0);
1328 vars[2] = 0.0; assert_eq!(expr.evaluate_vars(&vars), 0.0);
1330
1331 assert!(CalcExpression::parse("@@@").is_none());
1333 }
1334
1335 #[test]
1336 fn test_preset_trigger_count() {
1337 let mut cb = CircularBuffer::new(1, 1, TriggerCondition::External);
1338 cb.start(); cb.set_preset_trigger_count(2);
1340
1341 assert_eq!(cb.status(), BufferStatus::BufferFilling);
1345
1346 cb.push(make_array(1));
1347 assert_eq!(cb.status(), BufferStatus::BufferFilling);
1348
1349 cb.trigger();
1352 assert_eq!(cb.trigger_count(), 0);
1353 assert_eq!(cb.status(), BufferStatus::Flushing);
1354
1355 let done = cb.push(make_array(2));
1356 assert!(done.sequence_done);
1357 assert_eq!(cb.trigger_count(), 1); assert_eq!(cb.status(), BufferStatus::BufferFilling); cb.push(make_array(3));
1362
1363 cb.trigger();
1365 assert_eq!(cb.trigger_count(), 1);
1366 assert_eq!(cb.status(), BufferStatus::Flushing);
1367
1368 let done = cb.push(make_array(4));
1369 assert!(done.sequence_done);
1370 assert_eq!(cb.trigger_count(), 2);
1371 assert_eq!(cb.status(), BufferStatus::AcquisitionCompleted);
1372
1373 let done = cb.push(make_array(5));
1375 assert!(!done.sequence_done);
1376 assert_eq!(cb.status(), BufferStatus::AcquisitionCompleted);
1377
1378 cb.trigger();
1380 assert_eq!(cb.trigger_count(), 2); }
1382
1383 #[test]
1384 fn test_stop_resets_current_image_and_status() {
1385 use ad_core_rs::plugin::runtime::{ParamChangeValue, ParamUpdate, PluginParamSnapshot};
1388
1389 let mut processor = CircularBuffProcessor::new(2, 1, TriggerCondition::External, 100);
1390 processor.params.control = Some(10);
1391 processor.params.current_image = Some(11);
1392 processor.params.status = Some(12);
1393
1394 let snapshot = PluginParamSnapshot {
1395 enable_callbacks: true,
1396 reason: 10,
1397 addr: 0,
1398 value: ParamChangeValue::Int32(0), };
1400 let result = processor.on_param_change(10, &snapshot);
1401
1402 assert!(
1403 result.param_updates.iter().any(|u| matches!(
1404 u,
1405 ParamUpdate::Int32 {
1406 reason: 11,
1407 value: 0,
1408 ..
1409 }
1410 )),
1411 "stop must post CURRENT_IMAGE=0"
1412 );
1413 assert!(
1414 result.param_updates.iter().any(|u| matches!(
1415 u,
1416 ParamUpdate::Octet { reason: 12, value, .. } if value == "Acquisition Stopped"
1417 )),
1418 "stop must post STATUS=Acquisition Stopped"
1419 );
1420 }
1421
1422 #[test]
1423 fn test_pre_count_validation() {
1424 use ad_core_rs::plugin::runtime::{ParamChangeValue, ParamUpdate, PluginParamSnapshot};
1429
1430 let make_proc = || {
1431 let mut p = CircularBuffProcessor::new(3, 1, TriggerCondition::External, 10);
1432 p.params.pre_trigger = Some(20);
1433 p.params.status = Some(12);
1434 p
1435 };
1436 let write = |p: &mut CircularBuffProcessor, v: i32| {
1437 let snap = PluginParamSnapshot {
1438 enable_callbacks: true,
1439 reason: 20,
1440 addr: 0,
1441 value: ParamChangeValue::Int32(v),
1442 };
1443 p.on_param_change(20, &snap)
1444 };
1445
1446 let mut p = make_proc();
1448 p.buffer().start();
1449 let r = write(&mut p, 7);
1450 assert_eq!(
1451 p.buffer().pre_count,
1452 3,
1453 "reject while running, value unchanged"
1454 );
1455 assert!(r.param_updates.iter().any(|u| matches!(
1456 u,
1457 ParamUpdate::Octet { reason: 12, value, .. } if value == "Stop acquisition to set pre-count"
1458 )));
1459 assert!(r.param_updates.iter().any(|u| matches!(
1460 u,
1461 ParamUpdate::Int32 {
1462 reason: 20,
1463 value: 3,
1464 ..
1465 }
1466 )));
1467
1468 let mut p = make_proc();
1470 p.buffer().stop();
1471 let r = write(&mut p, -1);
1472 assert_eq!(
1473 p.buffer().pre_count,
1474 3,
1475 "negative rejected, value unchanged"
1476 );
1477 assert!(r.param_updates.iter().any(|u| matches!(
1478 u,
1479 ParamUpdate::Octet { reason: 12, value, .. } if value == "Invalid pre-count value"
1480 )));
1481
1482 let mut p = make_proc();
1484 p.buffer().stop();
1485 let r = write(&mut p, 10);
1486 assert_eq!(
1487 p.buffer().pre_count,
1488 3,
1489 "too-high rejected, value unchanged"
1490 );
1491 assert!(r.param_updates.iter().any(|u| matches!(
1492 u,
1493 ParamUpdate::Octet { reason: 12, value, .. } if value == "Pre-count too high"
1494 )));
1495 assert!(r.param_updates.iter().any(|u| matches!(
1496 u,
1497 ParamUpdate::Int32 {
1498 reason: 20,
1499 value: 3,
1500 ..
1501 }
1502 )));
1503
1504 let mut p = make_proc();
1506 p.buffer().stop();
1507 write(&mut p, 9);
1508 assert_eq!(p.buffer().pre_count, 9, "valid pre-count committed");
1509 }
1510
1511 #[test]
1512 fn test_frame_status_strings() {
1513 let mut cb = CircularBuffer::new(2, 2, TriggerCondition::External);
1517 cb.start(); assert_eq!(cb.push(make_array(1)).params.status, None);
1519 assert_eq!(
1522 cb.push(make_array(2)).params.status,
1523 Some("Buffer Wrapping")
1524 );
1525 assert_eq!(
1526 cb.push(make_array(3)).params.status,
1527 Some("Buffer Wrapping")
1528 );
1529 cb.trigger();
1531 assert_eq!(cb.push(make_array(4)).params.status, Some("Flushing"));
1532 assert_eq!(cb.push(make_array(5)).params.status, Some("Buffer filling"));
1534
1535 let mut cb = CircularBuffer::new(0, 1, TriggerCondition::External);
1538 cb.start(); assert_eq!(
1540 cb.push(make_array(1)).params.status,
1541 Some("Dropping frames")
1542 );
1543 cb.trigger();
1544 assert_eq!(
1545 cb.push(make_array(2)).params.status,
1546 Some("Dropping frames")
1547 );
1548
1549 let mut cb = CircularBuffer::new(2, 1, TriggerCondition::External);
1551 cb.start(); cb.set_preset_trigger_count(1);
1553 cb.trigger();
1554 assert_eq!(
1555 cb.push(make_array(1)).params.status,
1556 Some("Acquisition Completed")
1557 );
1558 }
1559
1560 #[test]
1561 fn test_post_count_posted_per_flushed_frame() {
1562 let mut cb = CircularBuffer::new(2, 3, TriggerCondition::External);
1568 cb.start(); assert_eq!(cb.push(make_array(1)).params.post_count, None);
1571 assert_eq!(cb.push(make_array(2)).params.post_count, None);
1572
1573 cb.trigger();
1574 assert_eq!(cb.push(make_array(3)).params.post_count, Some(1));
1575 assert_eq!(cb.push(make_array(4)).params.post_count, Some(2));
1576 assert_eq!(cb.push(make_array(5)).params.post_count, Some(0));
1579
1580 cb.trigger();
1582 assert_eq!(cb.push(make_array(6)).params.post_count, Some(1));
1583 }
1584
1585 #[test]
1586 fn test_post_count_survives_acquisition_completed() {
1587 let mut cb = CircularBuffer::new(1, 2, TriggerCondition::External);
1591 cb.start(); cb.set_preset_trigger_count(1);
1593 cb.trigger();
1594 assert_eq!(cb.push(make_array(1)).params.post_count, Some(1));
1595 let done = cb.push(make_array(2));
1596 assert_eq!(done.params.post_count, Some(2), "final count, not reset");
1597 assert_eq!(done.params.status, Some("Acquisition Completed"));
1598 assert_eq!(done.params.control, Some(0), "C turns acquisition off");
1599 }
1600
1601 #[test]
1602 fn test_current_image_frozen_during_flush() {
1603 let mut cb = CircularBuffer::new(3, 2, TriggerCondition::External);
1609 cb.start(); assert_eq!(cb.push(make_array(1)).params.current_image, Some(1));
1611 assert_eq!(cb.push(make_array(2)).params.current_image, Some(2));
1612
1613 cb.trigger();
1614 let r1 = cb.push(make_array(3));
1617 assert_eq!(cb.pre_buffer_len(), 0, "the flush drained the ring");
1618 assert_eq!(r1.params.current_image, None);
1619 assert_eq!(cb.push(make_array(4)).params.current_image, None);
1620
1621 assert_eq!(cb.push(make_array(5)).params.current_image, Some(1));
1623 }
1624
1625 #[test]
1626 fn test_actual_trigger_count_increments_at_sequence_completion() {
1627 let mut cb = CircularBuffer::new(1, 2, TriggerCondition::External);
1632 cb.start(); cb.push(make_array(1));
1634 assert_eq!(cb.trigger_count(), 0);
1635
1636 cb.trigger();
1637 assert_eq!(cb.trigger_count(), 0, "the trigger alone completes nothing");
1638
1639 let r1 = cb.push(make_array(2));
1641 assert_eq!(r1.params.actual_trigger_count, None);
1642 assert_eq!(cb.trigger_count(), 0);
1643
1644 let r2 = cb.push(make_array(3));
1647 assert!(r2.sequence_done);
1648 assert_eq!(r2.params.actual_trigger_count, Some(1));
1649 assert_eq!(cb.trigger_count(), 1);
1650 assert_eq!(r2.params.soft_trigger, Some(0), "C clears the soft latch");
1651 assert_eq!(r2.params.triggered, Some(0));
1652 assert_eq!(r2.params.control, Some(1), "still acquiring");
1653 }
1654
1655 #[test]
1656 fn test_processor_emits_the_frame_params() {
1657 use ad_core_rs::ndarray::{NDDataType, NDDimension};
1660 use ad_core_rs::plugin::runtime::ParamUpdate;
1661
1662 let mut p = CircularBuffProcessor::new(2, 2, TriggerCondition::External, 100);
1663 p.buffer().start(); p.params.current_image = Some(11);
1665 p.params.post_count = Some(13);
1666 p.params.actual_trigger_count = Some(16);
1667 let pool = NDArrayPool::new(0);
1668 let frame = || NDArray::new(vec![NDDimension::new(4)], NDDataType::UInt8);
1669 let int32s = |r: &ProcessResult| -> Vec<(usize, i32)> {
1670 r.param_updates
1671 .iter()
1672 .filter_map(|u| match u {
1673 ParamUpdate::Int32 { reason, value, .. } => Some((*reason, *value)),
1674 _ => None,
1675 })
1676 .collect()
1677 };
1678
1679 let r = p.process_array(&frame(), &pool);
1681 assert!(int32s(&r).contains(&(11, 1)));
1682 assert!(!int32s(&r).iter().any(|(reason, _)| *reason == 13));
1683
1684 p.trigger();
1687 let r = p.process_array(&frame(), &pool);
1688 assert!(int32s(&r).contains(&(13, 1)), "POST_COUNT posted per frame");
1689 assert!(
1690 !int32s(&r).iter().any(|(reason, _)| *reason == 11),
1691 "CURRENT_IMAGE frozen during the flush"
1692 );
1693 assert!(
1694 !int32s(&r).iter().any(|(reason, _)| *reason == 16),
1695 "ActualTriggerCount only moves at completion"
1696 );
1697
1698 let r = p.process_array(&frame(), &pool);
1700 assert!(int32s(&r).contains(&(16, 1)));
1701 assert!(int32s(&r).contains(&(13, 0)));
1702 }
1703
1704 #[test]
1716 fn test_soft_trigger_write_latches_only_for_a_nonzero_value() {
1717 use ad_core_rs::ndarray::{NDDataType, NDDimension};
1718 use ad_core_rs::plugin::runtime::{ParamChangeValue, ParamUpdate, PluginParamSnapshot};
1719
1720 let soft_trigger_write = |p: &mut CircularBuffProcessor, value: i32| {
1721 let reason = p.params.soft_trigger.unwrap();
1722 p.on_param_change(
1723 reason,
1724 &PluginParamSnapshot {
1725 enable_callbacks: true,
1726 reason,
1727 addr: 0,
1728 value: ParamChangeValue::Int32(value),
1729 },
1730 )
1731 };
1732 let processor = |flush_on_soft_trig: i32| {
1733 let mut p = CircularBuffProcessor::new(3, 2, TriggerCondition::External, 100);
1734 p.buffer().start(); p.params.soft_trigger = Some(20);
1736 p.params.triggered = Some(21);
1737 p.buffer().set_flush_on_soft_trigger(flush_on_soft_trig);
1738 let pool = NDArrayPool::new(0);
1739 for id in 1..=2 {
1741 let mut a = NDArray::new(vec![NDDimension::new(4)], NDDataType::UInt8);
1742 a.unique_id = id;
1743 p.process_array(&a, &pool);
1744 }
1745 assert_eq!(p.buffer().pre_buffer_len(), 2);
1746 p
1747 };
1748
1749 let latched = |r: &ad_core_rs::plugin::runtime::ParamChangeResult| {
1750 r.param_updates.iter().any(|u| {
1751 matches!(
1752 u,
1753 ParamUpdate::Int32 {
1754 reason: 21,
1755 value: 1,
1756 ..
1757 }
1758 )
1759 })
1760 };
1761
1762 let mut p = processor(1);
1765 let r = soft_trigger_write(&mut p, 0);
1766 assert!(!p.buffer().is_triggered(), "SoftTrigger 0 must not arm");
1767 assert!(!latched(&r), "SoftTrigger 0 must not post Triggered=1");
1768 assert!(r.output_arrays.is_empty(), "SoftTrigger 0 must not flush");
1769 assert_eq!(p.buffer().pre_buffer_len(), 2);
1770
1771 let mut p = processor(0);
1774 let r = soft_trigger_write(&mut p, 1);
1775 assert!(p.buffer().is_triggered());
1776 assert!(latched(&r));
1777 assert!(
1778 r.output_arrays.is_empty(),
1779 "no flush when FlushOnSoftTrig = 0"
1780 );
1781 assert_eq!(p.buffer().pre_buffer_len(), 2);
1782
1783 let mut p = processor(1);
1785 let r = soft_trigger_write(&mut p, 1);
1786 assert!(p.buffer().is_triggered());
1787 let ids: Vec<_> = r.output_arrays.iter().map(|a| a.unique_id).collect();
1788 assert_eq!(ids, vec![1, 2], "pre-buffer flushed from the write");
1789 assert_eq!(p.buffer().pre_buffer_len(), 0);
1790 assert!(latched(&r));
1791
1792 let pool = NDArrayPool::new(0);
1796 let mut a = NDArray::new(vec![NDDimension::new(4)], NDDataType::UInt8);
1797 a.unique_id = 3;
1798 let r = p.process_array(&a, &pool);
1799 let ids: Vec<_> = r.output_arrays.iter().map(|a| a.unique_id).collect();
1800 assert_eq!(ids, vec![3], "pre-buffer already flushed, not re-emitted");
1801 }
1802
1803 #[test]
1809 fn r11_63_flush_on_soft_trig_requires_a_positive_value() {
1810 use ad_core_rs::ndarray::{NDDataType, NDDimension};
1811 use ad_core_rs::plugin::runtime::{ParamChangeValue, PluginParamSnapshot};
1812
1813 const FLUSH_ON: usize = 22;
1814 const SOFT_TRIG: usize = 20;
1815
1816 let write = |p: &mut CircularBuffProcessor, reason: usize, value: i32| {
1817 p.on_param_change(
1818 reason,
1819 &PluginParamSnapshot {
1820 enable_callbacks: true,
1821 reason,
1822 addr: 0,
1823 value: ParamChangeValue::Int32(value),
1824 },
1825 )
1826 };
1827
1828 for (flush_on, expect_flush) in [(-1, false), (0, false), (1, true)] {
1829 let mut p = CircularBuffProcessor::new(3, 2, TriggerCondition::External, 100);
1830 p.buffer().start();
1831 p.params.soft_trigger = Some(SOFT_TRIG);
1832 p.params.triggered = Some(21);
1833 p.params.flush_on_soft_trigger = Some(FLUSH_ON);
1834
1835 write(&mut p, FLUSH_ON, flush_on);
1836 assert_eq!(
1837 p.buffer().flushes_on_soft_trigger(),
1838 expect_flush,
1839 "FlushOnSoftTrig = {flush_on}: C flushes only when > 0"
1840 );
1841
1842 let pool = NDArrayPool::new(0);
1843 for id in 1..=2 {
1844 let mut a = NDArray::new(vec![NDDimension::new(4)], NDDataType::UInt8);
1845 a.unique_id = id;
1846 p.process_array(&a, &pool);
1847 }
1848 assert_eq!(p.buffer().pre_buffer_len(), 2);
1849
1850 let r = write(&mut p, SOFT_TRIG, 1);
1851 assert!(p.buffer().is_triggered());
1852 if expect_flush {
1853 let ids: Vec<_> = r.output_arrays.iter().map(|a| a.unique_id).collect();
1854 assert_eq!(ids, vec![1, 2], "FlushOnSoftTrig = {flush_on}: flushed");
1855 assert_eq!(p.buffer().pre_buffer_len(), 0);
1856 } else {
1857 assert!(
1858 r.output_arrays.is_empty(),
1859 "FlushOnSoftTrig = {flush_on}: C does not flush from the write"
1860 );
1861 assert_eq!(p.buffer().pre_buffer_len(), 2);
1862 }
1863 }
1864 }
1865
1866 #[test]
1867 fn test_buffer_status_transitions() {
1868 let mut cb = CircularBuffer::new(2, 1, TriggerCondition::External);
1869
1870 assert_eq!(cb.status(), BufferStatus::Idle);
1872
1873 cb.start();
1876 assert_eq!(cb.status(), BufferStatus::BufferFilling);
1877
1878 cb.push(make_array(1));
1879 assert_eq!(cb.status(), BufferStatus::BufferFilling);
1880
1881 cb.push(make_array(2));
1882 assert_eq!(cb.status(), BufferStatus::BufferFilling);
1883
1884 cb.trigger();
1886 assert_eq!(cb.status(), BufferStatus::Flushing);
1887
1888 let done = cb.push(make_array(3));
1890 assert!(done.sequence_done);
1891 assert_eq!(cb.status(), BufferStatus::BufferFilling);
1892
1893 cb.reset();
1895 assert_eq!(cb.status(), BufferStatus::Idle);
1896 assert_eq!(cb.trigger_count(), 0);
1897 }
1898}