1use std::collections::HashMap;
17use std::sync::Arc;
18use std::time::{Duration, Instant, SystemTime};
19
20use std::any::Any;
21
22const AUTO_CONNECT_THROTTLE: Duration = Duration::from_secs(2);
28
29#[derive(Debug, Clone)]
31pub struct DeviceState {
32 pub connected: bool,
33 pub enabled: bool,
34 pub auto_connect: bool,
35 pub last_connect_disconnect: Option<Instant>,
41}
42
43impl Default for DeviceState {
44 fn default() -> Self {
45 Self {
46 connected: true,
47 enabled: true,
48 auto_connect: true,
49 last_connect_disconnect: None,
50 }
51 }
52}
53
54use crate::error::{AsynError, AsynResult, AsynStatus};
55use crate::exception::{AsynException, ExceptionEvent, ExceptionManager};
56use crate::interpose::{EomReason, OctetInterpose, OctetInterposeStack};
57use crate::interrupt::{InterruptManager, InterruptValue};
58use crate::param::{EnumEntry, InterruptReason, ParamList, ParamType};
59use crate::trace::TraceManager;
60use crate::user::AsynUser;
61
62#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Default)]
66pub enum QueuePriority {
67 Low = 0,
68 #[default]
69 Medium = 1,
70 High = 2,
71 Connect = 3,
73}
74
75#[derive(Debug, Clone, Copy)]
77pub struct PortFlags {
78 pub multi_device: bool,
80 pub can_block: bool,
89 pub destructible: bool,
91}
92
93impl Default for PortFlags {
94 fn default() -> Self {
95 Self {
103 multi_device: false,
104 can_block: false,
105 destructible: false,
106 }
107 }
108}
109
110pub struct PortDriverBase {
122 pub port_name: String,
123 pub max_addr: usize,
124 pub flags: PortFlags,
125 pub params: ParamList,
126 pub interrupts: InterruptManager,
127 pub connected: bool,
128 pub enabled: bool,
129 pub auto_connect: bool,
130 pub defunct: bool,
136 pub exception_sink: Option<Arc<ExceptionManager>>,
138 pub options: HashMap<String, String>,
139 pub input_eos: Vec<u8>,
141 pub output_eos: Vec<u8>,
143 pub interpose_octet: OctetInterposeStack,
144 pub trace: Option<Arc<TraceManager>>,
145 pub device_states: HashMap<i32, DeviceState>,
147 pub timestamp_source: Option<Arc<dyn Fn() -> SystemTime + Send + Sync>>,
149 pub last_connect_disconnect: Option<Instant>,
154}
155
156impl PortDriverBase {
157 pub fn new(port_name: &str, max_addr: usize, flags: PortFlags) -> Self {
158 Self {
159 port_name: port_name.to_string(),
160 max_addr: max_addr.max(1),
161 flags,
162 params: ParamList::new(max_addr, flags.multi_device),
163 interrupts: InterruptManager::new(256),
164 connected: true,
165 enabled: true,
166 auto_connect: true,
167 defunct: false,
168 exception_sink: None,
169 options: HashMap::new(),
170 input_eos: Vec::new(),
171 output_eos: Vec::new(),
172 interpose_octet: OctetInterposeStack::new(),
173 trace: None,
174 device_states: HashMap::new(),
175 timestamp_source: None,
176 last_connect_disconnect: None,
177 }
178 }
179
180 pub fn announce_exception(&self, exception: AsynException, addr: i32) {
182 if let Some(ref sink) = self.exception_sink {
183 sink.announce(&ExceptionEvent {
184 port_name: self.port_name.clone(),
185 exception,
186 addr,
187 });
188 }
189 }
190
191 pub fn is_connected(&self) -> bool {
193 self.connected
194 }
195
196 pub fn set_connected(&mut self, connected: bool) -> bool {
210 if self.connected == connected {
211 return false;
212 }
213 self.connected = connected;
214 if !connected {
215 self.last_connect_disconnect = Some(Instant::now());
219 }
220 self.announce_exception(AsynException::Connect, -1);
221 true
222 }
223
224 pub fn set_addr_connected(&mut self, addr: i32, connected: bool) -> bool {
227 let was = self.device_state(addr).connected;
228 if was == connected {
229 return false;
230 }
231 self.device_state(addr).connected = connected;
232 if !connected {
233 self.device_state(addr).last_connect_disconnect = Some(Instant::now());
236 }
237 self.announce_exception(AsynException::Connect, addr);
238 true
239 }
240
241 pub fn auto_connect_throttle_ok(&self, addr: i32, now: Instant) -> bool {
259 let last = if self.flags.multi_device {
260 self.device_states
261 .get(&addr)
262 .and_then(|d| d.last_connect_disconnect)
263 } else {
264 self.last_connect_disconnect
265 };
266 match last {
267 None => true,
268 Some(t) => now.saturating_duration_since(t) >= AUTO_CONNECT_THROTTLE,
269 }
270 }
271
272 pub fn stamp_auto_connect_attempt(&mut self, addr: i32, now: Instant) {
278 if self.flags.multi_device {
279 self.device_state(addr).last_connect_disconnect = Some(now);
280 } else {
281 self.last_connect_disconnect = Some(now);
282 }
283 }
284
285 pub fn is_enabled(&self) -> bool {
287 self.enabled
288 }
289
290 pub fn set_enabled(&mut self, enabled: bool) -> AsynResult<()> {
302 if self.defunct {
303 return Err(AsynError::Status {
304 status: AsynStatus::Disabled,
305 message: format!("port {} has been shut down (defunct)", self.port_name),
306 });
307 }
308 self.enabled = enabled;
309 self.announce_exception(AsynException::Enable, -1);
310 Ok(())
311 }
312
313 pub fn set_addr_enabled(&mut self, addr: i32, enabled: bool) -> AsynResult<()> {
318 if self.defunct {
319 return Err(AsynError::Status {
320 status: AsynStatus::Disabled,
321 message: format!("port {} has been shut down (defunct)", self.port_name),
322 });
323 }
324 self.device_state(addr).enabled = enabled;
325 self.announce_exception(AsynException::Enable, addr);
326 Ok(())
327 }
328
329 pub fn is_auto_connect(&self) -> bool {
331 self.auto_connect
332 }
333
334 pub fn set_auto_connect(&mut self, yes: bool) {
343 self.auto_connect = yes;
344 self.announce_exception(AsynException::AutoConnect, -1);
345 }
346
347 pub fn set_auto_connect_addr(&mut self, addr: i32, yes: bool) {
352 self.device_state(addr).auto_connect = yes;
353 self.announce_exception(AsynException::AutoConnect, addr);
354 }
355
356 pub fn is_defunct(&self) -> bool {
360 self.defunct
361 }
362
363 pub fn check_ready(&self) -> AsynResult<()> {
367 if self.defunct {
372 return Err(AsynError::Status {
373 status: AsynStatus::Disabled,
374 message: format!("port {} has been shut down (defunct)", self.port_name),
375 });
376 }
377 if !self.enabled {
378 return Err(AsynError::Status {
379 status: AsynStatus::Disabled,
380 message: format!("port {} is disabled", self.port_name),
381 });
382 }
383 if !self.connected {
384 return Err(AsynError::Status {
385 status: AsynStatus::Disconnected,
386 message: format!("port {} is disconnected", self.port_name),
387 });
388 }
389 Ok(())
390 }
391
392 pub fn shutdown_lifecycle(&mut self) -> AsynResult<()> {
408 if self.defunct {
409 return Ok(());
411 }
412 if !self.flags.destructible {
413 return Err(AsynError::Status {
414 status: AsynStatus::Error,
415 message: format!(
416 "port {} does not support shutting down (ASYN_DESTRUCTIBLE not set)",
417 self.port_name
418 ),
419 });
420 }
421 self.enabled = false;
422 self.defunct = true;
423 self.announce_exception(AsynException::Shutdown, -1);
424 Ok(())
425 }
426
427 pub fn check_ready_addr(&self, addr: i32) -> AsynResult<()> {
430 self.check_ready()?;
431 if self.flags.multi_device {
432 if let Some(ds) = self.device_states.get(&addr) {
433 if !ds.enabled {
434 return Err(AsynError::Status {
435 status: AsynStatus::Disabled,
436 message: format!("port {} addr {} is disabled", self.port_name, addr),
437 });
438 }
439 if !ds.connected {
440 return Err(AsynError::Status {
441 status: AsynStatus::Disconnected,
442 message: format!("port {} addr {} is disconnected", self.port_name, addr),
443 });
444 }
445 }
446 }
447 Ok(())
448 }
449
450 pub fn device_state(&mut self, addr: i32) -> &mut DeviceState {
452 self.device_states.entry(addr).or_default()
453 }
454
455 pub fn is_device_connected(&self, addr: i32) -> bool {
457 self.device_states
458 .get(&addr)
459 .map_or(true, |ds| ds.connected)
460 }
461
462 pub fn connect_addr(&mut self, addr: i32) {
472 self.set_addr_connected(addr, true);
473 }
474
475 pub fn disconnect_addr(&mut self, addr: i32) {
481 self.set_addr_connected(addr, false);
482 }
483
484 pub fn enable_addr(&mut self, addr: i32) {
487 let _ = self.set_addr_enabled(addr, true);
488 }
489
490 pub fn disable_addr(&mut self, addr: i32) {
493 let _ = self.set_addr_enabled(addr, false);
494 }
495
496 pub fn register_timestamp_source<F>(&mut self, source: F)
498 where
499 F: Fn() -> SystemTime + Send + Sync + 'static,
500 {
501 self.timestamp_source = Some(Arc::new(source));
502 }
503
504 pub fn current_timestamp(&self) -> SystemTime {
506 self.timestamp_source
507 .as_ref()
508 .map_or_else(SystemTime::now, |f| f())
509 }
510
511 pub fn create_param(&mut self, name: &str, param_type: ParamType) -> AsynResult<usize> {
512 self.params.create_param(name, param_type)
513 }
514
515 pub fn find_param(&self, name: &str) -> Option<usize> {
516 self.params.find_param(name)
517 }
518
519 pub fn set_int32_param(&mut self, index: usize, addr: i32, value: i32) -> AsynResult<()> {
522 self.params.set_int32(index, addr, value)
523 }
524
525 pub fn get_int32_param(&self, index: usize, addr: i32) -> AsynResult<i32> {
526 self.params.get_int32(index, addr)
527 }
528
529 pub fn get_int32_param_strict(&self, index: usize, addr: i32) -> AsynResult<i32> {
533 self.params.get_int32_strict(index, addr)
534 }
535
536 pub fn set_int64_param(&mut self, index: usize, addr: i32, value: i64) -> AsynResult<()> {
537 self.params.set_int64(index, addr, value)
538 }
539
540 pub fn get_int64_param(&self, index: usize, addr: i32) -> AsynResult<i64> {
541 self.params.get_int64(index, addr)
542 }
543
544 pub fn get_int64_param_strict(&self, index: usize, addr: i32) -> AsynResult<i64> {
546 self.params.get_int64_strict(index, addr)
547 }
548
549 pub fn set_float64_param(&mut self, index: usize, addr: i32, value: f64) -> AsynResult<()> {
550 self.params.set_float64(index, addr, value)
551 }
552
553 pub fn get_float64_param(&self, index: usize, addr: i32) -> AsynResult<f64> {
554 self.params.get_float64(index, addr)
555 }
556
557 pub fn get_float64_param_strict(&self, index: usize, addr: i32) -> AsynResult<f64> {
559 self.params.get_float64_strict(index, addr)
560 }
561
562 pub fn set_string_param(&mut self, index: usize, addr: i32, value: String) -> AsynResult<()> {
563 self.params.set_string(index, addr, value)
564 }
565
566 pub fn get_string_param(&self, index: usize, addr: i32) -> AsynResult<&str> {
567 self.params.get_string(index, addr)
568 }
569
570 pub fn get_string_param_strict(&self, index: usize, addr: i32) -> AsynResult<&str> {
572 self.params.get_string_strict(index, addr)
573 }
574
575 pub fn set_uint32_param(
581 &mut self,
582 index: usize,
583 addr: i32,
584 value: u32,
585 mask: u32,
586 interrupt_mask: u32,
587 ) -> AsynResult<()> {
588 self.params
589 .set_uint32(index, addr, value, mask, interrupt_mask)
590 }
591
592 pub fn get_uint32_param(&self, index: usize, addr: i32) -> AsynResult<u32> {
593 self.params.get_uint32(index, addr)
594 }
595
596 pub fn get_uint32_param_strict(&self, index: usize, addr: i32) -> AsynResult<u32> {
598 self.params.get_uint32_strict(index, addr)
599 }
600
601 pub fn get_enum_param(&self, index: usize, addr: i32) -> AsynResult<(usize, Arc<[EnumEntry]>)> {
602 self.params.get_enum(index, addr)
603 }
604
605 pub fn set_enum_index_param(
606 &mut self,
607 index: usize,
608 addr: i32,
609 value: usize,
610 ) -> AsynResult<()> {
611 self.params.set_enum_index(index, addr, value)
612 }
613
614 pub fn set_enum_choices_param(
615 &mut self,
616 index: usize,
617 addr: i32,
618 choices: Arc<[EnumEntry]>,
619 ) -> AsynResult<()> {
620 self.params.set_enum_choices(index, addr, choices)
621 }
622
623 pub fn get_generic_pointer_param(
624 &self,
625 index: usize,
626 addr: i32,
627 ) -> AsynResult<Arc<dyn Any + Send + Sync>> {
628 self.params.get_generic_pointer(index, addr)
629 }
630
631 pub fn set_generic_pointer_param(
632 &mut self,
633 index: usize,
634 addr: i32,
635 value: Arc<dyn Any + Send + Sync>,
636 ) -> AsynResult<()> {
637 self.params.set_generic_pointer(index, addr, value)
638 }
639
640 pub fn set_param_timestamp(
641 &mut self,
642 index: usize,
643 addr: i32,
644 ts: SystemTime,
645 ) -> AsynResult<()> {
646 self.params.set_timestamp(index, addr, ts)
647 }
648
649 pub fn set_param_status(
650 &mut self,
651 index: usize,
652 addr: i32,
653 status: AsynStatus,
654 alarm_status: u16,
655 alarm_severity: u16,
656 ) -> AsynResult<()> {
657 self.params
658 .set_param_status(index, addr, status, alarm_status, alarm_severity)
659 }
660
661 pub fn get_param_status(&self, index: usize, addr: i32) -> AsynResult<(AsynStatus, u16, u16)> {
662 self.params.get_param_status(index, addr)
663 }
664
665 pub fn report_params(&self, level: i32) {
667 eprintln!(" Number of parameters is {}", self.params.len());
668 if level < 1 {
669 return;
670 }
671 for i in 0..self.params.len() {
672 let name = self.params.param_name(i).unwrap_or("?");
673 let ptype = self
674 .params
675 .param_type(i)
676 .map(|t| format!("{t:?}"))
677 .unwrap_or("?".into());
678 if level >= 2 {
679 for addr in 0..self.max_addr.max(1) {
680 let val = self
681 .params
682 .get_value(i, addr as i32)
683 .map(|v| format!("{v:?}"))
684 .unwrap_or("undefined".into());
685 let (status, alarm_st, alarm_sev) = self
686 .params
687 .get_param_status(i, addr as i32)
688 .unwrap_or((AsynStatus::Success, 0, 0));
689 eprintln!(
690 " param[{i}] name={name} type={ptype} addr={addr} val={val} status={status:?} alarm=({alarm_st},{alarm_sev})"
691 );
692 }
693 } else {
694 eprintln!(" param[{i}] name={name} type={ptype}");
695 }
696 }
697 }
698
699 pub fn push_octet_interpose(&mut self, layer: Box<dyn OctetInterpose>) {
705 self.interpose_octet.push(layer);
706 }
707
708 pub fn call_param_callbacks(&mut self, addr: i32) -> AsynResult<()> {
711 let changed = self.params.take_changed(addr)?;
712 let now = self.current_timestamp();
713 for reason in changed {
714 let value = self.params.get_value(reason, addr)?.clone();
715 if !value.is_array() && !self.params.is_param_defined(reason, addr).unwrap_or(false) {
723 continue;
724 }
725 let ts = self.params.get_timestamp(reason, addr)?.unwrap_or(now);
726 let uint32_mask = self
731 .params
732 .take_uint32_interrupt_mask(reason, addr)
733 .unwrap_or(0);
734 let (aux_status, alarm_status, alarm_severity) = self
740 .params
741 .get_param_status(reason, addr)
742 .unwrap_or((AsynStatus::Success, 0, 0));
743 self.interrupts.notify(InterruptValue {
744 reason,
745 addr,
746 value,
747 timestamp: ts,
748 uint32_changed_mask: uint32_mask,
749 aux_status,
750 alarm_status,
751 alarm_severity,
752 });
753 }
754 Ok(())
755 }
756
757 pub fn call_param_callback(&mut self, addr: i32, reason: usize) -> AsynResult<()> {
761 if self.params.take_changed_single(reason, addr)? {
762 let value = self.params.get_value(reason, addr)?.clone();
763 if !value.is_array() && !self.params.is_param_defined(reason, addr).unwrap_or(false) {
767 return Ok(());
768 }
769 let now = self.current_timestamp();
770 let ts = self.params.get_timestamp(reason, addr)?.unwrap_or(now);
771 let uint32_mask = self
776 .params
777 .take_uint32_interrupt_mask(reason, addr)
778 .unwrap_or(0);
779 let (aux_status, alarm_status, alarm_severity) = self
781 .params
782 .get_param_status(reason, addr)
783 .unwrap_or((AsynStatus::Success, 0, 0));
784 self.interrupts.notify(InterruptValue {
785 reason,
786 addr,
787 value,
788 timestamp: ts,
789 uint32_changed_mask: uint32_mask,
790 aux_status,
791 alarm_status,
792 alarm_severity,
793 });
794 }
795 Ok(())
796 }
797
798 pub fn mark_param_changed(&mut self, index: usize, addr: i32) -> AsynResult<()> {
803 self.params.mark_changed(index, addr)
804 }
805}
806
807pub trait PortDriver: Send + Sync + 'static {
822 fn base(&self) -> &PortDriverBase;
823 fn base_mut(&mut self) -> &mut PortDriverBase;
824
825 fn connect(&mut self, _user: &AsynUser) -> AsynResult<()> {
828 self.base_mut().set_connected(true);
830 Ok(())
831 }
832
833 fn disconnect(&mut self, _user: &AsynUser) -> AsynResult<()> {
834 self.base_mut().set_connected(false);
835 Ok(())
836 }
837
838 fn enable(&mut self, _user: &AsynUser) -> AsynResult<()> {
839 self.base_mut().set_enabled(true)
842 }
843
844 fn disable(&mut self, _user: &AsynUser) -> AsynResult<()> {
845 self.base_mut().set_enabled(false)
846 }
847
848 fn connect_addr(&mut self, user: &AsynUser) -> AsynResult<()> {
849 self.base_mut().connect_addr(user.addr);
850 Ok(())
851 }
852
853 fn disconnect_addr(&mut self, user: &AsynUser) -> AsynResult<()> {
854 self.base_mut().disconnect_addr(user.addr);
855 Ok(())
856 }
857
858 fn enable_addr(&mut self, user: &AsynUser) -> AsynResult<()> {
859 self.base_mut().set_addr_enabled(user.addr, true)
861 }
862
863 fn disable_addr(&mut self, user: &AsynUser) -> AsynResult<()> {
864 self.base_mut().set_addr_enabled(user.addr, false)
865 }
866
867 fn get_option(&self, key: &str) -> AsynResult<String> {
868 self.base()
869 .options
870 .get(key)
871 .cloned()
872 .ok_or_else(|| AsynError::OptionNotFound(key.to_string()))
873 }
874
875 fn set_option(&mut self, key: &str, value: &str) -> AsynResult<()> {
876 self.base_mut()
877 .options
878 .insert(key.to_string(), value.to_string());
879 Ok(())
880 }
881
882 fn report(&self, level: i32) {
883 let base = self.base();
884 eprintln!("Port: {}", base.port_name);
885 eprintln!(
886 " connected: {}, max_addr: {}, params: {}, options: {}",
887 base.connected,
888 base.max_addr,
889 base.params.len(),
890 base.options.len()
891 );
892 if level >= 1 {
893 base.report_params(level.saturating_sub(1));
894 }
895 if level >= 2 {
896 for (k, v) in &base.options {
897 eprintln!(" option: {k} = {v}");
898 }
899 }
900 }
901
902 fn read_int32(&mut self, user: &AsynUser) -> AsynResult<i32> {
921 self.base().params.get_int32_strict(user.reason, user.addr)
922 }
923
924 fn write_int32(&mut self, user: &mut AsynUser, value: i32) -> AsynResult<()> {
925 self.base_mut()
926 .params
927 .set_int32(user.reason, user.addr, value)?;
928 self.base_mut().call_param_callbacks(user.addr)
929 }
930
931 fn read_int64(&mut self, user: &AsynUser) -> AsynResult<i64> {
932 self.base().params.get_int64_strict(user.reason, user.addr)
933 }
934
935 fn write_int64(&mut self, user: &mut AsynUser, value: i64) -> AsynResult<()> {
936 self.base_mut()
937 .params
938 .set_int64(user.reason, user.addr, value)?;
939 self.base_mut().call_param_callbacks(user.addr)
940 }
941
942 fn get_bounds_int32(&self, _user: &AsynUser) -> AsynResult<(i32, i32)> {
946 Ok((0, 0))
947 }
948
949 fn get_bounds_int64(&self, _user: &AsynUser) -> AsynResult<(i64, i64)> {
952 Ok((0, 0))
953 }
954
955 fn read_float64(&mut self, user: &AsynUser) -> AsynResult<f64> {
956 self.base()
957 .params
958 .get_float64_strict(user.reason, user.addr)
959 }
960
961 fn write_float64(&mut self, user: &mut AsynUser, value: f64) -> AsynResult<()> {
962 self.base_mut()
963 .params
964 .set_float64(user.reason, user.addr, value)?;
965 self.base_mut().call_param_callbacks(user.addr)
966 }
967
968 fn read_octet(&mut self, user: &AsynUser, buf: &mut [u8]) -> AsynResult<usize> {
969 let s = self
970 .base()
971 .params
972 .get_string_strict(user.reason, user.addr)?;
973 let bytes = s.as_bytes();
974 let n = bytes.len().min(buf.len());
975 buf[..n].copy_from_slice(&bytes[..n]);
976 Ok(n)
977 }
978
979 fn write_octet(&mut self, user: &mut AsynUser, data: &[u8]) -> AsynResult<()> {
980 let s = String::from_utf8_lossy(data).into_owned();
981 self.base_mut()
982 .params
983 .set_string(user.reason, user.addr, s)?;
984 self.base_mut().call_param_callbacks(user.addr)
985 }
986
987 fn read_uint32_digital(&mut self, user: &AsynUser, mask: u32) -> AsynResult<u32> {
988 let val = self
989 .base()
990 .params
991 .get_uint32_strict(user.reason, user.addr)?;
992 Ok(val & mask)
993 }
994
995 fn write_uint32_digital(
996 &mut self,
997 user: &mut AsynUser,
998 value: u32,
999 mask: u32,
1000 ) -> AsynResult<()> {
1001 self.base_mut()
1004 .params
1005 .set_uint32(user.reason, user.addr, value, mask, 0)?;
1006 self.base_mut().call_param_callbacks(user.addr)
1007 }
1008
1009 fn set_interrupt_uint32_digital(
1018 &mut self,
1019 user: &AsynUser,
1020 mask: u32,
1021 reason: InterruptReason,
1022 ) -> AsynResult<()> {
1023 self.base_mut()
1024 .params
1025 .set_uint32_interrupt(user.reason, user.addr, mask, reason)
1026 }
1027
1028 fn clear_interrupt_uint32_digital(&mut self, user: &AsynUser, mask: u32) -> AsynResult<()> {
1033 self.base_mut()
1034 .params
1035 .clear_uint32_interrupt(user.reason, user.addr, mask)
1036 }
1037
1038 fn get_interrupt_uint32_digital(
1042 &self,
1043 user: &AsynUser,
1044 reason: InterruptReason,
1045 ) -> AsynResult<u32> {
1046 self.base()
1047 .params
1048 .get_uint32_interrupt(user.reason, user.addr, reason)
1049 }
1050
1051 fn read_enum(&mut self, user: &AsynUser) -> AsynResult<(usize, Arc<[EnumEntry]>)> {
1054 self.base().params.get_enum(user.reason, user.addr)
1055 }
1056
1057 fn write_enum(&mut self, user: &mut AsynUser, index: usize) -> AsynResult<()> {
1058 self.base_mut()
1059 .params
1060 .set_enum_index(user.reason, user.addr, index)?;
1061 self.base_mut().call_param_callbacks(user.addr)
1062 }
1063
1064 fn write_enum_choices(
1065 &mut self,
1066 user: &mut AsynUser,
1067 choices: Arc<[EnumEntry]>,
1068 ) -> AsynResult<()> {
1069 self.base_mut()
1070 .params
1071 .set_enum_choices(user.reason, user.addr, choices)?;
1072 self.base_mut().call_param_callbacks(user.addr)
1073 }
1074
1075 fn read_generic_pointer(&mut self, user: &AsynUser) -> AsynResult<Arc<dyn Any + Send + Sync>> {
1078 self.base()
1079 .params
1080 .get_generic_pointer(user.reason, user.addr)
1081 }
1082
1083 fn write_generic_pointer(
1084 &mut self,
1085 user: &mut AsynUser,
1086 value: Arc<dyn Any + Send + Sync>,
1087 ) -> AsynResult<()> {
1088 self.base_mut()
1089 .params
1090 .set_generic_pointer(user.reason, user.addr, value)?;
1091 self.base_mut().call_param_callbacks(user.addr)
1092 }
1093
1094 fn read_float64_array(&mut self, _user: &AsynUser, _buf: &mut [f64]) -> AsynResult<usize> {
1097 Err(AsynError::InterfaceNotSupported("asynFloat64Array".into()))
1098 }
1099
1100 fn write_float64_array(&mut self, user: &AsynUser, data: &[f64]) -> AsynResult<()> {
1101 self.base_mut()
1102 .params
1103 .set_float64_array(user.reason, user.addr, data.to_vec())?;
1104 self.base_mut().call_param_callbacks(user.addr)
1105 }
1106
1107 fn read_int32_array(&mut self, _user: &AsynUser, _buf: &mut [i32]) -> AsynResult<usize> {
1108 Err(AsynError::InterfaceNotSupported("asynInt32Array".into()))
1109 }
1110
1111 fn write_int32_array(&mut self, user: &AsynUser, data: &[i32]) -> AsynResult<()> {
1112 self.base_mut()
1113 .params
1114 .set_int32_array(user.reason, user.addr, data.to_vec())?;
1115 self.base_mut().call_param_callbacks(user.addr)
1116 }
1117
1118 fn read_int8_array(&mut self, _user: &AsynUser, _buf: &mut [i8]) -> AsynResult<usize> {
1119 Err(AsynError::InterfaceNotSupported("asynInt8Array".into()))
1120 }
1121
1122 fn write_int8_array(&mut self, user: &AsynUser, data: &[i8]) -> AsynResult<()> {
1123 self.base_mut()
1124 .params
1125 .set_int8_array(user.reason, user.addr, data.to_vec())?;
1126 self.base_mut().call_param_callbacks(user.addr)
1127 }
1128
1129 fn read_int16_array(&mut self, _user: &AsynUser, _buf: &mut [i16]) -> AsynResult<usize> {
1130 Err(AsynError::InterfaceNotSupported("asynInt16Array".into()))
1131 }
1132
1133 fn write_int16_array(&mut self, user: &AsynUser, data: &[i16]) -> AsynResult<()> {
1134 self.base_mut()
1135 .params
1136 .set_int16_array(user.reason, user.addr, data.to_vec())?;
1137 self.base_mut().call_param_callbacks(user.addr)
1138 }
1139
1140 fn read_int64_array(&mut self, _user: &AsynUser, _buf: &mut [i64]) -> AsynResult<usize> {
1141 Err(AsynError::InterfaceNotSupported("asynInt64Array".into()))
1142 }
1143
1144 fn write_int64_array(&mut self, user: &AsynUser, data: &[i64]) -> AsynResult<()> {
1145 self.base_mut()
1146 .params
1147 .set_int64_array(user.reason, user.addr, data.to_vec())?;
1148 self.base_mut().call_param_callbacks(user.addr)
1149 }
1150
1151 fn read_float32_array(&mut self, _user: &AsynUser, _buf: &mut [f32]) -> AsynResult<usize> {
1152 Err(AsynError::InterfaceNotSupported("asynFloat32Array".into()))
1153 }
1154
1155 fn write_float32_array(&mut self, user: &AsynUser, data: &[f32]) -> AsynResult<()> {
1156 self.base_mut()
1157 .params
1158 .set_float32_array(user.reason, user.addr, data.to_vec())?;
1159 self.base_mut().call_param_callbacks(user.addr)
1160 }
1161
1162 fn io_read_octet(&mut self, user: &AsynUser, buf: &mut [u8]) -> AsynResult<usize> {
1167 self.read_octet(user, buf)
1168 }
1169
1170 fn io_read_octet_eom(
1180 &mut self,
1181 user: &AsynUser,
1182 buf: &mut [u8],
1183 ) -> AsynResult<(usize, EomReason)> {
1184 let cap = buf.len();
1185 let n = self.io_read_octet(user, buf)?;
1186 let eom = if n >= cap && cap > 0 {
1187 EomReason::CNT
1188 } else {
1189 EomReason::empty()
1190 };
1191 Ok((n, eom))
1192 }
1193
1194 fn io_write_octet(&mut self, user: &mut AsynUser, data: &[u8]) -> AsynResult<()> {
1195 self.write_octet(user, data)
1196 }
1197
1198 fn io_read_int32(&mut self, user: &AsynUser) -> AsynResult<i32> {
1199 self.read_int32(user)
1200 }
1201
1202 fn io_write_int32(&mut self, user: &mut AsynUser, value: i32) -> AsynResult<()> {
1203 self.write_int32(user, value)
1204 }
1205
1206 fn io_read_int64(&mut self, user: &AsynUser) -> AsynResult<i64> {
1207 self.read_int64(user)
1208 }
1209
1210 fn io_write_int64(&mut self, user: &mut AsynUser, value: i64) -> AsynResult<()> {
1211 self.write_int64(user, value)
1212 }
1213
1214 fn io_read_float64(&mut self, user: &AsynUser) -> AsynResult<f64> {
1215 self.read_float64(user)
1216 }
1217
1218 fn io_write_float64(&mut self, user: &mut AsynUser, value: f64) -> AsynResult<()> {
1219 self.write_float64(user, value)
1220 }
1221
1222 fn io_read_uint32_digital(&mut self, user: &AsynUser, mask: u32) -> AsynResult<u32> {
1223 self.read_uint32_digital(user, mask)
1224 }
1225
1226 fn io_write_uint32_digital(
1227 &mut self,
1228 user: &mut AsynUser,
1229 value: u32,
1230 mask: u32,
1231 ) -> AsynResult<()> {
1232 self.write_uint32_digital(user, value, mask)
1233 }
1234
1235 fn io_flush(&mut self, _user: &mut AsynUser) -> AsynResult<()> {
1236 Ok(())
1237 }
1238
1239 fn set_input_eos(&mut self, eos: &[u8]) -> AsynResult<()> {
1265 if eos.len() > 2 {
1266 return Err(AsynError::Status {
1267 status: AsynStatus::Error,
1268 message: format!("illegal eoslen {}", eos.len()),
1269 });
1270 }
1271 self.base_mut().input_eos = eos.to_vec();
1272 Ok(())
1273 }
1274
1275 fn get_input_eos(&self) -> Vec<u8> {
1276 self.base().input_eos.clone()
1277 }
1278
1279 fn set_output_eos(&mut self, eos: &[u8]) -> AsynResult<()> {
1280 if eos.len() > 2 {
1281 return Err(AsynError::Status {
1282 status: AsynStatus::Error,
1283 message: format!("illegal eoslen {}", eos.len()),
1284 });
1285 }
1286 self.base_mut().output_eos = eos.to_vec();
1287 Ok(())
1288 }
1289
1290 fn get_output_eos(&self) -> Vec<u8> {
1291 self.base().output_eos.clone()
1292 }
1293
1294 fn shutdown(&mut self) -> AsynResult<()> {
1299 Ok(())
1300 }
1301
1302 fn drv_user_create(&self, drv_info: &str) -> AsynResult<usize> {
1307 self.base()
1308 .params
1309 .find_param(drv_info)
1310 .ok_or_else(|| AsynError::ParamNotFound(drv_info.to_string()))
1311 }
1312
1313 fn capabilities(&self) -> Vec<crate::interfaces::Capability> {
1318 crate::interfaces::default_capabilities()
1319 }
1320
1321 fn supports(&self, cap: crate::interfaces::Capability) -> bool {
1323 self.capabilities().contains(&cap)
1324 }
1325
1326 fn init(&mut self) -> AsynResult<()> {
1327 Ok(())
1328 }
1329}
1330
1331#[cfg(test)]
1332mod tests {
1333 use super::*;
1334 struct TestDriver {
1335 base: PortDriverBase,
1336 }
1337
1338 impl TestDriver {
1339 fn new() -> Self {
1340 let mut base = PortDriverBase::new("test", 1, PortFlags::default());
1341 base.create_param("VAL", ParamType::Int32).unwrap();
1342 base.create_param("TEMP", ParamType::Float64).unwrap();
1343 base.create_param("MSG", ParamType::Octet).unwrap();
1344 base.create_param("BITS", ParamType::UInt32Digital).unwrap();
1345 Self { base }
1346 }
1347 }
1348
1349 impl PortDriver for TestDriver {
1350 fn base(&self) -> &PortDriverBase {
1351 &self.base
1352 }
1353 fn base_mut(&mut self) -> &mut PortDriverBase {
1354 &mut self.base
1355 }
1356 }
1357
1358 #[test]
1359 fn test_default_read_write_int32() {
1360 let mut drv = TestDriver::new();
1361 let mut user = AsynUser::new(0);
1362 drv.write_int32(&mut user, 42).unwrap();
1363 let user = AsynUser::new(0);
1364 assert_eq!(drv.read_int32(&user).unwrap(), 42);
1365 }
1366
1367 #[test]
1368 fn test_default_read_write_float64() {
1369 let mut drv = TestDriver::new();
1370 let mut user = AsynUser::new(1);
1371 drv.write_float64(&mut user, 3.14).unwrap();
1372 let user = AsynUser::new(1);
1373 assert!((drv.read_float64(&user).unwrap() - 3.14).abs() < 1e-10);
1374 }
1375
1376 #[test]
1377 fn test_default_read_write_octet() {
1378 let mut drv = TestDriver::new();
1379 let mut user = AsynUser::new(2);
1380 drv.write_octet(&mut user, b"hello").unwrap();
1381 let user = AsynUser::new(2);
1382 let mut buf = [0u8; 32];
1383 let n = drv.read_octet(&user, &mut buf).unwrap();
1384 assert_eq!(&buf[..n], b"hello");
1385 }
1386
1387 #[test]
1388 fn test_default_read_write_uint32() {
1389 let mut drv = TestDriver::new();
1390 let mut user = AsynUser::new(3);
1391 drv.write_uint32_digital(&mut user, 0xFF, 0x0F).unwrap();
1392 let user = AsynUser::new(3);
1393 assert_eq!(drv.read_uint32_digital(&user, 0xFF).unwrap(), 0x0F);
1394 }
1395
1396 #[test]
1397 fn test_connect_disconnect() {
1398 let mut drv = TestDriver::new();
1399 let user = AsynUser::default();
1400 assert!(drv.base().connected);
1401 drv.disconnect(&user).unwrap();
1402 assert!(!drv.base().connected);
1403 drv.connect(&user).unwrap();
1404 assert!(drv.base().connected);
1405 }
1406
1407 #[test]
1408 fn test_drv_user_create() {
1409 let drv = TestDriver::new();
1410 assert_eq!(drv.drv_user_create("VAL").unwrap(), 0);
1411 assert_eq!(drv.drv_user_create("TEMP").unwrap(), 1);
1412 assert!(drv.drv_user_create("NOPE").is_err());
1413 }
1414
1415 #[test]
1416 fn test_call_param_callbacks() {
1417 let mut drv = TestDriver::new();
1418 let mut rx = drv.base_mut().interrupts.subscribe_async();
1419
1420 drv.base_mut().set_int32_param(0, 0, 100).unwrap();
1421 drv.base_mut().set_float64_param(1, 0, 2.0).unwrap();
1422 drv.base_mut().call_param_callbacks(0).unwrap();
1423
1424 let v1 = rx.try_recv().unwrap();
1425 assert_eq!(v1.reason, 0);
1426 let v2 = rx.try_recv().unwrap();
1427 assert_eq!(v2.reason, 1);
1428 assert!(rx.try_recv().is_err());
1429 }
1430
1431 #[test]
1432 fn flush_skips_undefined_scalar_but_keeps_array_trigger() {
1433 let mut drv = TestDriver::new();
1439 let arr = drv
1440 .base_mut()
1441 .create_param("ARR", ParamType::Int32Array)
1442 .unwrap();
1443 let mut rx = drv.base_mut().interrupts.subscribe_async();
1444
1445 drv.base_mut()
1448 .params
1449 .set_param_status(0, 0, AsynStatus::Error, 0, 0)
1450 .unwrap();
1451 drv.base_mut().mark_param_changed(arr, 0).unwrap();
1453
1454 drv.base_mut().call_param_callbacks(0).unwrap();
1455
1456 let iv = rx.try_recv().unwrap();
1458 assert_eq!(
1459 iv.reason, arr,
1460 "array trigger must still fire while undefined"
1461 );
1462 assert!(
1463 rx.try_recv().is_err(),
1464 "undefined scalar must not emit an I/O Intr"
1465 );
1466
1467 drv.base_mut().set_int32_param(0, 0, 7).unwrap();
1469 drv.base_mut().call_param_callbacks(0).unwrap();
1470 let iv2 = rx.try_recv().unwrap();
1471 assert_eq!(iv2.reason, 0, "defined scalar must fire");
1472 }
1473
1474 #[test]
1475 fn uint32_callback_mask_does_not_leak_across_flushes() {
1476 let mut drv = TestDriver::new();
1480 let mut rx = drv.base_mut().interrupts.subscribe_async();
1481
1482 drv.base_mut()
1484 .params
1485 .set_uint32(3, 0, 0x01, 0x01, 0)
1486 .unwrap();
1487 drv.base_mut().call_param_callbacks(0).unwrap();
1488 let iv1 = rx.try_recv().unwrap();
1489 assert_eq!(iv1.reason, 3);
1490 assert_eq!(iv1.uint32_changed_mask, 0x01);
1491
1492 drv.base_mut()
1494 .params
1495 .set_uint32(3, 0, 0x02, 0x02, 0)
1496 .unwrap();
1497 drv.base_mut().call_param_callbacks(0).unwrap();
1498 let iv2 = rx.try_recv().unwrap();
1499 assert_eq!(
1500 iv2.uint32_changed_mask, 0x02,
1501 "second flush must not leak flush-1 bits via an un-reset mask"
1502 );
1503 assert_eq!(
1504 drv.base().params.get_uint32_interrupt_mask(3, 0).unwrap(),
1505 0,
1506 "the flush must consume (reset) the callback mask"
1507 );
1508 }
1509
1510 #[test]
1511 fn test_call_param_callbacks_propagates_aux_status_and_alarm() {
1512 let mut drv = TestDriver::new();
1517 let mut rx = drv.base_mut().interrupts.subscribe_async();
1518
1519 drv.base_mut().set_int32_param(0, 0, 99).unwrap();
1520 drv.base_mut()
1521 .params
1522 .set_param_status(0, 0, crate::error::AsynStatus::Timeout, 4, 2)
1523 .unwrap();
1524 drv.base_mut().call_param_callbacks(0).unwrap();
1525
1526 let iv = rx.try_recv().unwrap();
1527 assert_eq!(iv.reason, 0);
1528 assert!(matches!(iv.aux_status, crate::error::AsynStatus::Timeout));
1529 assert_eq!(iv.alarm_status, 4);
1530 assert_eq!(iv.alarm_severity, 2);
1531 }
1532
1533 #[test]
1534 fn test_call_param_callback_single_propagates_aux_status() {
1535 let mut drv = TestDriver::new();
1537 let mut rx = drv.base_mut().interrupts.subscribe_async();
1538
1539 drv.base_mut().set_int32_param(0, 0, 1).unwrap();
1540 drv.base_mut()
1541 .params
1542 .set_param_status(0, 0, crate::error::AsynStatus::Disconnected, 7, 3)
1543 .unwrap();
1544 drv.base_mut().call_param_callback(0, 0).unwrap();
1545
1546 let iv = rx.try_recv().unwrap();
1547 assert!(matches!(
1548 iv.aux_status,
1549 crate::error::AsynStatus::Disconnected
1550 ));
1551 assert_eq!(iv.alarm_status, 7);
1552 assert_eq!(iv.alarm_severity, 3);
1553 }
1554
1555 #[test]
1556 fn test_no_callback_for_unchanged() {
1557 let mut drv = TestDriver::new();
1558 let mut rx = drv.base_mut().interrupts.subscribe_async();
1559
1560 drv.base_mut().set_int32_param(0, 0, 5).unwrap();
1561 drv.base_mut().call_param_callbacks(0).unwrap();
1562 let _ = rx.try_recv().unwrap(); drv.base_mut().set_int32_param(0, 0, 5).unwrap();
1566 drv.base_mut().call_param_callbacks(0).unwrap();
1567 assert!(rx.try_recv().is_err());
1568 }
1569
1570 #[test]
1571 fn test_array_not_supported_by_default() {
1572 let mut drv = TestDriver::new();
1573 let user = AsynUser::new(0);
1574 let mut buf = [0f64; 10];
1575 assert!(drv.read_float64_array(&user, &mut buf).is_err());
1576 assert!(drv.write_float64_array(&user, &[1.0]).is_err());
1577 }
1578
1579 #[test]
1580 fn test_option_set_get() {
1581 let mut drv = TestDriver::new();
1582 drv.set_option("baud", "9600").unwrap();
1583 assert_eq!(drv.get_option("baud").unwrap(), "9600");
1584 drv.set_option("baud", "115200").unwrap();
1585 assert_eq!(drv.get_option("baud").unwrap(), "115200");
1586 }
1587
1588 #[test]
1589 fn test_option_not_found() {
1590 let drv = TestDriver::new();
1591 let err = drv.get_option("nonexistent").unwrap_err();
1592 assert!(matches!(err, AsynError::OptionNotFound(_)));
1593 }
1594
1595 #[test]
1596 fn test_report_no_panic() {
1597 let mut drv = TestDriver::new();
1598 drv.set_option("testkey", "testval").unwrap();
1599 drv.base_mut().set_int32_param(0, 0, 42).unwrap();
1600 for level in 0..=3 {
1601 drv.report(level);
1602 }
1603 }
1604
1605 #[test]
1606 fn test_callback_uses_param_timestamp() {
1607 let mut drv = TestDriver::new();
1608 let mut rx = drv.base_mut().interrupts.subscribe_async();
1609
1610 let custom_ts = SystemTime::UNIX_EPOCH + std::time::Duration::from_secs(1_000_000);
1611 drv.base_mut().set_int32_param(0, 0, 77).unwrap();
1612 drv.base_mut().set_param_timestamp(0, 0, custom_ts).unwrap();
1613 drv.base_mut().call_param_callbacks(0).unwrap();
1614
1615 let v = rx.try_recv().unwrap();
1616 assert_eq!(v.reason, 0);
1617 assert_eq!(v.timestamp, custom_ts);
1618 }
1619
1620 #[test]
1621 fn test_default_read_write_enum() {
1622 use crate::param::EnumEntry;
1623
1624 let mut base = PortDriverBase::new("test_enum", 1, PortFlags::default());
1625 base.create_param("MODE", ParamType::Enum).unwrap();
1626
1627 struct EnumDriver {
1628 base: PortDriverBase,
1629 }
1630 impl PortDriver for EnumDriver {
1631 fn base(&self) -> &PortDriverBase {
1632 &self.base
1633 }
1634 fn base_mut(&mut self) -> &mut PortDriverBase {
1635 &mut self.base
1636 }
1637 }
1638
1639 let mut drv = EnumDriver { base };
1640 let choices: Arc<[EnumEntry]> = Arc::from(vec![
1641 EnumEntry {
1642 string: "Off".into(),
1643 value: 0,
1644 severity: 0,
1645 },
1646 EnumEntry {
1647 string: "On".into(),
1648 value: 1,
1649 severity: 0,
1650 },
1651 ]);
1652 let mut user = AsynUser::new(0);
1653 drv.write_enum_choices(&mut user, choices).unwrap();
1654 drv.write_enum(&mut user, 1).unwrap();
1655 let (idx, ch) = drv.read_enum(&AsynUser::new(0)).unwrap();
1656 assert_eq!(idx, 1);
1657 assert_eq!(ch[1].string, "On");
1658 }
1659
1660 #[test]
1661 fn test_enum_callback() {
1662 use crate::param::{EnumEntry, ParamValue};
1663
1664 let mut base = PortDriverBase::new("test_enum_cb", 1, PortFlags::default());
1665 base.create_param("MODE", ParamType::Enum).unwrap();
1666 let mut rx = base.interrupts.subscribe_async();
1667
1668 struct EnumDriver {
1669 base: PortDriverBase,
1670 }
1671 impl PortDriver for EnumDriver {
1672 fn base(&self) -> &PortDriverBase {
1673 &self.base
1674 }
1675 fn base_mut(&mut self) -> &mut PortDriverBase {
1676 &mut self.base
1677 }
1678 }
1679
1680 let mut drv = EnumDriver { base };
1681 let choices: Arc<[EnumEntry]> = Arc::from(vec![
1682 EnumEntry {
1683 string: "A".into(),
1684 value: 0,
1685 severity: 0,
1686 },
1687 EnumEntry {
1688 string: "B".into(),
1689 value: 1,
1690 severity: 0,
1691 },
1692 ]);
1693 drv.base_mut()
1694 .set_enum_choices_param(0, 0, choices)
1695 .unwrap();
1696 drv.base_mut().set_enum_index_param(0, 0, 1).unwrap();
1697 drv.base_mut().call_param_callbacks(0).unwrap();
1698
1699 let v = rx.try_recv().unwrap();
1700 assert_eq!(v.reason, 0);
1701 assert!(matches!(v.value, ParamValue::Enum { index: 1, .. }));
1702 }
1703
1704 #[test]
1705 fn test_default_read_write_generic_pointer() {
1706 let mut base = PortDriverBase::new("test_gp", 1, PortFlags::default());
1707 base.create_param("PTR", ParamType::GenericPointer).unwrap();
1708
1709 struct GpDriver {
1710 base: PortDriverBase,
1711 }
1712 impl PortDriver for GpDriver {
1713 fn base(&self) -> &PortDriverBase {
1714 &self.base
1715 }
1716 fn base_mut(&mut self) -> &mut PortDriverBase {
1717 &mut self.base
1718 }
1719 }
1720
1721 let mut drv = GpDriver { base };
1722 let data: Arc<dyn std::any::Any + Send + Sync> = Arc::new(99i32);
1723 let mut user = AsynUser::new(0);
1724 drv.write_generic_pointer(&mut user, data).unwrap();
1725 let val = drv.read_generic_pointer(&AsynUser::new(0)).unwrap();
1726 assert_eq!(*val.downcast_ref::<i32>().unwrap(), 99);
1727 }
1728
1729 #[test]
1730 fn test_generic_pointer_callback() {
1731 use crate::param::ParamValue;
1732
1733 let mut base = PortDriverBase::new("test_gp_cb", 1, PortFlags::default());
1734 base.create_param("PTR", ParamType::GenericPointer).unwrap();
1735 let mut rx = base.interrupts.subscribe_async();
1736
1737 struct GpDriver {
1738 base: PortDriverBase,
1739 }
1740 impl PortDriver for GpDriver {
1741 fn base(&self) -> &PortDriverBase {
1742 &self.base
1743 }
1744 fn base_mut(&mut self) -> &mut PortDriverBase {
1745 &mut self.base
1746 }
1747 }
1748
1749 let mut drv = GpDriver { base };
1750 let data: Arc<dyn std::any::Any + Send + Sync> = Arc::new(vec![1, 2, 3]);
1751 drv.base_mut()
1752 .set_generic_pointer_param(0, 0, data)
1753 .unwrap();
1754 drv.base_mut().call_param_callbacks(0).unwrap();
1755
1756 let v = rx.try_recv().unwrap();
1757 assert_eq!(v.reason, 0);
1758 assert!(matches!(v.value, ParamValue::GenericPointer(_)));
1759 }
1760
1761 #[test]
1762 fn test_interpose_push_requires_lock() {
1763 use crate::interpose::{OctetInterpose, OctetNext, OctetReadResult};
1764 use parking_lot::Mutex;
1765 use std::sync::Arc;
1766
1767 struct NoopInterpose;
1768 impl OctetInterpose for NoopInterpose {
1769 fn read(
1770 &mut self,
1771 user: &AsynUser,
1772 buf: &mut [u8],
1773 next: &mut dyn OctetNext,
1774 ) -> AsynResult<OctetReadResult> {
1775 next.read(user, buf)
1776 }
1777 fn write(
1778 &mut self,
1779 user: &mut AsynUser,
1780 data: &[u8],
1781 next: &mut dyn OctetNext,
1782 ) -> AsynResult<usize> {
1783 next.write(user, data)
1784 }
1785 fn flush(&mut self, user: &mut AsynUser, next: &mut dyn OctetNext) -> AsynResult<()> {
1786 next.flush(user)
1787 }
1788 }
1789
1790 let port: Arc<Mutex<dyn PortDriver>> = Arc::new(Mutex::new(TestDriver::new()));
1791
1792 {
1793 let mut guard = port.lock();
1794 guard
1795 .base_mut()
1796 .push_octet_interpose(Box::new(NoopInterpose));
1797 assert_eq!(guard.base().interpose_octet.len(), 1);
1798 }
1799 }
1800
1801 #[test]
1802 fn test_default_read_write_int64() {
1803 let mut base = PortDriverBase::new("test_i64", 1, PortFlags::default());
1804 base.create_param("BIG", ParamType::Int64).unwrap();
1805
1806 struct I64Driver {
1807 base: PortDriverBase,
1808 }
1809 impl PortDriver for I64Driver {
1810 fn base(&self) -> &PortDriverBase {
1811 &self.base
1812 }
1813 fn base_mut(&mut self) -> &mut PortDriverBase {
1814 &mut self.base
1815 }
1816 }
1817
1818 let mut drv = I64Driver { base };
1819 let mut user = AsynUser::new(0);
1820 drv.write_int64(&mut user, i64::MAX).unwrap();
1821 assert_eq!(drv.read_int64(&AsynUser::new(0)).unwrap(), i64::MAX);
1822 }
1823
1824 #[test]
1825 fn test_get_bounds_int64_default() {
1826 let base = PortDriverBase::new("test_bounds", 1, PortFlags::default());
1827 struct BoundsDriver {
1828 base: PortDriverBase,
1829 }
1830 impl PortDriver for BoundsDriver {
1831 fn base(&self) -> &PortDriverBase {
1832 &self.base
1833 }
1834 fn base_mut(&mut self) -> &mut PortDriverBase {
1835 &mut self.base
1836 }
1837 }
1838 let drv = BoundsDriver { base };
1839 let (lo, hi) = drv.get_bounds_int64(&AsynUser::default()).unwrap();
1840 assert_eq!(lo, 0);
1843 assert_eq!(hi, 0);
1844 }
1845
1846 #[test]
1847 fn test_per_addr_device_state() {
1848 let mut base = PortDriverBase::new(
1849 "multi",
1850 4,
1851 PortFlags {
1852 multi_device: true,
1853 can_block: false,
1854 destructible: true,
1855 },
1856 );
1857 base.create_param("V", ParamType::Int32).unwrap();
1858
1859 assert!(base.is_device_connected(0));
1861 assert!(base.is_device_connected(1));
1862
1863 base.device_state(1).enabled = false;
1865 assert!(base.check_ready_addr(0).is_ok());
1866 let err = base.check_ready_addr(1).unwrap_err();
1867 assert!(format!("{err}").contains("disabled"));
1868
1869 base.device_state(2).connected = false;
1871 let err = base.check_ready_addr(2).unwrap_err();
1872 assert!(format!("{err}").contains("disconnected"));
1873 }
1874
1875 #[test]
1876 fn test_per_addr_single_device_ignored() {
1877 let mut base = PortDriverBase::new("single", 1, PortFlags::default());
1878 base.create_param("V", ParamType::Int32).unwrap();
1879 assert!(base.check_ready_addr(0).is_ok());
1881 }
1882
1883 #[test]
1884 fn test_timestamp_source() {
1885 let mut base = PortDriverBase::new("ts_test", 1, PortFlags::default());
1886 base.create_param("V", ParamType::Int32).unwrap();
1887
1888 let fixed_ts = SystemTime::UNIX_EPOCH + std::time::Duration::from_secs(999999);
1889 base.register_timestamp_source(move || fixed_ts);
1890
1891 assert_eq!(base.current_timestamp(), fixed_ts);
1892 }
1893
1894 #[test]
1895 fn test_timestamp_source_in_callbacks() {
1896 let mut base = PortDriverBase::new("ts_cb", 1, PortFlags::default());
1897 base.create_param("V", ParamType::Int32).unwrap();
1898 let mut rx = base.interrupts.subscribe_async();
1899
1900 let fixed_ts = SystemTime::UNIX_EPOCH + std::time::Duration::from_secs(123456);
1901 base.register_timestamp_source(move || fixed_ts);
1902
1903 struct TsDriver {
1904 base: PortDriverBase,
1905 }
1906 impl PortDriver for TsDriver {
1907 fn base(&self) -> &PortDriverBase {
1908 &self.base
1909 }
1910 fn base_mut(&mut self) -> &mut PortDriverBase {
1911 &mut self.base
1912 }
1913 }
1914 let mut drv = TsDriver { base };
1915 drv.base_mut().set_int32_param(0, 0, 42).unwrap();
1916 drv.base_mut().call_param_callbacks(0).unwrap();
1917
1918 let v = rx.try_recv().unwrap();
1919 assert_eq!(v.timestamp, fixed_ts);
1921 }
1922
1923 #[test]
1924 fn test_queue_priority_connect() {
1925 assert!(QueuePriority::Connect > QueuePriority::High);
1926 }
1927
1928 #[test]
1929 fn test_port_flags_destructible_default_is_opt_in() {
1930 let flags = PortFlags::default();
1936 assert!(
1937 !flags.destructible,
1938 "destructible must be opt-in (C parity)"
1939 );
1940 }
1941
1942 #[test]
1943 fn shutdown_lifecycle_refuses_non_destructible() {
1944 let mut base = PortDriverBase::new(
1945 "p_nondestr",
1946 1,
1947 PortFlags {
1948 multi_device: false,
1949 can_block: false,
1950 destructible: false,
1951 },
1952 );
1953 match base.shutdown_lifecycle() {
1954 Err(AsynError::Status { message, .. }) => {
1955 assert!(message.contains("ASYN_DESTRUCTIBLE"), "msg={message}");
1956 }
1957 other => panic!("expected ASYN_DESTRUCTIBLE refusal, got {other:?}"),
1958 }
1959 assert!(
1960 !base.is_defunct(),
1961 "non-destructible port must not flip defunct"
1962 );
1963 assert!(base.is_enabled(), "non-destructible port must stay enabled");
1964 }
1965
1966 #[test]
1967 fn shutdown_lifecycle_marks_destructible_defunct_and_idempotent() {
1968 let mut base = PortDriverBase::new(
1969 "p_destr",
1970 1,
1971 PortFlags {
1972 multi_device: false,
1973 can_block: false,
1974 destructible: true,
1975 },
1976 );
1977 assert!(base.is_enabled());
1978 assert!(!base.is_defunct());
1979 base.shutdown_lifecycle().unwrap();
1980 assert!(
1981 !base.is_enabled(),
1982 "shutdown_lifecycle must flip enabled=false"
1983 );
1984 assert!(
1985 base.is_defunct(),
1986 "shutdown_lifecycle must flip defunct=true"
1987 );
1988 base.shutdown_lifecycle().unwrap();
1990 assert!(base.is_defunct());
1991 match base.check_ready() {
1993 Err(AsynError::Status { message, .. }) => {
1994 assert!(message.contains("defunct"), "msg={message}");
1995 }
1996 other => panic!("expected defunct error, got {other:?}"),
1997 }
1998 }
1999
2000 #[test]
2003 fn test_connect_addr() {
2004 let mut base = PortDriverBase::new(
2005 "multi_conn",
2006 4,
2007 PortFlags {
2008 multi_device: true,
2009 can_block: false,
2010 destructible: true,
2011 },
2012 );
2013 base.create_param("V", ParamType::Int32).unwrap();
2014
2015 base.disconnect_addr(1);
2016 assert!(!base.is_device_connected(1));
2017 assert!(base.check_ready_addr(1).is_err());
2018
2019 base.connect_addr(1);
2020 assert!(base.is_device_connected(1));
2021 assert!(base.check_ready_addr(1).is_ok());
2022 }
2023
2024 #[test]
2025 fn test_enable_disable_addr() {
2026 let mut base = PortDriverBase::new(
2027 "multi_en",
2028 4,
2029 PortFlags {
2030 multi_device: true,
2031 can_block: false,
2032 destructible: true,
2033 },
2034 );
2035 base.create_param("V", ParamType::Int32).unwrap();
2036
2037 base.disable_addr(2);
2038 let err = base.check_ready_addr(2).unwrap_err();
2039 assert!(format!("{err}").contains("disabled"));
2040
2041 base.enable_addr(2);
2042 assert!(base.check_ready_addr(2).is_ok());
2043 }
2044
2045 #[test]
2046 fn test_port_level_overrides_addr() {
2047 let mut base = PortDriverBase::new(
2048 "multi_override",
2049 4,
2050 PortFlags {
2051 multi_device: true,
2052 can_block: false,
2053 destructible: true,
2054 },
2055 );
2056 base.create_param("V", ParamType::Int32).unwrap();
2057
2058 base.enabled = false;
2060 base.enable_addr(0); let err = base.check_ready_addr(0).unwrap_err();
2062 assert!(format!("{err}").contains("disabled"));
2063 }
2064
2065 #[test]
2066 fn test_per_addr_exception_announced() {
2067 use std::sync::atomic::{AtomicI32, Ordering};
2068
2069 let mut base = PortDriverBase::new(
2070 "multi_exc",
2071 4,
2072 PortFlags {
2073 multi_device: true,
2074 can_block: false,
2075 destructible: true,
2076 },
2077 );
2078 base.create_param("V", ParamType::Int32).unwrap();
2079
2080 let exc_mgr = Arc::new(crate::exception::ExceptionManager::new());
2081 base.exception_sink = Some(exc_mgr.clone());
2082
2083 let last_addr = Arc::new(AtomicI32::new(-99));
2084 let last_addr2 = last_addr.clone();
2085 exc_mgr.add_callback(move |event| {
2086 last_addr2.store(event.addr, Ordering::Relaxed);
2087 });
2088
2089 base.disconnect_addr(3);
2090 assert_eq!(last_addr.load(Ordering::Relaxed), 3);
2091
2092 base.enable_addr(2);
2093 assert_eq!(last_addr.load(Ordering::Relaxed), 2);
2094 }
2095
2096 #[test]
2103 fn test_connect_disconnect_announce_only_on_transition() {
2104 use std::sync::atomic::{AtomicUsize, Ordering};
2105
2106 let mut base = PortDriverBase::new(
2107 "edge",
2108 4,
2109 PortFlags {
2110 multi_device: true,
2111 can_block: false,
2112 destructible: true,
2113 },
2114 );
2115 base.create_param("V", ParamType::Int32).unwrap();
2116 let exc_mgr = Arc::new(crate::exception::ExceptionManager::new());
2117 base.exception_sink = Some(exc_mgr.clone());
2118
2119 let connect_hits = Arc::new(AtomicUsize::new(0));
2120 let hits2 = connect_hits.clone();
2121 exc_mgr.add_callback(move |event| {
2122 if event.exception == AsynException::Connect {
2123 hits2.fetch_add(1, Ordering::Relaxed);
2124 }
2125 });
2126
2127 base.connect_addr(2);
2130 assert_eq!(
2131 connect_hits.load(Ordering::Relaxed),
2132 0,
2133 "redundant connect_addr must not fan out"
2134 );
2135
2136 base.disconnect_addr(2);
2138 assert_eq!(connect_hits.load(Ordering::Relaxed), 1);
2139
2140 base.disconnect_addr(2);
2142 assert_eq!(
2143 connect_hits.load(Ordering::Relaxed),
2144 1,
2145 "redundant disconnect_addr must not fan out"
2146 );
2147
2148 base.connect_addr(2);
2150 assert_eq!(connect_hits.load(Ordering::Relaxed), 2);
2151 }
2152
2153 #[test]
2158 fn test_set_auto_connect_fires_unconditionally() {
2159 use std::sync::atomic::{AtomicUsize, Ordering};
2160
2161 let mut base = PortDriverBase::new("ac", 1, PortFlags::default());
2162 let exc_mgr = Arc::new(crate::exception::ExceptionManager::new());
2163 base.exception_sink = Some(exc_mgr.clone());
2164 let hits = Arc::new(AtomicUsize::new(0));
2165 let hits2 = hits.clone();
2166 exc_mgr.add_callback(move |event| {
2167 if event.exception == AsynException::AutoConnect {
2168 hits2.fetch_add(1, Ordering::Relaxed);
2169 }
2170 });
2171 base.set_auto_connect(true);
2174 base.set_auto_connect(false);
2175 base.set_auto_connect(false);
2176 assert_eq!(hits.load(Ordering::Relaxed), 3);
2177 }
2178
2179 #[test]
2180 fn auto_connect_throttle_gate_boundaries() {
2181 let mut base = PortDriverBase::new("thr", 1, PortFlags::default());
2184
2185 let t0 = Instant::now();
2188 assert!(base.auto_connect_throttle_ok(-1, t0));
2189
2190 base.last_connect_disconnect = Some(t0);
2193 assert!(!base.auto_connect_throttle_ok(-1, t0));
2195 assert!(!base.auto_connect_throttle_ok(-1, t0 + Duration::from_millis(1999)));
2197 assert!(base.auto_connect_throttle_ok(-1, t0 + Duration::from_secs(2)));
2199 assert!(base.auto_connect_throttle_ok(-1, t0 + Duration::from_secs(5)));
2201 }
2202
2203 #[test]
2204 fn auto_connect_throttle_stamps_on_disconnect_not_connect() {
2205 let mut base = PortDriverBase::new("thr", 1, PortFlags::default());
2208 assert!(base.last_connect_disconnect.is_none());
2210
2211 assert!(base.set_connected(false));
2213 assert!(base.last_connect_disconnect.is_some());
2214
2215 base.last_connect_disconnect = None;
2217 assert!(base.set_connected(true));
2218 assert!(base.last_connect_disconnect.is_none());
2219 }
2220
2221 #[test]
2222 fn auto_connect_throttle_per_device_anchor() {
2223 let flags = PortFlags {
2225 multi_device: true,
2226 ..PortFlags::default()
2227 };
2228 let mut base = PortDriverBase::new("thr", 4, flags);
2229 let t0 = Instant::now();
2230
2231 assert!(base.set_addr_connected(1, false));
2233 assert!(base.device_state(1).last_connect_disconnect.is_some());
2234 assert!(!base.auto_connect_throttle_ok(1, t0));
2236 assert!(base.auto_connect_throttle_ok(2, t0));
2237
2238 base.stamp_auto_connect_attempt(2, t0);
2240 assert!(!base.auto_connect_throttle_ok(2, t0));
2241 assert!(base.auto_connect_throttle_ok(2, t0 + Duration::from_secs(2)));
2242 }
2243
2244 #[test]
2245 fn set_enabled_refuses_defunct_port() {
2246 use std::sync::atomic::{AtomicUsize, Ordering};
2247 let flags = PortFlags {
2250 destructible: true,
2251 ..PortFlags::default()
2252 };
2253 let mut base = PortDriverBase::new("def", 1, flags);
2254 let exc_mgr = Arc::new(crate::exception::ExceptionManager::new());
2255 base.exception_sink = Some(exc_mgr.clone());
2256 let enable_hits = Arc::new(AtomicUsize::new(0));
2257 let h = enable_hits.clone();
2258 exc_mgr.add_callback(move |event| {
2259 if event.exception == AsynException::Enable {
2260 h.fetch_add(1, Ordering::Relaxed);
2261 }
2262 });
2263
2264 base.shutdown_lifecycle().unwrap();
2266 assert!(base.is_defunct());
2267 assert!(!base.is_enabled());
2268
2269 let err = base.set_enabled(true).unwrap_err();
2270 match err {
2271 AsynError::Status { status, .. } => assert_eq!(status, AsynStatus::Disabled),
2272 other => panic!("expected Disabled, got {other:?}"),
2273 }
2274 assert!(!base.is_enabled(), "defunct port must not re-enable");
2275 assert_eq!(
2276 enable_hits.load(Ordering::Relaxed),
2277 0,
2278 "no Enable exception may fire on a defunct port"
2279 );
2280 }
2281
2282 #[test]
2283 fn set_addr_enabled_refuses_defunct_port() {
2284 use std::sync::atomic::{AtomicUsize, Ordering};
2285 let flags = PortFlags {
2286 multi_device: true,
2287 destructible: true,
2288 ..PortFlags::default()
2289 };
2290 let mut base = PortDriverBase::new("def", 4, flags);
2291 let exc_mgr = Arc::new(crate::exception::ExceptionManager::new());
2292 base.exception_sink = Some(exc_mgr.clone());
2293 let enable_hits = Arc::new(AtomicUsize::new(0));
2294 let h = enable_hits.clone();
2295 exc_mgr.add_callback(move |event| {
2296 if event.exception == AsynException::Enable {
2297 h.fetch_add(1, Ordering::Relaxed);
2298 }
2299 });
2300
2301 base.shutdown_lifecycle().unwrap();
2302
2303 let err = base.set_addr_enabled(1, false).unwrap_err();
2304 match err {
2305 AsynError::Status { status, .. } => assert_eq!(status, AsynStatus::Disabled),
2306 other => panic!("expected Disabled, got {other:?}"),
2307 }
2308 assert!(
2311 !base.device_states.contains_key(&1),
2312 "refused per-device enable must not create device state"
2313 );
2314 base.disable_addr(1);
2316 assert!(!base.device_states.contains_key(&1));
2317 assert_eq!(
2318 enable_hits.load(Ordering::Relaxed),
2319 0,
2320 "no Enable exception may fire on a defunct port"
2321 );
2322 }
2323
2324 #[test]
2330 fn test_port_driver_uint32_interrupt_round_trip() {
2331 struct UInt32Drv {
2332 base: PortDriverBase,
2333 }
2334 impl PortDriver for UInt32Drv {
2335 fn base(&self) -> &PortDriverBase {
2336 &self.base
2337 }
2338 fn base_mut(&mut self) -> &mut PortDriverBase {
2339 &mut self.base
2340 }
2341 }
2342
2343 let mut base = PortDriverBase::new("uint32_int", 1, PortFlags::default());
2344 let idx = base
2345 .params
2346 .create_param("BITS", ParamType::UInt32Digital)
2347 .unwrap();
2348 let mut drv = UInt32Drv { base };
2349 let user = AsynUser::new(idx).with_addr(0);
2350
2351 drv.set_interrupt_uint32_digital(&user, 0xF0, InterruptReason::ZeroToOne)
2352 .unwrap();
2353 drv.set_interrupt_uint32_digital(&user, 0x0F, InterruptReason::OneToZero)
2354 .unwrap();
2355 assert_eq!(
2356 drv.get_interrupt_uint32_digital(&user, InterruptReason::Both)
2357 .unwrap(),
2358 0xFF
2359 );
2360 drv.clear_interrupt_uint32_digital(&user, 0x11).unwrap();
2361 assert_eq!(
2362 drv.get_interrupt_uint32_digital(&user, InterruptReason::ZeroToOne)
2363 .unwrap(),
2364 0xE0
2365 );
2366 assert_eq!(
2367 drv.get_interrupt_uint32_digital(&user, InterruptReason::OneToZero)
2368 .unwrap(),
2369 0x0E
2370 );
2371 }
2372
2373 #[test]
2384 fn default_scalar_reads_report_undefined_until_set() {
2385 struct AllTypesDrv {
2386 base: PortDriverBase,
2387 }
2388 impl PortDriver for AllTypesDrv {
2389 fn base(&self) -> &PortDriverBase {
2390 &self.base
2391 }
2392 fn base_mut(&mut self) -> &mut PortDriverBase {
2393 &mut self.base
2394 }
2395 }
2396
2397 let mut base = PortDriverBase::new("undef_read", 1, PortFlags::default());
2398 let i32_idx = base.params.create_param("I32", ParamType::Int32).unwrap();
2399 let i64_idx = base.params.create_param("I64", ParamType::Int64).unwrap();
2400 let f64_idx = base.params.create_param("F64", ParamType::Float64).unwrap();
2401 let oct_idx = base.params.create_param("OCT", ParamType::Octet).unwrap();
2402 let u32_idx = base
2403 .params
2404 .create_param("BITS", ParamType::UInt32Digital)
2405 .unwrap();
2406 let mut drv = AllTypesDrv { base };
2407
2408 assert!(matches!(
2410 drv.read_int32(&AsynUser::new(i32_idx).with_addr(0)),
2411 Err(AsynError::ParamUndefined(_))
2412 ));
2413 assert!(matches!(
2414 drv.read_int64(&AsynUser::new(i64_idx).with_addr(0)),
2415 Err(AsynError::ParamUndefined(_))
2416 ));
2417 assert!(matches!(
2418 drv.read_float64(&AsynUser::new(f64_idx).with_addr(0)),
2419 Err(AsynError::ParamUndefined(_))
2420 ));
2421 let mut buf = [0u8; 16];
2422 assert!(matches!(
2423 drv.read_octet(&AsynUser::new(oct_idx).with_addr(0), &mut buf),
2424 Err(AsynError::ParamUndefined(_))
2425 ));
2426 assert!(matches!(
2427 drv.read_uint32_digital(&AsynUser::new(u32_idx).with_addr(0), 0xFFFF_FFFF),
2428 Err(AsynError::ParamUndefined(_))
2429 ));
2430
2431 drv.base_mut().params.set_int32(i32_idx, 0, 7).unwrap();
2433 drv.base_mut().params.set_int64(i64_idx, 0, 9).unwrap();
2434 drv.base_mut().params.set_float64(f64_idx, 0, 1.5).unwrap();
2435 drv.base_mut()
2436 .params
2437 .set_string(oct_idx, 0, "hi".to_string())
2438 .unwrap();
2439 drv.base_mut()
2440 .params
2441 .set_uint32(u32_idx, 0, 0x05, 0xFFFF_FFFF, 0)
2442 .unwrap();
2443
2444 assert_eq!(
2445 drv.read_int32(&AsynUser::new(i32_idx).with_addr(0))
2446 .unwrap(),
2447 7
2448 );
2449 assert_eq!(
2450 drv.read_int64(&AsynUser::new(i64_idx).with_addr(0))
2451 .unwrap(),
2452 9
2453 );
2454 assert_eq!(
2455 drv.read_float64(&AsynUser::new(f64_idx).with_addr(0))
2456 .unwrap(),
2457 1.5
2458 );
2459 let n = drv
2460 .read_octet(&AsynUser::new(oct_idx).with_addr(0), &mut buf)
2461 .unwrap();
2462 assert_eq!(&buf[..n], b"hi");
2463 assert_eq!(
2464 drv.read_uint32_digital(&AsynUser::new(u32_idx).with_addr(0), 0xFFFF_FFFF)
2465 .unwrap(),
2466 0x05
2467 );
2468 }
2469}