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;
8
9#[derive(Debug, Clone)]
18pub struct CalcExpression {
19 compiled: calc::CompiledExpr,
20}
21
22impl CalcExpression {
23 pub fn parse(expr: &str) -> Option<CalcExpression> {
27 calc::compile(expr)
28 .ok()
29 .map(|compiled| CalcExpression { compiled })
30 }
31
32 pub fn evaluate(&self, a: f64, b: f64) -> f64 {
35 let mut inputs = calc::NumericInputs::new();
36 inputs.vars[0] = a; inputs.vars[1] = b; calc::eval(&self.compiled, &mut inputs).unwrap_or(0.0)
39 }
40
41 pub fn evaluate_vars(&self, vars: &[f64; calc::CALC_NARGS]) -> f64 {
45 let mut inputs = calc::NumericInputs::with_vars(*vars);
46 calc::eval(&self.compiled, &mut inputs).unwrap_or(0.0)
47 }
48}
49
50#[derive(Debug, Clone)]
52pub enum TriggerCondition {
53 AttributeThreshold { name: String, threshold: f64 },
55 External,
57 Calc {
62 attr_a: String,
63 attr_b: String,
64 expression: CalcExpression,
65 },
66}
67
68#[derive(Debug, Clone, Copy, PartialEq, Eq)]
70pub enum BufferStatus {
71 Idle,
72 BufferFilling,
73 Flushing,
74 AcquisitionCompleted,
75}
76
77#[derive(Debug, Clone, Copy)]
80pub struct TriggerValues {
81 pub a: f64,
83 pub b: f64,
85 pub calc: f64,
87}
88
89#[derive(Debug, Default)]
92pub struct PushResult {
93 pub forward: Vec<Arc<NDArray>>,
95 pub sequence_done: bool,
97 pub trigger_values: Option<TriggerValues>,
101}
102
103pub struct CircularBuffer {
105 pub(crate) pre_count: usize,
106 pub(crate) post_count: usize,
107 buffer: VecDeque<Arc<NDArray>>,
108 pub(crate) trigger_condition: TriggerCondition,
109 triggered: bool,
110 post_done: usize,
112 pre_flushed: bool,
114 captured: Vec<Arc<NDArray>>,
117 preset_trigger_count: usize,
119 trigger_count: usize,
121 flush_on_soft_trigger: bool,
123 pub(crate) status: BufferStatus,
125}
126
127impl CircularBuffer {
128 pub fn new(pre_count: usize, post_count: usize, condition: TriggerCondition) -> Self {
129 Self {
130 pre_count,
131 post_count,
132 buffer: VecDeque::with_capacity(pre_count + 1),
133 trigger_condition: condition,
134 triggered: false,
135 post_done: 0,
136 pre_flushed: false,
137 captured: Vec::new(),
138 preset_trigger_count: 0,
139 trigger_count: 0,
140 flush_on_soft_trigger: false,
141 status: BufferStatus::Idle,
142 }
143 }
144
145 pub fn set_preset_trigger_count(&mut self, count: usize) {
147 self.preset_trigger_count = count;
148 }
149
150 pub fn trigger_count(&self) -> usize {
152 self.trigger_count
153 }
154
155 pub fn status(&self) -> BufferStatus {
157 self.status
158 }
159
160 pub fn set_flush_on_soft_trigger(&mut self, flush: bool) {
162 self.flush_on_soft_trigger = flush;
163 }
164
165 pub fn push(&mut self, array: Arc<NDArray>) -> PushResult {
173 let mut result = PushResult::default();
174
175 if self.status == BufferStatus::AcquisitionCompleted {
177 return result;
178 }
179
180 if self.status == BufferStatus::Idle {
182 self.status = BufferStatus::BufferFilling;
183 }
184
185 if self.triggered {
186 if !self.pre_flushed {
189 self.pre_flushed = true;
190 let pre: Vec<_> = self.buffer.drain(..).collect();
191 self.captured.extend(pre.iter().cloned());
192 result.forward.extend(pre);
193 }
194 self.captured.push(Arc::clone(&array));
197 result.forward.push(array);
198 self.post_done += 1;
199 if self.post_done >= self.post_count {
200 self.complete_sequence(&mut result);
201 }
202 return result;
203 }
204
205 let trigger = match &self.trigger_condition {
208 TriggerCondition::AttributeThreshold { name, threshold } => array
209 .attributes
210 .get(name)
211 .and_then(|a| a.value.as_f64())
212 .map(|v| v >= *threshold)
213 .unwrap_or(false),
214 TriggerCondition::External => false,
215 TriggerCondition::Calc {
216 attr_a,
217 attr_b,
218 expression,
219 } => {
220 let a = array
221 .attributes
222 .get(attr_a)
223 .and_then(|a| a.value.as_f64())
224 .unwrap_or(f64::NAN);
225 let b = array
226 .attributes
227 .get(attr_b)
228 .and_then(|a| a.value.as_f64())
229 .unwrap_or(f64::NAN);
230 let mut vars = [0.0f64; calc::CALC_NARGS];
233 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);
240 result.trigger_values = Some(TriggerValues { a, b, calc });
243 calc.is_finite() && calc != 0.0
249 }
250 };
251
252 if trigger {
253 self.trigger();
256 self.pre_flushed = true;
260 let pre: Vec<_> = self.buffer.drain(..).collect();
261 self.captured.extend(pre.iter().cloned());
262 result.forward.extend(pre);
263 self.captured.push(Arc::clone(&array));
264 result.forward.push(array);
265 self.post_done += 1;
266 if self.post_done >= self.post_count {
267 self.complete_sequence(&mut result);
268 }
269 return result;
270 }
271
272 self.buffer.push_back(array);
274 if self.buffer.len() > self.pre_count {
275 self.buffer.pop_front();
276 }
277
278 result
279 }
280
281 fn complete_sequence(&mut self, result: &mut PushResult) {
285 self.triggered = false;
286 self.pre_flushed = false;
287 self.post_done = 0;
288 if self.preset_trigger_count > 0 && self.trigger_count >= self.preset_trigger_count {
289 self.status = BufferStatus::AcquisitionCompleted;
290 } else {
291 self.status = BufferStatus::BufferFilling;
292 }
293 result.sequence_done = true;
294 }
295
296 pub fn trigger(&mut self) {
298 if self.status == BufferStatus::AcquisitionCompleted {
300 return;
301 }
302
303 self.triggered = true;
304 self.post_done = 0;
305 self.pre_flushed = false;
306 self.trigger_count += 1;
307 self.status = BufferStatus::Flushing;
308 self.captured.clear();
311 }
312
313 pub fn take_captured(&mut self) -> Vec<Arc<NDArray>> {
315 std::mem::take(&mut self.captured)
316 }
317
318 pub fn is_triggered(&self) -> bool {
319 self.triggered
320 }
321
322 pub fn pre_buffer_len(&self) -> usize {
323 self.buffer.len()
324 }
325
326 pub fn reset(&mut self) {
327 self.buffer.clear();
328 self.captured.clear();
329 self.triggered = false;
330 self.post_done = 0;
331 self.pre_flushed = false;
332 self.trigger_count = 0;
333 self.status = BufferStatus::Idle;
334 }
335}
336
337fn flush_path_status(
349 sequence_done: bool,
350 acquisition_completed: bool,
351 forwarded: bool,
352 pre_buffer_full: bool,
353 pre_count: usize,
354) -> Option<&'static str> {
355 if sequence_done {
356 Some(if acquisition_completed {
357 "Acquisition Completed"
358 } else if pre_count > 0 {
359 "Buffer filling"
360 } else {
361 "Dropping frames"
362 })
363 } else if forwarded {
364 Some("Flushing")
365 } else if pre_buffer_full {
366 Some(if pre_count > 0 {
367 "Buffer Wrapping"
368 } else {
369 "Dropping frames"
370 })
371 } else {
372 None
373 }
374}
375
376#[derive(Default)]
378struct CBParamIndices {
379 control: Option<usize>,
380 status: Option<usize>,
381 trigger_a: Option<usize>,
382 trigger_b: Option<usize>,
383 trigger_a_val: Option<usize>,
384 trigger_b_val: Option<usize>,
385 trigger_calc: Option<usize>,
386 trigger_calc_val: Option<usize>,
387 pre_trigger: Option<usize>,
388 post_trigger: Option<usize>,
389 current_image: Option<usize>,
390 post_count: Option<usize>,
391 soft_trigger: Option<usize>,
392 triggered: Option<usize>,
393 preset_trigger_count: Option<usize>,
394 actual_trigger_count: Option<usize>,
395 flush_on_soft_trigger: Option<usize>,
396}
397
398pub struct CircularBuffProcessor {
399 buffer: CircularBuffer,
400 params: CBParamIndices,
401 max_buffers: usize,
405 trigger_a_name: String,
407 trigger_b_name: String,
408 trigger_calc_expr: String,
409}
410
411impl CircularBuffProcessor {
412 pub fn new(
413 pre_count: usize,
414 post_count: usize,
415 condition: TriggerCondition,
416 max_buffers: usize,
417 ) -> Self {
418 Self {
419 buffer: CircularBuffer::new(pre_count, post_count, condition),
420 params: CBParamIndices::default(),
421 max_buffers,
422 trigger_a_name: String::new(),
423 trigger_b_name: String::new(),
424 trigger_calc_expr: String::new(),
425 }
426 }
427
428 pub fn trigger(&mut self) {
429 self.buffer.trigger();
430 }
431
432 pub fn buffer(&self) -> &CircularBuffer {
433 &self.buffer
434 }
435
436 fn rebuild_trigger_condition(&mut self) {
438 if !self.trigger_calc_expr.is_empty() {
439 if let Some(expr) = CalcExpression::parse(&self.trigger_calc_expr) {
440 self.buffer.trigger_condition = TriggerCondition::Calc {
441 attr_a: self.trigger_a_name.clone(),
442 attr_b: self.trigger_b_name.clone(),
443 expression: expr,
444 };
445 return;
446 }
447 }
448 if !self.trigger_a_name.is_empty() {
449 self.buffer.trigger_condition = TriggerCondition::AttributeThreshold {
450 name: self.trigger_a_name.clone(),
451 threshold: 0.5,
452 };
453 } else {
454 self.buffer.trigger_condition = TriggerCondition::External;
455 }
456 }
457}
458
459impl NDPluginProcess for CircularBuffProcessor {
460 fn process_array(&mut self, array: &NDArray, _pool: &NDArrayPool) -> ProcessResult {
461 use ad_core_rs::plugin::runtime::ParamUpdate;
462
463 let push_result = self.buffer.push(Arc::new(array.clone()));
464
465 let mut updates = Vec::new();
466 if let Some(idx) = self.params.status {
467 let pre_count = self.buffer.pre_count;
471 if let Some(s) = flush_path_status(
472 push_result.sequence_done,
473 self.buffer.status() == BufferStatus::AcquisitionCompleted,
474 !push_result.forward.is_empty(),
475 self.buffer.pre_buffer_len() == pre_count,
476 pre_count,
477 ) {
478 updates.push(ParamUpdate::octet(idx, s.to_string()));
479 }
480 }
481 if let Some(idx) = self.params.current_image {
482 updates.push(ParamUpdate::int32(idx, self.buffer.pre_buffer_len() as i32));
483 }
484 if let Some(idx) = self.params.triggered {
485 updates.push(ParamUpdate::int32(
486 idx,
487 if self.buffer.is_triggered() { 1 } else { 0 },
488 ));
489 }
490 if let Some(idx) = self.params.actual_trigger_count {
491 updates.push(ParamUpdate::int32(idx, self.buffer.trigger_count() as i32));
492 }
493 if let Some(tv) = push_result.trigger_values {
496 if let Some(idx) = self.params.trigger_a_val {
497 updates.push(ParamUpdate::float64(idx, tv.a));
498 }
499 if let Some(idx) = self.params.trigger_b_val {
500 updates.push(ParamUpdate::float64(idx, tv.b));
501 }
502 if let Some(idx) = self.params.trigger_calc_val {
503 updates.push(ParamUpdate::float64(idx, tv.calc));
504 }
505 }
506
507 if push_result.forward.is_empty() {
511 ProcessResult::sink(updates)
512 } else {
513 let mut result = ProcessResult::arrays(push_result.forward);
514 result.param_updates = updates;
515 result
516 }
517 }
518
519 fn plugin_type(&self) -> &str {
520 "NDPluginCircularBuff"
521 }
522
523 fn register_params(
524 &mut self,
525 base: &mut asyn_rs::port::PortDriverBase,
526 ) -> asyn_rs::error::AsynResult<()> {
527 use asyn_rs::param::ParamType;
528 base.create_param("CIRC_BUFF_CONTROL", ParamType::Int32)?;
529 base.create_param("CIRC_BUFF_STATUS", ParamType::Octet)?;
532 base.create_param("CIRC_BUFF_TRIGGER_A", ParamType::Octet)?;
533 base.create_param("CIRC_BUFF_TRIGGER_B", ParamType::Octet)?;
534 base.create_param("CIRC_BUFF_TRIGGER_A_VAL", ParamType::Float64)?;
535 base.create_param("CIRC_BUFF_TRIGGER_B_VAL", ParamType::Float64)?;
536 base.create_param("CIRC_BUFF_TRIGGER_CALC", ParamType::Octet)?;
537 base.create_param("CIRC_BUFF_TRIGGER_CALC_VAL", ParamType::Float64)?;
538 base.create_param("CIRC_BUFF_PRE_TRIGGER", ParamType::Int32)?;
539 base.create_param("CIRC_BUFF_POST_TRIGGER", ParamType::Int32)?;
540 base.create_param("CIRC_BUFF_CURRENT_IMAGE", ParamType::Int32)?;
541 base.create_param("CIRC_BUFF_POST_COUNT", ParamType::Int32)?;
542 base.create_param("CIRC_BUFF_SOFT_TRIGGER", ParamType::Int32)?;
543 base.create_param("CIRC_BUFF_TRIGGERED", ParamType::Int32)?;
544 base.create_param("CIRC_BUFF_PRESET_TRIGGER_COUNT", ParamType::Int32)?;
545 base.create_param("CIRC_BUFF_ACTUAL_TRIGGER_COUNT", ParamType::Int32)?;
546 base.create_param("CIRC_BUFF_FLUSH_ON_SOFTTRIGGER", ParamType::Int32)?;
547
548 self.params.control = base.find_param("CIRC_BUFF_CONTROL");
549 self.params.status = base.find_param("CIRC_BUFF_STATUS");
550 self.params.trigger_a = base.find_param("CIRC_BUFF_TRIGGER_A");
551 self.params.trigger_b = base.find_param("CIRC_BUFF_TRIGGER_B");
552 self.params.trigger_a_val = base.find_param("CIRC_BUFF_TRIGGER_A_VAL");
553 self.params.trigger_b_val = base.find_param("CIRC_BUFF_TRIGGER_B_VAL");
554 self.params.trigger_calc = base.find_param("CIRC_BUFF_TRIGGER_CALC");
555 self.params.trigger_calc_val = base.find_param("CIRC_BUFF_TRIGGER_CALC_VAL");
556 self.params.pre_trigger = base.find_param("CIRC_BUFF_PRE_TRIGGER");
557 self.params.post_trigger = base.find_param("CIRC_BUFF_POST_TRIGGER");
558 self.params.current_image = base.find_param("CIRC_BUFF_CURRENT_IMAGE");
559 self.params.post_count = base.find_param("CIRC_BUFF_POST_COUNT");
560 self.params.soft_trigger = base.find_param("CIRC_BUFF_SOFT_TRIGGER");
561 self.params.triggered = base.find_param("CIRC_BUFF_TRIGGERED");
562 self.params.preset_trigger_count = base.find_param("CIRC_BUFF_PRESET_TRIGGER_COUNT");
563 self.params.actual_trigger_count = base.find_param("CIRC_BUFF_ACTUAL_TRIGGER_COUNT");
564 self.params.flush_on_soft_trigger = base.find_param("CIRC_BUFF_FLUSH_ON_SOFTTRIGGER");
565
566 if let Some(idx) = self.params.status {
569 base.set_string_param(idx, 0, "Idle".into())?;
570 }
571 Ok(())
572 }
573
574 fn on_param_change(
575 &mut self,
576 reason: usize,
577 params: &ad_core_rs::plugin::runtime::PluginParamSnapshot,
578 ) -> ad_core_rs::plugin::runtime::ParamChangeResult {
579 use ad_core_rs::plugin::runtime::{ParamChangeResult, ParamChangeValue, ParamUpdate};
580
581 let mut updates = Vec::new();
582 if Some(reason) == self.params.control {
583 let v = params.value.as_i32();
584 if v == 1 {
585 self.buffer.reset();
587 self.buffer.status = BufferStatus::BufferFilling;
588 if let Some(idx) = self.params.status {
591 let s = if self.buffer.pre_count > 0 {
592 "Buffer filling"
593 } else {
594 "Dropping frames"
595 };
596 updates.push(ParamUpdate::octet(idx, s.to_string()));
597 }
598 } else {
599 self.buffer.status = BufferStatus::Idle;
601 if let Some(idx) = self.params.current_image {
604 updates.push(ParamUpdate::int32(idx, 0));
605 }
606 if let Some(idx) = self.params.status {
609 updates.push(ParamUpdate::octet(idx, "Acquisition Stopped".to_string()));
610 }
611 }
612 } else if Some(reason) == self.params.pre_trigger {
613 let value = params.value.as_i32();
619 let running = matches!(
620 self.buffer.status(),
621 BufferStatus::BufferFilling | BufferStatus::Flushing
622 );
623 let reject_msg = if running {
624 Some("Stop acquisition to set pre-count")
625 } else if value > self.max_buffers as i32 - 1 {
626 Some("Pre-count too high")
628 } else if value < 0 {
629 Some("Invalid pre-count value")
630 } else {
631 None
632 };
633 if let Some(msg) = reject_msg {
634 if let Some(idx) = self.params.status {
635 updates.push(ParamUpdate::octet(idx, msg.to_string()));
636 }
637 if let Some(idx) = self.params.pre_trigger {
640 updates.push(ParamUpdate::int32(idx, self.buffer.pre_count as i32));
641 }
642 } else {
643 self.buffer.pre_count = value as usize;
644 }
645 } else if Some(reason) == self.params.post_trigger {
646 self.buffer.post_count = params.value.as_i32().max(0) as usize;
647 } else if Some(reason) == self.params.preset_trigger_count {
648 self.buffer
649 .set_preset_trigger_count(params.value.as_i32().max(0) as usize);
650 } else if Some(reason) == self.params.flush_on_soft_trigger {
651 self.buffer
652 .set_flush_on_soft_trigger(params.value.as_i32() != 0);
653 } else if Some(reason) == self.params.soft_trigger {
654 if params.value.as_i32() != 0 {
655 self.buffer.trigger();
656 }
657 } else if Some(reason) == self.params.trigger_a {
658 if let ParamChangeValue::Octet(s) = ¶ms.value {
659 self.trigger_a_name = s.clone();
660 self.rebuild_trigger_condition();
661 }
662 } else if Some(reason) == self.params.trigger_b {
663 if let ParamChangeValue::Octet(s) = ¶ms.value {
664 self.trigger_b_name = s.clone();
665 self.rebuild_trigger_condition();
666 }
667 } else if Some(reason) == self.params.trigger_calc {
668 if let ParamChangeValue::Octet(s) = ¶ms.value {
669 self.trigger_calc_expr = s.clone();
670 self.rebuild_trigger_condition();
671 }
672 }
673
674 ParamChangeResult::updates(updates)
675 }
676}
677
678#[cfg(test)]
679mod tests {
680 use super::*;
681 use ad_core_rs::attributes::{NDAttrSource, NDAttrValue, NDAttribute};
682 use ad_core_rs::ndarray::{NDDataType, NDDimension};
683
684 fn make_array(id: i32) -> Arc<NDArray> {
685 let mut arr = NDArray::new(vec![NDDimension::new(4)], NDDataType::UInt8);
686 arr.unique_id = id;
687 Arc::new(arr)
688 }
689
690 fn make_array_with_attr(id: i32, attr_val: f64) -> Arc<NDArray> {
691 let mut arr = NDArray::new(vec![NDDimension::new(4)], NDDataType::UInt8);
692 arr.unique_id = id;
693 arr.attributes.add(NDAttribute::new_static(
694 "trigger",
695 "",
696 NDAttrSource::Driver,
697 NDAttrValue::Float64(attr_val),
698 ));
699 Arc::new(arr)
700 }
701
702 fn make_array_with_attrs(id: i32, a_val: f64, b_val: f64) -> Arc<NDArray> {
703 let mut arr = NDArray::new(vec![NDDimension::new(4)], NDDataType::UInt8);
704 arr.unique_id = id;
705 arr.attributes.add(NDAttribute::new_static(
706 "attr_a",
707 "",
708 NDAttrSource::Driver,
709 NDAttrValue::Float64(a_val),
710 ));
711 arr.attributes.add(NDAttribute::new_static(
712 "attr_b",
713 "",
714 NDAttrSource::Driver,
715 NDAttrValue::Float64(b_val),
716 ));
717 Arc::new(arr)
718 }
719
720 #[test]
721 fn test_pre_trigger_buffering() {
722 let mut cb = CircularBuffer::new(3, 2, TriggerCondition::External);
723
724 for i in 0..5 {
725 cb.push(make_array(i));
726 }
727 assert_eq!(cb.pre_buffer_len(), 3);
729 }
730
731 #[test]
732 fn test_external_trigger() {
733 let mut cb = CircularBuffer::new(2, 2, TriggerCondition::External);
734
735 cb.push(make_array(1));
736 cb.push(make_array(2));
737 cb.push(make_array(3));
738 cb.trigger();
741 assert!(cb.is_triggered());
742
743 let r1 = cb.push(make_array(4));
745 assert!(!r1.sequence_done);
746 let ids1: Vec<_> = r1.forward.iter().map(|a| a.unique_id).collect();
747 assert_eq!(ids1, vec![2, 3, 4]); let r2 = cb.push(make_array(5));
751 assert!(r2.sequence_done);
752 let ids2: Vec<_> = r2.forward.iter().map(|a| a.unique_id).collect();
753 assert_eq!(ids2, vec![5]);
754
755 let captured = cb.take_captured();
756 assert_eq!(captured.len(), 4); assert_eq!(captured[0].unique_id, 2);
758 assert_eq!(captured[1].unique_id, 3);
759 assert_eq!(captured[2].unique_id, 4);
760 assert_eq!(captured[3].unique_id, 5);
761 }
762
763 #[test]
764 fn test_post_count_zero_no_underflow() {
765 let mut cb = CircularBuffer::new(2, 0, TriggerCondition::External);
768 cb.push(make_array(1));
769 cb.push(make_array(2));
770 cb.trigger();
771 assert!(cb.is_triggered());
772
773 let r = cb.push(make_array(3));
776 assert!(r.sequence_done);
777 let ids: Vec<_> = r.forward.iter().map(|a| a.unique_id).collect();
778 assert_eq!(ids, vec![1, 2, 3]);
779 assert!(!cb.is_triggered());
780 assert_eq!(cb.status(), BufferStatus::BufferFilling);
781
782 let r2 = cb.push(make_array(4));
784 assert!(!r2.sequence_done);
785 assert!(r2.forward.is_empty());
786 }
787
788 #[test]
789 fn test_attribute_trigger_post_count_zero() {
790 let mut cb = CircularBuffer::new(
793 1,
794 0,
795 TriggerCondition::AttributeThreshold {
796 name: "trigger".into(),
797 threshold: 5.0,
798 },
799 );
800 cb.push(make_array_with_attr(1, 1.0));
801 let r = cb.push(make_array_with_attr(2, 9.0));
802 assert!(r.sequence_done);
803 let ids: Vec<_> = r.forward.iter().map(|a| a.unique_id).collect();
804 assert_eq!(ids, vec![1, 2]); assert!(!cb.is_triggered());
806 }
807
808 #[test]
809 fn test_attribute_trigger() {
810 let mut cb = CircularBuffer::new(
811 1,
812 2,
813 TriggerCondition::AttributeThreshold {
814 name: "trigger".into(),
815 threshold: 5.0,
816 },
817 );
818
819 cb.push(make_array_with_attr(1, 1.0));
820 cb.push(make_array_with_attr(2, 2.0));
821 assert!(!cb.is_triggered());
822
823 let r3 = cb.push(make_array_with_attr(3, 5.0));
825 assert!(cb.is_triggered());
826 let ids3: Vec<_> = r3.forward.iter().map(|a| a.unique_id).collect();
828 assert_eq!(ids3, vec![2, 3]);
829
830 let r4 = cb.push(make_array(4));
831 assert!(r4.sequence_done);
832
833 let captured = cb.take_captured();
834 assert_eq!(captured.len(), 3);
836 assert_eq!(captured[0].unique_id, 2);
837 assert_eq!(captured[1].unique_id, 3);
838 assert_eq!(captured[2].unique_id, 4);
839 }
840
841 #[test]
844 fn test_calc_trigger() {
845 let expr = CalcExpression::parse("A>5").unwrap();
847 let mut cb = CircularBuffer::new(
848 1,
849 2,
850 TriggerCondition::Calc {
851 attr_a: "attr_a".into(),
852 attr_b: "attr_b".into(),
853 expression: expr,
854 },
855 );
856
857 cb.push(make_array_with_attrs(1, 3.0, 0.0));
859 assert!(!cb.is_triggered());
860
861 cb.push(make_array_with_attrs(2, 6.0, 0.0));
863 assert!(cb.is_triggered());
864
865 let done = cb.push(make_array(3));
866 assert!(done.sequence_done);
867
868 let captured = cb.take_captured();
869 assert_eq!(captured.len(), 3);
871 assert_eq!(captured[0].unique_id, 1);
872 assert_eq!(captured[1].unique_id, 2);
873 assert_eq!(captured[2].unique_id, 3);
874 }
875
876 #[test]
877 fn test_calc_trigger_values_surface() {
878 let expr = CalcExpression::parse("A+B").unwrap();
881 let mut cb = CircularBuffer::new(
884 2,
885 3,
886 TriggerCondition::Calc {
887 attr_a: "attr_a".into(),
888 attr_b: "attr_b".into(),
889 expression: expr,
890 },
891 );
892
893 let r = cb.push(make_array_with_attrs(1, 3.0, 4.0));
895 let tv = r.trigger_values.expect("calc path surfaces trigger values");
896 assert_eq!(tv.a, 3.0);
897 assert_eq!(tv.b, 4.0);
898 assert_eq!(tv.calc, 7.0);
899
900 let r2 = cb.push(make_array(2));
903 assert!(r2.trigger_values.is_none());
904 }
905
906 #[test]
907 fn test_calc_trigger_values_nan_when_attr_absent() {
908 let expr = CalcExpression::parse("A").unwrap();
911 let mut cb = CircularBuffer::new(
912 2,
913 1,
914 TriggerCondition::Calc {
915 attr_a: "missing_a".into(),
916 attr_b: "missing_b".into(),
917 expression: expr,
918 },
919 );
920 let r = cb.push(make_array(1));
921 let tv = r.trigger_values.expect("calc path surfaces trigger values");
922 assert!(tv.a.is_nan());
923 assert!(tv.b.is_nan());
924 assert!(tv.calc.is_nan());
925 }
926
927 #[test]
928 fn test_calc_trigger_skips_nan_and_inf_results() {
929 let push_calc = |val: f64| {
936 let expr = CalcExpression::parse("A").unwrap();
937 let mut cb = CircularBuffer::new(
938 2,
939 2,
940 TriggerCondition::Calc {
941 attr_a: "attr_a".into(),
942 attr_b: "attr_b".into(),
943 expression: expr,
944 },
945 );
946 cb.push(make_array_with_attrs(1, val, 0.0));
947 cb.is_triggered()
948 };
949 assert!(!push_calc(f64::NAN));
951 assert!(!push_calc(f64::INFINITY));
952 assert!(!push_calc(f64::NEG_INFINITY));
953 assert!(push_calc(1.0));
956 assert!(!push_calc(0.0));
957 }
958
959 #[test]
960 fn test_calc_expression_parse() {
961 let expr = CalcExpression::parse("A>5").unwrap();
963 assert_eq!(expr.evaluate(6.0, 0.0), 1.0);
964 assert_eq!(expr.evaluate(4.0, 0.0), 0.0);
965 assert_eq!(expr.evaluate(5.0, 0.0), 0.0); let expr = CalcExpression::parse("A>=5").unwrap();
969 assert_eq!(expr.evaluate(5.0, 0.0), 1.0);
970 assert_eq!(expr.evaluate(4.9, 0.0), 0.0);
971
972 let expr = CalcExpression::parse("A>3&&B<10").unwrap();
974 assert_eq!(expr.evaluate(4.0, 5.0), 1.0);
975 assert_eq!(expr.evaluate(2.0, 5.0), 0.0);
976 assert_eq!(expr.evaluate(4.0, 15.0), 0.0);
977
978 let expr = CalcExpression::parse("(A>10)||(B>10)").unwrap();
980 assert_eq!(expr.evaluate(11.0, 0.0), 1.0);
981 assert_eq!(expr.evaluate(0.0, 11.0), 1.0);
982 assert_eq!(expr.evaluate(0.0, 0.0), 0.0);
983
984 let expr = CalcExpression::parse("A!=0").unwrap();
986 assert_eq!(expr.evaluate(1.0, 0.0), 1.0);
987 assert_eq!(expr.evaluate(0.0, 0.0), 0.0);
988
989 let expr = CalcExpression::parse("A==B").unwrap();
991 assert_eq!(expr.evaluate(5.0, 5.0), 1.0);
992 assert_eq!(expr.evaluate(5.0, 6.0), 0.0);
993
994 let expr = CalcExpression::parse("!A").unwrap();
996 assert_eq!(expr.evaluate(0.0, 0.0), 1.0);
997 assert_eq!(expr.evaluate(1.0, 0.0), 0.0);
998
999 let expr = CalcExpression::parse("A=5").unwrap();
1002 assert_eq!(expr.evaluate(5.0, 0.0), 1.0);
1003 assert_eq!(expr.evaluate(4.0, 0.0), 0.0);
1004
1005 let expr = CalcExpression::parse("A&B").unwrap();
1006 assert_eq!(expr.evaluate(3.0, 1.0), 1.0);
1008
1009 let expr = CalcExpression::parse("ABS(A)").unwrap();
1011 assert_eq!(expr.evaluate(-5.0, 0.0), 5.0);
1012
1013 let expr = CalcExpression::parse("SQRT(A)").unwrap();
1014 assert!((expr.evaluate(9.0, 0.0) - 3.0).abs() < 1e-10);
1015
1016 let expr = CalcExpression::parse("A+B").unwrap();
1017 assert_eq!(expr.evaluate(3.0, 4.0), 7.0);
1018
1019 let expr = CalcExpression::parse("A-B").unwrap();
1020 assert_eq!(expr.evaluate(10.0, 3.0), 7.0);
1021
1022 let expr = CalcExpression::parse("A*B").unwrap();
1023 assert_eq!(expr.evaluate(3.0, 4.0), 12.0);
1024
1025 let expr = CalcExpression::parse("A/B").unwrap();
1026 assert_eq!(expr.evaluate(12.0, 4.0), 3.0);
1027
1028 let expr = CalcExpression::parse("A>5&&C>0").unwrap();
1030 let mut vars = [0.0f64; calc::CALC_NARGS];
1031 vars[0] = 6.0; vars[2] = 1.0; assert_eq!(expr.evaluate_vars(&vars), 1.0);
1034 vars[2] = 0.0; assert_eq!(expr.evaluate_vars(&vars), 0.0);
1036
1037 assert!(CalcExpression::parse("@@@").is_none());
1039 }
1040
1041 #[test]
1042 fn test_preset_trigger_count() {
1043 let mut cb = CircularBuffer::new(1, 1, TriggerCondition::External);
1044 cb.set_preset_trigger_count(2);
1045
1046 assert_eq!(cb.status(), BufferStatus::Idle);
1047
1048 cb.push(make_array(1));
1050 assert_eq!(cb.status(), BufferStatus::BufferFilling);
1051
1052 cb.trigger();
1054 assert_eq!(cb.trigger_count(), 1);
1055 assert_eq!(cb.status(), BufferStatus::Flushing);
1056
1057 let done = cb.push(make_array(2));
1058 assert!(done.sequence_done);
1059 assert_eq!(cb.status(), BufferStatus::BufferFilling); cb.take_captured();
1062
1063 cb.push(make_array(3));
1065
1066 cb.trigger();
1068 assert_eq!(cb.trigger_count(), 2);
1069 assert_eq!(cb.status(), BufferStatus::Flushing);
1070
1071 let done = cb.push(make_array(4));
1072 assert!(done.sequence_done);
1073 assert_eq!(cb.status(), BufferStatus::AcquisitionCompleted);
1074
1075 cb.take_captured();
1076
1077 let done = cb.push(make_array(5));
1079 assert!(!done.sequence_done);
1080 assert_eq!(cb.status(), BufferStatus::AcquisitionCompleted);
1081
1082 cb.trigger();
1084 assert_eq!(cb.trigger_count(), 2); }
1086
1087 #[test]
1088 fn test_stop_resets_current_image_and_status() {
1089 use ad_core_rs::plugin::runtime::{ParamChangeValue, ParamUpdate, PluginParamSnapshot};
1092
1093 let mut processor = CircularBuffProcessor::new(2, 1, TriggerCondition::External, 100);
1094 processor.params.control = Some(10);
1095 processor.params.current_image = Some(11);
1096 processor.params.status = Some(12);
1097
1098 let snapshot = PluginParamSnapshot {
1099 enable_callbacks: true,
1100 reason: 10,
1101 addr: 0,
1102 value: ParamChangeValue::Int32(0), };
1104 let result = processor.on_param_change(10, &snapshot);
1105
1106 assert!(
1107 result.param_updates.iter().any(|u| matches!(
1108 u,
1109 ParamUpdate::Int32 {
1110 reason: 11,
1111 value: 0,
1112 ..
1113 }
1114 )),
1115 "stop must post CURRENT_IMAGE=0"
1116 );
1117 assert!(
1118 result.param_updates.iter().any(|u| matches!(
1119 u,
1120 ParamUpdate::Octet { reason: 12, value, .. } if value == "Acquisition Stopped"
1121 )),
1122 "stop must post STATUS=Acquisition Stopped"
1123 );
1124 }
1125
1126 #[test]
1127 fn test_pre_count_validation() {
1128 use ad_core_rs::plugin::runtime::{ParamChangeValue, ParamUpdate, PluginParamSnapshot};
1133
1134 let make_proc = || {
1135 let mut p = CircularBuffProcessor::new(3, 1, TriggerCondition::External, 10);
1136 p.params.pre_trigger = Some(20);
1137 p.params.status = Some(12);
1138 p
1139 };
1140 let write = |p: &mut CircularBuffProcessor, v: i32| {
1141 let snap = PluginParamSnapshot {
1142 enable_callbacks: true,
1143 reason: 20,
1144 addr: 0,
1145 value: ParamChangeValue::Int32(v),
1146 };
1147 p.on_param_change(20, &snap)
1148 };
1149
1150 let mut p = make_proc();
1152 p.buffer.status = BufferStatus::BufferFilling;
1153 let r = write(&mut p, 7);
1154 assert_eq!(
1155 p.buffer.pre_count, 3,
1156 "reject while running, value unchanged"
1157 );
1158 assert!(r.param_updates.iter().any(|u| matches!(
1159 u,
1160 ParamUpdate::Octet { reason: 12, value, .. } if value == "Stop acquisition to set pre-count"
1161 )));
1162 assert!(r.param_updates.iter().any(|u| matches!(
1163 u,
1164 ParamUpdate::Int32 {
1165 reason: 20,
1166 value: 3,
1167 ..
1168 }
1169 )));
1170
1171 let mut p = make_proc();
1173 p.buffer.status = BufferStatus::Idle;
1174 let r = write(&mut p, -1);
1175 assert_eq!(p.buffer.pre_count, 3, "negative rejected, value unchanged");
1176 assert!(r.param_updates.iter().any(|u| matches!(
1177 u,
1178 ParamUpdate::Octet { reason: 12, value, .. } if value == "Invalid pre-count value"
1179 )));
1180
1181 let mut p = make_proc();
1183 p.buffer.status = BufferStatus::Idle;
1184 let r = write(&mut p, 10);
1185 assert_eq!(p.buffer.pre_count, 3, "too-high rejected, value unchanged");
1186 assert!(r.param_updates.iter().any(|u| matches!(
1187 u,
1188 ParamUpdate::Octet { reason: 12, value, .. } if value == "Pre-count too high"
1189 )));
1190 assert!(r.param_updates.iter().any(|u| matches!(
1191 u,
1192 ParamUpdate::Int32 {
1193 reason: 20,
1194 value: 3,
1195 ..
1196 }
1197 )));
1198
1199 let mut p = make_proc();
1201 p.buffer.status = BufferStatus::Idle;
1202 write(&mut p, 9);
1203 assert_eq!(p.buffer.pre_count, 9, "valid pre-count committed");
1204 }
1205
1206 #[test]
1207 fn test_flush_path_status_strings() {
1208 assert_eq!(
1212 flush_path_status(false, false, true, false, 5),
1213 Some("Flushing")
1214 );
1215 assert_eq!(
1217 flush_path_status(true, false, true, false, 5),
1218 Some("Buffer filling")
1219 );
1220 assert_eq!(
1222 flush_path_status(true, false, true, false, 0),
1223 Some("Dropping frames")
1224 );
1225 assert_eq!(
1227 flush_path_status(true, true, true, false, 5),
1228 Some("Acquisition Completed")
1229 );
1230 assert_eq!(
1232 flush_path_status(false, false, false, true, 5),
1233 Some("Buffer Wrapping")
1234 );
1235 assert_eq!(
1237 flush_path_status(false, false, false, true, 0),
1238 Some("Dropping frames")
1239 );
1240 assert_eq!(flush_path_status(false, false, false, false, 5), None);
1242 }
1243
1244 #[test]
1245 fn test_buffer_status_transitions() {
1246 let mut cb = CircularBuffer::new(2, 1, TriggerCondition::External);
1247
1248 assert_eq!(cb.status(), BufferStatus::Idle);
1250
1251 cb.push(make_array(1));
1253 assert_eq!(cb.status(), BufferStatus::BufferFilling);
1254
1255 cb.push(make_array(2));
1256 assert_eq!(cb.status(), BufferStatus::BufferFilling);
1257
1258 cb.trigger();
1260 assert_eq!(cb.status(), BufferStatus::Flushing);
1261
1262 let done = cb.push(make_array(3));
1264 assert!(done.sequence_done);
1265 assert_eq!(cb.status(), BufferStatus::BufferFilling);
1266
1267 cb.reset();
1269 assert_eq!(cb.status(), BufferStatus::Idle);
1270 assert_eq!(cb.trigger_count(), 0);
1271 }
1272}