pub mod device;
mod io_intr;
pub mod registry;
pub use device::{ASYN_RECORD_DTYP, AsynRecordDevice};
pub use registry::{
PortEntry, PortRegistry, asyn_record_factory, get_port, port_names, register_asyn_record_type,
register_port, unregister_port,
};
use io_intr::{IoIntrBinding, IoIntrSample, IoIntrScan};
use std::collections::HashMap;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use epics_base_rs::error::{CaError, CaResult};
use epics_base_rs::server::database::AsyncDbHandle;
use epics_base_rs::server::recgbl::{alarm_status, rec_gbl_set_sevr};
use epics_base_rs::server::record::{
AlarmSeverity, CommonFields, FieldDeclaration, ProcessOutcome, Record, RecordProcessResult,
ScanType,
};
use epics_base_rs::types::{DbFieldType, EpicsValue};
use crate::error::{AsynError, AsynResult, AsynStatus};
use crate::exception::{AsynException, ExceptionCallbackId, ExceptionManager};
use crate::interpose::EomReason;
use crate::port_handle::PortHandle;
use crate::request::{CancelToken, RequestOp, RequestResult};
use crate::trace::{TraceFile, TraceInfoMask, TraceIoMask, TraceManager, TraceMask};
use crate::user::AsynUser;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u16)]
enum TransferMode {
WriteRead = 0,
Write = 1,
Read = 2,
Flush = 3,
NoIo = 4,
}
impl TransferMode {
fn from_u16(v: u16) -> Self {
match v {
0 => Self::WriteRead,
1 => Self::Write,
2 => Self::Read,
3 => Self::Flush,
4 => Self::NoIo,
_ => Self::WriteRead,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u16)]
pub(crate) enum InterfaceType {
Octet = 0,
Int32 = 1,
UInt32Digital = 2,
Float64 = 3,
}
impl InterfaceType {
fn from_u16(v: u16) -> Self {
match v {
0 => Self::Octet,
1 => Self::Int32,
2 => Self::UInt32Digital,
3 => Self::Float64,
_ => Self::Octet,
}
}
fn as_asyn_iface(self) -> crate::interfaces::InterfaceType {
match self {
Self::Octet => crate::interfaces::InterfaceType::Octet,
Self::Int32 => crate::interfaces::InterfaceType::Int32,
Self::UInt32Digital => crate::interfaces::InterfaceType::UInt32Digital,
Self::Float64 => crate::interfaces::InterfaceType::Float64,
}
}
fn c_errs_name(self) -> &'static str {
match self {
Self::Octet => "Octet",
Self::Int32 => "Int32",
Self::UInt32Digital => "UInt32",
Self::Float64 => "Float64",
}
}
fn c_asyn_name(self) -> &'static str {
match self {
Self::Octet => "asynOctet",
Self::Int32 => "asynInt32",
Self::UInt32Digital => "asynUInt32Digital",
Self::Float64 => "asynFloat64",
}
}
fn registry_type(self) -> crate::interfaces::InterfaceType {
match self {
Self::Octet => crate::interfaces::InterfaceType::Octet,
Self::Int32 => crate::interfaces::InterfaceType::Int32,
Self::UInt32Digital => crate::interfaces::InterfaceType::UInt32Digital,
Self::Float64 => crate::interfaces::InterfaceType::Float64,
}
}
}
const BAUD_CHOICES: &[&str] = &[
"Unknown", "300", "600", "1200", "2400", "4800", "9600", "19200", "38400", "57600", "115200",
"230400", "460800", "576000", "921600", "1152000",
];
const PARITY_CHOICES: &[&str] = &["Unknown", "none", "even", "odd"];
const DBIT_CHOICES: &[&str] = &["Unknown", "5", "6", "7", "8"];
const SBIT_CHOICES: &[&str] = &["Unknown", "1", "2"];
const MCTL_CHOICES: &[&str] = &["Unknown", "Y", "N"];
const FCTL_CHOICES: &[&str] = &["Unknown", "N", "Y"];
const IX_CHOICES: &[&str] = &["Unknown", "N", "Y"];
const DRTO_CHOICES: &[&str] = &["Unknown", "N", "Y"];
fn menu_choice(choices: &'static [&'static str], index: i32) -> &'static str {
usize::try_from(index)
.ok()
.and_then(|i| choices.get(i).copied())
.unwrap_or(choices[0])
}
fn baud_choice_index(text: &str) -> i32 {
BAUD_CHOICES
.iter()
.position(|choice| *choice == text)
.unwrap_or(0) as i32
}
const ASYN_FMT_ASCII: i32 = 0;
const ASYN_FMT_HYBRID: i32 = 1;
const ASYN_FMT_BINARY: i32 = 2;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum AsynArrayField {
Bout,
Binp,
}
const AINP_SIZE: usize = 40;
pub(super) const TINP_SIZE: usize = 40;
pub(super) const EOS_SIZE: usize = 10;
pub(crate) fn translate_escape(s: &str) -> Vec<u8> {
translate_escape_bytes(s.as_bytes())
}
fn translate_escape_bytes(input: &[u8]) -> Vec<u8> {
let mut out = Vec::with_capacity(input.len());
let mut chars = input.iter().copied().peekable();
while let Some(c) = chars.next() {
if c != b'\\' {
out.push(c);
continue;
}
let Some(next) = chars.next() else {
out.push(b'\\');
break;
};
let decoded = match next {
b'r' => 0x0D,
b'n' => 0x0A,
b't' => 0x09,
b'\\' => b'\\',
b'"' => b'"',
b'\'' => b'\'',
b'0'..=b'7' => {
let mut val = u32::from(next - b'0');
for _ in 0..2 {
match chars.peek() {
Some(&d) if (b'0'..=b'7').contains(&d) => {
val = val * 8 + u32::from(d - b'0');
chars.next();
}
_ => break,
}
}
out.push((val & 0xFF) as u8);
continue;
}
b'a' => 0x07,
b'b' => 0x08,
b'f' => 0x0C,
b'v' => 0x0B,
other => {
out.push(b'\\');
out.push(other);
continue;
}
};
out.push(decoded);
}
out
}
fn open_trace_file(tfil: &str) -> std::io::Result<TraceFile> {
match tfil {
"" | "<stdout>" => Ok(TraceFile::Stdout),
"<stderr>" => Ok(TraceFile::Stderr),
"<errlog>" => Ok(TraceFile::Errlog),
path => std::fs::OpenOptions::new()
.append(true)
.create(true)
.open(path)
.map(|f| TraceFile::File(Arc::new(std::sync::Mutex::new(f)))),
}
}
pub struct AsynRecord {
pub port: String,
pub addr: i32,
pub pcnct: i32, pub drvinfo: String,
pub reason: i32,
pub tmod: i32, pub tmot: f64, pub iface: i32, pub octetiv: i32, pub optioniv: i32, pub gpibiv: i32, pub i32iv: i32, pub ui32iv: i32, pub f64iv: i32,
pub aout: String,
pub oeos: String,
pub bout: Vec<u8>,
pub omax: i32,
pub nowt: i32,
pub nawt: i32,
pub ofmt: i32,
pub ainp: String,
pub tinp: String,
pub ieos: String,
pub binp: Vec<u8>,
pub imax: i32,
pub nrrd: i32,
pub nord: i32,
pub ifmt: i32, pub eomr: i32,
pub i32inp: i32,
pub i32out: i32,
pub ui32inp: u32,
pub ui32out: u32,
pub ui32mask: u32,
pub f64inp: f64,
pub f64out: f64,
pub baud: i32,
pub lbaud: i32,
pub prty: i32,
pub dbit: i32,
pub sbit: i32,
pub mctl: i32,
pub fctl: i32,
pub ixon: i32,
pub ixoff: i32,
pub ixany: i32,
pub hostinfo: String,
pub drto: i32,
pub ucmd: i32,
pub acmd: i32,
pub spr: i32,
pub tmsk: i32,
pub tb0: i32,
pub tb1: i32,
pub tb2: i32,
pub tb3: i32,
pub tb4: i32,
pub tb5: i32,
pub tiom: i32,
pub tib0: i32,
pub tib1: i32,
pub tib2: i32,
pub tinm: i32,
pub tinb0: i32,
pub tinb1: i32,
pub tinb2: i32,
pub tinb3: i32,
pub tsiz: i32,
pub tfil: String,
pub auct: i32, pub cnct: i32, pub enbl: i32,
pub val: i32,
pub errs: String,
pub aqr: i32,
port_entry: Option<PortEntry>,
resolved_reason: usize,
async_ctx: Option<(String, AsyncDbHandle)>,
io_inflight: Option<IoInFlight>,
status_dirty: Arc<AtomicBool>,
except_cb: Option<(Arc<ExceptionManager>, ExceptionCallbackId)>,
io_alarm: Option<(u16, AlarmSeverity)>,
old_trace_file_id: Arc<Mutex<Option<usize>>>,
io_intr: Arc<IoIntrScan>,
}
impl Default for AsynRecord {
fn default() -> Self {
Self {
port: String::new(),
addr: 0,
pcnct: 0,
drvinfo: String::new(),
reason: 0,
tmod: 0,
tmot: 1.0,
iface: 0,
octetiv: 0,
optioniv: 0,
gpibiv: 0,
i32iv: 0,
ui32iv: 0,
f64iv: 0,
aout: String::new(),
oeos: String::new(),
bout: Vec::new(),
omax: 80,
nowt: 80,
nawt: 0,
ofmt: 0,
ainp: String::new(),
tinp: String::new(),
ieos: String::new(),
binp: Vec::new(),
imax: 80,
nrrd: 0,
nord: 0,
ifmt: 0,
eomr: 0,
i32inp: 0,
i32out: 0,
ui32inp: 0,
ui32out: 0,
ui32mask: 0xFFFFFFFF,
f64inp: 0.0,
f64out: 0.0,
baud: 0,
lbaud: 0,
prty: 0,
dbit: 0,
sbit: 0,
mctl: 0,
fctl: 0,
ixon: 0,
ixoff: 0,
ixany: 0,
hostinfo: String::new(),
drto: 0,
ucmd: 0,
acmd: 0,
spr: 0,
tmsk: 0,
tb0: 0,
tb1: 0,
tb2: 0,
tb3: 0,
tb4: 0,
tb5: 0,
tiom: 0,
tib0: 0,
tib1: 0,
tib2: 0,
tinm: 0,
tinb0: 0,
tinb1: 0,
tinb2: 0,
tinb3: 0,
tsiz: 80,
tfil: TFIL_UNKNOWN.to_string(),
auct: 1,
cnct: 0,
enbl: 1,
val: 0,
errs: String::new(),
aqr: 0,
port_entry: None,
resolved_reason: 0,
async_ctx: None,
io_inflight: None,
status_dirty: Arc::new(AtomicBool::new(false)),
except_cb: None,
io_alarm: None,
old_trace_file_id: Arc::new(Mutex::new(None)),
io_intr: Arc::new(IoIntrScan::new()),
}
}
}
struct IoInFlight {
cancel: CancelToken,
result: Arc<Mutex<Option<IoOutcome>>>,
}
#[derive(Clone, PartialEq, Eq, Debug)]
enum GpibCycle {
Universal(u8),
Addressed(Vec<u8>),
SerialPoll,
}
struct IoPlan {
tmod: TransferMode,
iface: InterfaceType,
gpib: Option<GpibCycle>,
reason: usize,
addr: i32,
timeout: std::time::Duration,
octet_out: Vec<u8>,
octet_out_len: usize,
ofmt: i32,
i32out: i32,
ui32out: u32,
ui32mask: u32,
f64out: f64,
octet_buf_size: usize,
in_len: usize,
ifmt: i32,
}
#[derive(Default)]
struct IoOutcome {
nawt: Option<i32>,
eomr: Option<i32>,
nord: Option<i32>,
tinp: Option<String>,
ainp: Option<String>,
binp: Option<Vec<u8>>,
i32inp: Option<i32>,
ui32inp: Option<u32>,
f64inp: Option<f64>,
spr: Option<i32>,
errs: Option<String>,
alarm: Option<(u16, AlarmSeverity)>,
}
impl IoOutcome {
fn report_error(&mut self, msg: String) {
self.errs = Some(msg);
}
fn report_canceled(&mut self) {
self.report_error(CANCELED_MSG.to_string());
raise_io_alarm(self, alarm_status::STATE_ALARM, AlarmSeverity::Major);
}
fn report_queue_timeout(&mut self) {
self.report_error(PROCESS_QUEUE_TIMEOUT_MSG.to_string());
raise_io_alarm(self, alarm_status::STATE_ALARM, AlarmSeverity::Major);
}
fn report_queue_refused(&mut self) {
self.report_error(PROCESS_QUEUE_REFUSED_MSG.to_string());
raise_io_alarm(self, alarm_status::STATE_ALARM, AlarmSeverity::Minor);
}
}
fn c_read_status_word(e: &crate::error::AsynError) -> &'static str {
use crate::error::AsynStatus;
match e.status() {
AsynStatus::Timeout => "timeout",
AsynStatus::Overflow => "overflow",
_ => "error",
}
}
fn raise_io_alarm(out: &mut IoOutcome, stat: u16, sevr: AlarmSeverity) {
let higher = match out.alarm {
Some((_, cur)) => (sevr as u16) > (cur as u16),
None => true,
};
if higher {
out.alarm = Some((stat, sevr));
}
}
fn io_user(plan: &IoPlan) -> AsynUser {
AsynUser::new(plan.reason)
.with_addr(plan.addr)
.with_timeout(plan.timeout)
.with_queue_timeout(QUEUE_TIMEOUT)
}
fn flush_user(plan: &IoPlan) -> AsynUser {
AsynUser::new(plan.reason)
.with_addr(plan.addr)
.with_queue_timeout(QUEUE_TIMEOUT)
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
enum IoPhase {
Flush,
Write,
Read,
GpibUniversal(u8),
GpibAddressed,
GpibSerialPollEnable,
GpibSerialPollRead,
GpibSerialPollDisable,
}
fn io_phases(plan: &IoPlan) -> Vec<IoPhase> {
if let Some(cycle) = &plan.gpib {
return match cycle {
GpibCycle::Universal(cmd) => vec![IoPhase::GpibUniversal(*cmd)],
GpibCycle::Addressed(_) => vec![IoPhase::GpibAddressed],
GpibCycle::SerialPoll => vec![
IoPhase::GpibSerialPollEnable,
IoPhase::GpibSerialPollRead,
IoPhase::GpibSerialPollDisable,
],
};
}
let mut phases = Vec::with_capacity(3);
if plan.iface == InterfaceType::Octet
&& matches!(plan.tmod, TransferMode::Flush | TransferMode::WriteRead)
{
phases.push(IoPhase::Flush);
}
if matches!(plan.tmod, TransferMode::Write | TransferMode::WriteRead) {
phases.push(IoPhase::Write);
}
if matches!(plan.tmod, TransferMode::Read | TransferMode::WriteRead) {
phases.push(IoPhase::Read);
}
phases
}
fn io_write_op(plan: &IoPlan) -> RequestOp {
match plan.iface {
InterfaceType::Octet => {
if plan.ofmt == ASYN_FMT_BINARY {
RequestOp::OctetWriteBinary {
data: plan.octet_out.clone(),
}
} else {
RequestOp::OctetWrite {
data: plan.octet_out.clone(),
}
}
}
InterfaceType::Int32 => RequestOp::Int32Write { value: plan.i32out },
InterfaceType::UInt32Digital => RequestOp::UInt32DigitalWrite {
value: plan.ui32out,
mask: plan.ui32mask,
},
InterfaceType::Float64 => RequestOp::Float64Write { value: plan.f64out },
}
}
fn io_read_op(plan: &IoPlan) -> RequestOp {
match plan.iface {
InterfaceType::Octet => {
if plan.ifmt == ASYN_FMT_BINARY {
RequestOp::OctetReadBinary {
buf_size: plan.octet_buf_size,
}
} else {
RequestOp::OctetRead {
buf_size: plan.octet_buf_size,
}
}
}
InterfaceType::Int32 => RequestOp::Int32Read,
InterfaceType::UInt32Digital => RequestOp::UInt32DigitalRead {
mask: plan.ui32mask,
},
InterfaceType::Float64 => RequestOp::Float64Read,
}
}
fn io_phase_op(plan: &IoPlan, phase: IoPhase) -> RequestOp {
match phase {
IoPhase::Flush => RequestOp::Flush,
IoPhase::Write => io_write_op(plan),
IoPhase::Read => io_read_op(plan),
IoPhase::GpibUniversal(cmd) => RequestOp::GpibUniversalCmd { cmd },
IoPhase::GpibAddressed => RequestOp::GpibAddressedCmd {
data: match &plan.gpib {
Some(GpibCycle::Addressed(frame)) => frame.clone(),
_ => Vec::new(),
},
},
IoPhase::GpibSerialPollEnable => RequestOp::GpibUniversalCmd {
cmd: crate::interfaces::gpib::IBSPE,
},
IoPhase::GpibSerialPollRead => RequestOp::OctetRead { buf_size: 1 },
IoPhase::GpibSerialPollDisable => RequestOp::GpibUniversalCmd {
cmd: crate::interfaces::gpib::IBSPD,
},
}
}
fn io_phase_user(plan: &IoPlan, phase: IoPhase) -> AsynUser {
match phase {
IoPhase::Flush => flush_user(plan),
_ => io_user(plan),
}
}
fn record_write_result(plan: &IoPlan, out: &mut IoOutcome, res: AsynResult<RequestResult>) {
if plan.iface != InterfaceType::Octet {
if let Err(e) = res {
out.report_error(format!(
"{} write error, {}",
plan.iface.c_errs_name(),
e.message()
));
raise_io_alarm(out, alarm_status::WRITE_ALARM, AlarmSeverity::Major);
}
return;
}
let (nawt, err) = match res {
Ok(result) => (result.nbytes, None),
Err(e) => (e.partial_write().unwrap_or(0), Some(e)),
};
out.nawt = Some(nawt as i32);
if err.is_some() || nawt != plan.octet_out_len {
let detail = match &err {
Some(e) => e.message(),
None => format!("wrote {} of {} chars", nawt, plan.octet_out_len),
};
out.report_error(format!("Write error, nout={nawt}, {detail}"));
}
}
fn record_read_result(plan: &IoPlan, out: &mut IoOutcome, res: AsynResult<RequestResult>) {
match plan.iface {
InterfaceType::Octet => {
let (data, eom, err) = match res {
Ok(result) => (
result.data.unwrap_or_default(),
result.eom_reason,
None::<crate::error::AsynError>,
),
Err(e) => {
let (data, eom) = match e.partial_read() {
Some(p) => (p.data.clone(), p.eom_reason.bits()),
None => (Vec::new(), 0),
};
(data, eom, Some(e))
}
};
let nread = data.len();
out.eomr = Some(eom as i32);
out.nord = Some(nread as i32);
let overflow = match plan.ifmt {
ASYN_FMT_BINARY => data.len() > plan.in_len,
_ => plan.in_len > 0 && data.len() >= plan.in_len,
};
out.tinp = Some(crate::escape::escaped_from_raw(&data, TINP_SIZE));
if plan.ifmt == ASYN_FMT_ASCII {
let bytes = if overflow {
&data[..plan.in_len - 1]
} else {
&data[..]
};
out.ainp = Some(String::from_utf8_lossy(bytes).to_string());
} else {
let mut data = data;
if overflow && plan.ifmt == ASYN_FMT_HYBRID {
data[plan.in_len - 1] = 0;
}
out.binp = Some(data);
}
if let Some(e) = &err {
out.report_error(format!(
"{} nread {} {}",
c_read_status_word(e),
nread,
e.message()
));
raise_io_alarm(out, alarm_status::READ_ALARM, AlarmSeverity::Major);
}
if overflow {
let tail = err.as_ref().map(|e| e.message()).unwrap_or_default();
out.report_error(format!("Overflow nread {nread} {tail}"));
raise_io_alarm(out, alarm_status::READ_ALARM, AlarmSeverity::Minor);
}
}
InterfaceType::Int32 => match res {
Ok(result) => {
if let Some(v) = result.int_val {
out.i32inp = Some(v);
}
}
Err(e) => {
out.report_error(format!(
"{} read error, {}",
plan.iface.c_errs_name(),
e.message()
));
raise_io_alarm(out, alarm_status::READ_ALARM, AlarmSeverity::Major);
}
},
InterfaceType::UInt32Digital => match res {
Ok(result) => {
if let Some(v) = result.uint_val {
out.ui32inp = Some(v);
}
}
Err(e) => {
out.report_error(format!(
"{} read error, {}",
plan.iface.c_errs_name(),
e.message()
));
raise_io_alarm(out, alarm_status::READ_ALARM, AlarmSeverity::Major);
}
},
InterfaceType::Float64 => match res {
Ok(result) => {
if let Some(v) = result.float_val {
out.f64inp = Some(v);
}
}
Err(e) => {
out.report_error(format!(
"{} read error, {}",
plan.iface.c_errs_name(),
e.message()
));
raise_io_alarm(out, alarm_status::READ_ALARM, AlarmSeverity::Major);
}
},
}
}
#[must_use]
enum PhaseFlow {
Continue,
Aborted,
}
fn record_phase_result(
plan: &IoPlan,
out: &mut IoOutcome,
phase: IoPhase,
res: AsynResult<RequestResult>,
) -> PhaseFlow {
if let Err(e) = &res {
if e.is_queue_timeout() {
out.report_queue_timeout();
return PhaseFlow::Aborted;
}
if e.is_queue_refused() {
out.report_queue_refused();
return PhaseFlow::Aborted;
}
}
match phase {
IoPhase::Flush => {
let _ = res;
}
IoPhase::Write => record_write_result(plan, out, res),
IoPhase::Read => record_read_result(plan, out, res),
IoPhase::GpibUniversal(_) => {
if let Err(e) = res {
out.report_error(format!("GPIB Universal command {}", e.message()));
raise_io_alarm(out, alarm_status::WRITE_ALARM, AlarmSeverity::Major);
}
}
IoPhase::GpibAddressed => {
if let Err(e) = res {
out.report_error(format!(
"Error in GPIB Addressed Command write, {}",
e.message()
));
raise_io_alarm(out, alarm_status::WRITE_ALARM, AlarmSeverity::Major);
}
}
IoPhase::GpibSerialPollEnable => {
if let Err(e) = res {
out.report_error(format!("Error in GPIB Serial Poll write, {}", e.message()));
raise_io_alarm(out, alarm_status::WRITE_ALARM, AlarmSeverity::Major);
}
}
IoPhase::GpibSerialPollRead => {
let (data, err) = match res {
Ok(result) => (result.data.unwrap_or_default(), None),
Err(e) => (
e.partial_read().map(|p| p.data.clone()).unwrap_or_default(),
Some(e),
),
};
if let Some(byte) = data.first() {
out.spr = Some(i32::from(*byte));
}
if err.is_some() || data.len() != 1 {
let detail = err.map(|e| e.message()).unwrap_or_default();
out.report_error(format!("Error in GPIB Serial Poll read, {detail}"));
raise_io_alarm(out, alarm_status::READ_ALARM, AlarmSeverity::Major);
}
}
IoPhase::GpibSerialPollDisable => {
if let Err(e) = res {
out.report_error(format!(
"Error in GPIB Serial Poll disable write, {}",
e.message()
));
raise_io_alarm(out, alarm_status::WRITE_ALARM, AlarmSeverity::Major);
}
}
}
PhaseFlow::Continue
}
async fn run_io_plan(handle: PortHandle, plan: IoPlan, cancel: CancelToken) -> IoOutcome {
let mut out = IoOutcome::default();
if cancel.is_cancelled() {
out.report_canceled();
return out;
}
for phase in io_phases(&plan) {
let res = handle
.submit_cancellable(
io_phase_op(&plan, phase),
io_phase_user(&plan, phase),
cancel.clone(),
)
.await;
if cancel.is_cancelled() {
out.report_canceled();
return out;
}
if let PhaseFlow::Aborted = record_phase_result(&plan, &mut out, phase, res) {
return out;
}
}
out
}
const CANCELED_MSG: &str = "I/O request canceled";
const QUEUE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
const PROCESS_QUEUE_TIMEOUT_MSG: &str = "process queueRequest timeout";
const PROCESS_QUEUE_REFUSED_MSG: &str = "queueRequest failed";
const SPECIAL_QUEUE_TIMEOUT_MSG: &str = "special queueRequest timeout";
const SPC_MOD_FIELDS: &[&str] = &[
"PORT", "ADDR", "PCNCT", "DRVINFO", "REASON", "IFACE", "OEOS", "IEOS", "UI32MASK", "BAUD",
"LBAUD", "PRTY", "DBIT", "SBIT", "MCTL", "FCTL", "IXON", "IXOFF", "IXANY", "HOSTINFO", "DRTO",
"TMSK", "TB0", "TB1", "TB2", "TB3", "TB4", "TB5", "TIOM", "TIB0", "TIB1", "TIB2", "TINM",
"TINB0", "TINB1", "TINB2", "TINB3", "TSIZ", "TFIL", "AUCT", "CNCT", "ENBL", "AQR",
];
const OPTION_READBACK_FIELDS: &[&str] = &[
"BAUD", "LBAUD", "PRTY", "DBIT", "SBIT", "MCTL", "FCTL", "IXON", "IXOFF", "IXANY", "HOSTINFO",
"DRTO",
];
const OPTION_READBACK_KEYS: &[&str] = &[
"baud",
"parity",
"bits",
"stop",
"crtscts",
"clocal",
"ixon",
"ixoff",
"ixany",
"hostinfo",
"disconnectOnReadTimeout",
];
const MONITOR_STATUS_FIELDS: &[&str] = &[
"TMSK", "TB0", "TB1", "TB2", "TB3", "TB4", "TB5", "TIOM", "TIB0", "TIB1", "TIB2", "TINM",
"TINB0", "TINB1", "TINB2", "TINB3", "TSIZ", "TFIL", "AUCT", "CNCT", "PCNCT", "REASON",
"DRVINFO", "ENBL", "OCTETIV", "OPTIONIV", "GPIBIV", "I32IV", "UI32IV", "F64IV",
];
const HOSTINFO_OPTION_KEY: &str = "hostinfo";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum OptionQueue {
Normal,
EvenIfNotConnected,
}
impl OptionQueue {
fn for_key(key: &str) -> Self {
if key == HOSTINFO_OPTION_KEY {
Self::EvenIfNotConnected
} else {
Self::Normal
}
}
}
const EOS_READBACK_FIELDS: &[&str] = &["IEOS", "OEOS"];
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
enum SpecialRan {
Yes,
No,
}
#[derive(Clone, Copy)]
struct TraceReadback {
trace_mask: u32,
io_mask: u32,
info_mask: u32,
truncate_size: i32,
file_changed: bool,
}
fn sample_trace_readback(
trace: &TraceManager,
port: &str,
addr: Option<i32>,
old_file_id: &Mutex<Option<usize>>,
) -> TraceReadback {
let snap = trace.snapshot(port, addr);
let file_changed = match old_file_id.lock() {
Ok(mut cached) => cached
.replace(snap.file_id)
.is_some_and(|old| old != snap.file_id),
Err(_) => false,
};
TraceReadback {
trace_mask: snap.trace_mask.bits(),
io_mask: snap.io_mask.bits(),
info_mask: snap.info_mask.bits(),
truncate_size: snap.io_truncate_size as i32,
file_changed,
}
}
const TFIL_UNKNOWN: &str = "Unknown";
fn trace_readback_fields(rb: &TraceReadback) -> Vec<(String, EpicsValue)> {
let TraceReadback {
trace_mask,
io_mask,
info_mask,
truncate_size,
file_changed,
} = *rb;
let bit = |mask: u32, flag: u32| EpicsValue::Enum(u16::from(mask & flag != 0));
let mut fields = vec![
("TMSK".to_string(), EpicsValue::Long(trace_mask as i32)),
("TB0".to_string(), bit(trace_mask, TraceMask::ERROR.bits())),
(
"TB1".to_string(),
bit(trace_mask, TraceMask::IO_DEVICE.bits()),
),
(
"TB2".to_string(),
bit(trace_mask, TraceMask::IO_FILTER.bits()),
),
(
"TB3".to_string(),
bit(trace_mask, TraceMask::IO_DRIVER.bits()),
),
("TB4".to_string(), bit(trace_mask, TraceMask::FLOW.bits())),
(
"TB5".to_string(),
bit(trace_mask, TraceMask::WARNING.bits()),
),
("TIOM".to_string(), EpicsValue::Long(io_mask as i32)),
("TIB0".to_string(), bit(io_mask, TraceIoMask::ASCII.bits())),
("TIB1".to_string(), bit(io_mask, TraceIoMask::ESCAPE.bits())),
("TIB2".to_string(), bit(io_mask, TraceIoMask::HEX.bits())),
("TINM".to_string(), EpicsValue::Long(info_mask as i32)),
(
"TINB0".to_string(),
bit(info_mask, TraceInfoMask::TIME.bits()),
),
(
"TINB1".to_string(),
bit(info_mask, TraceInfoMask::PORT.bits()),
),
(
"TINB2".to_string(),
bit(info_mask, TraceInfoMask::SOURCE.bits()),
),
(
"TINB3".to_string(),
bit(info_mask, TraceInfoMask::THREAD.bits()),
),
("TSIZ".to_string(), EpicsValue::Long(truncate_size)),
];
if file_changed {
fields.push(("TFIL".to_string(), EpicsValue::String(TFIL_UNKNOWN.into())));
}
fields
}
fn connect_readback_fields(
auto_connect: bool,
connected: bool,
enabled: bool,
) -> Vec<(String, EpicsValue)> {
vec![
(
"AUCT".to_string(),
EpicsValue::Enum(u16::from(auto_connect)),
),
("CNCT".to_string(), EpicsValue::Enum(u16::from(connected))),
("ENBL".to_string(), EpicsValue::Enum(u16::from(enabled))),
]
}
impl AsynRecord {
fn update_trace_bits_from_mask(&mut self) {
let mask = self.tmsk as u32;
self.tb0 = if mask & TraceMask::ERROR.bits() != 0 {
1
} else {
0
};
self.tb1 = if mask & TraceMask::IO_DEVICE.bits() != 0 {
1
} else {
0
};
self.tb2 = if mask & TraceMask::IO_FILTER.bits() != 0 {
1
} else {
0
};
self.tb3 = if mask & TraceMask::IO_DRIVER.bits() != 0 {
1
} else {
0
};
self.tb4 = if mask & TraceMask::FLOW.bits() != 0 {
1
} else {
0
};
self.tb5 = if mask & TraceMask::WARNING.bits() != 0 {
1
} else {
0
};
}
fn update_mask_from_trace_bits(&mut self) {
let mut mask: u32 = 0;
if self.tb0 != 0 {
mask |= TraceMask::ERROR.bits();
}
if self.tb1 != 0 {
mask |= TraceMask::IO_DEVICE.bits();
}
if self.tb2 != 0 {
mask |= TraceMask::IO_FILTER.bits();
}
if self.tb3 != 0 {
mask |= TraceMask::IO_DRIVER.bits();
}
if self.tb4 != 0 {
mask |= TraceMask::FLOW.bits();
}
if self.tb5 != 0 {
mask |= TraceMask::WARNING.bits();
}
self.tmsk = mask as i32;
}
fn update_io_bits_from_mask(&mut self) {
let mask = self.tiom as u32;
self.tib0 = if mask & TraceIoMask::ASCII.bits() != 0 {
1
} else {
0
};
self.tib1 = if mask & TraceIoMask::ESCAPE.bits() != 0 {
1
} else {
0
};
self.tib2 = if mask & TraceIoMask::HEX.bits() != 0 {
1
} else {
0
};
}
fn update_mask_from_io_bits(&mut self) {
let mut mask: u32 = 0;
if self.tib0 != 0 {
mask |= TraceIoMask::ASCII.bits();
}
if self.tib1 != 0 {
mask |= TraceIoMask::ESCAPE.bits();
}
if self.tib2 != 0 {
mask |= TraceIoMask::HEX.bits();
}
self.tiom = mask as i32;
}
fn update_info_bits_from_mask(&mut self) {
let mask = self.tinm as u32;
self.tinb0 = if mask & TraceInfoMask::TIME.bits() != 0 {
1
} else {
0
};
self.tinb1 = if mask & TraceInfoMask::PORT.bits() != 0 {
1
} else {
0
};
self.tinb2 = if mask & TraceInfoMask::SOURCE.bits() != 0 {
1
} else {
0
};
self.tinb3 = if mask & TraceInfoMask::THREAD.bits() != 0 {
1
} else {
0
};
}
fn update_mask_from_info_bits(&mut self) {
let mut mask: u32 = 0;
if self.tinb0 != 0 {
mask |= TraceInfoMask::TIME.bits();
}
if self.tinb1 != 0 {
mask |= TraceInfoMask::PORT.bits();
}
if self.tinb2 != 0 {
mask |= TraceInfoMask::SOURCE.bits();
}
if self.tinb3 != 0 {
mask |= TraceInfoMask::THREAD.bits();
}
self.tinm = mask as i32;
}
fn trace_addr_target(&self) -> Option<i32> {
match self.port_entry {
Some(ref entry) if self.addr >= 0 && entry.handle.is_multi_device() => Some(self.addr),
_ => None,
}
}
fn apply_trace_mask(&self) {
if let Some(ref entry) = self.port_entry {
let mask = TraceMask::from_bits_truncate(self.tmsk as u32);
match self.trace_addr_target() {
Some(addr) => entry.trace.set_device_trace_mask(&self.port, addr, mask),
None => entry.trace.set_trace_mask(Some(&self.port), mask),
}
}
}
fn apply_trace_io_mask(&self) {
if let Some(ref entry) = self.port_entry {
let mask = TraceIoMask::from_bits_truncate(self.tiom as u32);
match self.trace_addr_target() {
Some(addr) => entry.trace.set_device_trace_io_mask(&self.port, addr, mask),
None => entry.trace.set_trace_io_mask(Some(&self.port), mask),
}
}
}
fn apply_trace_info_mask(&self) {
if let Some(ref entry) = self.port_entry {
let mask = TraceInfoMask::from_bits_truncate(self.tinm as u32);
match self.trace_addr_target() {
Some(addr) => entry
.trace
.set_device_trace_info_mask(&self.port, addr, mask),
None => entry.trace.set_trace_info_mask(Some(&self.port), mask),
}
}
}
fn apply_trace_truncate_size(&self) {
if let Some(ref entry) = self.port_entry {
let size = self.tsiz as usize;
match self.trace_addr_target() {
Some(addr) => entry
.trace
.set_device_io_truncate_size(&self.port, addr, size),
None => entry.trace.set_io_truncate_size(Some(&self.port), size),
}
}
}
fn apply_trace_file(&mut self) {
let Some(entry) = self.port_entry.clone() else {
return;
};
let tfil = self.tfil.clone();
let file = match open_trace_file(&tfil) {
Ok(file) => file,
Err(_) => {
self.report_error(format!("Error opening trace file: {tfil}"));
return;
}
};
if let Ok(mut cached) = self.old_trace_file_id.lock() {
*cached = Some(file.id());
}
match self.trace_addr_target() {
Some(addr) => entry.trace.set_device_trace_file(&self.port, addr, file),
None => entry.trace.set_trace_file(Some(&self.port), file),
}
}
fn read_trace_state(&mut self) {
let Some(entry) = self.port_entry.clone() else {
return;
};
let addr = self.trace_addr_target();
let cache = Arc::clone(&self.old_trace_file_id);
let rb = sample_trace_readback(&entry.trace, &self.port, addr, &cache);
self.tmsk = rb.trace_mask as i32;
self.update_trace_bits_from_mask();
self.tiom = rb.io_mask as i32;
self.update_io_bits_from_mask();
self.tinm = rb.info_mask as i32;
self.update_info_bits_from_mask();
self.tsiz = rb.truncate_size;
if rb.file_changed {
self.tfil = TFIL_UNKNOWN.to_string();
}
}
fn register_exception_callback(&mut self) {
self.clear_exception_callback();
let Some(ref entry) = self.port_entry else {
return;
};
let Some(mgr) = entry.trace.exception_manager() else {
return;
};
let port = self.port.clone();
let trace = entry.trace.clone();
let handle = entry.handle.clone();
let dirty = Arc::clone(&self.status_dirty);
let immediate = match (
self.async_ctx.clone(),
tokio::runtime::Handle::try_current().ok(),
) {
(Some((name, db)), Some(rt)) => Some((name, db, rt)),
_ => None,
};
let old_file_id = Arc::clone(&self.old_trace_file_id);
let trace_addr = self.trace_addr_target();
let last_posted: Arc<Mutex<HashMap<String, EpicsValue>>> = Arc::new(Mutex::new(
trace_readback_fields(&sample_trace_readback(
&trace,
&port,
trace_addr,
&old_file_id,
))
.into_iter()
.chain(connect_readback_fields(
self.auct != 0,
self.cnct != 0,
self.enbl != 0,
))
.collect(),
));
let id = mgr.add_callback(move |ev| {
if ev.port_name != port {
return;
}
let Some((name, db, rt)) = immediate.clone() else {
dirty.store(true, Ordering::Release);
return;
};
let (port, trace, handle, last_posted, old_file_id) = (
port.clone(),
trace.clone(),
handle.clone(),
Arc::clone(&last_posted),
Arc::clone(&old_file_id),
);
rt.spawn(async move {
let auto = handle.is_auto_connect().await.unwrap_or(false);
let connected = handle.is_connected().await.unwrap_or(false);
let enabled = handle.is_enabled().await.unwrap_or(false);
let fields = trace_readback_fields(&sample_trace_readback(
&trace,
&port,
trace_addr,
&old_file_id,
))
.into_iter()
.chain(connect_readback_fields(auto, connected, enabled));
let changed: Vec<(String, EpicsValue)> = {
let mut cache = last_posted.lock().unwrap();
let mut changed = Vec::new();
for (field, value) in fields {
if cache.get(&field) != Some(&value) {
cache.insert(field.clone(), value.clone());
changed.push((field, value));
}
}
changed
};
if changed.is_empty() {
return;
}
let _ = db.post_fields(&name, changed);
});
});
self.except_cb = Some((mgr, id));
}
fn clear_exception_callback(&mut self) {
if let Some((mgr, id)) = self.except_cb.take() {
mgr.remove_callback(id);
}
}
fn read_options_from_driver(&mut self, handle: &PortHandle, queue: OptionQueue) {
if !handle.has_interface(crate::interfaces::InterfaceType::Option) {
return;
}
let Some(opts) = self.get_options(handle, OPTION_READBACK_KEYS, queue) else {
return;
};
let opt = |key: &str| opts.get(key).map(String::as_str).unwrap_or("");
let baud_text = opt("baud");
if let Some(rate) = crate::drivers::option_parse::sscanf_int(baud_text) {
self.lbaud = rate;
}
self.baud = baud_choice_index(baud_text);
self.prty = match opt("parity") {
"none" => 1,
"even" => 2,
"odd" => 3,
_ => 0, };
self.dbit = match opt("bits") {
"5" => 1,
"6" => 2,
"7" => 3,
"8" => 4,
_ => 0,
};
self.sbit = match opt("stop") {
"1" => 1,
"2" => 2,
_ => 0,
};
self.fctl = match opt("crtscts") {
"Y" | "Yes" => 2, "N" | "No" | "none" => 1, _ => 0,
};
self.mctl = match opt("clocal") {
"Y" | "Yes" => 1, "N" | "No" => 2, _ => 0,
};
self.ixon = match opt("ixon") {
"Y" | "Yes" => 2,
"N" | "No" => 1,
_ => 0,
};
self.ixoff = match opt("ixoff") {
"Y" | "Yes" => 2,
"N" | "No" => 1,
_ => 0,
};
self.ixany = match opt("ixany") {
"Y" | "Yes" => 2,
"N" | "No" => 1,
_ => 0,
};
self.hostinfo = opt("hostinfo").to_string();
self.drto = match opt("disconnectOnReadTimeout") {
"Y" | "Yes" => 2,
"N" | "No" => 1,
_ => 0,
};
}
fn get_options(
&mut self,
handle: &PortHandle,
keys: &[&str],
queue: OptionQueue,
) -> Option<HashMap<String, String>> {
let mut opts = HashMap::new();
for key in keys {
match handle.get_option_blocking(self.option_user_for(queue), key) {
Ok(val) => {
opts.insert((*key).to_string(), val);
}
Err(e) if e.is_queue_timeout() => {
self.report_special_queue_timeout();
return None;
}
Err(_) => {}
}
}
Some(opts)
}
fn read_eos_from_driver(&mut self, handle: &PortHandle) {
let mut ieos = String::new();
let mut oeos = String::new();
if self.octetiv != 0 {
let read = |res: AsynResult<Vec<u8>>| -> Result<String, bool> {
match res {
Ok(bytes) if !bytes.is_empty() => {
Ok(crate::escape::escaped_from_raw(&bytes, EOS_SIZE))
}
Ok(_) => Ok(String::new()),
Err(e) if e.is_queue_timeout() => Err(true),
Err(_) => Ok(String::new()),
}
};
match read(handle.get_input_eos_blocking(self.option_user())) {
Ok(v) => ieos = v,
Err(_) => {
self.report_special_queue_timeout();
return;
}
}
match read(handle.get_output_eos_blocking(self.option_user())) {
Ok(v) => oeos = v,
Err(_) => {
self.report_special_queue_timeout();
return;
}
}
}
self.ieos = ieos;
self.oeos = oeos;
}
fn field_snapshot(&self, names: &[&str]) -> Vec<(String, EpicsValue)> {
names
.iter()
.filter_map(|name| self.get_field(name).map(|v| ((*name).to_string(), v)))
.collect()
}
fn posting<T>(&mut self, fields: &[&str], readback: impl FnOnce(&mut Self) -> T) -> T {
let before = self.field_snapshot(fields);
let out = readback(self);
self.post_if_new(&before);
out
}
fn post_if_new(&self, before: &[(String, EpicsValue)]) {
let changed: Vec<(String, EpicsValue)> = before
.iter()
.filter_map(|(field, old)| {
let new = self.get_field(field)?;
(new != *old).then(|| (field.clone(), new))
})
.collect();
if changed.is_empty() {
return;
}
let (Some((name, db)), Ok(rt)) = (
self.async_ctx.clone(),
tokio::runtime::Handle::try_current(),
) else {
return;
};
rt.spawn(async move {
let _ = db.post_fields(&name, changed);
});
}
fn write_option(&mut self, key: &str, value: &str) {
let (key, value) = (key.to_string(), value.to_string());
self.special_callback(|this| this.write_option_body(&key, &value));
}
fn write_option_body(&mut self, key: &str, value: &str) -> SpecialRan {
let Some(entry) = self.port_entry.clone() else {
return SpecialRan::No;
};
if !self.port_has(crate::interfaces::InterfaceType::Option) {
self.report_no_interface("asynOption");
return SpecialRan::Yes;
}
let queue = OptionQueue::for_key(key);
if let Err(e) = entry
.handle
.set_option_blocking(self.option_user_for(queue), key, value)
{
if e.never_ran() {
return self.report_special_never_ran(&e);
}
self.report_error(format!("Error setting option, {}", e.message()));
}
self.posting(OPTION_READBACK_FIELDS, |this| {
this.read_options_from_driver(&entry.handle, queue);
});
SpecialRan::Yes
}
fn special_callback(&mut self, body: impl FnOnce(&mut Self) -> SpecialRan) {
let before = self.field_snapshot(MONITOR_STATUS_FIELDS);
if body(self) == SpecialRan::Yes {
self.monitor_status();
self.post_if_new(&before);
}
}
fn report_error(&mut self, msg: impl Into<String>) {
let before = self.field_snapshot(&["ERRS"]);
self.errs = msg.into();
self.post_if_new(&before);
}
fn reset_error(&mut self) {
self.report_error(String::new());
}
fn report_not_connected(&mut self) {
self.report_error("Not connect to a port");
self.io_alarm = Some((alarm_status::STATE_ALARM, AlarmSeverity::Minor));
}
fn report_no_interface(&mut self, asyn_name: &str) {
self.report_error(format!("No {asyn_name} interface"));
self.io_alarm = Some((alarm_status::COMM_ALARM, AlarmSeverity::Major));
}
fn report_special_queue_timeout(&mut self) {
self.report_error(SPECIAL_QUEUE_TIMEOUT_MSG);
}
fn report_special_never_ran(&mut self, e: &AsynError) -> SpecialRan {
if e.is_queue_refused() {
self.report_error(e.message());
} else {
self.report_special_queue_timeout();
}
SpecialRan::No
}
pub(crate) fn io_intr_scan(&self) -> Arc<IoIntrScan> {
self.io_intr.clone()
}
fn publish_io_intr_binding(&mut self) {
let binding = self.port_entry.as_ref().map(|entry| IoIntrBinding {
handle: entry.handle.clone(),
iface: InterfaceType::from_u16(self.iface as u16),
addr: self.addr,
reason: self.resolved_reason,
ui32mask: self.ui32mask,
});
if let Err(msg) = self.io_intr.rebind(binding) {
self.report_error(msg);
}
}
fn cancel_io_interrupt_scan(&mut self) {
if !self.io_intr.is_active() {
return;
}
let _ = self.io_intr.set_active(false);
let Some((name, db)) = self.async_ctx.clone() else {
return;
};
let passive = EpicsValue::Enum(ScanType::Passive.to_u16());
tokio::spawn(async move {
let _ = db.put_pv(&format!("{name}.SCAN"), passive).await;
});
}
fn apply_io_intr_sample(&mut self, sample: IoIntrSample) {
match sample {
IoIntrSample::Octet(s) => self.tinp = s,
IoIntrSample::Int32(v) => self.i32inp = v,
IoIntrSample::UInt32(v) => self.ui32inp = v,
IoIntrSample::Float64(v) => self.f64inp = v,
}
}
fn has_interface(&self, iface: InterfaceType) -> bool {
self.port_has(iface.registry_type())
}
fn port_has(&self, iface: crate::interfaces::InterfaceType) -> bool {
self.port_entry
.as_ref()
.is_some_and(|entry| entry.handle.has_interface(iface))
}
fn write_eos(&mut self, output: bool) {
self.special_callback(|this| this.write_eos_body(output));
}
fn write_eos_body(&mut self, output: bool) -> SpecialRan {
let Some(entry) = self.port_entry.clone() else {
return SpecialRan::No;
};
if !self.has_interface(InterfaceType::Octet) {
self.report_no_interface(InterfaceType::Octet.c_asyn_name());
return SpecialRan::Yes;
}
let field = if output { &self.oeos } else { &self.ieos };
let bytes = translate_escape(field);
let res = if output {
entry
.handle
.set_output_eos_blocking(self.option_user(), &bytes)
} else {
entry
.handle
.set_input_eos_blocking(self.option_user(), &bytes)
};
if let Err(e) = res {
if e.never_ran() {
return self.report_special_never_ran(&e);
}
let which = if output { "output" } else { "input" };
self.report_error(format!("Error setting {which} eos, {}", e.message()));
}
self.posting(EOS_READBACK_FIELDS, |this| {
this.read_eos_from_driver(&entry.handle);
});
SpecialRan::Yes
}
fn refresh_connected_state(&mut self) {
let connected = match self.port_entry {
Some(ref entry) => entry.handle.is_connected_blocking().unwrap_or(false),
None => false,
};
self.cnct = i32::from(connected);
}
fn monitor_status(&mut self) {
self.read_trace_state();
let (enabled, auto) = match self.port_entry {
Some(ref entry) => (
entry.handle.is_enabled_blocking().unwrap_or(false),
entry.handle.is_auto_connect_blocking().unwrap_or(false),
),
None => (false, false),
};
self.enbl = i32::from(enabled);
self.auct = i32::from(auto);
self.refresh_connected_state();
}
fn connect_device(&mut self) -> Result<(), String> {
self.reset_error();
let before_status = self.field_snapshot(MONITOR_STATUS_FIELDS);
if self.port.is_empty() {
self.pcnct = 0;
self.port_entry = None;
self.clear_exception_callback();
self.monitor_status();
self.post_if_new(&before_status);
self.publish_io_intr_binding();
return Err(self.report_connect_error(
"asynManager:connectDevice no port name provided".to_string(),
));
}
match registry::get_port(&self.port) {
Some(entry) => {
if entry
.handle
.has_interface(crate::interfaces::InterfaceType::DrvUser)
{
if !self.drvinfo.is_empty() {
let req = crate::port::DrvUserRequest::new(&self.drvinfo, self.addr)
.with_iface(InterfaceType::from_u16(self.iface as u16).as_asyn_iface());
match entry.handle.drv_user_create_blocking(&req) {
Ok(info) => {
self.resolved_reason = info.reason;
self.reason = info.reason as i32;
}
Err(_) => {
self.report_error("Error in asynDrvUser->create()");
self.resolved_reason = 0;
}
}
} else {
self.resolved_reason = self.reason as usize;
}
} else {
self.reason = 0;
self.resolved_reason = 0;
if !self.drvinfo.is_empty() {
self.report_error("asynDrvUser not supported but drvInfo not blank");
}
}
let has = |iface| i32::from(entry.handle.has_interface(iface));
self.octetiv = has(crate::interfaces::InterfaceType::Octet);
self.i32iv = has(crate::interfaces::InterfaceType::Int32);
self.ui32iv = has(crate::interfaces::InterfaceType::UInt32Digital);
self.f64iv = has(crate::interfaces::InterfaceType::Float64);
self.optioniv = has(crate::interfaces::InterfaceType::Option);
self.gpibiv = has(crate::interfaces::InterfaceType::Gpib);
self.port_entry = Some(entry.clone());
self.posting(OPTION_READBACK_FIELDS, |this| {
this.read_options_from_driver(&entry.handle, OptionQueue::EvenIfNotConnected);
});
self.posting(EOS_READBACK_FIELDS, |this| {
this.read_eos_from_driver(&entry.handle);
});
self.pcnct = 1;
self.monitor_status();
self.post_if_new(&before_status);
self.register_exception_callback();
self.publish_io_intr_binding();
Ok(())
}
None => {
self.pcnct = 0;
self.port_entry = None;
self.clear_exception_callback();
self.monitor_status();
self.post_if_new(&before_status);
self.publish_io_intr_binding();
Err(self.report_connect_error(format!(
"asynManager:connectDevice port {} not found",
self.port
)))
}
}
}
fn report_connect_error(&mut self, manager_message: String) -> String {
self.report_error(format!(
"Connect error, status={}, {manager_message}",
AsynStatus::Error as i32
));
manager_message
}
fn io_timeout(&self) -> std::time::Duration {
crate::user::timeout_from_secs(self.tmot)
}
fn option_user(&self) -> AsynUser {
AsynUser::new(self.resolved_reason)
.with_addr(self.addr)
.with_timeout(self.io_timeout())
.with_queue_timeout(QUEUE_TIMEOUT)
}
fn option_user_for(&self, queue: OptionQueue) -> AsynUser {
match queue {
OptionQueue::Normal => self.option_user(),
OptionQueue::EvenIfNotConnected => self.option_user().queue_even_if_not_connected(),
}
}
fn put_array_field(&mut self, field: AsynArrayField, data: Vec<u8>) {
let n_new = data.len().min(i32::MAX as usize) as i32;
match field {
AsynArrayField::Bout => {
self.bout = data;
self.nowt = n_new;
}
AsynArrayField::Binp => {
self.binp = data;
self.nord = n_new;
}
}
}
fn clamp_transfer_sizes(&mut self, in_len: usize) {
if self.ofmt == ASYN_FMT_BINARY && self.nowt > self.omax {
self.nowt = self.omax;
}
let in_len = in_len.min(i32::MAX as usize) as i32;
if self.nrrd > in_len {
self.nrrd = in_len;
}
}
fn octet_output_buffer(&self) -> Vec<u8> {
match self.ofmt {
ASYN_FMT_BINARY => {
let nowt = self.nowt.max(0) as usize;
self.bout[..nowt.min(self.bout.len())].to_vec()
}
ASYN_FMT_HYBRID => {
let end = self
.bout
.iter()
.position(|&b| b == 0)
.unwrap_or(self.bout.len());
translate_escape_bytes(&self.bout[..end])
}
_ => translate_escape(&self.aout),
}
}
fn build_io_plan(&mut self) -> IoPlan {
self.build_io_plan_for(None)
}
fn build_io_plan_for(&mut self, gpib: Option<GpibCycle>) -> IoPlan {
let iface = InterfaceType::from_u16(self.iface as u16);
let in_len = if self.ifmt == ASYN_FMT_ASCII {
AINP_SIZE
} else {
self.imax.max(0) as usize
};
let (octet_out, octet_out_len) = if gpib.is_none() && iface == InterfaceType::Octet {
self.clamp_transfer_sizes(in_len);
let out = self.octet_output_buffer();
let len = out.len();
(out, len)
} else {
(Vec::new(), 0)
};
let octet_buf_size = if self.nrrd > 0 {
(self.nrrd as usize).min(in_len)
} else {
in_len
};
let timeout = self.io_timeout();
IoPlan {
tmod: TransferMode::from_u16(self.tmod as u16),
iface,
gpib,
reason: self.resolved_reason,
addr: self.addr,
timeout,
octet_out,
octet_out_len,
ofmt: self.ofmt,
i32out: self.i32out,
ui32out: self.ui32out,
ui32mask: self.ui32mask,
f64out: self.f64out,
octet_buf_size,
in_len,
ifmt: self.ifmt,
}
}
fn take_gpib_cycle(&mut self) -> Option<GpibCycle> {
use crate::interfaces::gpib::{
GpibAddressedRequest, addressed_request, universal_cmd_byte,
};
if self.ucmd != 0 {
let cmd = universal_cmd_byte(self.ucmd);
self.ucmd = 0;
return Some(GpibCycle::Universal(cmd));
}
if self.acmd != 0 {
let request = addressed_request(self.acmd, self.addr);
self.acmd = 0;
return Some(match request {
GpibAddressedRequest::Frame(frame) => GpibCycle::Addressed(frame),
GpibAddressedRequest::SerialPoll => GpibCycle::SerialPoll,
});
}
None
}
fn apply_io_outcome(&mut self, out: IoOutcome) {
if let Some(v) = out.nawt {
self.nawt = v;
}
if let Some(v) = out.eomr {
self.eomr = v;
}
if let Some(v) = out.nord {
self.nord = v;
}
if let Some(v) = out.tinp {
self.tinp = v;
}
if let Some(v) = out.ainp {
self.ainp = v;
}
if let Some(v) = out.binp {
self.binp = v;
}
if let Some(v) = out.i32inp {
self.i32inp = v;
}
if let Some(v) = out.ui32inp {
self.ui32inp = v;
}
if let Some(v) = out.f64inp {
self.f64inp = v;
}
if let Some(v) = out.spr {
self.spr = v;
}
if let Some(v) = out.errs {
self.report_error(v);
}
if let Some(a) = out.alarm {
self.io_alarm = Some(a);
}
}
fn perform_io(&mut self, plan: IoPlan) -> CaResult<()> {
let entry = match &self.port_entry {
Some(e) => e.clone(),
None => {
self.report_not_connected();
return Ok(());
}
};
let mut out = IoOutcome::default();
for phase in io_phases(&plan) {
let res = entry
.handle
.submit_blocking(io_phase_op(&plan, phase), io_phase_user(&plan, phase));
if let PhaseFlow::Aborted = record_phase_result(&plan, &mut out, phase, res) {
break;
}
}
self.apply_io_outcome(out);
Ok(())
}
fn spawn_async_io(
&mut self,
handle: PortHandle,
name: String,
db: AsyncDbHandle,
plan: IoPlan,
) -> ProcessOutcome {
let cancel = CancelToken::new();
let slot: Arc<Mutex<Option<IoOutcome>>> = Arc::new(Mutex::new(None));
let cancel_task = cancel.clone();
let slot_task = slot.clone();
tokio::spawn(async move {
let outcome = run_io_plan(handle, plan, cancel_task).await;
*slot_task.lock().unwrap() = Some(outcome);
if let Some(token) = db.mint_async_token(&name) {
let (waitset, completion) = AsyncDbHandle::new_put_notify();
waitset.leave();
let _ = db.reprocess_on_notify(token, completion);
}
});
self.io_inflight = Some(IoInFlight {
cancel,
result: slot,
});
ProcessOutcome::async_pending()
}
}
const MENU_ASYN_TMOD: &[&str] = &["Write/Read", "Write", "Read", "Flush", "NoI/O"];
const MENU_ASYN_INTERFACE: &[&str] =
&["asynOctet", "asynInt32", "asynUInt32Digital", "asynFloat64"];
const MENU_ASYN_FMT: &[&str] = &["ASCII", "Hybrid", "Binary"];
const MENU_ASYN_TRACE: &[&str] = &["Off", "On"];
const MENU_ASYN_AUTOCONNECT: &[&str] = &["noAutoConnect", "autoConnect"];
const MENU_ASYN_CONNECT: &[&str] = &["Disconnect", "Connect"];
const MENU_ASYN_ENABLE: &[&str] = &["Disable", "Enable"];
const MENU_ASYN_EOMREASON: &[&str] = &[
"None",
"Count",
"Eos",
"Count Eos",
"End",
"Count End",
"Eos End",
"Count Eos End",
];
const MENU_SERIAL_BAUD: &[&str] = &[
"Unknown", "300", "600", "1200", "2400", "4800", "9600", "19200", "38400", "57600", "115200",
"230400", "460800", "576000", "921600", "1152000",
];
const MENU_SERIAL_PRTY: &[&str] = &["Unknown", "None", "Even", "Odd"];
const MENU_SERIAL_DBIT: &[&str] = &["Unknown", "5", "6", "7", "8"];
const MENU_SERIAL_SBIT: &[&str] = &["Unknown", "1", "2"];
const MENU_SERIAL_MCTL: &[&str] = &["Unknown", "CLOCAL", "YES"];
const MENU_SERIAL_FCTL: &[&str] = &["Unknown", "None", "Hardware"];
const MENU_SERIAL_IX: &[&str] = &["Unknown", "No", "Yes"];
const MENU_IP_DRTO: &[&str] = &["Unknown", "No", "Yes"];
const MENU_GPIB_UCMD: &[&str] = &[
"None",
"Device Clear (DCL)",
"Local Lockout (LL0)",
"Serial Poll Disable (SPD)",
"Serial Poll Enable (SPE)",
"Unlisten (UNL)",
"Untalk (UNT)",
];
const MENU_GPIB_ACMD: &[&str] = &[
"None",
"Group Execute Trig. (GET)",
"Go To Local (GTL)",
"Selected Dev. Clear (SDC)",
"Take Control (TCT)",
"Serial Poll",
];
impl Record for AsynRecord {
fn record_type(&self) -> &'static str {
"asyn"
}
fn property_support(&self) -> epics_base_rs::server::snapshot::PropertySupport {
use epics_base_rs::server::snapshot::PropertySupport as P;
P {
precision: true,
..P::NONE
}
}
fn menu_field_choices(&self, field: &str) -> Option<&'static [&'static str]> {
match field {
"TMOD" => Some(MENU_ASYN_TMOD),
"IFACE" => Some(MENU_ASYN_INTERFACE),
"OFMT" | "IFMT" => Some(MENU_ASYN_FMT),
"TB0" | "TB1" | "TB2" | "TB3" | "TB4" | "TB5" | "TIB0" | "TIB1" | "TIB2" | "TINB0"
| "TINB1" | "TINB2" | "TINB3" => Some(MENU_ASYN_TRACE),
"AUCT" => Some(MENU_ASYN_AUTOCONNECT),
"CNCT" | "PCNCT" => Some(MENU_ASYN_CONNECT),
"ENBL" => Some(MENU_ASYN_ENABLE),
"EOMR" => Some(MENU_ASYN_EOMREASON),
"BAUD" => Some(MENU_SERIAL_BAUD),
"PRTY" => Some(MENU_SERIAL_PRTY),
"DBIT" => Some(MENU_SERIAL_DBIT),
"SBIT" => Some(MENU_SERIAL_SBIT),
"MCTL" => Some(MENU_SERIAL_MCTL),
"FCTL" => Some(MENU_SERIAL_FCTL),
"IXON" | "IXOFF" | "IXANY" => Some(MENU_SERIAL_IX),
"DRTO" => Some(MENU_IP_DRTO),
"UCMD" => Some(MENU_GPIB_UCMD),
"ACMD" => Some(MENU_GPIB_ACMD),
_ => None,
}
}
fn set_async_context(&mut self, name: String, db: AsyncDbHandle) {
self.async_ctx = Some((name, db));
}
fn check_alarms(&mut self, common: &mut CommonFields) {
if let Some((stat, sevr)) = self.io_alarm.take() {
rec_gbl_set_sevr(common, stat, sevr);
}
}
fn get_field(&self, name: &str) -> Option<EpicsValue> {
match name {
"PORT" => Some(EpicsValue::String(self.port.clone().into())),
"ADDR" => Some(EpicsValue::Long(self.addr)),
"PCNCT" => Some(EpicsValue::Enum(self.pcnct as u16)),
"DRVINFO" => Some(EpicsValue::String(self.drvinfo.clone().into())),
"REASON" => Some(EpicsValue::Long(self.reason)),
"TMOD" => Some(EpicsValue::Enum(self.tmod as u16)),
"TMOT" => Some(EpicsValue::Double(self.tmot)),
"IFACE" => Some(EpicsValue::Enum(self.iface as u16)),
"OCTETIV" => Some(EpicsValue::Long(self.octetiv)),
"OPTIONIV" => Some(EpicsValue::Long(self.optioniv)),
"GPIBIV" => Some(EpicsValue::Long(self.gpibiv)),
"I32IV" => Some(EpicsValue::Long(self.i32iv)),
"UI32IV" => Some(EpicsValue::Long(self.ui32iv)),
"F64IV" => Some(EpicsValue::Long(self.f64iv)),
"AOUT" => Some(EpicsValue::String(self.aout.clone().into())),
"OEOS" => Some(EpicsValue::String(self.oeos.clone().into())),
"BOUT" => Some(EpicsValue::CharArray(self.bout.clone())),
"OMAX" => Some(EpicsValue::Long(self.omax)),
"NOWT" => Some(EpicsValue::Long(self.nowt)),
"NAWT" => Some(EpicsValue::Long(self.nawt)),
"OFMT" => Some(EpicsValue::Enum(self.ofmt as u16)),
"AINP" => Some(EpicsValue::String(self.ainp.clone().into())),
"TINP" => Some(EpicsValue::String(self.tinp.clone().into())),
"IEOS" => Some(EpicsValue::String(self.ieos.clone().into())),
"BINP" => Some(EpicsValue::CharArray(self.binp.clone())),
"IMAX" => Some(EpicsValue::Long(self.imax)),
"NRRD" => Some(EpicsValue::Long(self.nrrd)),
"NORD" => Some(EpicsValue::Long(self.nord)),
"IFMT" => Some(EpicsValue::Enum(self.ifmt as u16)),
"EOMR" => Some(EpicsValue::Enum(self.eomr as u16)),
"I32INP" => Some(EpicsValue::Long(self.i32inp)),
"I32OUT" => Some(EpicsValue::Long(self.i32out)),
"UI32INP" => Some(EpicsValue::ULong(self.ui32inp)),
"UI32OUT" => Some(EpicsValue::ULong(self.ui32out)),
"UI32MASK" => Some(EpicsValue::ULong(self.ui32mask)),
"F64INP" => Some(EpicsValue::Double(self.f64inp)),
"F64OUT" => Some(EpicsValue::Double(self.f64out)),
"BAUD" => Some(EpicsValue::Enum(self.baud as u16)),
"LBAUD" => Some(EpicsValue::Long(self.lbaud)),
"PRTY" => Some(EpicsValue::Enum(self.prty as u16)),
"DBIT" => Some(EpicsValue::Enum(self.dbit as u16)),
"SBIT" => Some(EpicsValue::Enum(self.sbit as u16)),
"MCTL" => Some(EpicsValue::Enum(self.mctl as u16)),
"FCTL" => Some(EpicsValue::Enum(self.fctl as u16)),
"IXON" => Some(EpicsValue::Enum(self.ixon as u16)),
"IXOFF" => Some(EpicsValue::Enum(self.ixoff as u16)),
"IXANY" => Some(EpicsValue::Enum(self.ixany as u16)),
"HOSTINFO" => Some(EpicsValue::String(self.hostinfo.clone().into())),
"DRTO" => Some(EpicsValue::Enum(self.drto as u16)),
"UCMD" => Some(EpicsValue::Enum(self.ucmd as u16)),
"ACMD" => Some(EpicsValue::Enum(self.acmd as u16)),
"SPR" => Some(EpicsValue::Char(self.spr as u8)),
"TMSK" => Some(EpicsValue::Long(self.tmsk)),
"TB0" => Some(EpicsValue::Enum(self.tb0 as u16)),
"TB1" => Some(EpicsValue::Enum(self.tb1 as u16)),
"TB2" => Some(EpicsValue::Enum(self.tb2 as u16)),
"TB3" => Some(EpicsValue::Enum(self.tb3 as u16)),
"TB4" => Some(EpicsValue::Enum(self.tb4 as u16)),
"TB5" => Some(EpicsValue::Enum(self.tb5 as u16)),
"TIOM" => Some(EpicsValue::Long(self.tiom)),
"TIB0" => Some(EpicsValue::Enum(self.tib0 as u16)),
"TIB1" => Some(EpicsValue::Enum(self.tib1 as u16)),
"TIB2" => Some(EpicsValue::Enum(self.tib2 as u16)),
"TINM" => Some(EpicsValue::Long(self.tinm)),
"TINB0" => Some(EpicsValue::Enum(self.tinb0 as u16)),
"TINB1" => Some(EpicsValue::Enum(self.tinb1 as u16)),
"TINB2" => Some(EpicsValue::Enum(self.tinb2 as u16)),
"TINB3" => Some(EpicsValue::Enum(self.tinb3 as u16)),
"TSIZ" => Some(EpicsValue::Long(self.tsiz)),
"TFIL" => Some(EpicsValue::String(self.tfil.clone().into())),
"AUCT" => Some(EpicsValue::Enum(self.auct as u16)),
"CNCT" => Some(EpicsValue::Enum(self.cnct as u16)),
"ENBL" => Some(EpicsValue::Enum(self.enbl as u16)),
"VAL" => Some(EpicsValue::Long(self.val)),
"ERRS" => Some(EpicsValue::String(self.errs.clone().into())),
"AQR" => Some(EpicsValue::Char(self.aqr as u8)),
_ => None,
}
}
fn put_field(&mut self, name: &str, value: EpicsValue) -> CaResult<()> {
let to_i32 = |v: &EpicsValue| -> i32 { v.to_f64().unwrap_or(0.0) as i32 };
let to_u32 = |v: &EpicsValue| -> u32 { v.to_f64().unwrap_or(0.0) as u32 };
let to_f64 = |v: &EpicsValue| -> f64 { v.to_f64().unwrap_or(0.0) };
let to_str = |v: &EpicsValue| -> String { format!("{v}") };
let to_bytes = |v: &EpicsValue| -> Vec<u8> {
match v {
EpicsValue::CharArray(b) => b.clone(),
EpicsValue::String(s) => s.as_bytes().to_vec(),
_ => Vec::new(),
}
};
match name {
"PORT" => {
self.port = to_str(&value);
}
"ADDR" => {
self.addr = to_i32(&value);
}
"PCNCT" => {
self.pcnct = to_i32(&value);
}
"DRVINFO" => {
self.drvinfo = to_str(&value);
}
"REASON" => {
self.reason = to_i32(&value);
}
"TMOD" => {
self.tmod = to_i32(&value);
}
"TMOT" => {
self.tmot = to_f64(&value);
}
"IFACE" => {
self.iface = to_i32(&value);
}
"OCTETIV" => {
self.octetiv = to_i32(&value);
}
"OPTIONIV" => {
self.optioniv = to_i32(&value);
}
"GPIBIV" => {
self.gpibiv = to_i32(&value);
}
"I32IV" => {
self.i32iv = to_i32(&value);
}
"UI32IV" => {
self.ui32iv = to_i32(&value);
}
"F64IV" => {
self.f64iv = to_i32(&value);
}
"AOUT" => {
self.aout = to_str(&value);
}
"OEOS" => {
self.oeos = to_str(&value);
}
"BOUT" => {
self.put_array_field(AsynArrayField::Bout, to_bytes(&value));
}
"OMAX" => {
self.omax = to_i32(&value);
}
"NOWT" => {
self.nowt = to_i32(&value);
}
"NAWT" => {
self.nawt = to_i32(&value);
}
"OFMT" => {
self.ofmt = to_i32(&value);
}
"AINP" => {
self.ainp = to_str(&value);
}
"TINP" => {
self.tinp = to_str(&value);
}
"IEOS" => {
self.ieos = to_str(&value);
}
"BINP" => {
self.put_array_field(AsynArrayField::Binp, to_bytes(&value));
}
"IMAX" => {
self.imax = to_i32(&value);
}
"NRRD" => {
self.nrrd = to_i32(&value);
}
"NORD" => {
self.nord = to_i32(&value);
}
"IFMT" => {
self.ifmt = to_i32(&value);
}
"EOMR" => {
self.eomr = to_i32(&value);
}
"I32INP" => {
self.i32inp = to_i32(&value);
}
"I32OUT" => {
self.i32out = to_i32(&value);
}
"UI32INP" => {
self.ui32inp = to_u32(&value);
}
"UI32OUT" => {
self.ui32out = to_u32(&value);
}
"UI32MASK" => {
self.ui32mask = to_u32(&value);
}
"F64INP" => {
self.f64inp = to_f64(&value);
}
"F64OUT" => {
self.f64out = to_f64(&value);
}
"BAUD" => {
self.baud = to_i32(&value);
}
"LBAUD" => {
self.lbaud = to_i32(&value);
}
"PRTY" => {
self.prty = to_i32(&value);
}
"DBIT" => {
self.dbit = to_i32(&value);
}
"SBIT" => {
self.sbit = to_i32(&value);
}
"MCTL" => {
self.mctl = to_i32(&value);
}
"FCTL" => {
self.fctl = to_i32(&value);
}
"IXON" => {
self.ixon = to_i32(&value);
}
"IXOFF" => {
self.ixoff = to_i32(&value);
}
"IXANY" => {
self.ixany = to_i32(&value);
}
"HOSTINFO" => {
self.hostinfo = to_str(&value);
}
"DRTO" => {
self.drto = to_i32(&value);
}
"UCMD" => {
self.ucmd = to_i32(&value);
}
"ACMD" => {
self.acmd = to_i32(&value);
}
"SPR" => {
self.spr = to_i32(&value);
}
"TMSK" => {
self.tmsk = to_i32(&value);
}
"TB0" => {
self.tb0 = to_i32(&value);
}
"TB1" => {
self.tb1 = to_i32(&value);
}
"TB2" => {
self.tb2 = to_i32(&value);
}
"TB3" => {
self.tb3 = to_i32(&value);
}
"TB4" => {
self.tb4 = to_i32(&value);
}
"TB5" => {
self.tb5 = to_i32(&value);
}
"TIOM" => {
self.tiom = to_i32(&value);
}
"TIB0" => {
self.tib0 = to_i32(&value);
}
"TIB1" => {
self.tib1 = to_i32(&value);
}
"TIB2" => {
self.tib2 = to_i32(&value);
}
"TINM" => {
self.tinm = to_i32(&value);
}
"TINB0" => {
self.tinb0 = to_i32(&value);
}
"TINB1" => {
self.tinb1 = to_i32(&value);
}
"TINB2" => {
self.tinb2 = to_i32(&value);
}
"TINB3" => {
self.tinb3 = to_i32(&value);
}
"TSIZ" => {
self.tsiz = to_i32(&value);
}
"TFIL" => {
self.tfil = to_str(&value);
}
"AUCT" => {
self.auct = to_i32(&value);
}
"CNCT" => {
self.cnct = to_i32(&value);
}
"ENBL" => {
self.enbl = to_i32(&value);
}
"VAL" => {
self.val = to_i32(&value);
}
"ERRS" => {
self.errs = to_str(&value);
}
"AQR" => {
self.aqr = to_i32(&value);
}
_ => {
return Err(CaError::InvalidValue(format!("unknown field: {name}")));
}
}
Ok(())
}
fn init_record(&mut self, pass: u8) -> CaResult<()> {
if pass == 1 && !self.port.is_empty() {
let _ = self.connect_device();
}
Ok(())
}
fn special(&mut self, field: &str, after: bool) -> CaResult<()> {
if !after {
return Ok(());
}
if !SPC_MOD_FIELDS.contains(&field) {
return Ok(());
}
self.reset_error();
match field {
"PORT" | "ADDR" | "DRVINFO" => {
if let Err(manager_message) = self.connect_device() {
self.report_error(format!("connectDevice failed: {manager_message}"));
}
}
"TMSK" => {
self.update_trace_bits_from_mask();
self.apply_trace_mask();
}
"TB0" | "TB1" | "TB2" | "TB3" | "TB4" | "TB5" => {
self.update_mask_from_trace_bits();
self.apply_trace_mask();
}
"TIOM" => {
self.update_io_bits_from_mask();
self.apply_trace_io_mask();
}
"TIB0" | "TIB1" | "TIB2" => {
self.update_mask_from_io_bits();
self.apply_trace_io_mask();
}
"TINM" => {
self.update_info_bits_from_mask();
self.apply_trace_info_mask();
}
"TINB0" | "TINB1" | "TINB2" | "TINB3" => {
self.update_mask_from_info_bits();
self.apply_trace_info_mask();
}
"TSIZ" => {
self.apply_trace_truncate_size();
}
"TFIL" => {
self.apply_trace_file();
}
"ENBL" => {
if let Some(ref entry) = self.port_entry {
let _ = entry.handle.set_enable_blocking(self.enbl != 0);
}
}
"AUCT" => {
if let Some(ref entry) = self.port_entry {
let _ = entry.handle.set_auto_connect_blocking(self.auct != 0);
}
}
"CNCT" => self.special_callback(|this| {
let want = this.cnct != 0;
match this.port_entry {
Some(ref entry) => {
let handle = entry.handle.clone();
let cnct_user = || AsynUser::new(0).with_addr(this.addr);
let res = match (want, handle.is_connected_blocking()) {
(true, Ok(false)) => Some((
"connect",
handle
.submit_blocking(RequestOp::Connect, cnct_user())
.map(|_| ()),
)),
(false, Ok(true)) => Some((
"disconnect",
handle
.submit_blocking(RequestOp::Disconnect, cnct_user())
.map(|_| ()),
)),
(_, Ok(_)) => None,
(_, Err(e)) => {
this.report_error(format!(
"asynCallbackSpecial isConnected error: {e}"
));
None
}
};
if let Some((what, Err(e))) = res {
if e.never_ran() {
return this.report_special_never_ran(&e);
}
this.report_error(format!(
"asynCallbackSpecial callbackConnect {what}: {e}"
));
}
}
None => {
this.report_error("asynCallbackSpecial isConnected error");
}
}
SpecialRan::Yes
}),
"PCNCT" => {
if self.pcnct != 0 {
let _ = self.connect_device();
} else {
self.port_entry = None;
self.clear_exception_callback();
self.posting(&["CNCT"], |this| this.refresh_connected_state());
self.cancel_io_interrupt_scan();
self.publish_io_intr_binding();
}
}
"IFACE" => {
self.cancel_io_interrupt_scan();
self.publish_io_intr_binding();
}
"REASON" => self.special_callback(|this| {
this.resolved_reason = this.reason as usize;
this.drvinfo.clear();
this.cancel_io_interrupt_scan();
this.publish_io_intr_binding();
SpecialRan::Yes
}),
"BAUD" => {
let val = menu_choice(BAUD_CHOICES, self.baud);
self.write_option("baud", val);
}
"LBAUD" => {
self.write_option("baud", &self.lbaud.to_string());
}
"PRTY" => {
let val = menu_choice(PARITY_CHOICES, self.prty);
self.write_option("parity", val);
}
"DBIT" => {
let val = menu_choice(DBIT_CHOICES, self.dbit);
self.write_option("bits", val);
}
"SBIT" => {
let val = menu_choice(SBIT_CHOICES, self.sbit);
self.write_option("stop", val);
}
"MCTL" => {
let val = menu_choice(MCTL_CHOICES, self.mctl);
self.write_option("clocal", val);
}
"FCTL" => {
let val = menu_choice(FCTL_CHOICES, self.fctl);
self.write_option("crtscts", val);
}
"IXON" => {
let val = menu_choice(IX_CHOICES, self.ixon);
self.write_option("ixon", val);
}
"IXOFF" => {
let val = menu_choice(IX_CHOICES, self.ixoff);
self.write_option("ixoff", val);
}
"IXANY" => {
let val = menu_choice(IX_CHOICES, self.ixany);
self.write_option("ixany", val);
}
"HOSTINFO" => {
self.write_option("hostinfo", &self.hostinfo.clone());
}
"DRTO" => {
let val = menu_choice(DRTO_CHOICES, self.drto);
self.write_option("disconnectOnReadTimeout", val);
}
"AQR" => {
if let Some(inflight) = &self.io_inflight {
inflight.cancel.cancel();
}
}
"OEOS" => self.write_eos(true),
"IEOS" => self.write_eos(false),
"UI32MASK" => {
self.cancel_io_interrupt_scan();
self.publish_io_intr_binding();
}
_ => {}
}
Ok(())
}
fn process(&mut self) -> CaResult<ProcessOutcome> {
let outcome = self.process_cycle();
let went_async = matches!(
&outcome,
Ok(o) if matches!(
o.result,
RecordProcessResult::AsyncPending | RecordProcessResult::AsyncPendingNotify(_)
)
);
if !went_async {
self.io_intr.clear_sample();
}
outcome
}
fn set_io_intr_scan(&mut self, active: bool) {
if let Err(msg) = self.io_intr.set_active(active) {
self.report_error(msg);
}
}
fn as_any_mut(&mut self) -> Option<&mut dyn std::any::Any> {
Some(self)
}
fn clears_udf(&self) -> bool {
true
}
fn init_resets_alarms(&self) -> bool {
true
}
fn field_native_count(&self, field: &str) -> Option<u32> {
match field {
"BOUT" => Some(self.omax.max(0) as u32),
"BINP" => Some(self.imax.max(0) as u32),
_ => None,
}
}
}
impl AsynRecord {
fn process_cycle(&mut self) -> CaResult<ProcessOutcome> {
if let Some(inflight) = self.io_inflight.take() {
let outcome = inflight.result.lock().unwrap().take().unwrap_or_default();
self.apply_io_outcome(outcome);
return Ok(ProcessOutcome::complete());
}
if self.status_dirty.swap(false, Ordering::AcqRel) {
self.monitor_status();
}
self.reset_error();
let Some(entry) = self.port_entry.clone() else {
self.report_not_connected();
return Ok(ProcessOutcome::complete());
};
if let Some(sample) = self.io_intr.take_sample() {
self.apply_io_intr_sample(sample);
return Ok(ProcessOutcome::complete());
}
let plan = if let Some(cycle) = self.take_gpib_cycle() {
if !self.port_has(crate::interfaces::InterfaceType::Gpib) {
self.report_no_interface(crate::interfaces::InterfaceType::Gpib.asyn_name());
return Ok(ProcessOutcome::complete());
}
self.build_io_plan_for(Some(cycle))
} else {
if TransferMode::from_u16(self.tmod as u16) == TransferMode::NoIo {
return Ok(ProcessOutcome::complete());
}
let iface = InterfaceType::from_u16(self.iface as u16);
if !self.has_interface(iface) {
self.report_no_interface(iface.c_asyn_name());
return Ok(ProcessOutcome::complete());
}
self.build_io_plan()
};
let blocking_handle = entry.handle.can_block().then(|| entry.handle.clone());
if let (Some(handle), Some((name, db))) = (blocking_handle, self.async_ctx.clone()) {
return Ok(self.spawn_async_io(handle, name, db, plan));
}
self.perform_io(plan)?;
Ok(ProcessOutcome::complete())
}
}
impl Drop for AsynRecord {
fn drop(&mut self) {
self.clear_exception_callback();
}
}
#[cfg(test)]
#[allow(clippy::field_reassign_with_default)]
mod tests {
use super::*;
use epics_base_rs::server::record::RecordProcessResult;
#[test]
fn test_default_fields() {
let rec = AsynRecord::default();
assert_eq!(rec.record_type(), "asyn");
assert_eq!(rec.cnct, 0);
assert_eq!(rec.tmot, 1.0);
assert_eq!(rec.omax, 80);
assert_eq!(rec.imax, 80);
assert_eq!(rec.tsiz, 80);
assert_eq!(rec.ui32mask, 0xFFFFFFFF);
assert_eq!(rec.auct, 1);
assert_eq!(rec.enbl, 1);
}
#[test]
fn test_field_list_count() {
let rec = AsynRecord::default();
assert_eq!(rec.field_list().len(), 76);
}
#[test]
fn menu_fields_are_declared_and_served_as_enum() {
let rec = AsynRecord::default();
const MENU_FIELDS: &[&str] = &[
"TMOD", "IFACE", "OFMT", "IFMT", "EOMR", "TB0", "TB1", "TB2", "TB3", "TB4", "TB5",
"TIB0", "TIB1", "TIB2", "TINB0", "TINB1", "TINB2", "TINB3", "AUCT", "CNCT", "PCNCT",
"ENBL", "BAUD", "PRTY", "DBIT", "SBIT", "MCTL", "FCTL", "IXON", "IXOFF", "IXANY",
"DRTO", "UCMD", "ACMD",
];
assert_eq!(MENU_FIELDS.len(), 34, "asynRecord has 34 DBF_MENU fields");
for f in MENU_FIELDS {
let choices = rec.menu_field_choices(f);
assert!(choices.is_some(), "{f} must serve menu choices");
assert!(
!choices.unwrap().is_empty(),
"{f} choices must be non-empty"
);
let desc = rec
.field_list()
.iter()
.find(|d| d.name == *f)
.unwrap_or_else(|| panic!("{f} must be declared"));
assert_eq!(
desc.dbf_type,
DbFieldType::Enum,
"{f} is DBF_MENU: its native type must be Enum, as in C"
);
assert!(
matches!(rec.get_field(f), Some(EpicsValue::Enum(_))),
"{f} must be served as an Enum so a client reads the label"
);
}
for desc in rec.field_list() {
if rec.menu_field_choices(desc.name).is_some() {
assert_eq!(desc.dbf_type, DbFieldType::Enum, "{}", desc.name);
} else {
assert_ne!(desc.dbf_type, DbFieldType::Enum, "{}", desc.name);
}
}
}
#[test]
fn menu_choice_strings_match_dbd() {
let rec = AsynRecord::default();
assert_eq!(rec.menu_field_choices("TB0"), Some(&["Off", "On"][..]));
assert_eq!(
rec.menu_field_choices("IFACE"),
Some(&["asynOctet", "asynInt32", "asynUInt32Digital", "asynFloat64"][..])
);
assert_eq!(
rec.menu_field_choices("CNCT"),
Some(&["Disconnect", "Connect"][..])
);
assert_eq!(
rec.menu_field_choices("SBIT"),
Some(&["Unknown", "1", "2"][..])
);
let baud = rec.menu_field_choices("BAUD").unwrap();
assert_eq!(baud.len(), 16);
assert_eq!(baud[0], "Unknown");
assert_eq!(baud[15], "1152000");
}
#[test]
fn ui32_fields_are_unsigned_long() {
let mut rec = AsynRecord::default();
for f in ["UI32INP", "UI32OUT", "UI32MASK"] {
let desc = rec.field_list().iter().find(|d| d.name == f).unwrap();
assert_eq!(desc.dbf_type, DbFieldType::ULong, "{f}");
}
assert_eq!(
rec.get_field("UI32MASK"),
Some(EpicsValue::ULong(0xFFFF_FFFF))
);
rec.put_field("UI32OUT", EpicsValue::ULong(0x8000_0001))
.unwrap();
assert_eq!(
rec.get_field("UI32OUT"),
Some(EpicsValue::ULong(0x8000_0001))
);
rec.put_field("UI32MASK", EpicsValue::Double(4294967295.0))
.unwrap();
assert_eq!(
rec.get_field("UI32MASK"),
Some(EpicsValue::ULong(0xFFFF_FFFF))
);
}
#[test]
fn non_menu_fields_have_no_choices() {
let rec = AsynRecord::default();
for f in [
"PORT", "ADDR", "REASON", "TMSK", "TIOM", "TINM", "TSIZ", "LBAUD", "ERRS",
] {
assert_eq!(rec.menu_field_choices(f), None, "{f} is not a menu field");
}
}
#[test]
fn test_get_put_roundtrip() {
let mut rec = AsynRecord::default();
rec.put_field("PORT", EpicsValue::String("SIM1".into()))
.unwrap();
assert_eq!(
rec.get_field("PORT"),
Some(EpicsValue::String("SIM1".into()))
);
rec.put_field("ADDR", EpicsValue::Long(3)).unwrap();
assert_eq!(rec.get_field("ADDR"), Some(EpicsValue::Long(3)));
rec.put_field("TMOT", EpicsValue::Double(2.5)).unwrap();
assert_eq!(rec.get_field("TMOT"), Some(EpicsValue::Double(2.5)));
rec.put_field("F64OUT", EpicsValue::Double(3.14)).unwrap();
assert_eq!(rec.get_field("F64OUT"), Some(EpicsValue::Double(3.14)));
}
#[test]
fn test_trace_bit_sync() {
let mut rec = AsynRecord::default();
rec.tmsk = (TraceMask::ERROR | TraceMask::FLOW).bits() as i32;
rec.update_trace_bits_from_mask();
assert_eq!(rec.tb0, 1); assert_eq!(rec.tb4, 1); assert_eq!(rec.tb1, 0);
assert_eq!(rec.tb2, 0);
assert_eq!(rec.tb3, 0);
assert_eq!(rec.tb5, 0);
rec.tb0 = 1;
rec.tb1 = 1;
rec.tb2 = 0;
rec.tb3 = 0;
rec.tb4 = 0;
rec.tb5 = 1;
rec.update_mask_from_trace_bits();
let expected = TraceMask::ERROR | TraceMask::IO_DEVICE | TraceMask::WARNING;
assert_eq!(rec.tmsk, expected.bits() as i32);
}
#[test]
fn test_io_bit_sync() {
let mut rec = AsynRecord::default();
rec.tiom = (TraceIoMask::ASCII | TraceIoMask::HEX).bits() as i32;
rec.update_io_bits_from_mask();
assert_eq!(rec.tib0, 1); assert_eq!(rec.tib1, 0); assert_eq!(rec.tib2, 1); }
#[test]
fn test_info_bit_sync() {
let mut rec = AsynRecord::default();
rec.tinm = (TraceInfoMask::TIME | TraceInfoMask::THREAD).bits() as i32;
rec.update_info_bits_from_mask();
assert_eq!(rec.tinb0, 1); assert_eq!(rec.tinb1, 0); assert_eq!(rec.tinb2, 0); assert_eq!(rec.tinb3, 1); }
#[test]
fn test_connect_nonexistent_port() {
let mut rec = AsynRecord::default();
rec.port = "NONEXISTENT".to_string();
let _ = rec.connect_device();
assert_eq!(rec.cnct, 0);
assert!(rec.errs.contains("not found"));
}
#[test]
fn test_connect_empty_port() {
let mut rec = AsynRecord::default();
let _ = rec.connect_device();
assert_eq!(rec.cnct, 0);
assert!(rec.port_entry.is_none());
}
#[test]
fn test_process_no_io_mode() {
let mut rec = AsynRecord::default();
rec.tmod = TransferMode::NoIo as i32;
let result = rec.process().unwrap();
assert_eq!(result.result, RecordProcessResult::Complete);
}
#[test]
fn test_process_not_connected() {
let mut rec = AsynRecord::default();
rec.tmod = TransferMode::Read as i32;
rec.process().unwrap();
assert_eq!(rec.errs, "Not connect to a port");
}
#[test]
fn a_process_with_no_port_raises_state_minor() {
let mut rec = AsynRecord::default();
rec.tmod = TransferMode::Read as i32;
rec.process().unwrap();
assert_eq!(rec.errs, "Not connect to a port");
assert_eq!(
read_alarm(&mut rec),
(alarm_status::STATE_ALARM, AlarmSeverity::Minor),
"C asynRecord.c:361 alarms the stateNoDevice refusal STATE/MINOR"
);
rec.process().unwrap();
assert_eq!(
read_alarm(&mut rec),
(alarm_status::STATE_ALARM, AlarmSeverity::Minor),
"the refusal alarm is re-raised on the next process, not one-shot"
);
}
#[tokio::test]
async fn a_canceled_queued_request_raises_state_major() {
let entry = canblock_int32_entry(11);
let cancel = CancelToken::new();
assert!(cancel.cancel(), "a queued request cancels");
let mut rec = AsynRecord::default();
rec.tmod = TransferMode::Read as i32;
rec.iface = InterfaceType::Int32 as i32;
let plan = rec.build_io_plan();
let out = run_io_plan(entry.handle.clone(), plan, cancel).await;
rec.apply_io_outcome(out);
assert_eq!(rec.errs, "I/O request canceled");
assert_eq!(
read_alarm(&mut rec),
(alarm_status::STATE_ALARM, AlarmSeverity::Major),
"C asynRecord.c:399 alarms the canceled request STATE/MAJOR"
);
assert_eq!(
rec.i32inp, 0,
"the canceled request performed no I/O — the device value never landed"
);
}
#[test]
fn connect_failure_errs_texts_are_c_texts() {
let mut rec = AsynRecord::default();
rec.port = "NO_SUCH_PORT_R9_51".to_string();
rec.init_record(1).unwrap();
assert_eq!(
rec.errs,
"Connect error, status=3, asynManager:connectDevice port NO_SUCH_PORT_R9_51 not found"
);
let mut rec = AsynRecord::default();
rec.port = "NO_SUCH_PORT_R9_51".to_string();
rec.special("PORT", true).unwrap();
assert_eq!(
rec.errs,
"connectDevice failed: asynManager:connectDevice port NO_SUCH_PORT_R9_51 not found"
);
let mut rec = AsynRecord::default();
rec.special("PORT", true).unwrap();
assert_eq!(
rec.errs,
"connectDevice failed: asynManager:connectDevice no port name provided"
);
}
#[test]
fn spc_mod_fields_match_the_dbd() {
let dbd = [
"PORT", "ADDR", "PCNCT", "DRVINFO", "REASON", "IFACE", "OEOS", "IEOS", "UI32MASK",
"BAUD", "LBAUD", "PRTY", "DBIT", "SBIT", "MCTL", "FCTL", "IXON", "IXOFF", "IXANY",
"HOSTINFO", "DRTO", "TMSK", "TB0", "TB1", "TB2", "TB3", "TB4", "TB5", "TIOM", "TIB0",
"TIB1", "TIB2", "TINM", "TINB0", "TINB1", "TINB2", "TINB3", "TSIZ", "TFIL", "AUCT",
"CNCT", "ENBL", "AQR",
];
let mut ours = SPC_MOD_FIELDS.to_vec();
let mut theirs = dbd.to_vec();
ours.sort_unstable();
theirs.sort_unstable();
assert_eq!(ours, theirs);
}
#[test]
fn a_noio_process_clears_a_stale_errs() {
let calls = Arc::new(Mutex::new(GpibCalls::default()));
let (mut rec, _rt) = gpib_record(
"r10_47_noio",
GpibSpyPort::new("r10_47_noio", calls.clone()),
);
rec.tmod = TransferMode::NoIo as i32;
rec.errs = "Read error, timeout".to_string();
rec.process().unwrap();
assert_eq!(
rec.errs, "",
"C process() resetError (asynRecord.c:339) runs above the TMOD check"
);
}
#[test]
fn a_special_put_clears_a_stale_errs() {
let mut rec = AsynRecord::default();
rec.errs = "Write error, nout=0, timeout".to_string();
rec.tmsk = 0;
rec.special("TMSK", true).unwrap();
assert_eq!(
rec.errs, "",
"C special() resetError (asynRecord.c:390) runs before the field dispatch"
);
}
#[test]
fn a_non_spc_mod_put_keeps_errs() {
let mut rec = AsynRecord::default();
rec.errs = "Read error, timeout".to_string();
rec.special("VAL", true).unwrap();
rec.special("ERRS", true).unwrap();
assert_eq!(
rec.errs, "Read error, timeout",
"a non-SPC_MOD put does not reach C special() and does not reset ERRS"
);
}
#[test]
fn connect_device_entry_clears_the_previous_attempts_errs() {
use crate::interrupt::InterruptManager;
use crate::param::ParamType;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use tokio::sync::mpsc;
struct ParamDriver(PortDriverBase);
impl PortDriver for ParamDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
}
let port_name = "r10_47_reason0";
let mut base = PortDriverBase::new(port_name, 1, PortFlags::default());
assert_eq!(
base.params.create_param("FIRST", ParamType::Int32).unwrap(),
0
);
let (tx, rx) = mpsc::channel(16);
let actor = PortActor::new(Box::new(ParamDriver(base)), rx);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(16)),
actor_id,
);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.drvinfo = "FIRST".to_string();
rec.errs = "Connect error, status=3, port down".to_string();
rec.connect_device().unwrap();
assert_eq!(rec.resolved_reason, 0, "DRVINFO resolved to parameter 0");
assert_eq!(
rec.errs, "",
"a successful connect leaves no diagnostic behind, whatever the reason resolves to"
);
}
#[test]
fn hostinfo_put_and_option_readback_run_on_a_disconnected_port() {
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use tokio::sync::mpsc;
#[derive(Default)]
struct Log {
option_sets: Vec<(String, String)>,
eos_sets: Vec<Vec<u8>>,
}
struct DeadLineDriver {
base: PortDriverBase,
log: Arc<Mutex<Log>>,
host: Arc<Mutex<String>>,
}
impl PortDriver for DeadLineDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn set_option(
&mut self,
_user: &mut AsynUser,
key: &str,
value: &str,
) -> crate::error::AsynResult<()> {
self.log
.lock()
.unwrap()
.option_sets
.push((key.to_string(), value.to_string()));
if key == "hostinfo" {
*self.host.lock().unwrap() = value.to_string();
}
Ok(())
}
fn get_option(&self, key: &str) -> crate::error::AsynResult<String> {
match key {
"baud" => Ok("9600".to_string()),
"parity" => Ok("even".to_string()),
"hostinfo" => Ok(self.host.lock().unwrap().clone()),
_ => Err(crate::error::AsynError::Status {
status: crate::error::AsynStatus::Error,
message: format!("unsupported option {key}"),
}),
}
}
fn set_input_eos(
&mut self,
_user: &AsynUser,
eos: &[u8],
) -> crate::error::AsynResult<()> {
self.log.lock().unwrap().eos_sets.push(eos.to_vec());
Ok(())
}
}
let port_name = "r12_47_dead_line";
let log = Arc::new(Mutex::new(Log::default()));
let host = Arc::new(Mutex::new("oldhost:5000".to_string()));
let mut base = PortDriverBase::new(port_name, 1, PortFlags::default());
base.auto_connect = false;
base.set_connected(false);
let (tx, rx) = mpsc::channel(64);
let actor = PortActor::new(
Box::new(DeadLineDriver {
base,
log: log.clone(),
host: host.clone(),
}),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(16)),
actor_id,
);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
assert_eq!(
rec.baud,
baud_choice_index("9600"),
"C queues getOptions with QUEUE_EVEN_IF_NOT_CONNECTED, so BAUD is the \
port's real rate, not Unknown"
);
assert_eq!(rec.prty, 2, "…and PRTY likewise");
assert_eq!(rec.hostinfo, "oldhost:5000", "…and HOSTINFO likewise");
log.lock().unwrap().option_sets.clear();
rec.hostinfo = "newhost:5000".to_string();
rec.special("HOSTINFO", true).unwrap();
assert_eq!(
log.lock().unwrap().option_sets,
vec![("hostinfo".to_string(), "newhost:5000".to_string())],
"the only route to repoint a dead IP port must not be refused"
);
assert_eq!(rec.errs, "", "a waived put reports no error");
assert_eq!(rec.hostinfo, "newhost:5000");
log.lock().unwrap().option_sets.clear();
rec.baud = baud_choice_index("19200");
rec.special("BAUD", true).unwrap();
assert!(
log.lock().unwrap().option_sets.is_empty(),
"a BAUD put on a disconnected port is asynDisconnected in C, not a driver call"
);
assert!(
rec.errs.contains("not connected"),
"the gate's refusal message is what reaches ERRS: {}",
rec.errs
);
assert!(
!rec.errs.starts_with("Error setting option"),
"a refused put never entered setOption: {}",
rec.errs
);
rec.ieos = "\\r\\n".to_string();
rec.special("IEOS", true).unwrap();
assert!(
log.lock().unwrap().eos_sets.is_empty(),
"the EOS put must not inherit the HOSTINFO waiver"
);
}
#[test]
fn an_unknown_option_put_reaches_the_driver_and_refreshes_the_readbacks() {
use crate::error::{AsynError, AsynStatus};
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use tokio::sync::mpsc;
struct LoggingDriver {
base: PortDriverBase,
sets: Arc<Mutex<Vec<(String, String)>>>,
}
impl PortDriver for LoggingDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn set_option(
&mut self,
_user: &mut AsynUser,
key: &str,
value: &str,
) -> crate::error::AsynResult<()> {
self.sets
.lock()
.unwrap()
.push((key.to_string(), value.to_string()));
if value == "Unknown" {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: format!("Bad {key}"),
});
}
Ok(())
}
fn get_option(&self, key: &str) -> crate::error::AsynResult<String> {
match key {
"baud" => Ok("9600".to_string()),
"parity" => Ok("even".to_string()),
_ => Err(AsynError::Status {
status: AsynStatus::Error,
message: format!("unsupported option {key}"),
}),
}
}
}
let port_name = "r10_48_unknown_option";
let sets = Arc::new(Mutex::new(Vec::new()));
let (tx, rx) = mpsc::channel(64);
let actor = PortActor::new(
Box::new(LoggingDriver {
base: PortDriverBase::new(port_name, 1, PortFlags::default()),
sets: sets.clone(),
}),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(16)),
actor_id,
);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
sets.lock().unwrap().clear();
assert_eq!(rec.baud, baud_choice_index("9600"));
assert_eq!(rec.prty, 2);
rec.baud = 0;
rec.special("BAUD", true).unwrap();
assert_eq!(
*sets.lock().unwrap(),
vec![("baud".to_string(), "Unknown".to_string())],
"C sends baud_choices[0] = \"Unknown\" to the driver (asynRecord.c:1780)"
);
assert_eq!(
rec.errs, "Error setting option, Bad baud",
"the driver's rejection is C's \"Error setting option, %s\" (asynRecord.c:1828-1830)"
);
assert_eq!(
rec.baud,
baud_choice_index("9600"),
"the getOptions fall-through snaps BAUD back to the port's real rate"
);
assert_eq!(rec.lbaud, 9600, "…and LBAUD with it");
assert_eq!(
rec.prty, 2,
"every option readback refreshes, not just BAUD"
);
sets.lock().unwrap().clear();
rec.hostinfo = String::new();
rec.special("HOSTINFO", true).unwrap();
assert_eq!(
*sets.lock().unwrap(),
vec![("hostinfo".to_string(), String::new())],
"C gates the HOSTINFO set on nothing — the driver decides"
);
}
#[test]
fn trace_file_open_failure_errs_text_is_the_c_text() {
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use tokio::sync::mpsc;
struct TfilDriver(PortDriverBase);
impl PortDriver for TfilDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
}
let port_name = "r9_51_tfil";
let (tx, rx) = mpsc::channel(16);
let actor = PortActor::new(
Box::new(TfilDriver(PortDriverBase::new(
port_name,
1,
PortFlags::default(),
))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(16)),
actor_id,
);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
rec.tfil = "/nonexistent-dir-r9-51/trace.log".to_string();
rec.special("TFIL", true).unwrap();
assert_eq!(
rec.errs,
"Error opening trace file: /nonexistent-dir-r9-51/trace.log"
);
}
use crate::port::{PortDriver, PortDriverBase, PortFlags};
#[derive(Default)]
struct GpibCalls {
universal: Vec<u8>,
addressed: Vec<Vec<u8>>,
reads: usize,
}
struct GpibSpyPort {
base: PortDriverBase,
calls: Arc<Mutex<GpibCalls>>,
poll_byte: Option<u8>,
fail: Option<String>,
}
impl GpibSpyPort {
fn new(name: &str, calls: Arc<Mutex<GpibCalls>>) -> Self {
Self {
base: PortDriverBase::new(name, 1, PortFlags::default()),
calls,
poll_byte: None,
fail: None,
}
}
fn check(&self) -> AsynResult<()> {
match &self.fail {
Some(msg) => Err(crate::error::AsynError::Status {
status: crate::error::AsynStatus::Error,
message: msg.clone(),
}),
None => Ok(()),
}
}
}
impl PortDriver for GpibSpyPort {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn capabilities(&self) -> Vec<crate::interfaces::Capability> {
crate::interfaces::gpib::gpib_port_capabilities()
}
fn gpib_universal_cmd(&mut self, _user: &mut AsynUser, cmd: u8) -> AsynResult<()> {
self.calls.lock().unwrap().universal.push(cmd);
self.check()
}
fn gpib_addressed_cmd(&mut self, _user: &mut AsynUser, data: &[u8]) -> AsynResult<()> {
self.calls.lock().unwrap().addressed.push(data.to_vec());
self.check()
}
fn read_octet(&mut self, _user: &AsynUser, buf: &mut [u8]) -> AsynResult<usize> {
self.calls.lock().unwrap().reads += 1;
match self.poll_byte {
Some(b) if !buf.is_empty() => {
buf[0] = b;
Ok(1)
}
_ => Ok(0),
}
}
}
fn gpib_record(
port_name: &str,
port: GpibSpyPort,
) -> (AsynRecord, crate::runtime::PortRuntimeHandle) {
use crate::runtime::{RuntimeConfig, create_port_runtime};
let (rt, _jh) = create_port_runtime(port, RuntimeConfig::default());
register_port(
port_name,
rt.port_handle().clone(),
Arc::new(TraceManager::new()),
)
.unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.connect_device().unwrap();
(rec, rt)
}
#[test]
fn process_ucmd_sends_the_universal_command_to_the_port() {
let calls = Arc::new(Mutex::new(GpibCalls::default()));
let (mut rec, _rt) = gpib_record(
"r10_55_ucmd",
GpibSpyPort::new("r10_55_ucmd", calls.clone()),
);
rec.tmod = TransferMode::NoIo as i32;
rec.ucmd = 1; rec.process().unwrap();
assert_eq!(
calls.lock().unwrap().universal,
vec![crate::interfaces::gpib::IBDCL]
);
assert_eq!(rec.errs, "");
assert_eq!(rec.ucmd, 0, "UCMD resets to None after dispatch");
}
#[test]
fn process_acmd_sends_the_addressed_frame_to_the_port() {
use crate::interfaces::gpib::{IBGET, IBUNL, IBUNT, LADBASE};
let calls = Arc::new(Mutex::new(GpibCalls::default()));
let (mut rec, _rt) = gpib_record(
"r10_55_acmd",
GpibSpyPort::new("r10_55_acmd", calls.clone()),
);
rec.addr = 7;
rec.acmd = 1; rec.process().unwrap();
assert_eq!(
calls.lock().unwrap().addressed,
vec![vec![IBUNT, IBUNL, 7 + LADBASE, IBGET, IBUNT, IBUNL]]
);
assert_eq!(rec.errs, "");
assert_eq!(rec.acmd, 0, "ACMD resets to None after dispatch");
}
#[test]
fn process_acmd_serial_poll_runs_spe_read_spd() {
use crate::interfaces::gpib::{IBSPD, IBSPE};
let calls = Arc::new(Mutex::new(GpibCalls::default()));
let mut port = GpibSpyPort::new("r10_55_poll", calls.clone());
port.poll_byte = Some(0x41);
let (mut rec, _rt) = gpib_record("r10_55_poll", port);
rec.acmd = 5; rec.process().unwrap();
let seen = calls.lock().unwrap();
assert_eq!(seen.universal, vec![IBSPE, IBSPD]);
assert_eq!(seen.reads, 1, "one status-byte read between SPE and SPD");
assert!(seen.addressed.is_empty(), "serial poll sends no ACMD frame");
drop(seen);
assert_eq!(rec.spr, 0x41, "the status byte lands in SPR");
assert_eq!(rec.errs, "");
}
#[test]
fn a_failing_universal_command_reports_the_driver_text() {
let calls = Arc::new(Mutex::new(GpibCalls::default()));
let mut port = GpibSpyPort::new("r10_55_ucmd_fail", calls.clone());
port.fail = Some("prologixUniversalCmd unimplemented".to_string());
let (mut rec, _rt) = gpib_record("r10_55_ucmd_fail", port);
rec.ucmd = 1;
rec.process().unwrap();
assert_eq!(
rec.errs,
"GPIB Universal command prologixUniversalCmd unimplemented"
);
let mut c = CommonFields::default();
rec.check_alarms(&mut c);
assert_eq!(c.nsta, alarm_status::WRITE_ALARM);
assert_eq!(c.nsev, AlarmSeverity::Major);
}
#[test]
fn a_port_without_asyngpib_refuses_ucmd_and_acmd() {
use crate::interfaces::octet_transport_capabilities;
use crate::runtime::{RuntimeConfig, create_port_runtime};
struct OctetTransport(PortDriverBase);
impl PortDriver for OctetTransport {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn capabilities(&self) -> Vec<crate::interfaces::Capability> {
octet_transport_capabilities()
}
}
let port_name = "r10_55_no_gpib";
let (rt, _jh) = create_port_runtime(
OctetTransport(PortDriverBase::new(port_name, 1, PortFlags::default())),
RuntimeConfig::default(),
);
register_port(
port_name,
rt.port_handle().clone(),
Arc::new(TraceManager::new()),
)
.unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.connect_device().unwrap();
assert_eq!(rec.gpibiv, 0, "an octet transport has no asynGpib");
rec.tmod = TransferMode::NoIo as i32;
rec.ucmd = 1;
rec.process().unwrap();
assert_eq!(rec.errs, "No asynGpib interface");
assert_eq!(rec.ucmd, 0, "UCMD is consumed even when refused");
let mut c = CommonFields::default();
rec.check_alarms(&mut c);
assert_eq!(c.nsta, alarm_status::COMM_ALARM);
assert_eq!(c.nsev, AlarmSeverity::Major);
rec.acmd = 1;
rec.process().unwrap();
assert_eq!(rec.errs, "No asynGpib interface");
assert_eq!(rec.acmd, 0, "ACMD is consumed even when refused");
}
#[test]
fn a_record_with_no_port_leaves_ucmd_pending() {
let mut rec = AsynRecord::default();
rec.ucmd = 1;
rec.process().unwrap();
assert_eq!(rec.errs, "Not connect to a port");
assert_eq!(
rec.ucmd, 1,
"nothing was dispatched, so nothing is consumed"
);
let mut c = CommonFields::default();
rec.check_alarms(&mut c);
assert_eq!(c.nsta, alarm_status::STATE_ALARM);
assert_eq!(c.nsev, AlarmSeverity::Minor);
}
#[test]
fn process_ucmd_takes_priority_over_acmd() {
let calls = Arc::new(Mutex::new(GpibCalls::default()));
let (mut rec, _rt) = gpib_record(
"r10_55_priority",
GpibSpyPort::new("r10_55_priority", calls.clone()),
);
rec.ucmd = 1;
rec.acmd = 1;
rec.process().unwrap();
assert_eq!(rec.ucmd, 0, "UCMD consumed first");
assert_eq!(rec.acmd, 1, "ACMD left pending while UCMD was set");
assert_eq!(
calls.lock().unwrap().universal,
vec![crate::interfaces::gpib::IBDCL],
"only the universal command went to the bus"
);
assert!(calls.lock().unwrap().addressed.is_empty());
}
#[test]
fn test_special_trace_mask() {
let mut rec = AsynRecord::default();
rec.tmsk = (TraceMask::ERROR | TraceMask::WARNING | TraceMask::FLOW).bits() as i32;
rec.special("TMSK", true).unwrap();
assert_eq!(rec.tb0, 1); assert_eq!(rec.tb4, 1); assert_eq!(rec.tb5, 1); }
#[test]
fn test_special_trace_bits() {
let mut rec = AsynRecord::default();
rec.tb0 = 1;
rec.tb3 = 1;
rec.special("TB0", true).unwrap();
assert_eq!(
rec.tmsk as u32,
(TraceMask::ERROR | TraceMask::IO_DRIVER).bits()
);
}
#[test]
fn tfil_serves_unknown_with_no_live_port() {
let rec = AsynRecord::default();
assert_eq!(rec.tfil, "Unknown");
assert_eq!(
rec.get_field("TFIL"),
Some(EpicsValue::String("Unknown".into())),
"get(TFIL) must serve C's init_record seed with no live port"
);
}
#[tokio::test]
async fn init_resets_alarms_to_defined_no_alarm() {
use epics_base_rs::server::ioc_builder::IocBuilder;
use std::collections::HashMap;
let macros = HashMap::new();
let (name, factory) = crate::asyn_record::asyn_record_factory();
let (db, _autosave) = IocBuilder::new()
.register_record_type(name, factory)
.db_string("record(asyn, \"TEST:ASYN:UDF\") {}\n", ¯os)
.unwrap()
.build()
.await
.unwrap();
let rec = db.get_record("TEST:ASYN:UDF").expect("asyn record loaded");
let inst = rec.read();
assert_eq!(
inst.get_common_field("UDF"),
Some(EpicsValue::UChar(0)),
"C init_record pass 0: udf=0 (a device-config record is defined at load)"
);
assert_eq!(
inst.get_common_field("STAT"),
Some(EpicsValue::Short(0)),
"recGblResetAlarms → STAT=NO_ALARM"
);
assert_eq!(
inst.get_common_field("SEVR"),
Some(EpicsValue::Short(0)),
"recGblResetAlarms → SEVR=NO_ALARM"
);
}
#[test]
fn monitor_status_refreshes_tsiz_and_flags_a_foreign_trace_file() {
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use tokio::sync::mpsc;
struct TsizDriver(PortDriverBase);
impl PortDriver for TsizDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
}
let port_name = "r9_52_tsiz";
let (tx, rx) = mpsc::channel(16);
let actor = PortActor::new(
Box::new(TsizDriver(PortDriverBase::new(
port_name,
1,
PortFlags::default(),
))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(16)),
actor_id,
);
let trace = Arc::new(TraceManager::new());
register_port(port_name, handle, trace.clone()).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
assert_eq!(rec.tsiz, 80);
assert_eq!(rec.tfil, "Unknown");
trace.set_io_truncate_size(Some(port_name), 17);
rec.monitor_status();
assert_eq!(rec.tsiz, 17);
assert_eq!(rec.tfil, "Unknown", "the trace file did not change");
rec.tfil = "<stdout>".to_string();
rec.special("TFIL", true).unwrap();
rec.monitor_status();
assert_eq!(rec.tfil, "<stdout>");
trace.set_trace_file(Some(port_name), TraceFile::Stderr);
rec.monitor_status();
assert_eq!(rec.tfil, "Unknown");
}
#[test]
fn test_register_and_get_port() {
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use tokio::sync::mpsc;
struct TestDriver(PortDriverBase);
impl TestDriver {
fn new() -> Self {
Self(PortDriverBase::new(
"test_asyn_rec",
1,
PortFlags::default(),
))
}
}
impl PortDriver for TestDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
}
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(Box::new(TestDriver::new()), rx);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, "test_asyn_rec".into(), interrupts, actor_id);
let trace = Arc::new(TraceManager::new());
register_port("test_asyn_rec", handle, trace).unwrap();
let entry = registry::get_port("test_asyn_rec");
assert!(entry.is_some());
assert_eq!(entry.unwrap().handle.port_name(), "test_asyn_rec");
}
#[test]
fn trace_controls_route_to_device_on_multi_device_port() {
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use tokio::sync::mpsc;
struct MdDriver(PortDriverBase);
impl PortDriver for MdDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
}
let port_name = "test_trace_addr_md";
let flags = PortFlags {
multi_device: true,
..PortFlags::default()
};
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(MdDriver(PortDriverBase::new(port_name, 4, flags))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let mut handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
handle.set_capabilities(true, 4);
let trace = Arc::new(TraceManager::new());
register_port(port_name, handle, trace.clone()).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.addr = 3;
let _ = rec.connect_device();
assert_eq!(rec.cnct, 1);
assert!(rec.trace_addr_target() == Some(3));
rec.tmsk = (TraceMask::ERROR | TraceMask::FLOW).bits() as i32;
rec.apply_trace_mask();
rec.tsiz = 17;
rec.apply_trace_truncate_size();
assert!(trace.is_enabled_device(port_name, 3, TraceMask::FLOW));
assert!(!trace.is_enabled_device(port_name, 4, TraceMask::FLOW));
assert!(!trace.is_enabled(port_name, TraceMask::FLOW));
let mut rec0 = AsynRecord::default();
rec0.port = port_name.to_string();
rec0.addr = -1; let _ = rec0.connect_device();
assert!(rec0.trace_addr_target().is_none());
}
#[test]
fn read_trace_state_imports_info_mask_on_connect() {
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use tokio::sync::mpsc;
struct D(PortDriverBase);
impl PortDriver for D {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
}
let port_name = "test_trace_info_sync";
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(D(PortDriverBase::new(port_name, 1, PortFlags::default()))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
let trace = Arc::new(TraceManager::new());
trace.set_trace_info_mask(
Some(port_name),
TraceInfoMask::SOURCE | TraceInfoMask::THREAD,
);
register_port(port_name, handle, trace).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
assert_eq!(rec.cnct, 1);
assert_eq!(
rec.tinm as u32,
(TraceInfoMask::SOURCE | TraceInfoMask::THREAD).bits()
);
assert_eq!(rec.tinb0, 0); assert_eq!(rec.tinb1, 0); assert_eq!(rec.tinb2, 1); assert_eq!(rec.tinb3, 1); }
#[test]
fn external_trace_info_mask_reflected_after_process() {
use crate::exception::ExceptionManager;
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use tokio::sync::mpsc;
struct D(PortDriverBase);
impl PortDriver for D {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
}
let port_name = "test_trace_info_live";
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(D(PortDriverBase::new(port_name, 1, PortFlags::default()))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
let trace = Arc::new(TraceManager::new());
trace.set_exception_sink(Arc::new(ExceptionManager::new()));
trace.set_trace_info_mask(Some(port_name), TraceInfoMask::TIME);
register_port(port_name, handle, trace.clone()).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.tmod = TransferMode::NoIo as i32; let _ = rec.connect_device();
assert_eq!(rec.cnct, 1);
assert_eq!(rec.tinm as u32, TraceInfoMask::TIME.bits());
assert_eq!(rec.tinb0, 1); assert_eq!(rec.tinb1, 0);
trace.set_trace_info_mask(Some(port_name), TraceInfoMask::PORT | TraceInfoMask::THREAD);
assert_eq!(rec.tinm as u32, TraceInfoMask::TIME.bits());
rec.process().unwrap();
assert_eq!(
rec.tinm as u32,
(TraceInfoMask::PORT | TraceInfoMask::THREAD).bits()
);
assert_eq!(rec.tinb0, 0); assert_eq!(rec.tinb1, 1); assert_eq!(rec.tinb3, 1); }
#[test]
fn octet_read_updates_only_ifmt_selected_field() {
use crate::interpose::EomReason;
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::user::AsynUser;
use tokio::sync::mpsc;
const PAYLOAD: &[u8] = &[0xFF, b'A', b'B'];
struct OctetDriver(PortDriverBase);
impl PortDriver for OctetDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn io_read_octet_eom(
&mut self,
_user: &AsynUser,
buf: &mut [u8],
) -> crate::error::AsynResult<(usize, EomReason)> {
let n = PAYLOAD.len().min(buf.len());
buf[..n].copy_from_slice(&PAYLOAD[..n]);
Ok((n, EomReason::END))
}
}
let port_name = "test_ifmt_octet_read";
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(OctetDriver(PortDriverBase::new(
port_name,
1,
PortFlags::default(),
))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut ascii = AsynRecord::default();
ascii.port = port_name.to_string();
let _ = ascii.connect_device();
ascii.iface = 0; ascii.tmod = TransferMode::Read as i32;
ascii.imax = 256;
ascii.ifmt = ASYN_FMT_ASCII;
ascii.binp = b"SENTINEL".to_vec();
ascii.process().unwrap();
assert_eq!(ascii.errs, "");
assert_eq!(ascii.nord, PAYLOAD.len() as i32);
assert_eq!(ascii.ainp, String::from_utf8_lossy(PAYLOAD));
assert_eq!(
ascii.binp,
b"SENTINEL".to_vec(),
"ASCII read must not touch BINP"
);
assert!(!ascii.tinp.is_empty(), "TINP is posted for every read mode");
let mut binary = AsynRecord::default();
binary.port = port_name.to_string();
let _ = binary.connect_device();
binary.iface = 0;
binary.tmod = TransferMode::Read as i32;
binary.imax = 256;
binary.ifmt = ASYN_FMT_BINARY;
binary.ainp = "SENTINEL".to_string();
binary.process().unwrap();
assert_eq!(binary.errs, "");
assert_eq!(binary.nord, PAYLOAD.len() as i32);
assert_eq!(binary.binp, PAYLOAD.to_vec());
assert_eq!(binary.ainp, "SENTINEL", "Binary read must not touch AINP");
}
#[test]
fn escaped_record_fields_are_cut_at_their_c_buffer_size() {
use crate::interpose::EomReason;
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::user::AsynUser;
use tokio::sync::mpsc;
struct LongDriver(PortDriverBase);
impl PortDriver for LongDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn io_read_octet_eom(
&mut self,
_user: &AsynUser,
buf: &mut [u8],
) -> crate::error::AsynResult<(usize, EomReason)> {
let payload: Vec<u8> = b"\r\n".repeat(100);
let n = payload.len().min(buf.len());
buf[..n].copy_from_slice(&payload[..n]);
Ok((n, EomReason::END))
}
fn set_input_eos(
&mut self,
_user: &AsynUser,
_eos: &[u8],
) -> crate::error::AsynResult<()> {
Ok(())
}
fn get_input_eos(&self, _user: &AsynUser) -> Vec<u8> {
b"\r\n\r\n\r\n".to_vec()
}
}
let port_name = "r17_48_escape_bounds";
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(LongDriver(PortDriverBase::new(
port_name,
1,
PortFlags::default(),
))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
rec.iface = InterfaceType::Octet as i32;
rec.tmod = TransferMode::Read as i32;
rec.ifmt = ASYN_FMT_BINARY;
rec.imax = 256;
rec.process().unwrap();
assert_eq!(rec.errs, "");
assert_eq!(rec.nord, 200, "NORD is the raw transfer count, unbounded");
assert_eq!(
rec.tinp.len(),
TINP_SIZE - 1,
"TINP is a 40-byte DBF_STRING"
);
assert!(
rec.tinp.ends_with(r"\r\n\r\"),
"C cuts mid escape pair and leaves the backslash: {}",
rec.tinp
);
rec.ieos = r"\r".to_string();
rec.special("IEOS", true).unwrap();
assert_eq!(rec.errs, "");
assert_eq!(rec.ieos, r"\r\n\r\n\", "IEOS is cut at EOS_SIZE - 1");
assert_eq!(rec.ieos.len(), EOS_SIZE - 1);
}
#[test]
fn octet_error_resets_transfer_fields() {
use crate::error::{AsynError, AsynStatus};
use crate::interpose::EomReason;
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::user::AsynUser;
use tokio::sync::mpsc;
struct FailingOctetDriver(PortDriverBase);
impl PortDriver for FailingOctetDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn io_read_octet_eom(
&mut self,
_user: &AsynUser,
_buf: &mut [u8],
) -> crate::error::AsynResult<(usize, EomReason)> {
Err(AsynError::Status {
status: AsynStatus::Timeout,
message: "read boom".into(),
})
}
fn io_write_octet(
&mut self,
_user: &mut AsynUser,
_data: &[u8],
) -> crate::error::AsynResult<usize> {
Err(AsynError::Status {
status: AsynStatus::Timeout,
message: "write boom".into(),
})
}
}
let port_name = "test_octet_error_reset";
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(FailingOctetDriver(PortDriverBase::new(
port_name,
1,
PortFlags::default(),
))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut ascii = AsynRecord::default();
ascii.port = port_name.to_string();
let _ = ascii.connect_device();
ascii.iface = 0; ascii.tmod = TransferMode::Read as i32;
ascii.imax = 256;
ascii.ifmt = ASYN_FMT_ASCII;
ascii.nord = 5;
ascii.eomr = 2;
ascii.ainp = "STALE".to_string();
ascii.tinp = "STALE".to_string();
ascii.binp = b"KEEP".to_vec();
ascii.process().unwrap();
assert!(!ascii.errs.is_empty(), "read error must set ERRS");
assert_eq!(ascii.nord, 0, "failed read must reset NORD to 0");
assert_eq!(ascii.eomr, 0, "failed read must reset EOMR to 0");
assert_eq!(ascii.ainp, "", "failed ASCII read must clear AINP");
assert_eq!(ascii.tinp, "", "failed read must clear TINP");
assert_eq!(
ascii.binp,
b"KEEP".to_vec(),
"ASCII path must not touch BINP"
);
let mut binary = AsynRecord::default();
binary.port = port_name.to_string();
let _ = binary.connect_device();
binary.iface = 0;
binary.tmod = TransferMode::Read as i32;
binary.imax = 256;
binary.ifmt = ASYN_FMT_BINARY;
binary.nord = 7;
binary.eomr = 3;
binary.binp = b"STALE".to_vec();
binary.ainp = "KEEP".to_string();
binary.tinp = "STALE".to_string();
binary.process().unwrap();
assert!(!binary.errs.is_empty(), "read error must set ERRS");
assert_eq!(binary.nord, 0, "failed read must reset NORD to 0");
assert_eq!(binary.eomr, 0, "failed read must reset EOMR to 0");
assert_eq!(
binary.binp,
Vec::<u8>::new(),
"failed binary read must clear BINP"
);
assert_eq!(binary.tinp, "", "failed read must clear TINP");
assert_eq!(binary.ainp, "KEEP", "Binary path must not touch AINP");
let mut writer = AsynRecord::default();
writer.port = port_name.to_string();
let _ = writer.connect_device();
writer.iface = 0;
writer.tmod = TransferMode::Write as i32;
writer.ofmt = ASYN_FMT_ASCII;
writer.aout = "hello".to_string();
writer.nawt = 9;
writer.process().unwrap();
assert!(!writer.errs.is_empty(), "write error must set ERRS");
assert_eq!(writer.nawt, 0, "failed write must reset NAWT to 0");
}
#[test]
fn register_and_option_errors_use_the_c_errs_text() {
use crate::error::{AsynError, AsynStatus};
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::user::AsynUser;
use tokio::sync::mpsc;
fn boom(what: &str) -> AsynError {
AsynError::Status {
status: AsynStatus::Error,
message: format!("{what} boom"),
}
}
struct AllFailDriver(PortDriverBase);
impl PortDriver for AllFailDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn write_int32(&mut self, _u: &mut AsynUser, _v: i32) -> crate::error::AsynResult<()> {
Err(boom("i32w"))
}
fn read_int32(&mut self, _u: &AsynUser) -> crate::error::AsynResult<i32> {
Err(boom("i32r"))
}
fn write_uint32_digital(
&mut self,
_u: &mut AsynUser,
_v: u32,
_m: u32,
) -> crate::error::AsynResult<()> {
Err(boom("u32w"))
}
fn read_uint32_digital(
&mut self,
_u: &AsynUser,
_m: u32,
) -> crate::error::AsynResult<u32> {
Err(boom("u32r"))
}
fn write_float64(
&mut self,
_u: &mut AsynUser,
_v: f64,
) -> crate::error::AsynResult<()> {
Err(boom("f64w"))
}
fn read_float64(&mut self, _u: &AsynUser) -> crate::error::AsynResult<f64> {
Err(boom("f64r"))
}
fn set_option(
&mut self,
_user: &mut AsynUser,
_k: &str,
_v: &str,
) -> crate::error::AsynResult<()> {
Err(boom("opt"))
}
fn set_input_eos(
&mut self,
_user: &AsynUser,
_eos: &[u8],
) -> crate::error::AsynResult<()> {
Err(boom("ieos"))
}
fn set_output_eos(
&mut self,
_user: &AsynUser,
_eos: &[u8],
) -> crate::error::AsynResult<()> {
Err(boom("oeos"))
}
}
let port_name = "test_register_errs_text";
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(AllFailDriver(PortDriverBase::new(
port_name,
1,
PortFlags::default(),
))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let reg = |iface: InterfaceType, tmod: TransferMode| -> String {
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
rec.iface = iface as i32;
rec.tmod = tmod as i32;
rec.process().unwrap();
rec.errs.clone()
};
assert_eq!(
reg(InterfaceType::Int32, TransferMode::Write),
"Int32 write error, i32w boom"
);
assert_eq!(
reg(InterfaceType::Int32, TransferMode::Read),
"Int32 read error, i32r boom"
);
assert_eq!(
reg(InterfaceType::UInt32Digital, TransferMode::Write),
"UInt32 write error, u32w boom",
"C writes UInt32, not UInt32Digital"
);
assert_eq!(
reg(InterfaceType::UInt32Digital, TransferMode::Read),
"UInt32 read error, u32r boom"
);
assert_eq!(
reg(InterfaceType::Float64, TransferMode::Write),
"Float64 write error, f64w boom"
);
assert_eq!(
reg(InterfaceType::Float64, TransferMode::Read),
"Float64 read error, f64r boom"
);
let mut opt = AsynRecord::default();
opt.port = port_name.to_string();
let _ = opt.connect_device();
opt.lbaud = 9600;
opt.special("LBAUD", true).unwrap();
assert_eq!(opt.errs, "Error setting option, opt boom");
opt.ieos = "\\n".to_string();
opt.special("IEOS", true).unwrap();
assert_eq!(opt.errs, "Error setting input eos, ieos boom");
opt.oeos = "\\r".to_string();
opt.special("OEOS", true).unwrap();
assert_eq!(opt.errs, "Error setting output eos, oeos boom");
assert_eq!(opt.ieos, "");
assert_eq!(opt.oeos, "");
}
#[test]
fn option_and_eos_puts_fall_through_to_a_driver_re_read() {
use crate::error::{AsynError, AsynStatus};
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use tokio::sync::mpsc;
struct StubbornDriver(PortDriverBase);
impl PortDriver for StubbornDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn set_option(
&mut self,
_user: &mut AsynUser,
_k: &str,
_v: &str,
) -> crate::error::AsynResult<()> {
Ok(())
}
fn get_option(&self, key: &str) -> crate::error::AsynResult<String> {
match key {
"baud" => Ok("9600".to_string()),
"parity" => Ok("even".to_string()),
_ => Err(AsynError::Status {
status: AsynStatus::Error,
message: format!("unsupported option {key}"),
}),
}
}
fn set_input_eos(
&mut self,
_user: &AsynUser,
_eos: &[u8],
) -> crate::error::AsynResult<()> {
Err(AsynError::Status {
status: AsynStatus::Error,
message: "ieos boom".to_string(),
})
}
fn get_input_eos(&self, _user: &AsynUser) -> Vec<u8> {
b"\r\n".to_vec()
}
fn set_output_eos(
&mut self,
_user: &AsynUser,
_eos: &[u8],
) -> crate::error::AsynResult<()> {
Err(AsynError::Status {
status: AsynStatus::Error,
message: "oeos boom".to_string(),
})
}
fn get_output_eos(&self, _user: &AsynUser) -> Vec<u8> {
b"\n".to_vec()
}
}
let port_name = "test_option_eos_reread";
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(StubbornDriver(PortDriverBase::new(
port_name,
1,
PortFlags::default(),
))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
rec.lbaud = 115_200;
rec.special("LBAUD", true).unwrap();
assert_eq!(rec.lbaud, 9600, "LBAUD is a driver readback, not a latch");
assert_eq!(rec.baud, baud_choice_index("9600"));
assert_eq!(rec.errs, "", "the set succeeded — no ERRS");
assert_eq!(rec.prty, 2, "parity readback: even");
rec.ieos = "\\n".to_string();
rec.special("IEOS", true).unwrap();
assert_eq!(rec.errs, "Error setting input eos, ieos boom");
assert_eq!(rec.ieos, "\\r\\n", "IEOS snaps back to the driver's EOS");
assert_eq!(
rec.oeos, "\\n",
"an IEOS put re-reads the output EOS too (asynRecord.c:2009)"
);
}
#[test]
fn flush_failure_is_discarded_and_does_not_reach_errs() {
use crate::error::{AsynError, AsynStatus};
use crate::interpose::EomReason;
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::user::AsynUser;
use tokio::sync::mpsc;
struct FlushFailsDriver(PortDriverBase);
impl PortDriver for FlushFailsDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn io_flush(&mut self, _user: &mut AsynUser) -> crate::error::AsynResult<()> {
Err(AsynError::Status {
status: AsynStatus::Error,
message: "flush boom".into(),
})
}
fn io_write_octet(
&mut self,
_user: &mut AsynUser,
data: &[u8],
) -> crate::error::AsynResult<usize> {
Ok(data.len())
}
fn io_read_octet_eom(
&mut self,
_user: &AsynUser,
buf: &mut [u8],
) -> crate::error::AsynResult<(usize, EomReason)> {
buf[..2].copy_from_slice(b"OK");
Ok((2, EomReason::END))
}
}
let port_name = "test_flush_failure_discarded";
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(FlushFailsDriver(PortDriverBase::new(
port_name,
1,
PortFlags::default(),
))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
rec.iface = InterfaceType::Octet as i32;
rec.tmod = TransferMode::WriteRead as i32;
rec.ofmt = ASYN_FMT_ASCII;
rec.ifmt = ASYN_FMT_ASCII;
rec.aout = "CMD".to_string();
rec.process().unwrap();
assert_eq!(rec.ainp, "OK", "the write/read after the flush still ran");
assert_eq!(
rec.errs, "",
"C discards the flush status; a good transaction reports nothing"
);
assert_eq!(read_alarm(&mut rec).1, AlarmSeverity::NoAlarm);
rec.tmod = TransferMode::Flush as i32;
rec.process().unwrap();
assert_eq!(rec.errs, "", "a flush-only cycle reports nothing either");
}
#[test]
fn octet_read_error_errs_matches_c_report_error() {
use crate::error::{AsynError, AsynStatus};
use crate::interpose::{EomReason, PartialOctetRead};
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::user::AsynUser;
use tokio::sync::mpsc;
struct StatusDriver {
base: PortDriverBase,
n: usize,
}
impl PortDriver for StatusDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn io_read_octet_eom(
&mut self,
_user: &AsynUser,
_buf: &mut [u8],
) -> crate::error::AsynResult<(usize, EomReason)> {
self.n += 1;
match self.n {
1 => Err(AsynError::Status {
status: AsynStatus::Timeout,
message: "read timeout".into(),
}),
2 => Err(AsynError::Status {
status: AsynStatus::Overflow,
message: "buffer full".into(),
}),
_ => Err(AsynError::Status {
status: AsynStatus::Error,
message: "device fault".into(),
}
.with_partial_read(PartialOctetRead {
data: b"xy".to_vec(),
eom_reason: EomReason::empty(),
})),
}
}
}
let port_name = "test_octet_read_errs_text";
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(StatusDriver {
base: PortDriverBase::new(port_name, 1, PortFlags::default()),
n: 0,
}),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
rec.iface = InterfaceType::Octet as i32;
rec.tmod = TransferMode::Read as i32;
rec.ifmt = ASYN_FMT_ASCII;
rec.process().unwrap();
assert_eq!(rec.errs, "timeout nread 0 read timeout");
rec.process().unwrap();
assert_eq!(rec.errs, "overflow nread 0 buffer full");
rec.process().unwrap();
assert_eq!(
rec.errs, "error nread 2 device fault",
"the count is the bytes the failing read did deliver"
);
}
#[test]
fn octet_partial_read_delivers_bytes_with_the_timeout() {
use crate::error::{AsynError, AsynStatus};
use crate::interpose::{EomReason, PartialOctetRead};
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::user::AsynUser;
use epics_base_rs::server::recgbl::alarm_status;
use epics_base_rs::server::record::{AlarmSeverity, CommonFields};
use tokio::sync::mpsc;
struct PartialThenTimeout(PortDriverBase);
impl PortDriver for PartialThenTimeout {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn io_read_octet_eom(
&mut self,
_user: &AsynUser,
buf: &mut [u8],
) -> crate::error::AsynResult<(usize, EomReason)> {
buf[..3].copy_from_slice(b"abc");
buf[3] = 0;
Err(AsynError::Status {
status: AsynStatus::Timeout,
message: "read timeout".into(),
}
.with_partial_read(PartialOctetRead {
data: b"abc".to_vec(),
eom_reason: EomReason::empty(),
}))
}
}
let port_name = "test_octet_partial_read";
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(PartialThenTimeout(PortDriverBase::new(
port_name,
1,
PortFlags::default(),
))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
rec.iface = 0; rec.tmod = TransferMode::Read as i32;
rec.imax = 256;
rec.ifmt = ASYN_FMT_ASCII;
rec.process().unwrap();
assert_eq!(rec.ainp, "abc", "the partial line reaches AINP (C :1503)");
assert_eq!(rec.nord, 3, "NORD = nbytesTransfered (C :1627)");
assert_eq!(rec.eomr, 0, "no EOS matched, buffer never filled (C :1591)");
assert_eq!(rec.tinp, "abc", "TINP is posted for the failed read too");
assert_eq!(
rec.errs, "timeout nread 3 read timeout",
"C :1593-1598 \"%s nread %d %s\" with the transferred count"
);
let mut c = CommonFields::default();
rec.check_alarms(&mut c);
assert_eq!(c.nsta, alarm_status::READ_ALARM, "C :1599 recGblSetSevr");
assert_eq!(c.nsev, AlarmSeverity::Major, "READ_ALARM is MAJOR");
let mut bin = AsynRecord::default();
bin.port = port_name.to_string();
let _ = bin.connect_device();
bin.iface = 0;
bin.tmod = TransferMode::Read as i32;
bin.imax = 256;
bin.ifmt = ASYN_FMT_BINARY;
bin.process().unwrap();
assert_eq!(bin.binp, b"abc".to_vec(), "partial bytes reach BINP too");
assert_eq!(bin.nord, 3);
}
#[test]
fn cnct_drives_the_transport_and_pcnct_the_attachment() {
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::user::AsynUser;
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::sync::mpsc;
struct CountingDriver {
base: PortDriverBase,
connects: Arc<AtomicUsize>,
disconnects: Arc<AtomicUsize>,
}
impl PortDriver for CountingDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn connect(&mut self, _user: &AsynUser) -> crate::error::AsynResult<()> {
self.connects.fetch_add(1, Ordering::SeqCst);
self.base.set_connected(true);
Ok(())
}
fn disconnect(&mut self, _user: &AsynUser) -> crate::error::AsynResult<()> {
self.disconnects.fetch_add(1, Ordering::SeqCst);
self.base.set_connected(false);
Ok(())
}
}
let port_name = "test_cnct_transport";
let connects = Arc::new(AtomicUsize::new(0));
let disconnects = Arc::new(AtomicUsize::new(0));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(CountingDriver {
base: PortDriverBase::new(port_name, 1, PortFlags::default()),
connects: connects.clone(),
disconnects: disconnects.clone(),
}),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(256)),
actor_id,
);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
assert_eq!(rec.cnct, 1, "CNCT is the port's transport state");
assert_eq!(rec.pcnct, 1, "PCNCT is the attachment");
let base_connects = connects.load(Ordering::SeqCst);
rec.cnct = 0;
rec.special("CNCT", true).unwrap();
assert_eq!(
disconnects.load(Ordering::SeqCst),
1,
"CNCT=0 must disconnect the driver's transport"
);
assert_eq!(rec.cnct, 0, "readback follows the wire");
assert_eq!(
rec.pcnct, 1,
"CNCT must not detach the record (that is PCNCT)"
);
assert!(
rec.port_entry.is_some(),
"CNCT must not drop the port binding"
);
rec.cnct = 1;
rec.special("CNCT", true).unwrap();
assert_eq!(
connects.load(Ordering::SeqCst),
base_connects,
"a device-addressed CNCT=1 put must be refused on a disconnected \
port (C checkPortConnect)"
);
assert!(
rec.errs.contains("not connected"),
"the gate's refusal message is what reaches ERRS: {}",
rec.errs
);
assert!(
!rec.errs.contains("callbackConnect"),
"a refused request never entered the callback: {}",
rec.errs
);
assert_eq!(
rec.cnct, 1,
"no callback ran, so no monitorStatus tail snapped CNCT back \
(asynRecord.c:571-576 vs :897)"
);
rec.port_entry
.as_ref()
.unwrap()
.handle
.connect_blocking()
.unwrap();
rec.cnct = 1;
rec.special("CNCT", true).unwrap();
assert_eq!(rec.cnct, 1);
rec.cnct = 1;
rec.special("CNCT", true).unwrap();
assert_eq!(
connects.load(Ordering::SeqCst),
base_connects + 1,
"an already-connected port must not be re-connected"
);
let disconnects_before = disconnects.load(Ordering::SeqCst);
rec.pcnct = 0;
rec.special("PCNCT", true).unwrap();
assert!(rec.port_entry.is_none(), "PCNCT=0 detaches the record");
assert_eq!(
disconnects.load(Ordering::SeqCst),
disconnects_before,
"PCNCT must not touch the driver's transport"
);
assert_eq!(
rec.cnct, 0,
"detached: isConnected has no device to report on (C :1091)"
);
}
#[test]
fn a_gate_refused_request_reports_the_refusal_and_runs_no_callback() {
use crate::error::{AsynError, AsynStatus};
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::sync::mpsc;
#[derive(Default)]
struct Log {
option_sets: Vec<(String, String)>,
eos_sets: Vec<Vec<u8>>,
reads: usize,
}
struct GateDriver {
base: PortDriverBase,
log: Arc<Mutex<Log>>,
connects: Arc<AtomicUsize>,
}
impl PortDriver for GateDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn set_option(
&mut self,
_user: &mut AsynUser,
key: &str,
value: &str,
) -> crate::error::AsynResult<()> {
self.log
.lock()
.unwrap()
.option_sets
.push((key.to_string(), value.to_string()));
if value == "Unknown" {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: format!("Bad {key}"),
});
}
Ok(())
}
fn get_option(&self, key: &str) -> crate::error::AsynResult<String> {
match key {
"baud" => Ok("9600".to_string()),
_ => Err(AsynError::Status {
status: AsynStatus::Error,
message: format!("unsupported option {key}"),
}),
}
}
fn set_input_eos(
&mut self,
_user: &AsynUser,
eos: &[u8],
) -> crate::error::AsynResult<()> {
self.log.lock().unwrap().eos_sets.push(eos.to_vec());
Ok(())
}
fn read_octet(
&mut self,
_user: &AsynUser,
buf: &mut [u8],
) -> crate::error::AsynResult<usize> {
self.log.lock().unwrap().reads += 1;
let data = b"hi";
buf[..data.len()].copy_from_slice(data);
Ok(data.len())
}
fn connect(&mut self, _user: &AsynUser) -> crate::error::AsynResult<()> {
self.connects.fetch_add(1, Ordering::SeqCst);
self.base.set_connected(true);
Ok(())
}
}
let port_name = "r14_46_gate_refusal";
let log = Arc::new(Mutex::new(Log::default()));
let connects = Arc::new(AtomicUsize::new(0));
let (tx, rx) = mpsc::channel(64);
let actor = PortActor::new(
Box::new(GateDriver {
base: PortDriverBase::new(port_name, 1, PortFlags::default()),
log: log.clone(),
connects: connects.clone(),
}),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(16)),
actor_id,
);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
assert_eq!(rec.baud, baud_choice_index("9600"), "the port runs at 9600");
let entry_handle = rec.port_entry.as_ref().unwrap().handle.clone();
log.lock().unwrap().option_sets.clear();
rec.baud = 0; rec.special("BAUD", true).unwrap();
assert_eq!(
log.lock().unwrap().option_sets,
vec![("baud".to_string(), "Unknown".to_string())],
"control: the request ran, so the driver saw the put"
);
assert_eq!(rec.errs, "Error setting option, Bad baud");
assert_eq!(
rec.baud,
baud_choice_index("9600"),
"control: the callback's readback tail ran"
);
entry_handle.set_enable_blocking(false).unwrap();
log.lock().unwrap().option_sets.clear();
rec.baud = baud_choice_index("19200");
rec.special("BAUD", true).unwrap();
assert!(
log.lock().unwrap().option_sets.is_empty(),
"a refused put must not reach the driver"
);
assert!(
rec.errs.contains("disabled"),
"ERRS is the gate's errorMessage (asynRecord.c:575): {}",
rec.errs
);
assert_eq!(
rec.baud,
baud_choice_index("19200"),
"no callback ran, so no getOptions readback snapped BAUD back"
);
rec.tmod = TransferMode::Read as i32;
rec.nord = 0;
rec.process().unwrap();
assert_eq!(
log.lock().unwrap().reads,
0,
"a refused process request must not reach the driver"
);
assert_eq!(rec.errs, "queueRequest failed");
assert_eq!(rec.nord, 0, "no transfer, so no NORD to publish");
assert_eq!(
read_alarm(&mut rec),
(alarm_status::STATE_ALARM, AlarmSeverity::Minor),
"C alarms the refusal STATE/MINOR (:361), not COMM/MAJOR"
);
entry_handle.set_enable_blocking(true).unwrap();
entry_handle.disconnect_blocking().unwrap();
log.lock().unwrap().eos_sets.clear();
rec.ieos = "\\n".to_string();
rec.special("IEOS", true).unwrap();
assert!(
log.lock().unwrap().eos_sets.is_empty(),
"a refused EOS put must not reach the driver"
);
assert!(
rec.errs.contains("not connected"),
"ERRS is the gate's errorMessage: {}",
rec.errs
);
let connects_before = connects.load(Ordering::SeqCst);
rec.cnct = 1;
rec.special("CNCT", true).unwrap();
assert_eq!(
connects.load(Ordering::SeqCst),
connects_before,
"a refused CNCT put must not touch the wire"
);
assert!(
rec.errs.contains("not connected"),
"ERRS is the gate's errorMessage, not a callback-shaped text: {}",
rec.errs
);
assert_eq!(rec.cnct, 1, "no callback ran, so no monitorStatus tail");
}
#[test]
fn octet_write_reports_the_transferred_count_on_both_arms() {
use crate::error::{AsynError, AsynStatus};
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::user::AsynUser;
use epics_base_rs::server::record::{AlarmSeverity, CommonFields};
use tokio::sync::mpsc;
struct ShortWriteDriver {
base: PortDriverBase,
accept: usize,
fail: bool,
}
impl PortDriver for ShortWriteDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn io_write_octet(
&mut self,
_user: &mut AsynUser,
_data: &[u8],
) -> crate::error::AsynResult<usize> {
if self.fail {
Err(AsynError::Status {
status: AsynStatus::Timeout,
message: "serial write timeout".into(),
}
.with_partial_write(self.accept))
} else {
Ok(self.accept)
}
}
}
fn writer(port_name: &'static str, accept: usize, fail: bool) -> AsynRecord {
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(ShortWriteDriver {
base: PortDriverBase::new(port_name, 1, PortFlags::default()),
accept,
fail,
}),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
rec.iface = 0; rec.tmod = TransferMode::Write as i32;
rec.ofmt = ASYN_FMT_ASCII;
rec.aout = "hello".to_string(); rec.nawt = 99; rec
}
let mut partial = writer("test_short_write_err", 3, true);
partial.process().unwrap();
assert_eq!(partial.nawt, 3, "NAWT = nbytesTransfered (C :1547)");
assert!(
partial.errs.starts_with("Write error, nout=3,"),
"C :1552 reportError text, got {:?}",
partial.errs
);
let mut c = CommonFields::default();
partial.check_alarms(&mut c);
assert_eq!(
c.nsev,
AlarmSeverity::NoAlarm,
"octet write error raises no record severity"
);
let mut short_ok = writer("test_short_write_ok", 2, false);
short_ok.process().unwrap();
assert_eq!(short_ok.nawt, 2, "NAWT comes from the reply, not from OPTR");
assert!(
short_ok.errs.starts_with("Write error, nout=2,"),
"C :1551 fires on a short write even with asynSuccess, got {:?}",
short_ok.errs
);
let mut full = writer("test_full_write", 5, false);
full.process().unwrap();
assert_eq!(full.nawt, 5);
assert!(full.errs.is_empty(), "a complete write reports nothing");
}
#[test]
fn io_errors_raise_record_alarm() {
use crate::error::{AsynError, AsynStatus};
use crate::interpose::EomReason;
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::user::AsynUser;
use epics_base_rs::server::recgbl::alarm_status;
use epics_base_rs::server::record::{AlarmSeverity, CommonFields};
use tokio::sync::mpsc;
fn boom() -> AsynError {
AsynError::Status {
status: AsynStatus::Timeout,
message: "boom".into(),
}
}
struct FailDriver(PortDriverBase);
impl PortDriver for FailDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn io_read_octet_eom(
&mut self,
_user: &AsynUser,
_buf: &mut [u8],
) -> crate::error::AsynResult<(usize, EomReason)> {
Err(boom())
}
fn io_write_octet(
&mut self,
_user: &mut AsynUser,
_data: &[u8],
) -> crate::error::AsynResult<usize> {
Err(boom())
}
fn io_read_int32(&mut self, _user: &AsynUser) -> crate::error::AsynResult<i32> {
Err(boom())
}
fn io_write_int32(
&mut self,
_user: &mut AsynUser,
_value: i32,
) -> crate::error::AsynResult<()> {
Err(boom())
}
}
let port_name = "test_io_alarm_fail";
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(FailDriver(PortDriverBase::new(
port_name,
1,
PortFlags::default(),
))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mk = |iface: i32, tmod: TransferMode| {
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
rec.iface = iface;
rec.tmod = tmod as i32;
rec.imax = 256;
rec
};
let mut octet_rd = mk(0, TransferMode::Read);
octet_rd.ifmt = ASYN_FMT_ASCII;
octet_rd.process().unwrap();
let mut c = CommonFields::default();
octet_rd.check_alarms(&mut c);
assert_eq!(c.nsta, alarm_status::READ_ALARM, "octet read err -> READ");
assert_eq!(c.nsev, AlarmSeverity::Major, "octet read err -> MAJOR");
let mut octet_wr = mk(0, TransferMode::Write);
octet_wr.ofmt = ASYN_FMT_ASCII;
octet_wr.aout = "x".to_string();
octet_wr.process().unwrap();
let mut c = CommonFields::default();
octet_wr.check_alarms(&mut c);
assert_eq!(
c.nsev,
AlarmSeverity::NoAlarm,
"octet write err raises no record alarm (C parity)"
);
let mut int_rd = mk(1, TransferMode::Read);
int_rd.process().unwrap();
let mut c = CommonFields::default();
int_rd.check_alarms(&mut c);
assert_eq!(c.nsta, alarm_status::READ_ALARM, "int32 read err -> READ");
assert_eq!(c.nsev, AlarmSeverity::Major, "int32 read err -> MAJOR");
let mut int_wr = mk(1, TransferMode::Write);
int_wr.process().unwrap();
let mut c = CommonFields::default();
int_wr.check_alarms(&mut c);
assert_eq!(
c.nsta,
alarm_status::WRITE_ALARM,
"int32 write err -> WRITE"
);
assert_eq!(c.nsev, AlarmSeverity::Major, "int32 write err -> MAJOR");
}
fn spawn_fill_port(
port_name: &'static str,
fill: usize,
eom: crate::interpose::EomReason,
) -> Arc<Mutex<Option<usize>>> {
use crate::interpose::EomReason;
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::user::AsynUser;
use tokio::sync::mpsc;
struct FillDriver {
base: PortDriverBase,
fill: usize,
eom: EomReason,
requested: Arc<Mutex<Option<usize>>>,
}
impl PortDriver for FillDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn io_read_octet_eom(
&mut self,
_user: &AsynUser,
buf: &mut [u8],
) -> crate::error::AsynResult<(usize, EomReason)> {
*self.requested.lock().unwrap() = Some(buf.len());
let n = self.fill.min(buf.len());
for b in buf[..n].iter_mut() {
*b = b'Z';
}
Ok((n, self.eom))
}
}
let requested = Arc::new(Mutex::new(None));
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(FillDriver {
base: PortDriverBase::new(port_name, 1, PortFlags::default()),
fill,
eom,
requested: requested.clone(),
}),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
requested
}
#[test]
fn tmot_is_passed_through_verbatim_when_non_negative() {
use std::time::Duration;
let plan_timeout = |tmot: f64| {
let mut rec = AsynRecord::default();
rec.tmot = tmot;
rec.build_io_plan().timeout
};
assert_eq!(
plan_timeout(0.0),
Duration::ZERO,
"TMOT=0 is C's non-blocking poll"
);
assert_eq!(plan_timeout(2.5), Duration::from_millis(2500));
assert_eq!(plan_timeout(0.001), Duration::from_millis(1));
assert_eq!(
plan_timeout(-1.0),
Duration::from_secs(1),
"negative TMOT falls back to the DRV-42 bounded default"
);
assert_eq!(plan_timeout(f64::NAN), Duration::from_secs(1));
assert_eq!(plan_timeout(f64::INFINITY), Duration::from_secs(1));
}
#[test]
fn tmot_zero_reaches_the_driver_as_a_zero_timeout() {
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::user::AsynUser;
use tokio::sync::mpsc;
struct TimeoutSpy {
base: PortDriverBase,
seen: Arc<Mutex<Option<std::time::Duration>>>,
}
impl PortDriver for TimeoutSpy {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn io_read_octet(
&mut self,
user: &AsynUser,
buf: &mut [u8],
) -> crate::error::AsynResult<usize> {
*self.seen.lock().unwrap() = Some(user.timeout);
buf[0] = b'x';
Ok(1)
}
}
let port_name = "test_tmot_zero_to_driver";
let seen = Arc::new(Mutex::new(None));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(TimeoutSpy {
base: PortDriverBase::new(port_name, 1, PortFlags::default()),
seen: seen.clone(),
}),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(256)),
actor_id,
);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let mut rec = read_rec(port_name, 0, 40, 0);
rec.tmot = 0.0;
rec.process().unwrap();
assert_eq!(
*seen.lock().unwrap(),
Some(std::time::Duration::ZERO),
"the driver's asynUser must carry the operator's TMOT=0 poll"
);
}
fn read_rec(port_name: &str, ifmt: i32, imax: i32, nrrd: i32) -> AsynRecord {
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
let _ = rec.connect_device();
rec.iface = 0;
rec.tmod = TransferMode::Read as i32;
rec.ifmt = ifmt;
rec.imax = imax;
rec.nrrd = nrrd;
rec
}
fn read_alarm(rec: &mut AsynRecord) -> (u16, epics_base_rs::server::record::AlarmSeverity) {
use epics_base_rs::server::record::CommonFields;
let mut c = CommonFields::default();
rec.check_alarms(&mut c);
(c.nsta, c.nsev)
}
fn spawn_phase_log_port(port_name: &'static str) -> Arc<Mutex<Vec<&'static str>>> {
use crate::interpose::EomReason;
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::user::AsynUser;
use tokio::sync::mpsc;
struct PhaseLogDriver {
base: PortDriverBase,
log: Arc<Mutex<Vec<&'static str>>>,
}
impl PortDriver for PhaseLogDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn io_flush(&mut self, _user: &mut AsynUser) -> crate::error::AsynResult<()> {
self.log.lock().unwrap().push("flush");
Ok(())
}
fn io_write_octet(
&mut self,
_user: &mut AsynUser,
data: &[u8],
) -> crate::error::AsynResult<usize> {
self.log.lock().unwrap().push("write");
Ok(data.len())
}
fn io_read_octet_eom(
&mut self,
_user: &AsynUser,
buf: &mut [u8],
) -> crate::error::AsynResult<(usize, EomReason)> {
self.log.lock().unwrap().push("read");
let resp = b"OK";
let n = resp.len().min(buf.len());
buf[..n].copy_from_slice(&resp[..n]);
Ok((n, EomReason::EOS))
}
}
let log = Arc::new(Mutex::new(Vec::new()));
let interrupts = Arc::new(InterruptManager::new(256));
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(PhaseLogDriver {
base: PortDriverBase::new(port_name, 1, PortFlags::default()),
log: log.clone(),
}),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(tx, port_name.into(), interrupts, actor_id);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
log
}
#[test]
fn write_read_flushes_the_input_before_the_write() {
let port = "test_tmod_writeread_flush";
let log = spawn_phase_log_port(port);
let mut rec = AsynRecord::default();
rec.port = port.to_string();
let _ = rec.connect_device();
rec.iface = InterfaceType::Octet as i32;
rec.tmod = TransferMode::WriteRead as i32;
rec.ofmt = ASYN_FMT_ASCII;
rec.ifmt = ASYN_FMT_ASCII;
rec.aout = "CMD".to_string();
rec.process().unwrap();
assert_eq!(
*log.lock().unwrap(),
vec!["flush", "write", "read"],
"Write/Read must flush the input before the write"
);
assert_eq!(rec.ainp, "OK");
}
#[test]
fn io_phase_plan_matches_c_perform_io() {
let plan = |tmod: TransferMode, iface: InterfaceType| -> Vec<IoPhase> {
let mut rec = AsynRecord::default();
rec.tmod = tmod as i32;
rec.iface = iface as i32;
io_phases(&rec.build_io_plan())
};
use InterfaceType::{Int32, Octet};
use IoPhase::{Flush, Read, Write};
assert_eq!(plan(TransferMode::WriteRead, Octet), [Flush, Write, Read]);
assert_eq!(plan(TransferMode::Write, Octet), [Write]);
assert_eq!(plan(TransferMode::Read, Octet), [Read]);
assert_eq!(plan(TransferMode::Flush, Octet), [Flush]);
assert_eq!(plan(TransferMode::NoIo, Octet), []);
assert_eq!(plan(TransferMode::WriteRead, Int32), [Write, Read]);
assert_eq!(plan(TransferMode::Flush, Int32), []);
}
#[test]
fn ascii_read_is_sized_by_ainp_size_not_imax() {
use crate::interpose::EomReason;
use epics_base_rs::server::recgbl::alarm_status;
use epics_base_rs::server::record::AlarmSeverity;
let port = "test_ascii_inlen_is_ainp";
let requested = spawn_fill_port(port, 100, EomReason::CNT);
let mut rec = read_rec(port, ASYN_FMT_ASCII, 80, 0);
rec.process().unwrap();
assert_eq!(
*requested.lock().unwrap(),
Some(AINP_SIZE),
"the ASCII read length is sizeof(ainp), not IMAX"
);
assert_eq!(rec.nord, AINP_SIZE as i32, "NORD is the raw transfer count");
assert_eq!(rec.ainp.len(), AINP_SIZE - 1, "AINP NUL-truncated at 40");
assert_eq!(rec.errs, "Overflow nread 40 ");
assert_eq!(
read_alarm(&mut rec),
(alarm_status::READ_ALARM, AlarmSeverity::Minor)
);
}
#[test]
fn overflow_read_reports_errs_in_every_ifmt() {
use crate::interpose::EomReason;
use epics_base_rs::server::recgbl::alarm_status;
use epics_base_rs::server::record::AlarmSeverity;
let port = "test_overflow_errs_ascii";
spawn_fill_port(port, 100, EomReason::CNT);
let mut ascii = read_rec(port, ASYN_FMT_ASCII, 80, 0);
ascii.process().unwrap();
assert_eq!(ascii.errs, "Overflow nread 40 ");
assert_eq!(
read_alarm(&mut ascii),
(alarm_status::READ_ALARM, AlarmSeverity::Minor)
);
let hport = "test_overflow_errs_hybrid";
spawn_fill_port(hport, 100, EomReason::CNT);
let mut hybrid = read_rec(hport, ASYN_FMT_HYBRID, 16, 0);
hybrid.process().unwrap();
assert_eq!(hybrid.errs, "Overflow nread 16 ");
assert_eq!(
read_alarm(&mut hybrid),
(alarm_status::READ_ALARM, AlarmSeverity::Minor)
);
let sport = "test_overflow_errs_short";
spawn_fill_port(sport, 8, EomReason::END);
let mut short = read_rec(sport, ASYN_FMT_ASCII, 80, 0);
short.process().unwrap();
assert_eq!(short.errs, "", "a read inside the buffer reports nothing");
}
#[test]
fn ascii_read_below_ainp_size_is_not_overflow() {
use crate::interpose::EomReason;
use epics_base_rs::server::record::AlarmSeverity;
let port = "test_ascii_inlen_under";
spawn_fill_port(port, AINP_SIZE - 1, EomReason::CNT);
let mut rec = read_rec(port, ASYN_FMT_ASCII, 80, 0);
rec.process().unwrap();
assert_eq!(rec.nord, (AINP_SIZE - 1) as i32);
assert_eq!(rec.ainp.len(), AINP_SIZE - 1, "all 39 bytes land in AINP");
assert_eq!(read_alarm(&mut rec).1, AlarmSeverity::NoAlarm);
}
#[test]
fn ascii_overflow_fires_even_when_the_read_ended_on_eos() {
use crate::interpose::EomReason;
use epics_base_rs::server::recgbl::alarm_status;
use epics_base_rs::server::record::AlarmSeverity;
let port = "test_ascii_overflow_eos";
spawn_fill_port(port, 100, EomReason::EOS);
let mut rec = read_rec(port, ASYN_FMT_ASCII, 80, 0);
rec.process().unwrap();
assert_eq!(rec.nord, AINP_SIZE as i32);
assert_eq!(
read_alarm(&mut rec),
(alarm_status::READ_ALARM, AlarmSeverity::Minor)
);
}
#[test]
fn nrrd_limited_ascii_read_is_never_overflow() {
use crate::interpose::EomReason;
use epics_base_rs::server::record::AlarmSeverity;
let port = "test_ascii_nrrd_short";
let requested = spawn_fill_port(port, 100, EomReason::CNT);
let mut rec = read_rec(port, ASYN_FMT_ASCII, 80, 20);
rec.process().unwrap();
assert_eq!(*requested.lock().unwrap(), Some(20), "NRRD sizes the read");
assert_eq!(rec.nord, 20);
assert_eq!(rec.ainp.len(), 20, "no truncation below the threshold");
assert_eq!(read_alarm(&mut rec).1, AlarmSeverity::NoAlarm);
}
#[test]
fn hybrid_read_is_sized_by_imax_and_nuls_the_buffer_end_on_overflow() {
use crate::interpose::EomReason;
use epics_base_rs::server::recgbl::alarm_status;
use epics_base_rs::server::record::AlarmSeverity;
let port = "test_hybrid_inlen_is_imax";
let requested = spawn_fill_port(port, 100, EomReason::CNT);
let mut rec = read_rec(port, ASYN_FMT_HYBRID, 4, 0);
rec.process().unwrap();
assert_eq!(
*requested.lock().unwrap(),
Some(4),
"IMAX sizes a Hybrid read"
);
assert_eq!(rec.nord, 4);
assert_eq!(rec.binp, vec![b'Z', b'Z', b'Z', 0], "BINP NUL at IMAX-1");
assert_eq!(rec.ainp, "", "Hybrid must not touch AINP");
assert_eq!(
read_alarm(&mut rec),
(alarm_status::READ_ALARM, AlarmSeverity::Minor)
);
}
#[test]
fn test_register_asyn_record_type() {
register_asyn_record_type();
let rec = epics_base_rs::server::db_loader::create_record("asyn").unwrap();
assert_eq!(rec.record_type(), "asyn");
assert!(rec.field_list().len() > 3);
}
#[test]
fn test_translate_escape_standard_sequences() {
assert_eq!(translate_escape("\\r\\n"), vec![0x0D, 0x0A]);
assert_eq!(translate_escape("\\t"), vec![0x09]);
assert_eq!(translate_escape("\\\\"), vec![b'\\']);
assert_eq!(translate_escape("\\0"), vec![0x00]);
assert_eq!(translate_escape("abc"), vec![b'a', b'b', b'c']);
assert_eq!(translate_escape("\\x"), vec![b'\\', b'x']);
assert_eq!(translate_escape("a\\"), vec![b'a', b'\\']);
}
#[test]
fn test_translate_escape_octal() {
assert_eq!(translate_escape("\\033"), vec![0x1B]); assert_eq!(translate_escape("\\7"), vec![0x07]); assert_eq!(translate_escape("\\101"), vec![b'A']); assert_eq!(translate_escape("\\0119"), vec![0x09, b'9']);
assert_eq!(translate_escape("\\0"), vec![0x00]);
assert_eq!(translate_escape("\\015\\012"), vec![0x0D, 0x0A]);
}
#[test]
fn test_octet_output_buffer_by_ofmt() {
let mut rec = AsynRecord::default();
rec.ofmt = ASYN_FMT_ASCII;
rec.aout = "hi\\r\\n".to_string();
rec.nowt = 2; assert_eq!(
rec.octet_output_buffer(),
vec![b'h', b'i', 0x0D, 0x0A],
"ASCII must escape-translate AOUT and ignore NOWT"
);
rec.ofmt = ASYN_FMT_HYBRID;
rec.bout = b"x\\t".to_vec();
assert_eq!(
rec.octet_output_buffer(),
vec![b'x', 0x09],
"Hybrid must escape-translate the BOUT buffer"
);
rec.bout = b"ab\0cd".to_vec();
assert_eq!(rec.octet_output_buffer(), vec![b'a', b'b']);
rec.ofmt = ASYN_FMT_BINARY;
rec.bout = vec![b'\\', b'r', 0x00, 0x01, 0x02];
rec.omax = 80;
rec.nowt = 4;
assert_eq!(
rec.octet_output_buffer(),
vec![b'\\', b'r', 0x00, 0x01],
"Binary writes raw BOUT untranslated, NOWT bytes"
);
rec.omax = 3;
rec.nowt = 10;
rec.clamp_transfer_sizes(AINP_SIZE);
assert_eq!(rec.nowt, 3);
assert_eq!(rec.octet_output_buffer(), vec![b'\\', b'r', 0x00]);
}
#[test]
fn an_array_put_sets_its_element_count() {
let mut rec = AsynRecord::default();
assert_eq!(rec.nowt, 80, "default NOWT (asynRecord.dbd)");
let payload: Vec<u8> = (0..120u32).map(|i| (i % 251) as u8).collect();
rec.omax = 1000;
rec.ofmt = ASYN_FMT_BINARY;
rec.put_field("BOUT", EpicsValue::CharArray(payload.clone()))
.unwrap();
assert_eq!(rec.nowt, 120, "BOUT put must set NOWT = nNew");
assert_eq!(
rec.octet_output_buffer(),
payload,
"every byte the client put must reach the device"
);
rec.put_field("BINP", EpicsValue::CharArray(vec![1, 2, 3]))
.unwrap();
assert_eq!(rec.nord, 3, "BINP put must set NORD = nNew");
rec.put_field("BOUT", EpicsValue::CharArray(vec![9, 9]))
.unwrap();
assert_eq!(rec.nowt, 2);
assert_eq!(rec.octet_output_buffer(), vec![9, 9]);
}
#[test]
fn bout_binp_native_count_is_the_buffer_capacity() {
use epics_base_rs::server::record::Record;
let rec = AsynRecord::default();
assert_eq!(rec.field_native_count("BOUT"), Some(80));
assert_eq!(rec.field_native_count("BINP"), Some(80));
assert_eq!(rec.get_field("BOUT"), Some(EpicsValue::CharArray(vec![])));
assert_eq!(rec.get_field("BINP"), Some(EpicsValue::CharArray(vec![])));
assert_eq!(rec.field_native_count("AOUT"), None);
assert_eq!(rec.field_native_count("PORT"), None);
let mut rec = AsynRecord::default();
rec.omax = 256;
rec.imax = 512;
assert_eq!(rec.field_native_count("BOUT"), Some(256));
assert_eq!(rec.field_native_count("BINP"), Some(512));
}
#[test]
fn nowt_and_nrrd_are_clamped_back_into_the_record() {
let mut rec = AsynRecord::default();
rec.iface = InterfaceType::Octet as i32;
rec.tmod = TransferMode::WriteRead as i32;
rec.ofmt = ASYN_FMT_BINARY;
rec.ifmt = ASYN_FMT_BINARY;
rec.omax = 8;
rec.nowt = 500;
rec.imax = 16;
rec.nrrd = 900;
rec.build_io_plan();
assert_eq!(rec.nowt, 8, "NOWT clamps back to OMAX");
assert_eq!(rec.nrrd, 16, "NRRD clamps back to IMAX for a Binary read");
let mut ascii = AsynRecord::default();
ascii.iface = InterfaceType::Octet as i32;
ascii.tmod = TransferMode::Read as i32;
ascii.ifmt = ASYN_FMT_ASCII;
ascii.imax = 80;
ascii.nrrd = 900;
ascii.build_io_plan();
assert_eq!(
ascii.nrrd, AINP_SIZE as i32,
"ASCII NRRD clamps to AINP_SIZE"
);
let mut read_only = AsynRecord::default();
read_only.iface = InterfaceType::Octet as i32;
read_only.tmod = TransferMode::Read as i32;
read_only.ofmt = ASYN_FMT_BINARY;
read_only.omax = 4;
read_only.nowt = 99;
read_only.build_io_plan();
assert_eq!(read_only.nowt, 4, "TMOD=Read still clamps NOWT");
let mut ascii_out = AsynRecord::default();
ascii_out.iface = InterfaceType::Octet as i32;
ascii_out.tmod = TransferMode::Write as i32;
ascii_out.ofmt = ASYN_FMT_ASCII;
ascii_out.omax = 4;
ascii_out.nowt = 99;
ascii_out.build_io_plan();
assert_eq!(ascii_out.nowt, 99, "ASCII output leaves NOWT alone");
let mut i32_rec = AsynRecord::default();
i32_rec.iface = InterfaceType::Int32 as i32;
i32_rec.tmod = TransferMode::WriteRead as i32;
i32_rec.ofmt = ASYN_FMT_BINARY;
i32_rec.omax = 8;
i32_rec.nowt = 500;
i32_rec.imax = 16;
i32_rec.nrrd = 900;
i32_rec.build_io_plan();
assert_eq!(i32_rec.nowt, 500, "Int32 interface leaves NOWT alone");
assert_eq!(i32_rec.nrrd, 900, "Int32 interface leaves NRRD alone");
}
#[test]
fn test_tfil_special_targets() {
assert!(matches!(open_trace_file("").unwrap(), TraceFile::Stdout));
assert!(matches!(
open_trace_file("<stdout>").unwrap(),
TraceFile::Stdout
));
assert!(matches!(
open_trace_file("<stderr>").unwrap(),
TraceFile::Stderr
));
assert!(matches!(
open_trace_file("<errlog>").unwrap(),
TraceFile::Errlog
));
}
#[test]
fn test_tfil_bare_names_are_file_paths() {
let dir = std::env::temp_dir().join(format!("asynrec_tfil_bare_{}", std::process::id()));
std::fs::create_dir_all(&dir).unwrap();
for name in ["stdout", "stderr"] {
let path = dir.join(name);
let p = path.to_str().unwrap();
assert!(
matches!(open_trace_file(p).unwrap(), TraceFile::File(_)),
"bare {name} must resolve to a file path, not a console sink"
);
}
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn test_tfil_path_appends_not_truncates() {
let path = std::env::temp_dir().join(format!("asynrec_tfil_append_{}", std::process::id()));
let p = path.to_str().unwrap();
let _ = std::fs::remove_file(&path);
open_trace_file(p).unwrap().write_line("first\n");
open_trace_file(p).unwrap().write_line("second\n");
let contents = std::fs::read_to_string(&path).unwrap();
assert_eq!(
contents, "first\nsecond\n",
"re-opening a trace file must append, not truncate"
);
let _ = std::fs::remove_file(&path);
}
fn canblock_int32_entry(value: i32) -> super::registry::PortEntry {
use crate::interrupt::InterruptManager;
use crate::param::ParamType;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::trace::TraceManager;
use tokio::sync::mpsc;
struct ReadDriver {
base: PortDriverBase,
}
impl PortDriver for ReadDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
}
let mut base = PortDriverBase::new("ASYNIO", 1, PortFlags::default());
let val = base.create_param("VAL", ParamType::Int32).unwrap();
base.set_int32_param(val, 0, value).unwrap();
let (tx, rx) = mpsc::channel(16);
let actor = PortActor::new(Box::new(ReadDriver { base }), rx);
let actor_id = actor.id();
std::thread::Builder::new()
.name("asynio-test-actor".into())
.spawn(move || actor.run())
.unwrap();
let mut handle = PortHandle::new(
tx,
"ASYNIO".into(),
Arc::new(InterruptManager::new(16)),
actor_id,
);
handle.set_can_block(true);
super::registry::PortEntry {
handle,
trace: Arc::new(TraceManager::new()),
}
}
#[tokio::test]
async fn nonblocking_canblock_port_defers_then_applies_on_reentry() {
use epics_base_rs::server::database::PvDatabase;
let mut rec = AsynRecord::default();
rec.port_entry = Some(canblock_int32_entry(7));
rec.tmod = TransferMode::Read as i32;
rec.iface = InterfaceType::Int32 as i32;
rec.resolved_reason = 0;
let db = PvDatabase::new();
rec.async_ctx = Some(("ASYNIO_REC".to_string(), db.async_handle()));
let out = rec.process().unwrap();
assert_eq!(
out.result,
RecordProcessResult::AsyncPending,
"a can_block port with async context must defer, not run inline"
);
assert!(
rec.io_inflight.is_some(),
"the deferred request is held in io_inflight until completion"
);
assert_eq!(
rec.i32inp, 0,
"the scan thread returned before the read value landed"
);
let slot = rec.io_inflight.as_ref().unwrap().result.clone();
let mut filled = false;
for _ in 0..2000 {
if slot.lock().unwrap().is_some() {
filled = true;
break;
}
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
}
assert!(filled, "the orchestration must fill the result slot");
let out2 = rec.process().unwrap();
assert_eq!(out2.result, RecordProcessResult::Complete);
assert!(
rec.io_inflight.is_none(),
"completion re-entry clears the in-flight slot"
);
assert_eq!(rec.i32inp, 7, "the read value is applied on re-entry");
}
#[tokio::test]
async fn aqr_after_driver_dequeue_runs_to_completion() {
use crate::interrupt::InterruptManager;
use crate::param::ParamType;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::trace::TraceManager;
use epics_base_rs::server::database::PvDatabase;
use std::sync::Barrier;
use std::sync::atomic::AtomicBool;
use tokio::sync::mpsc;
struct BlockingReadDriver {
base: PortDriverBase,
value: i32,
entered: Arc<AtomicBool>,
release: Arc<Barrier>,
}
impl PortDriver for BlockingReadDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn read_int32(&mut self, _user: &AsynUser) -> AsynResult<i32> {
self.entered.store(true, Ordering::SeqCst);
self.release.wait();
Ok(self.value)
}
}
let entered = Arc::new(AtomicBool::new(false));
let release = Arc::new(Barrier::new(2));
let mut base = PortDriverBase::new("ASYNAQR", 1, PortFlags::default());
base.create_param("VAL", ParamType::Int32).unwrap();
let driver = BlockingReadDriver {
base,
value: 7,
entered: entered.clone(),
release: release.clone(),
};
let (tx, rx) = mpsc::channel(16);
let actor = PortActor::new(Box::new(driver), rx);
let actor_id = actor.id();
std::thread::Builder::new()
.name("asynaqr-test-actor".into())
.spawn(move || actor.run())
.unwrap();
let mut handle = PortHandle::new(
tx,
"ASYNAQR".into(),
Arc::new(InterruptManager::new(16)),
actor_id,
);
handle.set_can_block(true);
let entry = super::registry::PortEntry {
handle,
trace: Arc::new(TraceManager::new()),
};
let mut rec = AsynRecord::default();
rec.port_entry = Some(entry);
rec.tmod = TransferMode::Read as i32;
rec.iface = InterfaceType::Int32 as i32;
rec.resolved_reason = 0;
let db = PvDatabase::new();
rec.async_ctx = Some(("ASYNAQR_REC".to_string(), db.async_handle()));
rec.special("AQR", true).unwrap();
assert!(rec.errs.is_empty(), "AQR with no in-flight request is idle");
let out = rec.process().unwrap();
assert_eq!(out.result, RecordProcessResult::AsyncPending);
for _ in 0..2000 {
if entered.load(Ordering::SeqCst) {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
}
assert!(
entered.load(Ordering::SeqCst),
"the off-thread read must reach the driver before AQR"
);
rec.special("AQR", true).unwrap();
release.wait();
let slot = rec.io_inflight.as_ref().unwrap().result.clone();
let mut filled = false;
for _ in 0..2000 {
if slot.lock().unwrap().is_some() {
filled = true;
break;
}
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
}
assert!(filled, "the running request still produces an outcome");
let out2 = rec.process().unwrap();
assert_eq!(out2.result, RecordProcessResult::Complete);
assert!(rec.io_inflight.is_none(), "completion re-entry leaves idle");
assert!(
rec.errs.is_empty(),
"a cancel that lost the race does not report CANCELED (wasQueued==0)"
);
assert_eq!(
rec.i32inp, 7,
"the device read value applies normally when the cancel loses the race"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn trace_change_posts_readback_fields_immediately() {
use crate::exception::ExceptionManager;
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::trace::TraceManager;
use epics_base_rs::server::database::PvDatabase;
use tokio::sync::mpsc;
struct TraceDriver(PortDriverBase);
impl PortDriver for TraceDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
}
let port_name = "test_trace_immediate_post";
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(TraceDriver(PortDriverBase::new(
port_name,
1,
PortFlags::default(),
))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(256)),
actor_id,
);
let trace = Arc::new(TraceManager::new());
trace.set_exception_sink(Arc::new(ExceptionManager::new()));
trace.set_trace_mask(Some(port_name), TraceMask::empty());
super::registry::register_port(port_name, handle, trace.clone()).unwrap();
let db = PvDatabase::new();
let rec_name = "TRACE_IMM_REC";
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.set_async_context(rec_name.to_string(), db.async_handle());
let _ = rec.connect_device();
assert_eq!(rec.cnct, 1, "record must connect to the registered port");
assert_eq!(rec.tmsk, 0, "baseline trace mask is empty");
db.add_record(rec_name, Box::new(rec)).await.unwrap();
let new_mask = TraceMask::ERROR | TraceMask::FLOW;
{
let tm = trace.clone();
let pn = port_name.to_string();
std::thread::spawn(move || tm.set_trace_mask(Some(&pn), new_mask))
.join()
.unwrap();
}
let want = new_mask.bits() as i32;
let mut posted = false;
for _ in 0..2000 {
let inst = db.get_record(rec_name).unwrap();
let got = inst.read().record.get_field("TMSK");
if got == Some(EpicsValue::Long(want)) {
posted = true;
break;
}
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
}
assert!(
posted,
"trace change must post TMSK immediately, no process()"
);
let inst = db.get_record(rec_name).unwrap();
let g = inst.read();
assert_eq!(g.record.get_field("TMSK"), Some(EpicsValue::Long(want)));
assert_eq!(
g.record.get_field("TB0"),
Some(EpicsValue::Enum(1)),
"ERROR bit posted"
);
assert_eq!(
g.record.get_field("TB4"),
Some(EpicsValue::Enum(1)),
"FLOW bit posted"
);
assert_eq!(
g.record.get_field("TB1"),
Some(EpicsValue::Enum(0)),
"IO_DEVICE bit stays clear"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn unchanged_trace_field_is_not_reposted() {
use crate::exception::ExceptionManager;
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::trace::{TraceIoMask, TraceManager};
use epics_base_rs::server::database::PvDatabase;
use epics_base_rs::server::database::db_access::DbSubscription;
use std::time::Duration;
use tokio::sync::mpsc;
struct TraceDriver(PortDriverBase);
impl PortDriver for TraceDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
}
let port_name = "test_trace_post_if_new";
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(
Box::new(TraceDriver(PortDriverBase::new(
port_name,
1,
PortFlags::default(),
))),
rx,
);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(256)),
actor_id,
);
let trace = Arc::new(TraceManager::new());
trace.set_exception_sink(Arc::new(ExceptionManager::new()));
trace.set_trace_mask(Some(port_name), TraceMask::ERROR);
trace.set_trace_io_mask(Some(port_name), TraceIoMask::empty());
super::registry::register_port(port_name, handle, trace.clone()).unwrap();
let db = PvDatabase::new();
let rec_name = "TRACE_POSTIFNEW_REC";
let rec = AsynRecord::default();
db.add_record(rec_name, Box::new(rec)).await.unwrap();
{
let inst = db.get_record(rec_name).unwrap();
let mut g = inst.write();
let rec = g
.record
.as_any_mut()
.unwrap()
.downcast_mut::<AsynRecord>()
.unwrap();
rec.tmsk = TraceMask::ERROR.bits() as i32;
rec.update_trace_bits_from_mask();
rec.port = port_name.to_string();
rec.special("PORT", true).unwrap();
assert_eq!(rec.cnct, 1, "record must connect to the registered port");
}
let mut tmsk_sub = DbSubscription::subscribe(&db, &format!("{rec_name}.TMSK"))
.await
.expect("subscribe TMSK");
let mut tiom_sub = DbSubscription::subscribe(&db, &format!("{rec_name}.TIOM"))
.await
.expect("subscribe TIOM");
let new_tmsk = TraceMask::ERROR | TraceMask::FLOW;
{
let (tm, pn) = (trace.clone(), port_name.to_string());
std::thread::spawn(move || tm.set_trace_mask(Some(&pn), new_tmsk))
.join()
.unwrap();
}
assert_eq!(
tmsk_sub.recv().await,
Some(EpicsValue::Long(new_tmsk.bits() as i32)),
"a changed TMSK is posted to its monitor"
);
{
let (tm, pn) = (trace.clone(), port_name.to_string());
std::thread::spawn(move || tm.set_trace_io_mask(Some(&pn), TraceIoMask::ASCII))
.join()
.unwrap();
}
assert_eq!(
tiom_sub.recv().await,
Some(EpicsValue::Long(TraceIoMask::ASCII.bits() as i32)),
"a changed TIOM is posted to its monitor"
);
let reposted = tokio::time::timeout(Duration::from_millis(500), tmsk_sub.recv()).await;
assert!(
reposted.is_err(),
"unchanged TMSK must not be re-posted on an IO-mask-only change, got {reposted:?}"
);
}
fn errs_on_the_wire(text: &str) -> EpicsValue {
EpicsValue::CharArray(text.bytes().collect())
}
#[tokio::test]
async fn every_errs_write_posts_and_an_unchanged_one_does_not() {
use crate::exception::ExceptionManager;
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::trace::TraceManager;
use epics_base_rs::server::database::PvDatabase;
use epics_base_rs::server::database::db_access::DbSubscription;
use std::time::Duration;
use tokio::sync::mpsc;
struct DownDriver(PortDriverBase);
impl PortDriver for DownDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
}
let port_name = "test_errs_post";
let mut base = PortDriverBase::new(port_name, 1, PortFlags::default());
base.init_connected(false);
base.auto_connect = false;
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(Box::new(DownDriver(base)), rx);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(256)),
actor_id,
);
let trace = Arc::new(TraceManager::new());
trace.set_exception_sink(Arc::new(ExceptionManager::new()));
super::registry::register_port(port_name, handle, trace).unwrap();
let db = PvDatabase::new();
let rec_name = "ERRS_POST_REC";
db.add_record(rec_name, Box::new(AsynRecord::default()))
.await
.unwrap();
{
let inst = db.get_record(rec_name).unwrap();
let mut g = inst.write();
let rec = g
.record
.as_any_mut()
.unwrap()
.downcast_mut::<AsynRecord>()
.unwrap();
rec.port = port_name.to_string();
rec.special("PORT", true).unwrap();
rec.report_error(String::new()); }
let mut errs_sub = DbSubscription::subscribe(&db, &format!("{rec_name}.ERRS"))
.await
.expect("subscribe ERRS");
{
let inst = db.get_record(rec_name).unwrap();
let mut g = inst.write();
let rec = g
.record
.as_any_mut()
.unwrap()
.downcast_mut::<AsynRecord>()
.unwrap();
rec.baud = 9600;
rec.special("BAUD", true).unwrap();
assert_eq!(
rec.errs,
format!("port {port_name} not connected"),
"precondition: the queue gate refused the option put"
);
}
assert_eq!(
errs_sub.recv().await,
Some(errs_on_the_wire(&format!("port {port_name} not connected"))),
"a refused put must fire the ERRS monitor with the refusal text"
);
{
let inst = db.get_record(rec_name).unwrap();
let mut g = inst.write();
let rec = g
.record
.as_any_mut()
.unwrap()
.downcast_mut::<AsynRecord>()
.unwrap();
rec.report_error(format!("port {port_name} not connected"));
}
let reposted = tokio::time::timeout(Duration::from_millis(300), errs_sub.recv()).await;
assert!(
reposted.is_err(),
"an unchanged ERRS must not re-post, got {reposted:?}"
);
{
let inst = db.get_record(rec_name).unwrap();
let mut g = inst.write();
let rec = g
.record
.as_any_mut()
.unwrap()
.downcast_mut::<AsynRecord>()
.unwrap();
rec.reset_error();
}
assert_eq!(
errs_sub.recv().await,
Some(errs_on_the_wire("")),
"resetError's clear must post"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn connect_device_posts_every_readback_it_refreshes() {
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::trace::TraceManager;
use epics_base_rs::server::database::PvDatabase;
use epics_base_rs::server::database::db_access::DbSubscription;
use tokio::sync::mpsc;
struct Serialish(PortDriverBase);
impl PortDriver for Serialish {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn capabilities(&self) -> Vec<crate::interfaces::Capability> {
crate::interfaces::octet_transport_capabilities()
}
fn get_option(&self, key: &str) -> crate::error::AsynResult<String> {
match key {
"baud" => Ok("19200".to_string()),
_ => Err(crate::error::AsynError::Status {
status: AsynStatus::Error,
message: format!("unsupported option {key}"),
}),
}
}
}
let port_name = "r14_47_connect_posts";
let mut drv = Serialish(PortDriverBase::new(port_name, 1, PortFlags::default()));
drv.set_input_eos(&AsynUser::default(), b"\r\n").unwrap();
let (tx, rx) = mpsc::channel(64);
let actor = PortActor::new(Box::new(drv), rx);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(64)),
actor_id,
);
register_port(port_name, handle, Arc::new(TraceManager::new())).unwrap();
let db = PvDatabase::new();
let rec_name = "R14_47_REC";
let rec = AsynRecord::default();
db.add_record(rec_name, Box::new(rec)).await.unwrap();
let mut baud_sub = DbSubscription::subscribe(&db, &format!("{rec_name}.BAUD"))
.await
.expect("subscribe BAUD");
let mut ieos_sub = DbSubscription::subscribe(&db, &format!("{rec_name}.IEOS"))
.await
.expect("subscribe IEOS");
let mut pcnct_sub = DbSubscription::subscribe(&db, &format!("{rec_name}.PCNCT"))
.await
.expect("subscribe PCNCT");
{
let inst = db.get_record(rec_name).unwrap();
let mut g = inst.write();
let rec = g
.record
.as_any_mut()
.unwrap()
.downcast_mut::<AsynRecord>()
.unwrap();
rec.port = port_name.to_string();
rec.special("PORT", true).unwrap();
assert_eq!(rec.baud, baud_choice_index("19200"), "the field refreshed");
assert_eq!(rec.ieos, "\\r\\n", "…and the EOS with it");
assert_eq!(rec.pcnct, 1, "…and the record attached");
}
assert_eq!(
baud_sub.recv().await,
Some(EpicsValue::Enum(baud_choice_index("19200") as u16)),
"getOptions posts BAUD (asynRecord.c:1926)"
);
assert_eq!(
ieos_sub.recv().await,
Some(EpicsValue::String("\\r\\n".into())),
"getEos posts IEOS (asynRecord.c:2016-2020)"
);
assert_eq!(
pcnct_sub.recv().await,
Some(EpicsValue::Enum(1)),
"monitorStatus posts PCNCT (asynRecord.c:1128)"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn any_port_exception_refreshes_the_connect_state_immediately() {
use crate::exception::ExceptionManager;
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use crate::trace::TraceManager;
use epics_base_rs::server::database::PvDatabase;
use tokio::sync::mpsc;
struct PlainDriver(PortDriverBase);
impl PortDriver for PlainDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
}
let port_name = "test_except_connect_state";
let exceptions = Arc::new(ExceptionManager::new());
let trace = Arc::new(TraceManager::new());
trace.set_exception_sink(exceptions.clone());
let mut base = PortDriverBase::new(port_name, 1, PortFlags::default());
base.bind_exception_sink(exceptions.clone());
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(Box::new(PlainDriver(base)), rx);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(256)),
actor_id,
);
super::registry::register_port(port_name, handle.clone(), trace.clone()).unwrap();
let db = PvDatabase::new();
let rec_name = "EXCEPT_CNCT_REC";
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.set_async_context(rec_name.to_string(), db.async_handle());
let _ = rec.connect_device();
assert_eq!(rec.cnct, 1, "the port starts connected");
assert_eq!(rec.enbl, 1, "the port starts enabled");
db.add_record(rec_name, Box::new(rec)).await.unwrap();
{
let h = handle.clone();
std::thread::spawn(move || h.disconnect_blocking().unwrap())
.join()
.unwrap();
}
let mut posted = false;
for _ in 0..2000 {
let inst = db.get_record(rec_name).unwrap();
let got = inst.read().record.get_field("CNCT");
if got == Some(EpicsValue::Enum(0)) {
posted = true;
break;
}
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
}
assert!(
posted,
"a disconnect must refresh CNCT immediately, with no process()"
);
let inst = db.get_record(rec_name).unwrap();
let g = inst.read();
assert_eq!(g.record.get_field("ENBL"), Some(EpicsValue::Enum(1)));
assert_eq!(g.record.get_field("AUCT"), Some(EpicsValue::Enum(1)));
}
#[test]
fn a_port_exception_refreshes_the_connect_state_on_the_next_process() {
use crate::exception::ExceptionManager;
use crate::interrupt::InterruptManager;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::port_actor::PortActor;
use tokio::sync::mpsc;
struct PlainDriver(PortDriverBase);
impl PortDriver for PlainDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
}
let port_name = "test_except_connect_state_dirty";
let exceptions = Arc::new(ExceptionManager::new());
let trace = Arc::new(TraceManager::new());
trace.set_exception_sink(exceptions.clone());
let mut base = PortDriverBase::new(port_name, 1, PortFlags::default());
base.bind_exception_sink(exceptions.clone());
let (tx, rx) = mpsc::channel(256);
let actor = PortActor::new(Box::new(PlainDriver(base)), rx);
let actor_id = actor.id();
std::thread::spawn(move || actor.run());
let handle = PortHandle::new(
tx,
port_name.into(),
Arc::new(InterruptManager::new(256)),
actor_id,
);
register_port(port_name, handle.clone(), trace).unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.tmod = TransferMode::NoIo as i32;
let _ = rec.connect_device();
assert_eq!(rec.cnct, 1);
handle.disconnect_blocking().unwrap();
rec.process().unwrap();
assert_eq!(
rec.cnct, 0,
"the dropped link reaches CNCT on the next scan"
);
assert_eq!(rec.enbl, 1, "ENBL is re-read and unchanged");
assert_eq!(rec.auct, 1, "AUCT is re-read and unchanged");
}
#[test]
fn a_gpib_port_reads_back_gpibiv_and_i32iv() {
use crate::drivers::prologix::DrvAsynPrologixPort;
use crate::drivers::vxi11::DrvVxi11Port;
use crate::runtime::{RuntimeConfig, create_port_runtime};
let vxi_name = "r10_55_vxi11";
let (vxi_rt, _vxi_jh) = create_port_runtime(
DrvVxi11Port::configure(vxi_name, "192.0.2.1", 0, "", "gpib0", 0, true).unwrap(),
RuntimeConfig::default(),
);
register_port(
vxi_name,
vxi_rt.port_handle().clone(),
Arc::new(TraceManager::new()),
)
.unwrap();
let mut rec = AsynRecord::default();
rec.port = vxi_name.to_string();
rec.connect_device().unwrap();
assert_eq!(rec.gpibiv, 1, "asynGpib (C :1228-1234)");
assert_eq!(rec.i32iv, 1, "asynInt32, registered by asynGpib (C :140)");
assert_eq!(rec.octetiv, 1, "asynOctet");
assert_eq!(rec.optioniv, 1, "vxi11 registers asynOption (:1777)");
let prologix_name = "r10_55_prologix";
let (p_rt, _p_jh) = create_port_runtime(
DrvAsynPrologixPort::new(prologix_name, "192.0.2.1:1234", true).unwrap(),
RuntimeConfig::default(),
);
register_port(
prologix_name,
p_rt.port_handle().clone(),
Arc::new(TraceManager::new()),
)
.unwrap();
let mut rec = AsynRecord::default();
rec.port = prologix_name.to_string();
rec.connect_device().unwrap();
assert_eq!(rec.gpibiv, 1, "asynGpib (C :1228-1234)");
assert_eq!(rec.i32iv, 1, "asynInt32, registered by asynGpib (C :140)");
assert_eq!(rec.octetiv, 1, "asynOctet");
assert_eq!(
rec.optioniv, 0,
"drvPrologixGPIB registers no asynOption (:592)"
);
}
#[test]
fn connect_device_reads_the_ports_interface_registry() {
use crate::interfaces::octet_transport_capabilities;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::runtime::{RuntimeConfig, create_port_runtime};
struct OctetTransport(PortDriverBase);
impl PortDriver for OctetTransport {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn capabilities(&self) -> Vec<crate::interfaces::Capability> {
octet_transport_capabilities()
}
}
let port_name = "r9_53_octet_transport";
let (rt, _jh) = create_port_runtime(
OctetTransport(PortDriverBase::new(port_name, 1, PortFlags::default())),
RuntimeConfig::default(),
);
register_port(
port_name,
rt.port_handle().clone(),
Arc::new(TraceManager::new()),
)
.unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.connect_device().unwrap();
assert_eq!(rec.octetiv, 1, "the transport carries asynOctet");
assert_eq!(rec.optioniv, 1, "and asynOption");
assert_eq!(rec.i32iv, 0, "but no asynInt32 (C :1204-1210)");
assert_eq!(rec.ui32iv, 0, "no asynUInt32Digital (C :1212-1218)");
assert_eq!(rec.f64iv, 0, "no asynFloat64 (C :1220-1226)");
assert_eq!(rec.gpibiv, 0, "no asynGpib (C :1228-1234)");
rec.iface = InterfaceType::Int32 as i32;
rec.tmod = TransferMode::Read as i32;
rec.process().unwrap();
assert_eq!(rec.errs, "No asynInt32 interface");
assert_eq!(rec.i32inp, 0, "no read ran");
let mut c = CommonFields::default();
rec.check_alarms(&mut c);
assert_eq!(c.nsta, alarm_status::COMM_ALARM);
assert_eq!(c.nsev, AlarmSeverity::Major);
rec.iface = InterfaceType::Octet as i32;
rec.errs.clear();
rec.imax = 16;
rec.ifmt = ASYN_FMT_ASCII;
rec.process().unwrap();
assert_ne!(
rec.errs, "No asynOctet interface",
"asynOctet is implemented and must not be refused"
);
}
#[test]
fn an_option_put_reaches_the_driver_with_the_records_tmot() {
use crate::interfaces::{Capability, octet_transport_capabilities};
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::runtime::{RuntimeConfig, create_port_runtime};
use std::sync::Mutex as StdMutex;
use std::time::Duration;
static SEEN: StdMutex<Option<Duration>> = StdMutex::new(None);
struct OptionSpy(PortDriverBase);
impl PortDriver for OptionSpy {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn capabilities(&self) -> Vec<Capability> {
octet_transport_capabilities()
}
fn set_option(
&mut self,
user: &mut AsynUser,
_key: &str,
_value: &str,
) -> crate::error::AsynResult<()> {
*SEEN.lock().unwrap() = Some(user.timeout);
Ok(())
}
}
let port_name = "r9_55_option_user";
let (rt, _jh) = create_port_runtime(
OptionSpy(PortDriverBase::new(port_name, 1, PortFlags::default())),
RuntimeConfig::default(),
);
register_port(
port_name,
rt.port_handle().clone(),
Arc::new(TraceManager::new()),
)
.unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.connect_device().unwrap();
rec.tmot = 3.5;
rec.lbaud = 19200;
rec.special("LBAUD", true).unwrap();
assert_eq!(
SEEN.lock().unwrap().take(),
Some(Duration::from_millis(3500)),
"the driver's setOption runs under the record's TMOT"
);
}
#[test]
fn option_and_eos_are_refused_on_a_port_that_lacks_the_interface() {
use crate::interfaces::Capability;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::runtime::{RuntimeConfig, create_port_runtime};
struct OctetNoOption(PortDriverBase);
impl PortDriver for OctetNoOption {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn capabilities(&self) -> Vec<Capability> {
vec![Capability::OctetRead, Capability::OctetWrite]
}
}
let port_name = "r9_53_no_option";
let (rt, _jh) = create_port_runtime(
OctetNoOption(PortDriverBase::new(port_name, 1, PortFlags::default())),
RuntimeConfig::default(),
);
register_port(
port_name,
rt.port_handle().clone(),
Arc::new(TraceManager::new()),
)
.unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.connect_device().unwrap();
assert_eq!(rec.optioniv, 0, "no asynOption interface");
assert_eq!(rec.octetiv, 1);
rec.lbaud = 9600;
rec.special("LBAUD", true).unwrap();
assert_eq!(rec.errs, "No asynOption interface");
let mut c = CommonFields::default();
rec.check_alarms(&mut c);
assert_eq!(c.nsta, alarm_status::COMM_ALARM);
assert_eq!(c.nsev, AlarmSeverity::Major);
rec.errs.clear();
rec.ieos = "\\n".to_string();
rec.special("IEOS", true).unwrap();
assert_ne!(rec.errs, "No asynOctet interface");
}
#[test]
fn every_record_request_carries_c_queue_timeout() {
use std::time::Duration;
assert_eq!(
QUEUE_TIMEOUT,
Duration::from_secs(10),
"C `#define QUEUE_TIMEOUT 10.0`"
);
let mut rec = AsynRecord::default();
let plan = rec.build_io_plan();
assert_eq!(io_user(&plan).queue_timeout, Some(QUEUE_TIMEOUT));
assert_eq!(flush_user(&plan).queue_timeout, Some(QUEUE_TIMEOUT));
assert_eq!(rec.option_user().queue_timeout, Some(QUEUE_TIMEOUT));
assert_eq!(AsynUser::default().queue_timeout, None);
assert_eq!(AsynUser::new(3).with_addr(1).queue_timeout, None);
}
#[test]
fn a_process_queue_timeout_reports_the_c_text_and_state_major() {
let mut rec = AsynRecord::default();
rec.tmod = TransferMode::WriteRead as i32;
let plan = rec.build_io_plan();
let mut out = IoOutcome::default();
let flow = record_phase_result(
&plan,
&mut out,
IoPhase::Write,
Err(crate::error::AsynError::QueueTimeout { port: "p".into() }),
);
assert!(
matches!(flow, PhaseFlow::Aborted),
"the request never ran: C never entered performIO, so no read follows"
);
assert_eq!(out.errs.as_deref(), Some("process queueRequest timeout"));
assert_eq!(
out.alarm,
Some((alarm_status::STATE_ALARM, AlarmSeverity::Major))
);
assert_eq!(out.nawt, None, "nothing was written — nothing to publish");
for phase in [IoPhase::Read, IoPhase::Flush] {
let mut out = IoOutcome::default();
let flow = record_phase_result(
&plan,
&mut out,
phase,
Err(crate::error::AsynError::QueueTimeout { port: "p".into() }),
);
assert!(matches!(flow, PhaseFlow::Aborted), "{phase:?}");
assert_eq!(out.errs.as_deref(), Some("process queueRequest timeout"));
assert_eq!(
out.alarm,
Some((alarm_status::STATE_ALARM, AlarmSeverity::Major)),
"{phase:?}"
);
}
}
#[test]
fn an_io_timeout_is_not_a_queue_timeout() {
let mut rec = AsynRecord::default();
rec.tmod = TransferMode::WriteRead as i32;
let plan = rec.build_io_plan();
let mut out = IoOutcome::default();
let flow = record_phase_result(
&plan,
&mut out,
IoPhase::Read,
Err(crate::error::AsynError::Status {
status: AsynStatus::Timeout,
message: "no response".into(),
}),
);
assert!(
matches!(flow, PhaseFlow::Continue),
"a phase that ran and failed does not abort the cycle"
);
let errs = out.errs.clone().unwrap();
assert!(
errs.contains("timeout") && !errs.contains("queueRequest"),
"the driver's read-error text, not the queue timeout's: {errs}"
);
}
#[test]
fn a_special_queue_timeout_reports_the_c_text_and_no_severity() {
let mut rec = AsynRecord::default();
rec.report_special_queue_timeout();
assert_eq!(rec.errs, "special queueRequest timeout");
let mut c = CommonFields::default();
rec.check_alarms(&mut c);
assert_eq!(
c.nsev,
AlarmSeverity::NoAlarm,
"C raises no recGblSetSevr in queueTimeoutCallbackSpecial"
);
}
fn register_baud_text_port(port_name: &'static str, baud_text: &'static str) {
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::runtime::{RuntimeConfig, create_port_runtime};
struct BaudTextPort(PortDriverBase, &'static str);
impl PortDriver for BaudTextPort {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn get_option(&self, key: &str) -> crate::error::AsynResult<String> {
if key == "baud" {
Ok(self.1.to_string())
} else {
Err(crate::error::AsynError::OptionNotFound(key.to_string()))
}
}
}
let (rt, jh) = create_port_runtime(
BaudTextPort(
PortDriverBase::new(port_name, 1, PortFlags::default()),
baud_text,
),
RuntimeConfig::default(),
);
std::mem::forget(jh);
register_port(
port_name,
rt.port_handle().clone(),
Arc::new(TraceManager::new()),
)
.unwrap();
std::mem::forget(rt);
}
#[test]
fn a_baud_readback_with_no_number_leaves_lbaud_alone() {
let port_name = "r10_52_no_number";
register_baud_text_port(port_name, "unknown");
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.lbaud = 115200;
rec.dbit = 4;
rec.connect_device().unwrap();
assert_eq!(
rec.lbaud, 115200,
"C's sscanf leaves LBAUD untouched when the text carries no number"
);
assert_eq!(rec.baud, 0, "BAUD still lands on its Unknown choice");
assert_eq!(
rec.dbit, 0,
"a key the port refuses is an empty readback, and C zeroes the enum before the choice walk"
);
}
#[test]
fn a_baud_readback_with_a_number_reaches_lbaud() {
let port_name = "r10_52_number";
register_baud_text_port(port_name, "9600 baud");
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.lbaud = 115200;
rec.connect_device().unwrap();
assert_eq!(rec.lbaud, 9600);
assert_eq!(
rec.baud, 0,
"BAUD matches the choice *text*, which \"9600 baud\" is not"
);
}
fn register_octet_transport(port_name: &'static str) {
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::runtime::{RuntimeConfig, create_port_runtime};
struct OctetTransport(PortDriverBase);
impl PortDriver for OctetTransport {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn capabilities(&self) -> Vec<crate::interfaces::Capability> {
crate::interfaces::octet_transport_capabilities()
}
}
let (rt, jh) = create_port_runtime(
OctetTransport(PortDriverBase::new(port_name, 1, PortFlags::default())),
RuntimeConfig::default(),
);
std::mem::forget(jh);
register_port(
port_name,
rt.port_handle().clone(),
Arc::new(TraceManager::new()),
)
.unwrap();
std::mem::forget(rt);
}
fn register_param_port(port_name: &'static str, param: &str) {
use crate::param::ParamType;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::runtime::{RuntimeConfig, create_port_runtime};
struct ParamDriver(PortDriverBase);
impl PortDriver for ParamDriver {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
}
let mut base = PortDriverBase::new(port_name, 1, PortFlags::default());
base.params
.create_param("FILLER", ParamType::Int32)
.unwrap();
base.params.create_param(param, ParamType::Int32).unwrap();
let (rt, jh) = create_port_runtime(ParamDriver(base), RuntimeConfig::default());
std::mem::forget(jh);
register_port(
port_name,
rt.port_handle().clone(),
Arc::new(TraceManager::new()),
)
.unwrap();
std::mem::forget(rt);
}
#[test]
fn a_port_without_asyn_drv_user_forces_reason_zero_and_reports_drvinfo() {
let port_name = "r11_48_no_drvuser";
register_octet_transport(port_name);
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.reason = 7;
rec.drvinfo = "SOME_PARAM".to_string();
rec.connect_device().unwrap();
assert_eq!(
rec.reason, 0,
"C zeroes REASON on a port with no asynDrvUser"
);
assert_eq!(
rec.resolved_reason, 0,
"and no I/O may carry the stale reason"
);
assert_eq!(
rec.errs, "asynDrvUser not supported but drvInfo not blank",
"C reportError text (asynRecord.c:1264-1265)"
);
}
#[test]
fn connect_device_reads_the_eos_back_from_the_driver() {
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::runtime::{RuntimeConfig, create_port_runtime};
struct EosPort(PortDriverBase);
impl PortDriver for EosPort {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn capabilities(&self) -> Vec<crate::interfaces::Capability> {
crate::interfaces::octet_transport_capabilities()
}
}
for (port_name, connected) in [("w10_d3_eos_up", true), ("w10_d3_eos_down", false)] {
let mut base = PortDriverBase::new(port_name, 1, PortFlags::default());
base.auto_connect = false;
base.init_connected(connected);
let mut drv = EosPort(base);
drv.set_input_eos(&AsynUser::default(), b"\r\n").unwrap();
drv.set_output_eos(&AsynUser::default(), b"\n").unwrap();
let (rt, jh) = create_port_runtime(drv, RuntimeConfig::default());
std::mem::forget(jh);
register_port(
port_name,
rt.port_handle().clone(),
Arc::new(TraceManager::new()),
)
.unwrap();
std::mem::forget(rt);
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.connect_device().unwrap();
if connected {
assert_eq!(
rec.ieos, "\\r\\n",
"connectDevice queues callbackGetEos (asynRecord.c:1289-1300)"
);
assert_eq!(rec.oeos, "\\n");
} else {
assert_eq!(
rec.ieos, "",
"a disconnected port's EOS readback is refused (asynRecord.c:1296)"
);
assert_eq!(rec.oeos, "");
}
}
}
#[test]
fn every_special_callback_arm_ends_in_monitor_status() {
for (field, port_name) in [
("BAUD", "w10_d2_option"),
("IEOS", "w10_d2_eos"),
("CNCT", "w10_d2_cnct"),
] {
register_octet_transport(port_name);
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.connect_device().unwrap();
assert_eq!(rec.enbl, 1, "{field}: the port starts enabled");
assert_eq!(rec.cnct, 1, "{field}: …and connected");
assert_eq!(rec.auct, 1, "{field}: …and auto-connecting");
let handle = rec.port_entry.as_ref().unwrap().handle.clone();
handle.set_auto_connect_blocking(false).unwrap();
rec.special(field, true).unwrap();
assert_eq!(
rec.auct, 0,
"{field}: monitorStatus re-imports AUCT from isAutoConnect \
(asynRecord.c:1084-1088) at the tail of every arm (:897)"
);
}
}
#[test]
fn a_reason_put_blanks_drvinfo_so_a_reconnect_cannot_undo_it() {
let port_name = "r11_c8_reason_put";
register_param_port(port_name, "GAIN");
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.drvinfo = "GAIN".to_string();
rec.connect_device().unwrap();
assert_eq!(rec.reason, 1, "DRVINFO resolved GAIN to parameter 1");
assert_eq!(rec.resolved_reason, 1);
rec.reason = 3;
rec.special("REASON", true).unwrap();
assert_eq!(
rec.resolved_reason, 3,
"C assigns pasynUser->reason from the field (asynRecord.c:488)"
);
assert_eq!(
rec.drvinfo, "",
"C blanks DRVINFO in the same arm (asynRecord.c:489)"
);
rec.connect_device().unwrap();
assert_eq!(
rec.reason, 3,
"a stale DRVINFO would have re-resolved REASON back to 1"
);
assert_eq!(rec.resolved_reason, 3, "and the I/O user with it");
}
#[test]
fn a_port_without_asyn_drv_user_zeroes_reason_even_with_blank_drvinfo() {
let port_name = "r11_48_no_drvuser_blank";
register_octet_transport(port_name);
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.reason = 7;
rec.connect_device().unwrap();
assert_eq!(rec.reason, 0);
assert_eq!(rec.resolved_reason, 0);
assert_eq!(rec.errs, "", "a blank DRVINFO is not an error");
}
#[test]
fn a_port_with_asyn_drv_user_still_resolves_and_still_reports_create_failure() {
let port_name = "r11_48_drvuser";
register_param_port(port_name, "GAIN");
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.drvinfo = "GAIN".to_string();
rec.connect_device().unwrap();
assert_eq!(rec.reason, 1, "DRVINFO resolved through the driver");
assert_eq!(rec.resolved_reason, 1);
assert_eq!(rec.errs, "");
rec.drvinfo = "NO_SUCH_PARAM".to_string();
rec.connect_device().unwrap();
assert_eq!(rec.errs, "Error in asynDrvUser->create()");
assert_eq!(rec.resolved_reason, 0);
rec.drvinfo.clear();
rec.reason = 5;
rec.connect_device().unwrap();
assert_eq!(rec.reason, 5);
assert_eq!(rec.resolved_reason, 5);
assert_eq!(rec.errs, "");
}
struct IoIntrDriver {
base: PortDriverBase,
reads: Arc<std::sync::atomic::AtomicUsize>,
}
impl PortDriver for IoIntrDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn read_int32(&mut self, _user: &AsynUser) -> AsynResult<i32> {
self.reads.fetch_add(1, Ordering::SeqCst);
Ok(42)
}
}
fn io_intr_port(
name: &str,
) -> (
crate::runtime::PortRuntimeHandle,
Arc<std::sync::atomic::AtomicUsize>,
) {
use crate::port::PortFlags;
use crate::runtime::{RuntimeConfig, create_port_runtime};
let reads = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let (rt, _jh) = create_port_runtime(
IoIntrDriver {
base: PortDriverBase::new(name, 1, PortFlags::default()),
reads: reads.clone(),
},
RuntimeConfig::default(),
);
register_port(
name,
rt.port_handle().clone(),
Arc::new(TraceManager::new()),
)
.unwrap();
(rt, reads)
}
fn fire_interrupt(
rt: &crate::runtime::PortRuntimeHandle,
reason: usize,
value: crate::param::ParamValue,
iface: crate::interfaces::InterfaceType,
changed_mask: u32,
) {
rt.port_handle()
.interrupts()
.notify(crate::interrupt::InterruptValue {
reason,
addr: 0,
value,
uint32_changed_mask: changed_mask,
iface: Some(iface),
..Default::default()
});
}
#[test]
fn an_int32_interrupt_lands_in_i32inp_and_that_cycle_does_no_io() {
use crate::param::ParamValue;
use epics_base_rs::server::device_support::DeviceSupport;
let (rt, reads) = io_intr_port("r11_46_int32");
let mut rec = AsynRecord::default();
rec.port = "r11_46_int32".to_string();
rec.iface = InterfaceType::Int32 as i32;
rec.tmod = TransferMode::Read as i32;
rec.connect_device().unwrap();
let mut dev = AsynRecordDevice::new();
dev.init(&mut rec).unwrap();
let mut wakeups = dev.io_intr_receiver().expect("the DSET owns a scan list");
rec.set_io_intr_scan(true);
assert_eq!(rec.errs, "", "the port implements asynInt32");
fire_interrupt(
&rt,
0,
ParamValue::Int32(7),
crate::interfaces::InterfaceType::Int32,
0,
);
wakeups
.try_recv()
.expect("C `scanIoRequest` — the callback asks for a process");
rec.process().unwrap();
assert_eq!(rec.i32inp, 7, "the pushed value, not a read");
assert_eq!(
reads.load(Ordering::SeqCst),
0,
"C `goto done`: the interrupt-driven cycle queues no I/O"
);
rec.process().unwrap();
assert_eq!(reads.load(Ordering::SeqCst), 1, "gotValue was cleared");
assert_eq!(rec.i32inp, 42, "…and the driver's read landed");
}
#[test]
fn a_driver_value_is_ignored_while_the_record_is_not_on_the_io_intr_list() {
use crate::param::ParamValue;
let (rt, reads) = io_intr_port("r11_46_not_armed");
let mut rec = AsynRecord::default();
rec.port = "r11_46_not_armed".to_string();
rec.iface = InterfaceType::Int32 as i32;
rec.tmod = TransferMode::Read as i32;
rec.connect_device().unwrap();
let scan = rec.io_intr_scan();
let _wakeups = scan.take_receiver().unwrap();
fire_interrupt(
&rt,
0,
ParamValue::Int32(7),
crate::interfaces::InterfaceType::Int32,
0,
);
rec.process().unwrap();
assert_eq!(
reads.load(Ordering::SeqCst),
1,
"SCAN is Passive: a real read"
);
assert_eq!(rec.i32inp, 42, "the driver's read, not the interrupt value");
}
#[test]
fn a_second_interrupt_before_the_record_processes_is_dropped() {
use crate::param::ParamValue;
let (rt, _reads) = io_intr_port("r11_46_drop");
let mut rec = AsynRecord::default();
rec.port = "r11_46_drop".to_string();
rec.iface = InterfaceType::Int32 as i32;
rec.tmod = TransferMode::Read as i32;
rec.connect_device().unwrap();
let scan = rec.io_intr_scan();
let _wakeups = scan.take_receiver().unwrap();
rec.set_io_intr_scan(true);
fire_interrupt(
&rt,
0,
ParamValue::Int32(7),
crate::interfaces::InterfaceType::Int32,
0,
);
fire_interrupt(
&rt,
0,
ParamValue::Int32(9),
crate::interfaces::InterfaceType::Int32,
0,
);
rec.process().unwrap();
assert_eq!(rec.i32inp, 7, "C keeps the first unprocessed value");
fire_interrupt(
&rt,
0,
ParamValue::Int32(11),
crate::interfaces::InterfaceType::Int32,
0,
);
rec.process().unwrap();
assert_eq!(rec.i32inp, 11);
}
#[test]
fn an_interrupt_sample_does_not_survive_a_cycle_with_no_port() {
use crate::param::ParamValue;
let (rt, _reads) = io_intr_port("r12_46_no_port");
let mut rec = AsynRecord::default();
rec.port = "r12_46_no_port".to_string();
rec.iface = InterfaceType::Int32 as i32;
rec.tmod = TransferMode::Read as i32;
rec.connect_device().unwrap();
let scan = rec.io_intr_scan();
let _wakeups = scan.take_receiver().unwrap();
rec.set_io_intr_scan(true);
fire_interrupt(
&rt,
0,
ParamValue::Int32(7),
crate::interfaces::InterfaceType::Int32,
0,
);
rec.port = "R12_46_NO_SUCH_PORT".to_string();
let _ = rec.connect_device();
assert!(rec.port_entry.is_none(), "the record now has no port");
rec.process().unwrap();
assert_eq!(
rec.errs, "Not connect to a port",
"C takes the stateNoDevice arm and never looks at gotValue"
);
assert_eq!(
read_alarm(&mut rec),
(alarm_status::STATE_ALARM, AlarmSeverity::Minor),
"C asynRecord.c:361 alarms the refusal STATE/MINOR"
);
assert_eq!(
rec.i32inp, 0,
"the stale interrupt value was never published — a portless record \
must not look healthy and freshly updated"
);
rec.process().unwrap();
assert_eq!(
rec.i32inp, 0,
"the discarded sample does not resurface on the next cycle"
);
}
#[test]
fn a_re_armed_scan_does_not_inherit_the_previous_registrations_value() {
use crate::param::ParamValue;
let (rt, _reads) = io_intr_port("r12_46_rearm");
let mut rec = AsynRecord::default();
rec.port = "r12_46_rearm".to_string();
rec.iface = InterfaceType::Int32 as i32;
rec.tmod = TransferMode::Read as i32;
rec.connect_device().unwrap();
let scan = rec.io_intr_scan();
let _wakeups = scan.take_receiver().unwrap();
rec.set_io_intr_scan(true);
fire_interrupt(
&rt,
0,
ParamValue::Int32(7),
crate::interfaces::InterfaceType::Int32,
0,
);
rec.set_io_intr_scan(false);
rec.set_io_intr_scan(true);
rec.process().unwrap();
assert_eq!(
rec.i32inp, 42,
"the re-armed scan performed a real read — it did not inherit the 7 \
pushed under the previous registration"
);
}
#[test]
fn the_interrupt_gate_still_fires_when_the_record_has_a_port() {
use crate::param::ParamValue;
let (rt, reads) = io_intr_port("r12_46_with_port");
let mut rec = AsynRecord::default();
rec.port = "r12_46_with_port".to_string();
rec.iface = InterfaceType::Int32 as i32;
rec.tmod = TransferMode::Read as i32;
rec.connect_device().unwrap();
let scan = rec.io_intr_scan();
let _wakeups = scan.take_receiver().unwrap();
rec.set_io_intr_scan(true);
fire_interrupt(
&rt,
0,
ParamValue::Int32(7),
crate::interfaces::InterfaceType::Int32,
0,
);
rec.process().unwrap();
assert_eq!(rec.i32inp, 7, "the interrupt value, not a read");
assert_eq!(
reads.load(Ordering::SeqCst),
0,
"C `goto done` queues no I/O on an interrupt-driven cycle"
);
assert_eq!(rec.errs, "", "an interrupt-driven cycle raises nothing");
}
#[test]
fn each_iface_interrupt_writes_its_own_input_field() {
use crate::interfaces::InterfaceType as RegIface;
use crate::param::ParamValue;
let (rt, _reads) = io_intr_port("r11_46_octet");
let mut rec = AsynRecord::default();
rec.port = "r11_46_octet".to_string();
rec.iface = InterfaceType::Octet as i32;
rec.connect_device().unwrap();
let scan = rec.io_intr_scan();
let _wakeups = scan.take_receiver().unwrap();
rec.set_io_intr_scan(true);
fire_interrupt(&rt, 0, ParamValue::Octet("ab\n".into()), RegIface::Octet, 0);
rec.process().unwrap();
assert_eq!(rec.tinp, "ab\\n", "C `epicsStrSnPrintEscaped` into TINP");
let (rt, _reads) = io_intr_port("r11_46_f64");
let mut rec = AsynRecord::default();
rec.port = "r11_46_f64".to_string();
rec.iface = InterfaceType::Float64 as i32;
rec.connect_device().unwrap();
let scan = rec.io_intr_scan();
let _wakeups = scan.take_receiver().unwrap();
rec.set_io_intr_scan(true);
fire_interrupt(&rt, 0, ParamValue::Float64(2.5), RegIface::Float64, 0);
rec.process().unwrap();
assert_eq!(rec.f64inp, 2.5);
let (rt, _reads) = io_intr_port("r11_46_ui32");
let mut rec = AsynRecord::default();
rec.port = "r11_46_ui32".to_string();
rec.iface = InterfaceType::UInt32Digital as i32;
rec.ui32mask = 0x0F;
rec.connect_device().unwrap();
let scan = rec.io_intr_scan();
let _wakeups = scan.take_receiver().unwrap();
rec.set_io_intr_scan(true);
fire_interrupt(
&rt,
0,
ParamValue::UInt32Digital(0xF0),
RegIface::UInt32Digital,
0xF0,
);
rec.process().unwrap();
assert_eq!(rec.ui32inp, 0, "no bit of the change is in UI32MASK");
fire_interrupt(
&rt,
0,
ParamValue::UInt32Digital(0x03),
RegIface::UInt32Digital,
0x03,
);
rec.process().unwrap();
assert_eq!(rec.ui32inp, 0x03);
}
#[test]
fn arming_io_intr_against_a_port_without_the_iface_reports_the_c_text() {
use crate::interfaces::octet_transport_capabilities;
use crate::param::ParamValue;
use crate::port::PortFlags;
use crate::runtime::{RuntimeConfig, create_port_runtime};
struct OctetOnly(PortDriverBase);
impl PortDriver for OctetOnly {
fn base(&self) -> &PortDriverBase {
&self.0
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.0
}
fn capabilities(&self) -> Vec<crate::interfaces::Capability> {
octet_transport_capabilities()
}
}
let port_name = "r11_46_no_int32";
let (rt, _jh) = create_port_runtime(
OctetOnly(PortDriverBase::new(port_name, 1, PortFlags::default())),
RuntimeConfig::default(),
);
register_port(
port_name,
rt.port_handle().clone(),
Arc::new(TraceManager::new()),
)
.unwrap();
let mut rec = AsynRecord::default();
rec.port = port_name.to_string();
rec.iface = InterfaceType::Int32 as i32;
rec.connect_device().unwrap();
let scan = rec.io_intr_scan();
let _wakeups = scan.take_receiver().unwrap();
rec.set_io_intr_scan(true);
assert_eq!(rec.errs, "No asynInt32 interface");
fire_interrupt(
&rt,
0,
ParamValue::Int32(7),
crate::interfaces::InterfaceType::Int32,
0,
);
rec.process().unwrap();
assert_eq!(rec.i32inp, 0, "the refused registration pushed nothing");
}
#[test]
fn reason_iface_ui32mask_and_pcnct_puts_cancel_the_io_intr_scan() {
use crate::param::ParamValue;
let (rt, reads) = io_intr_port("r11_46_cancel");
let mut rec = AsynRecord::default();
rec.port = "r11_46_cancel".to_string();
rec.iface = InterfaceType::Int32 as i32;
rec.tmod = TransferMode::Read as i32;
rec.connect_device().unwrap();
let scan = rec.io_intr_scan();
let _wakeups = scan.take_receiver().unwrap();
let mut reason = 0usize;
for field in ["REASON", "IFACE", "UI32MASK", "PCNCT"] {
rec.set_io_intr_scan(true);
assert!(scan.is_active(), "{field}: armed");
let before = reads.load(Ordering::SeqCst);
fire_interrupt(
&rt,
reason,
ParamValue::Int32(7),
crate::interfaces::InterfaceType::Int32,
0,
);
rec.process().unwrap();
assert_eq!(
rec.i32inp, 7,
"{field}: armed, the interrupt drove the cycle"
);
assert_eq!(
reads.load(Ordering::SeqCst),
before,
"{field}: and did no I/O"
);
match field {
"REASON" => {
rec.reason = 3;
reason = 3;
}
"IFACE" => rec.iface = InterfaceType::Int32 as i32,
"UI32MASK" => rec.ui32mask = 0x0F,
"PCNCT" => rec.pcnct = 0,
_ => unreachable!(),
}
rec.special(field, true).unwrap();
assert!(
!scan.is_active(),
"{field}: C forces SCAN back to Passive, cancelling the registration"
);
if field == "PCNCT" {
rec.pcnct = 1;
rec.special("PCNCT", true).unwrap();
assert!(!scan.is_active(), "a reconnect does not re-arm I/O Intr");
}
let before = reads.load(Ordering::SeqCst);
fire_interrupt(
&rt,
reason,
ParamValue::Int32(7),
crate::interfaces::InterfaceType::Int32,
0,
);
rec.i32inp = 0;
rec.process().unwrap();
assert_eq!(
reads.load(Ordering::SeqCst),
before + 1,
"{field}: no interrupt value gated the cycle"
);
assert_eq!(rec.i32inp, 42, "{field}: the value came from the read");
}
}
}