use std::sync::{Arc, Mutex, OnceLock};
use std::time::Duration;
use epics_base_rs::server::ioc_app::IocApplication;
use epics_base_rs::server::iocsh::registry::{
ArgDesc, ArgType, ArgValue, CommandContext, CommandDef, CommandOutcome, CommandResult,
};
use crate::drivers::ftdi::DrvAsynFtdiPort;
use crate::drivers::ip_port::DrvAsynIPPort;
use crate::drivers::ip_server_port::{DrvAsynIPServerPort, IpServerConfig};
use crate::drivers::prologix::DrvAsynPrologixPort;
use crate::drivers::serial_port::DrvAsynSerialPort;
use crate::drivers::usbtmc::DrvAsynUsbtmcPort;
use crate::drivers::vxi11::DrvVxi11Port;
use crate::error::AsynResult;
use crate::escape::escaped_from_raw;
use crate::manager::PortManager;
use crate::port::PortDriver;
use crate::runtime::config::RuntimeConfig;
use crate::runtime::port::{PortRuntimeHandle, create_port_runtime};
use crate::services::PortServices;
use crate::trace::{TraceFile, TraceInfoMask, TraceIoMask, TraceManager, TraceMask};
use crate::user::AsynUser;
pub fn register_asyn_commands(mut app: IocApplication, mgr: Arc<PortManager>) -> IocApplication {
for def in build_asyn_commands(mgr) {
app = app.register_startup_command(def.clone());
app = app.register_shell_command(def);
}
app
}
fn arg_int(args: &[ArgValue], i: usize) -> Option<i64> {
match args.get(i) {
Some(ArgValue::Int(v)) => Some(*v),
Some(ArgValue::Double(v)) => Some(*v as i64),
Some(ArgValue::String(s)) => s.parse::<i64>().ok(),
_ => None,
}
}
fn arg_f64(args: &[ArgValue], i: usize) -> Option<f64> {
match args.get(i) {
Some(ArgValue::Double(v)) => Some(*v),
Some(ArgValue::Int(v)) => Some(*v as f64),
Some(ArgValue::String(s)) => s.parse::<f64>().ok(),
_ => None,
}
}
fn arg_str(args: &[ArgValue], i: usize) -> Option<String> {
match args.get(i) {
Some(ArgValue::String(s)) => Some(s.clone()),
Some(ArgValue::Int(v)) => Some(v.to_string()),
Some(ArgValue::Double(v)) => Some(v.to_string()),
_ => None,
}
}
fn raw_from_escaped(src: &str) -> Vec<u8> {
let b = src.as_bytes();
let mut out = Vec::with_capacity(b.len());
let mut i = 0;
while i < b.len() {
let c = b[i];
i += 1;
if c == 0 {
break;
}
if c != b'\\' {
out.push(c);
continue;
}
if i >= b.len() {
break;
}
let e = b[i];
i += 1;
if e == 0 {
break;
}
match e {
b'a' => out.push(0x07),
b'b' => out.push(0x08),
b'f' => out.push(0x0C),
b'n' => out.push(b'\n'),
b'r' => out.push(b'\r'),
b't' => out.push(b'\t'),
b'v' => out.push(0x0B),
b'\\' => out.push(b'\\'),
b'\'' => out.push(b'\''),
b'"' => out.push(b'"'),
b'0' => out.push(0),
b'x' => {
let mut u: u32 = 0;
let mut n = 0;
while n < 2
&& i < b.len()
&& b[i] != 0
&& let Some(d) = (b[i] as char).to_digit(16)
{
u = (u << 4) | d;
i += 1;
n += 1;
}
if n > 0 {
out.push(u as u8);
}
}
other => out.push(other),
}
}
out
}
const SHOW_EOS_BUF_SIZE: usize = 4 * 10 + 2;
const SHELL_IO_TIMEOUT: Duration = Duration::from_secs(2);
fn shell_eos_user(addr: i32) -> AsynUser {
AsynUser::default()
.with_addr(addr)
.with_timeout(SHELL_IO_TIMEOUT)
.queue_even_if_not_connected()
}
fn asyn_set_eos(
mgr: &Arc<PortManager>,
ctx: &CommandContext,
args: &[ArgValue],
set_input: bool,
) -> CommandResult {
let cmd = if set_input {
"asynOctetSetInputEos"
} else {
"asynOctetSetOutputEos"
};
let port = arg_str(args, 0).ok_or_else(|| "portName required".to_string())?;
let addr = arg_int(args, 1).unwrap_or(0) as i32;
let eos = raw_from_escaped(&arg_str(args, 2).unwrap_or_default());
match mgr.find_port_handle(&port) {
Ok(handle) => {
let user = shell_eos_user(addr);
let res = if set_input {
handle.set_input_eos_blocking(user, &eos)
} else {
handle.set_output_eos_blocking(user, &eos)
};
if let Err(e) = res {
ctx.println(&format!("{cmd}: {e}"));
}
}
Err(e) => ctx.println(&format!("{cmd}: {e}")),
}
Ok(CommandOutcome::Continue)
}
fn asyn_show_eos(
mgr: &Arc<PortManager>,
ctx: &CommandContext,
args: &[ArgValue],
get_input: bool,
) -> CommandResult {
let cmd = if get_input {
"asynOctetGetInputEos"
} else {
"asynOctetGetOutputEos"
};
let port = arg_str(args, 0).ok_or_else(|| "portName required".to_string())?;
let addr = arg_int(args, 1).unwrap_or(0) as i32;
match mgr.find_port_handle(&port) {
Ok(handle) => {
let user = shell_eos_user(addr);
let res = if get_input {
handle.get_input_eos_blocking(user)
} else {
handle.get_output_eos_blocking(user)
};
match res {
Ok(eos) => ctx.println(&format!(
"\"{}\"",
escaped_from_raw(&eos, SHOW_EOS_BUF_SIZE)
)),
Err(e) => ctx.println(&format!("Get EOS failed: {e}")),
}
}
Err(e) => ctx.println(&format!("{cmd}: {e}")),
}
Ok(CommandOutcome::Continue)
}
fn duration_from_secs(secs: f64) -> Duration {
if secs.is_finite() && secs > 0.0 {
Duration::from_secs_f64(secs)
} else {
Duration::ZERO
}
}
fn enable_style_args() -> Vec<ArgDesc> {
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: false,
},
ArgDesc {
name: "yesNo",
arg_type: ArgType::Int,
optional: false,
},
]
}
fn addresses_a_device(handle: &crate::port_handle::PortHandle, addr: i32) -> bool {
addr >= 0 && handle.is_multi_device()
}
fn shell_enable_target(
mgr: &Arc<PortManager>,
ctx: &CommandContext,
cmd: &str,
args: &[ArgValue],
) -> Option<(crate::port_handle::PortHandle, i32, bool)> {
let port = arg_str(args, 0).filter(|s| !s.is_empty())?;
let addr = arg_int(args, 1).unwrap_or(0) as i32;
let yes = arg_int(args, 2).unwrap_or(0) != 0;
match mgr.find_port_handle(&port) {
Ok(handle) => Some((handle, addr, yes)),
Err(e) => {
ctx.println(&format!("{cmd}: {e}"));
None
}
}
}
fn report_ports(mgr: &Arc<PortManager>, level: i32, port: Option<&str>) {
if let Some(name) = port {
match mgr.find_runtime_handle(name) {
Ok(handle) => {
let _ = handle
.port_handle()
.report_blocking(level)
.map_err(|e| eprintln!("asynReport {name}: {e}"));
}
Err(e) => eprintln!("asynReport: {e}"),
}
} else {
for name in mgr.list_port_names() {
if let Ok(handle) = mgr.find_runtime_handle(&name) {
let _ = handle
.port_handle()
.report_blocking(level)
.map_err(|e| eprintln!("asynReport {name}: {e}"));
}
}
}
}
pub fn register_asyn_commands_on_shell(
shell: &epics_base_rs::server::iocsh::IocShell,
mgr: Arc<PortManager>,
) {
for def in build_asyn_commands(mgr) {
shell.register(def);
}
}
pub fn build_asyn_commands(mgr: Arc<PortManager>) -> Vec<CommandDef> {
let mut out = Vec::new();
let services = mgr.services().clone();
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynReport",
vec![
ArgDesc {
name: "level",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "port",
arg_type: ArgType::String,
optional: true,
},
],
"asynReport [level] [portName] - Report registered ports",
move |args: &[ArgValue], _ctx: &CommandContext| {
let level = arg_int(args, 0).unwrap_or(0) as i32;
let port = arg_str(args, 1);
report_ports(&mgr_r, level, port.as_deref());
Ok(CommandOutcome::Continue)
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynSetOption",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: false,
},
ArgDesc {
name: "key",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "value",
arg_type: ArgType::String,
optional: false,
},
],
"asynSetOption portName addr key value",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0).ok_or_else(|| "portName required".to_string())?;
let addr = arg_int(args, 1).unwrap_or(0) as i32;
let key = arg_str(args, 2).ok_or_else(|| "key required".to_string())?;
let value = arg_str(args, 3).unwrap_or_default();
let user = AsynUser::default()
.with_addr(addr)
.with_timeout(SHELL_IO_TIMEOUT)
.queue_even_if_not_connected();
match mgr_r.find_port_handle(&port) {
Ok(handle) => match handle.set_option_blocking(user, &key, &value) {
Ok(()) => Ok(CommandOutcome::Continue),
Err(e) => {
ctx.println(&format!("asynSetOption: {e}"));
Ok(CommandOutcome::Continue)
}
},
Err(e) => {
ctx.println(&format!("asynSetOption: {e}"));
Ok(CommandOutcome::Continue)
}
}
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynOctetSetInputEos",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: false,
},
ArgDesc {
name: "eos",
arg_type: ArgType::String,
optional: false,
},
],
"asynOctetSetInputEos portName addr eos - set the port input EOS (e.g. \"\\r\\n\")",
move |args: &[ArgValue], ctx: &CommandContext| asyn_set_eos(&mgr_r, ctx, args, true),
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynOctetSetOutputEos",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: false,
},
ArgDesc {
name: "eos",
arg_type: ArgType::String,
optional: false,
},
],
"asynOctetSetOutputEos portName addr eos - set the port output EOS (e.g. \"\\r\\n\")",
move |args: &[ArgValue], ctx: &CommandContext| asyn_set_eos(&mgr_r, ctx, args, false),
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynOctetGetInputEos",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: false,
},
],
"asynOctetGetInputEos portName addr - print the device's input EOS",
move |args: &[ArgValue], ctx: &CommandContext| asyn_show_eos(&mgr_r, ctx, args, true),
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynOctetGetOutputEos",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: false,
},
],
"asynOctetGetOutputEos portName addr - print the device's output EOS",
move |args: &[ArgValue], ctx: &CommandContext| asyn_show_eos(&mgr_r, ctx, args, false),
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynInterposeEcho",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: true,
},
],
"asynInterposeEcho portName [addr] - install the echo interpose \
(half-duplex devices that echo each char)",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0)
.filter(|s| !s.is_empty())
.ok_or_else(|| "portName required".to_string())?;
let addr = arg_int(args, 1).unwrap_or(0) as i32;
match mgr_r.find_port_handle(&port) {
Ok(handle) => {
if let Err(e) = handle.push_echo_interpose_blocking(addr) {
ctx.println(&format!("{port} interposeInterface failed: {e}"));
}
}
Err(_) => ctx.println(&format!("{port} interposeInterface failed.")),
}
Ok(CommandOutcome::Continue)
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynInterposeDelay",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: false,
},
ArgDesc {
name: "delay(sec)",
arg_type: ArgType::Double,
optional: false,
},
],
"asynInterposeDelay portName addr delay(sec) - install the delay \
interpose (one write per character, delay after each)",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0)
.filter(|s| !s.is_empty())
.ok_or_else(|| "portName required".to_string())?;
let addr = arg_int(args, 1).unwrap_or(0) as i32;
let delay =
crate::interpose::delay::delay_from_secs(arg_f64(args, 2).unwrap_or(0.0));
match mgr_r.find_port_handle(&port) {
Ok(handle) => {
if let Err(e) = handle.push_delay_interpose_blocking(addr, delay) {
ctx.println(&format!(
"{port} interposeInterface asynOctetType failed: {e}"
));
}
}
Err(_) => {
ctx.println(&format!("{port} interposeInterface asynOctetType failed."))
}
}
Ok(CommandOutcome::Continue)
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynShowOption",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: false,
},
ArgDesc {
name: "key",
arg_type: ArgType::String,
optional: false,
},
],
"asynShowOption portName addr key - print one driver option",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0).ok_or_else(|| "portName required".to_string())?;
let addr = arg_int(args, 1).unwrap_or(0) as i32;
let Some(key) = arg_str(args, 2) else {
ctx.println("Missing key argument");
return Ok(CommandOutcome::Continue);
};
let user = AsynUser::default()
.with_addr(addr)
.with_timeout(SHELL_IO_TIMEOUT)
.queue_even_if_not_connected();
match mgr_r.find_port_handle(&port) {
Ok(handle) => match handle.get_option_blocking(user, &key) {
Ok(value) => ctx.println(&format!("{key}={value}")),
Err(e) => ctx.println(&format!("getOption failed {e}")),
},
Err(e) => ctx.println(&format!("asynShowOption: {e}")),
}
Ok(CommandOutcome::Continue)
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynSetTraceIOTruncateSize",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: false,
},
ArgDesc {
name: "size",
arg_type: ArgType::Int,
optional: false,
},
],
"asynSetTraceIOTruncateSize portName addr size - bytes of each I/O to trace",
move |args: &[ArgValue], _ctx: &CommandContext| {
let port = arg_str(args, 0).filter(|s| !s.is_empty());
let addr = arg_int(args, 1).unwrap_or(-1) as i32;
let size = arg_int(args, 2).unwrap_or(0).max(0) as usize;
let trace = mgr_r.trace_manager();
match port.as_deref() {
Some(p) if addr >= 0 => trace.set_device_io_truncate_size(p, addr, size),
Some(p) => trace.set_io_truncate_size(Some(p), size),
None => trace.set_io_truncate_size(None, size),
}
Ok(CommandOutcome::Continue)
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynRegisterTimeStampSource",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "functionName",
arg_type: ArgType::String,
optional: false,
},
],
"asynRegisterTimeStampSource portName functionName - stamp this port's values with \
the named source",
move |args: &[ArgValue], ctx: &CommandContext| {
let (Some(port), Some(function)) = (
arg_str(args, 0).filter(|s| !s.is_empty()),
arg_str(args, 1).filter(|s| !s.is_empty()),
) else {
ctx.println("Usage: asynRegisterTimeStampSource portName functionName");
return Ok(CommandOutcome::Continue);
};
match mgr_r.find_port_handle(&port) {
Ok(handle) => {
if let Err(e) = handle.set_time_stamp_source_blocking(Some(&function)) {
ctx.println(&format!("asynRegisterTimeStampSource: {e}"));
}
}
Err(_) => ctx.println(&format!(
"asynRegisterTimeStampSource, cannot connect to port {port}"
)),
}
Ok(CommandOutcome::Continue)
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynUnregisterTimeStampSource",
vec![ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
}],
"asynUnregisterTimeStampSource portName - back to the port's default clock",
move |args: &[ArgValue], ctx: &CommandContext| {
let Some(port) = arg_str(args, 0).filter(|s| !s.is_empty()) else {
ctx.println("Usage: asynUnregisterTimeStampSource portName");
return Ok(CommandOutcome::Continue);
};
match mgr_r.find_port_handle(&port) {
Ok(handle) => {
if let Err(e) = handle.set_time_stamp_source_blocking(None) {
ctx.println(&format!("asynUnregisterTimeStampSource: {e}"));
}
}
Err(_) => ctx.println(&format!(
"asynUnregisterTimeStampSource, cannot connect to port {port}"
)),
}
Ok(CommandOutcome::Continue)
},
));
}
out.push(CommandDef::new(
"asynSetMinTimerPeriod",
vec![ArgDesc {
name: "minimum period",
arg_type: ArgType::Double,
optional: false,
}],
"asynSetMinTimerPeriod period - Windows-only timer resolution (no effect here)",
move |_args: &[ArgValue], ctx: &CommandContext| {
ctx.println("asynSetMinTimerPeriod is not currently supported on this OS");
Ok(CommandOutcome::Continue)
},
));
{
let mgr_r = mgr.clone();
let services_wait = services.clone();
out.push(CommandDef::new(
"asynWaitConnect",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "timeout",
arg_type: ArgType::Double,
optional: false,
},
],
"asynWaitConnect portName timeout - block until the port is connected",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0)
.filter(|s| !s.is_empty())
.ok_or_else(|| "portName required".to_string())?;
let timeout = duration_from_secs(arg_f64(args, 1).unwrap_or(0.0));
let handle = match mgr_r.find_port_handle(&port) {
Ok(h) => h,
Err(e) => {
ctx.println(&format!("asynWaitConnect: {e}"));
return Ok(CommandOutcome::Continue);
}
};
let waiter = crate::runtime::port::ConnectWaiter::arm(&services_wait, &port);
if handle.is_connected_blocking().unwrap_or(false) {
return Ok(CommandOutcome::Continue);
}
if !waiter.wait(timeout) {
ctx.println(&format!("asynWaitConnect: {port} not connected"));
}
Ok(CommandOutcome::Continue)
},
));
}
out.push(CommandDef::new(
"asynSetAutoConnectTimeout",
vec![ArgDesc {
name: "timeout",
arg_type: ArgType::Double,
optional: false,
}],
"asynSetAutoConnectTimeout timeout - seconds a new port waits for its first connect \
(C default 0.5)",
move |args: &[ArgValue], _ctx: &CommandContext| {
crate::runtime::config::set_auto_connect_timeout(duration_from_secs(
arg_f64(args, 0).unwrap_or(0.0),
));
Ok(CommandOutcome::Continue)
},
));
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynInterposeEosConfig",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: false,
},
ArgDesc {
name: "processIn",
arg_type: ArgType::Int,
optional: false,
},
ArgDesc {
name: "processOut",
arg_type: ArgType::Int,
optional: false,
},
],
"asynInterposeEosConfig portName addr processIn processOut - install the EOS \
interpose (terminator handling) on the port",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0)
.filter(|s| !s.is_empty())
.ok_or_else(|| "portName required".to_string())?;
let addr = arg_int(args, 1).unwrap_or(0) as i32;
let process_in = arg_int(args, 2).unwrap_or(0) != 0;
let process_out = arg_int(args, 3).unwrap_or(0) != 0;
match mgr_r.find_port_handle(&port) {
Ok(handle) => {
if let Err(e) =
handle.push_eos_interpose_blocking(addr, process_in, process_out)
{
ctx.println(&format!("{port} interposeInterface failed: {e}"));
}
}
Err(_) => ctx.println(&format!("{port} interposeInterface failed.")),
}
Ok(CommandOutcome::Continue)
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynInterposeFlushConfig",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: false,
},
ArgDesc {
name: "timeout(msec)",
arg_type: ArgType::Double,
optional: false,
},
],
"asynInterposeFlushConfig portName addr timeout(msec) - install the flush \
interpose (a read with this timeout drains the input before each write)",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0)
.filter(|s| !s.is_empty())
.ok_or_else(|| "portName required".to_string())?;
let addr = arg_int(args, 1).unwrap_or(0) as i32;
let ms = arg_f64(args, 2).unwrap_or(0.0) as i64;
let ms = if ms <= 0 { 1 } else { ms };
let timeout = Duration::from_millis(ms as u64);
match mgr_r.find_port_handle(&port) {
Ok(handle) => {
if let Err(e) = handle.push_flush_interpose_blocking(addr, timeout) {
ctx.println(&format!("{port} interposeInterface failed: {e}"));
}
}
Err(_) => ctx.println(&format!("{port} interposeInterface failed.")),
}
Ok(CommandOutcome::Continue)
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynSetTraceMask",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: true,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "mask",
arg_type: ArgType::String,
optional: false,
},
],
"asynSetTraceMask [portName] [addr] mask",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0).filter(|s| !s.is_empty());
let addr = arg_int(args, 1).unwrap_or(-1) as i32;
let mask_str = arg_str(args, 2).ok_or_else(|| "mask required".to_string())?;
match TraceMask::from_symbolic(&mask_str) {
Ok(m) => {
let trace = mgr_r.trace_manager();
if let Some(p) = port.as_deref() {
if addr >= 0 {
trace.set_device_trace_mask(p, addr, m);
} else {
trace.set_trace_mask(Some(p), m);
}
} else {
trace.set_trace_mask(None, m);
}
Ok(CommandOutcome::Continue)
}
Err(e) => {
ctx.println(&format!("asynSetTraceMask: {e}"));
Ok(CommandOutcome::Continue)
}
}
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynSetTraceIOMask",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: true,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "mask",
arg_type: ArgType::String,
optional: false,
},
],
"asynSetTraceIOMask [portName] [addr] mask",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0).filter(|s| !s.is_empty());
let addr = arg_int(args, 1).unwrap_or(-1) as i32;
let mask_str = arg_str(args, 2).ok_or_else(|| "mask required".to_string())?;
match TraceIoMask::from_symbolic(&mask_str) {
Ok(m) => {
let trace = mgr_r.trace_manager();
if let Some(p) = port.as_deref() {
if addr >= 0 {
trace.set_device_trace_io_mask(p, addr, m);
} else {
trace.set_trace_io_mask(Some(p), m);
}
} else {
trace.set_trace_io_mask(None, m);
}
Ok(CommandOutcome::Continue)
}
Err(e) => {
ctx.println(&format!("asynSetTraceIOMask: {e}"));
Ok(CommandOutcome::Continue)
}
}
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynSetTraceInfoMask",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: true,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "mask",
arg_type: ArgType::String,
optional: false,
},
],
"asynSetTraceInfoMask [portName] [addr] mask",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0).filter(|s| !s.is_empty());
let addr = arg_int(args, 1).unwrap_or(-1) as i32;
let mask_str = arg_str(args, 2).ok_or_else(|| "mask required".to_string())?;
match TraceInfoMask::from_symbolic(&mask_str) {
Ok(m) => {
let trace = mgr_r.trace_manager();
if let Some(p) = port.as_deref() {
if addr >= 0 {
trace.set_device_trace_info_mask(p, addr, m);
} else {
trace.set_trace_info_mask(Some(p), m);
}
} else {
trace.set_trace_info_mask(None, m);
}
Ok(CommandOutcome::Continue)
}
Err(e) => {
ctx.println(&format!("asynSetTraceInfoMask: {e}"));
Ok(CommandOutcome::Continue)
}
}
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynSetTraceFile",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: true,
},
ArgDesc {
name: "addr",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "filename",
arg_type: ArgType::String,
optional: true,
},
],
"asynSetTraceFile [portName] [addr] [filename]",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0).filter(|s| !s.is_empty());
let addr = arg_int(args, 1).unwrap_or(-1) as i32;
let filename = arg_str(args, 2).unwrap_or_default();
let target = match filename.as_str() {
"" | "stderr" => TraceFile::Stderr,
"stdout" => TraceFile::Stdout,
path => match std::fs::File::create(path) {
Ok(f) => TraceFile::File(Arc::new(std::sync::Mutex::new(f))),
Err(e) => {
ctx.println(&format!("asynSetTraceFile: fopen failed: {e}"));
return Ok(CommandOutcome::Continue);
}
},
};
let trace = mgr_r.trace_manager();
if let Some(p) = port.as_deref() {
if addr >= 0 {
trace.set_device_trace_file(p, addr, target);
} else {
trace.set_trace_file(Some(p), target);
}
} else {
trace.set_trace_file(None, target);
}
Ok(CommandOutcome::Continue)
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynEnable",
enable_style_args(),
"asynEnable portName addr yesNo - enable (1) or disable (0) a port or one of its devices",
move |args: &[ArgValue], ctx: &CommandContext| {
let (handle, addr, yes) = match shell_enable_target(&mgr_r, ctx, "asynEnable", args)
{
Some(t) => t,
None => return Ok(CommandOutcome::Continue),
};
let r = match (addresses_a_device(&handle, addr), yes) {
(true, true) => handle.enable_addr_blocking(addr),
(true, false) => handle.disable_addr_blocking(addr),
(false, yes) => handle.set_enable_blocking(yes),
};
if let Err(e) = r {
ctx.println(&format!("asynEnable: {e}"));
}
Ok(CommandOutcome::Continue)
},
));
}
{
let mgr_r = mgr.clone();
out.push(CommandDef::new(
"asynAutoConnect",
enable_style_args(),
"asynAutoConnect portName addr yesNo - turn auto-connect on (1) or off (0) for a port \
or one of its devices",
move |args: &[ArgValue], ctx: &CommandContext| {
let (handle, addr, yes) =
match shell_enable_target(&mgr_r, ctx, "asynAutoConnect", args) {
Some(t) => t,
None => return Ok(CommandOutcome::Continue),
};
let r = if addresses_a_device(&handle, addr) {
handle.set_auto_connect_addr_blocking(addr, yes)
} else {
handle.set_auto_connect_blocking(yes)
};
if let Err(e) = r {
ctx.println(&format!("asynAutoConnect: {e}"));
}
Ok(CommandOutcome::Continue)
},
));
}
out.push(drv_asyn_ip_port_configure_command(services.clone()));
out.push(drv_asyn_ip_server_port_configure_command(services.clone()));
out.push(drv_asyn_serial_port_configure_command(services.clone()));
out.push(drv_asyn_ftdi_port_configure_command(services.clone()));
out.push(vxi11_configure_command(services.clone()));
out.push(usbtmc_configure_command(services.clone()));
out.push(drv_asyn_prologix_port_configure_command(services));
out
}
pub(crate) fn build_configured_ip_port(
port: &str,
host: &str,
no_auto_connect: bool,
no_process_eos: bool,
) -> AsynResult<DrvAsynIPPort> {
DrvAsynIPPort::new_configured(port, host, no_auto_connect, no_process_eos)
}
pub fn drv_asyn_ip_port_configure_command(services: PortServices) -> CommandDef {
CommandDef::new(
"drvAsynIPPortConfigure",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "hostInfo",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "priority",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "noAutoConnect",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "noProcessEos",
arg_type: ArgType::Int,
optional: true,
},
],
"drvAsynIPPortConfigure portName hostInfo [priority] [noAutoConnect] [noProcessEos] \
- create an IP octet port",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0)
.filter(|s| !s.is_empty())
.ok_or_else(|| "portName required".to_string())?;
let host = arg_str(args, 1)
.filter(|s| !s.is_empty())
.ok_or_else(|| "hostInfo required".to_string())?;
let no_auto_connect = arg_int(args, 3).unwrap_or(0) != 0;
let no_process_eos = arg_int(args, 4).unwrap_or(0) != 0;
let driver =
match build_configured_ip_port(&port, &host, no_auto_connect, no_process_eos) {
Ok(d) => d,
Err(e) => {
ctx.println(&format!("drvAsynIPPortConfigure: {e}"));
return Ok(CommandOutcome::Continue);
}
};
let config = RuntimeConfig {
services: services.clone(),
..RuntimeConfig::default()
};
let (handle, _jh) = create_port_runtime(driver, config);
if let Err(e) = crate::asyn_record::register_port(
&port,
handle.port_handle().clone(),
services.trace().clone(),
) {
ctx.println(&format!("drvAsynIPPortConfigure: {e}"));
handle.shutdown();
return Ok(CommandOutcome::Continue);
}
drop(handle);
ctx.println(&format!(
"drvAsynIPPortConfigure: octet port '{port}' -> {host}"
));
Ok(CommandOutcome::Continue)
},
)
}
pub fn drv_asyn_ip_server_port_configure_command(services: PortServices) -> CommandDef {
CommandDef::new(
"drvAsynIPServerPortConfigure",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "serverInfo",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "maxClients",
arg_type: ArgType::Int,
optional: false,
},
ArgDesc {
name: "priority",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "noAutoConnect",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "noProcessEos",
arg_type: ArgType::Int,
optional: true,
},
],
"drvAsynIPServerPortConfigure portName serverInfo maxClients [priority] \
[noAutoConnect] [noProcessEos] - create an IP server (listening) port",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0)
.filter(|s| !s.is_empty())
.ok_or_else(|| "portName required".to_string())?;
let server_info = arg_str(args, 1)
.filter(|s| !s.is_empty())
.ok_or_else(|| "serverInfo required".to_string())?;
let max_clients = arg_int(args, 2).unwrap_or(0);
let no_auto_connect = arg_int(args, 4).unwrap_or(0) != 0;
let no_process_eos = arg_int(args, 5).unwrap_or(0) != 0;
let mut config = match IpServerConfig::parse(&server_info) {
Ok(c) => c,
Err(e) => {
ctx.println(&format!("drvAsynIPServerPortConfigure: {e}"));
return Ok(CommandOutcome::Continue);
}
};
config.max_clients = max_clients.max(0) as usize;
config.no_process_eos = no_process_eos;
let mut driver = match DrvAsynIPServerPort::with_config(&port, config) {
Ok(d) => d,
Err(e) => {
ctx.println(&format!("drvAsynIPServerPortConfigure: {e}"));
return Ok(CommandOutcome::Continue);
}
};
if no_auto_connect {
driver.base_mut().auto_connect = false;
}
let n_children = driver.child_port_names().len();
let children = (0..n_children)
.map(|i| driver.make_subport(i))
.collect::<AsynResult<Vec<_>>>();
let children = match children {
Ok(c) => c,
Err(e) => {
ctx.println(&format!("drvAsynIPServerPortConfigure: {e}"));
return Ok(CommandOutcome::Continue);
}
};
let config = RuntimeConfig {
services: services.clone(),
..RuntimeConfig::default()
};
let (handle, _jh) = create_port_runtime(driver, config.clone());
if let Err(e) = crate::asyn_record::register_port(
&port,
handle.port_handle().clone(),
services.trace().clone(),
) {
ctx.println(&format!("drvAsynIPServerPortConfigure: {e}"));
handle.shutdown();
return Ok(CommandOutcome::Continue);
}
if let Err(e) = handle.port_handle().connect_blocking() {
ctx.println(&format!(
"drvAsynIPServerPortConfigure: cannot listen on {server_info}: {e}"
));
crate::asyn_record::unregister_port(&port);
handle.shutdown();
return Ok(CommandOutcome::Continue);
}
drop(handle);
for child in children {
let name = child.base().port_name.clone();
let (child_handle, _cjh) = create_port_runtime(child, config.clone());
if let Err(e) = crate::asyn_record::register_port(
&name,
child_handle.port_handle().clone(),
services.trace().clone(),
) {
ctx.println(&format!("drvAsynIPServerPortConfigure: {name}: {e}"));
child_handle.shutdown();
continue;
}
drop(child_handle);
}
ctx.println(&format!(
"drvAsynIPServerPortConfigure: server port '{port}' -> {server_info}"
));
Ok(CommandOutcome::Continue)
},
)
}
pub fn drv_asyn_serial_port_configure_command(services: PortServices) -> CommandDef {
CommandDef::new(
"drvAsynSerialPortConfigure",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "ttyName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "priority",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "noAutoConnect",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "noProcessEos",
arg_type: ArgType::Int,
optional: true,
},
],
"drvAsynSerialPortConfigure portName ttyName [priority] [noAutoConnect] [noProcessEos] \
- create a serial octet port",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0)
.filter(|s| !s.is_empty())
.ok_or_else(|| "portName required".to_string())?;
let tty = arg_str(args, 1)
.filter(|s| !s.is_empty())
.ok_or_else(|| "ttyName required".to_string())?;
let no_auto_connect = arg_int(args, 3).unwrap_or(0) != 0;
let no_process_eos = arg_int(args, 4).unwrap_or(0) != 0;
let driver =
match DrvAsynSerialPort::configure(&port, &tty, no_auto_connect, no_process_eos) {
Ok(d) => d,
Err(e) => {
ctx.println(&format!("drvAsynSerialPortConfigure: {e}"));
return Ok(CommandOutcome::Continue);
}
};
let config = RuntimeConfig {
services: services.clone(),
..RuntimeConfig::default()
};
let (handle, _jh) = create_port_runtime(driver, config);
if let Err(e) = crate::asyn_record::register_port(
&port,
handle.port_handle().clone(),
services.trace().clone(),
) {
ctx.println(&format!("drvAsynSerialPortConfigure: {e}"));
handle.shutdown();
return Ok(CommandOutcome::Continue);
}
drop(handle);
ctx.println(&format!(
"drvAsynSerialPortConfigure: octet port '{port}' -> {tty}"
));
Ok(CommandOutcome::Continue)
},
)
}
pub fn drv_asyn_prologix_port_configure_command(services: PortServices) -> CommandDef {
CommandDef::new(
"prologixGPIBConfigure",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "host",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "priority",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "noAutoConnect",
arg_type: ArgType::Int,
optional: true,
},
],
"prologixGPIBConfigure portName host [priority] [noAutoConnect] \
- create a Prologix GPIB-Ethernet port",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0)
.filter(|s| !s.is_empty())
.ok_or_else(|| "portName required".to_string())?;
let host = arg_str(args, 1)
.filter(|s| !s.is_empty())
.ok_or_else(|| "host required".to_string())?;
let no_auto_connect = arg_int(args, 3).unwrap_or(0) != 0;
let driver = match DrvAsynPrologixPort::new(&port, &host, no_auto_connect) {
Ok(d) => d,
Err(e) => {
ctx.println(&format!("prologixGPIBConfigure: {e}"));
return Ok(CommandOutcome::Continue);
}
};
let config = RuntimeConfig {
services: services.clone(),
..RuntimeConfig::default()
};
let (handle, _jh) = create_port_runtime(driver, config);
if let Err(e) = crate::asyn_record::register_port(
&port,
handle.port_handle().clone(),
services.trace().clone(),
) {
ctx.println(&format!("prologixGPIBConfigure: {e}"));
handle.shutdown();
return Ok(CommandOutcome::Continue);
}
drop(handle);
ctx.println(&format!(
"prologixGPIBConfigure: GPIB port '{port}' -> {host}"
));
Ok(CommandOutcome::Continue)
},
)
}
fn publish_configured_port<D: PortDriver>(
command: &str,
port: &str,
driver: D,
services: &PortServices,
ctx: &CommandContext,
) -> bool {
let config = RuntimeConfig {
services: services.clone(),
..RuntimeConfig::default()
};
let (handle, _jh) = create_port_runtime(driver, config);
if let Err(e) = crate::asyn_record::register_port(
port,
handle.port_handle().clone(),
services.trace().clone(),
) {
ctx.println(&format!("{command}: {e}"));
handle.shutdown();
return false;
}
drop(handle);
true
}
pub fn drv_asyn_ftdi_port_configure_command(services: PortServices) -> CommandDef {
let int_args = [
"vendorID",
"productID",
"baudrate",
"latency",
"priority",
"noAutoConnect",
"noProcessEos",
"mode",
];
let mut arg_descs = vec![ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
}];
arg_descs.extend(int_args.into_iter().map(|name| ArgDesc {
name,
arg_type: ArgType::Int,
optional: true,
}));
CommandDef::new(
"drvAsynFTDIPortConfigure",
arg_descs,
"drvAsynFTDIPortConfigure portName vendorID productID baudrate latency [priority] \
[noAutoConnect] [noProcessEos] [mode] - create an FTDI octet port",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0)
.filter(|s| !s.is_empty())
.ok_or_else(|| "portName required".to_string())?;
let driver = match DrvAsynFtdiPort::configure(
&port,
arg_int(args, 1).unwrap_or(0) as i32,
arg_int(args, 2).unwrap_or(0) as i32,
arg_int(args, 3).unwrap_or(0) as i32,
arg_int(args, 4).unwrap_or(0) as i32,
arg_int(args, 5).unwrap_or(0) as u32,
arg_int(args, 6).unwrap_or(0) != 0,
arg_int(args, 7).unwrap_or(0) != 0,
arg_int(args, 8).unwrap_or(0) as i32,
) {
Ok(d) => d,
Err(e) => {
ctx.println(&format!("drvAsynFTDIPortConfigure: {e}"));
return Ok(CommandOutcome::Continue);
}
};
if publish_configured_port("drvAsynFTDIPortConfigure", &port, driver, &services, ctx) {
ctx.println(&format!(
"drvAsynFTDIPortConfigure: FTDI port '{port}' created"
));
}
Ok(CommandOutcome::Continue)
},
)
}
pub fn vxi11_configure_command(services: PortServices) -> CommandDef {
CommandDef::new(
"vxi11Configure",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "hostName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "flags",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "defaultTimeout",
arg_type: ArgType::String,
optional: true,
},
ArgDesc {
name: "vxiName",
arg_type: ArgType::String,
optional: true,
},
ArgDesc {
name: "priority",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "noAutoConnect",
arg_type: ArgType::Int,
optional: true,
},
],
"vxi11Configure portName hostName [flags] [defaultTimeout] [vxiName] [priority] \
[noAutoConnect] - create a VXI-11 port",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0)
.filter(|s| !s.is_empty())
.ok_or_else(|| "portName required".to_string())?;
let host = arg_str(args, 1)
.filter(|s| !s.is_empty())
.ok_or_else(|| "hostName required".to_string())?;
let driver = match DrvVxi11Port::configure(
&port,
&host,
arg_int(args, 2).unwrap_or(0) as i32,
&arg_str(args, 3).unwrap_or_default(),
&arg_str(args, 4).unwrap_or_default(),
arg_int(args, 5).unwrap_or(0) as i32,
arg_int(args, 6).unwrap_or(0) != 0,
) {
Ok(d) => d,
Err(e) => {
ctx.println(&format!("vxi11Configure: {e}"));
return Ok(CommandOutcome::Continue);
}
};
if publish_configured_port("vxi11Configure", &port, driver, &services, ctx) {
ctx.println(&format!("vxi11Configure: VXI-11 port '{port}' -> {host}"));
}
Ok(CommandOutcome::Continue)
},
)
}
pub fn usbtmc_configure_command(services: PortServices) -> CommandDef {
CommandDef::new(
"usbtmcConfigure",
vec![
ArgDesc {
name: "portName",
arg_type: ArgType::String,
optional: false,
},
ArgDesc {
name: "vendorID",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "productID",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "serialNumber",
arg_type: ArgType::String,
optional: true,
},
ArgDesc {
name: "priority",
arg_type: ArgType::Int,
optional: true,
},
ArgDesc {
name: "flags",
arg_type: ArgType::Int,
optional: true,
},
],
"usbtmcConfigure portName [vendorID] [productID] [serialNumber] [priority] [flags] \
- create a USBTMC port",
move |args: &[ArgValue], ctx: &CommandContext| {
let port = arg_str(args, 0)
.filter(|s| !s.is_empty())
.ok_or_else(|| "portName required".to_string())?;
let driver = match DrvAsynUsbtmcPort::configure(
&port,
arg_int(args, 1).unwrap_or(0) as i32,
arg_int(args, 2).unwrap_or(0) as i32,
&arg_str(args, 3).unwrap_or_default(),
arg_int(args, 4).unwrap_or(0) as i32,
arg_int(args, 5).unwrap_or(0) as i32,
) {
Ok(d) => d,
Err(e) => {
ctx.println(&format!("usbtmcConfigure: {e}"));
return Ok(CommandOutcome::Continue);
}
};
if publish_configured_port("usbtmcConfigure", &port, driver, &services, ctx) {
ctx.println(&format!("usbtmcConfigure: USBTMC port '{port}' created"));
}
Ok(CommandOutcome::Continue)
},
)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::error::AsynResult;
use crate::exception::AsynException;
use crate::param::ParamType;
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::user::AsynUser;
use std::sync::Mutex;
struct DummyDriver {
base: PortDriverBase,
}
impl DummyDriver {
fn new(name: &str) -> Self {
let mut base = PortDriverBase::new(name, 1, PortFlags::default());
base.create_param("VAL", ParamType::Int32).unwrap();
Self { base }
}
}
impl PortDriver for DummyDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
}
impl DummyDriver {
fn disconnected(name: &str) -> Self {
let mut d = Self::new(name);
d.base.init_connected(false);
d.base.set_auto_connect(false);
d
}
fn multi_device(name: &str, max_addr: usize) -> Self {
let mut base = PortDriverBase::new(
name,
max_addr,
PortFlags {
multi_device: true,
..PortFlags::default()
},
);
base.create_param("VAL", ParamType::Int32).unwrap();
Self { base }
}
}
struct OctetDriver {
base: PortDriverBase,
input: Vec<u8>,
pos: usize,
written: Arc<Mutex<Vec<u8>>>,
}
impl OctetDriver {
fn new(name: &str, input: &[u8], written: Arc<Mutex<Vec<u8>>>) -> Self {
Self {
base: PortDriverBase::new(name, 1, PortFlags::default()),
input: input.to_vec(),
pos: 0,
written,
}
}
}
impl PortDriver for OctetDriver {
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]) -> AsynResult<usize> {
let n = (self.input.len() - self.pos).min(buf.len());
buf[..n].copy_from_slice(&self.input[self.pos..self.pos + n]);
self.pos += n;
Ok(n)
}
fn io_write_octet(&mut self, _user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
self.written.lock().unwrap().extend_from_slice(data);
Ok(data.len())
}
}
fn fresh_mgr_with_port(name: &str) -> Arc<PortManager> {
let mgr = Arc::new(PortManager::new());
let _ = mgr.register_port(DummyDriver::new(name)).unwrap();
mgr
}
#[test]
fn iocsh_set_trace_mask_updates_port_mask() {
let mgr = fresh_mgr_with_port("trace_mask_port");
let observed: Arc<Mutex<Vec<AsynException>>> = Arc::new(Mutex::new(Vec::new()));
let observed_clone = observed.clone();
mgr.exception_manager().add_callback(move |ev| {
observed_clone.lock().unwrap().push(ev.exception);
});
let cmds = build_asyn_commands(mgr.clone());
let set_trace_mask = cmds
.iter()
.find(|c| c.name == "asynSetTraceMask")
.expect("asynSetTraceMask must be registered");
let trace = mgr.trace_manager().clone();
let mask = TraceMask::from_symbolic("ERROR+WARNING").unwrap();
trace.set_trace_mask(Some("trace_mask_port"), mask);
assert!(set_trace_mask.args.len() == 3);
let evs = observed.lock().unwrap();
assert!(
evs.iter().any(|e| matches!(e, AsynException::TraceMask)),
"set_trace_mask must fire asynExceptionTraceMask"
);
assert!(trace.is_enabled("trace_mask_port", TraceMask::ERROR));
assert!(trace.is_enabled("trace_mask_port", TraceMask::WARNING));
}
#[test]
fn iocsh_registers_c_parity_commands() {
let mgr = Arc::new(PortManager::new());
let cmds = build_asyn_commands(mgr);
let mut names: Vec<&str> = cmds.iter().map(|c| c.name.as_str()).collect();
names.sort_unstable();
let mut expected = vec![
"asynReport",
"asynSetOption",
"asynSetTraceMask",
"asynSetTraceIOMask",
"asynSetTraceInfoMask",
"asynSetTraceFile",
"asynEnable",
"asynAutoConnect",
"asynWaitConnect",
"asynSetAutoConnectTimeout",
"asynShowOption",
"asynSetTraceIOTruncateSize",
"asynRegisterTimeStampSource",
"asynUnregisterTimeStampSource",
"asynSetMinTimerPeriod",
"asynOctetSetInputEos",
"asynOctetGetInputEos",
"asynOctetSetOutputEos",
"asynOctetGetOutputEos",
"asynInterposeEcho",
"asynInterposeDelay",
"asynInterposeEosConfig",
"asynInterposeFlushConfig",
"drvAsynIPPortConfigure",
"drvAsynIPServerPortConfigure",
"drvAsynSerialPortConfigure",
"drvAsynFTDIPortConfigure",
"vxi11Configure",
"usbtmcConfigure",
"prologixGPIBConfigure",
];
expected.sort_unstable();
assert_eq!(names, expected);
}
#[test]
fn iocsh_set_trace_io_truncate_size_routes_addr_to_the_device() {
let mgr = fresh_mgr_with_port("trunc_port");
let trace = mgr.trace_manager().clone();
let cmds = build_asyn_commands(mgr);
let ctx = make_ctx();
let set = |addr: i64, size: i64| {
cmds.iter()
.find(|c| c.name == "asynSetTraceIOTruncateSize")
.expect("asynSetTraceIOTruncateSize must be registered")
.handler
.call(
&[
ArgValue::String("trunc_port".into()),
ArgValue::Int(addr),
ArgValue::Int(size),
],
&ctx,
)
.unwrap();
};
set(-1, 40);
assert_eq!(trace.snapshot("trunc_port", None).io_truncate_size, 40);
set(2, 4096);
assert_eq!(trace.snapshot("trunc_port", Some(2)).io_truncate_size, 4096);
assert_eq!(trace.snapshot("trunc_port", None).io_truncate_size, 40);
}
#[test]
fn iocsh_show_option_reads_back_what_set_option_wrote() {
let mgr = fresh_mgr_with_port("show_opt");
let handle = mgr.find_port_handle("show_opt").unwrap();
handle
.set_option_blocking(AsynUser::default(), "baud", "115200")
.unwrap();
let cmds = build_asyn_commands(mgr);
let ctx = make_ctx();
cmds.iter()
.find(|c| c.name == "asynShowOption")
.expect("asynShowOption must be registered")
.handler
.call(
&[
ArgValue::String("show_opt".into()),
ArgValue::Int(0),
ArgValue::String("baud".into()),
],
&ctx,
)
.unwrap();
assert_eq!(
handle
.get_option_blocking(AsynUser::default(), "baud")
.unwrap(),
"115200"
);
}
#[test]
fn iocsh_time_stamp_source_is_installed_by_name_and_refused_when_unknown() {
use crate::timestamp::{find_time_stamp_source, register_time_stamp_source};
use std::time::{Duration as StdDuration, UNIX_EPOCH};
let fixed = UNIX_EPOCH + StdDuration::from_secs(1_234_567);
register_time_stamp_source("iocsh_fixed_clock", move || fixed);
let mgr = fresh_mgr_with_port("ts_port");
let handle = mgr.find_port_handle("ts_port").unwrap();
let cmds = build_asyn_commands(mgr);
let ctx = make_ctx();
let register = |name: &str| {
cmds.iter()
.find(|c| c.name == "asynRegisterTimeStampSource")
.expect("asynRegisterTimeStampSource must be registered")
.handler
.call(
&[
ArgValue::String("ts_port".into()),
ArgValue::String(name.into()),
],
&ctx,
)
.unwrap();
};
assert!(find_time_stamp_source("no_such_clock").is_none());
register("no_such_clock");
assert!(
handle
.set_time_stamp_source_blocking(Some("no_such_clock"))
.is_err()
);
register("iocsh_fixed_clock");
assert!(
handle
.set_time_stamp_source_blocking(Some("iocsh_fixed_clock"))
.is_ok()
);
cmds.iter()
.find(|c| c.name == "asynUnregisterTimeStampSource")
.expect("asynUnregisterTimeStampSource must be registered")
.handler
.call(&[ArgValue::String("ts_port".into())], &ctx)
.unwrap();
assert!(handle.set_time_stamp_source_blocking(None).is_ok());
}
#[test]
fn iocsh_wait_connect_covers_before_during_and_never() {
use std::time::Instant;
let mgr = Arc::new(PortManager::new());
let cfg = RuntimeConfig {
auto_connect: false,
services: mgr.services().clone(),
..RuntimeConfig::default()
};
let _ = mgr
.register_port_with_config(DummyDriver::disconnected("wc_late"), cfg.clone())
.unwrap();
let _ = mgr
.register_port_with_config(DummyDriver::disconnected("wc_never"), cfg)
.unwrap();
let cmds = build_asyn_commands(mgr.clone());
let ctx = make_ctx();
let wait = |port: &str, timeout: f64| {
let t0 = Instant::now();
cmds.iter()
.find(|c| c.name == "asynWaitConnect")
.expect("asynWaitConnect must be registered")
.handler
.call(
&[ArgValue::String(port.into()), ArgValue::Double(timeout)],
&ctx,
)
.unwrap();
t0.elapsed()
};
let elapsed = wait("wc_never", 0.1);
assert!(
elapsed >= Duration::from_millis(100),
"must have waited the full timeout, waited {elapsed:?}"
);
assert!(
!mgr.find_port_handle("wc_never")
.unwrap()
.is_connected_blocking()
.unwrap()
);
let late = mgr.find_port_handle("wc_late").unwrap();
let connector = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(50));
late.connect_blocking().unwrap();
});
let elapsed = wait("wc_late", 5.0);
connector.join().unwrap();
assert!(
elapsed < Duration::from_secs(4),
"must have returned on the connect, not on the timeout ({elapsed:?})"
);
let late = mgr.find_port_handle("wc_late").unwrap();
assert!(late.is_connected_blocking().unwrap());
let elapsed = wait("wc_late", 5.0);
assert!(
elapsed < Duration::from_secs(1),
"an already-connected port must not wait ({elapsed:?})"
);
}
#[test]
fn iocsh_set_auto_connect_timeout_rewrites_the_registration_wait() {
assert_eq!(
RuntimeConfig::default().auto_connect_timeout,
Duration::from_millis(500),
"C DEFAULT_AUTOCONNECT_TIMEOUT (asynManager.c:49)"
);
let mgr = Arc::new(PortManager::new());
let cmds = build_asyn_commands(mgr);
let ctx = make_ctx();
cmds.iter()
.find(|c| c.name == "asynSetAutoConnectTimeout")
.expect("asynSetAutoConnectTimeout must be registered")
.handler
.call(&[ArgValue::Double(2.5)], &ctx)
.unwrap();
assert_eq!(
RuntimeConfig::default().auto_connect_timeout,
Duration::from_millis(2500)
);
cmds.iter()
.find(|c| c.name == "asynSetAutoConnectTimeout")
.unwrap()
.handler
.call(&[ArgValue::Double(-1.0)], &ctx)
.unwrap();
assert_eq!(
RuntimeConfig::default().auto_connect_timeout,
Duration::ZERO
);
}
#[test]
fn iocsh_interpose_eos_config_gates_each_direction_on_its_flag() {
let mgr = Arc::new(PortManager::new());
let written = Arc::new(Mutex::new(Vec::new()));
let _ = mgr
.register_port(OctetDriver::new(
"eos_cfg",
b"line1\nline2\n",
written.clone(),
))
.unwrap();
let cmds = build_asyn_commands(mgr.clone());
let ctx = make_ctx();
cmds.iter()
.find(|c| c.name == "asynInterposeEosConfig")
.expect("asynInterposeEosConfig must be registered")
.handler
.call(
&[
ArgValue::String("eos_cfg".into()),
ArgValue::Int(0),
ArgValue::Int(1), ArgValue::Int(0), ],
&ctx,
)
.unwrap();
let handle = mgr.find_port_handle("eos_cfg").unwrap();
handle
.set_input_eos_blocking(shell_eos_user(0), b"\n")
.unwrap();
handle
.set_output_eos_blocking(shell_eos_user(0), b"\n")
.unwrap();
let user = AsynUser::default()
.with_addr(0)
.with_timeout(SHELL_IO_TIMEOUT);
let first = handle
.submit_blocking(crate::request::RequestOp::OctetRead { buf_size: 32 }, user)
.unwrap();
assert_eq!(first.data.as_deref(), Some(&b"line1"[..]));
let user = AsynUser::default()
.with_addr(0)
.with_timeout(SHELL_IO_TIMEOUT);
handle
.submit_blocking(
crate::request::RequestOp::OctetWrite {
data: b"CMD".to_vec(),
},
user,
)
.unwrap();
assert_eq!(written.lock().unwrap().as_slice(), b"CMD");
}
#[test]
fn iocsh_interpose_flush_config_installs_the_layer() {
let mgr = Arc::new(PortManager::new());
let written = Arc::new(Mutex::new(Vec::new()));
let _ = mgr
.register_port(OctetDriver::new("flush_cfg", b"stale", written.clone()))
.unwrap();
let cmds = build_asyn_commands(mgr.clone());
let ctx = make_ctx();
cmds.iter()
.find(|c| c.name == "asynInterposeFlushConfig")
.expect("asynInterposeFlushConfig must be registered")
.handler
.call(
&[
ArgValue::String("flush_cfg".into()),
ArgValue::Int(0),
ArgValue::Double(0.0),
],
&ctx,
)
.unwrap();
let handle = mgr.find_port_handle("flush_cfg").unwrap();
let user = AsynUser::default()
.with_addr(0)
.with_timeout(SHELL_IO_TIMEOUT);
handle
.submit_blocking(crate::request::RequestOp::Flush, user)
.unwrap();
let user = AsynUser::default()
.with_addr(0)
.with_timeout(SHELL_IO_TIMEOUT);
let after = handle
.submit_blocking(crate::request::RequestOp::OctetRead { buf_size: 16 }, user)
.unwrap();
assert_eq!(
after.nbytes, 0,
"the flush layer must have drained the driver's stale input"
);
let user = AsynUser::default()
.with_addr(0)
.with_timeout(SHELL_IO_TIMEOUT);
handle
.submit_blocking(
crate::request::RequestOp::OctetWrite {
data: b"GO".to_vec(),
},
user,
)
.unwrap();
assert_eq!(written.lock().unwrap().as_slice(), b"GO");
}
#[test]
fn iocsh_enable_and_autoconnect_follow_c_find_dp_common() {
let mgr = Arc::new(PortManager::new());
let _ = mgr.register_port(DummyDriver::new("edp_single")).unwrap();
let _ = mgr
.register_port(DummyDriver::multi_device("edp_multi", 2))
.unwrap();
let seen: Arc<Mutex<Vec<(String, AsynException, i32)>>> = Arc::new(Mutex::new(Vec::new()));
let seen_cb = seen.clone();
mgr.exception_manager().add_callback(move |ev| {
seen_cb
.lock()
.unwrap()
.push((ev.port_name.clone(), ev.exception, ev.addr));
});
let cmds = build_asyn_commands(mgr.clone());
let ctx = make_ctx();
let call = |name: &str, port: &str, addr: i64, yes: i64| {
cmds.iter()
.find(|c| c.name == name)
.unwrap_or_else(|| panic!("{name} not registered"))
.handler
.call(
&[
ArgValue::String(port.into()),
ArgValue::Int(addr),
ArgValue::Int(yes),
],
&ctx,
)
.unwrap();
};
let single = mgr.find_port_handle("edp_single").unwrap();
let multi = mgr.find_port_handle("edp_multi").unwrap();
call("asynEnable", "edp_single", 0, 0);
assert!(!single.is_enabled_blocking().unwrap());
call("asynEnable", "edp_multi", 1, 0);
assert!(multi.is_enabled_blocking().unwrap());
assert!(
seen.lock()
.unwrap()
.contains(&("edp_multi".to_string(), AsynException::Enable, 1)),
"the device-level enable must announce at its addr"
);
call("asynEnable", "edp_multi", -1, 0);
assert!(!multi.is_enabled_blocking().unwrap());
call("asynAutoConnect", "edp_multi", 1, 0);
assert!(
multi.is_auto_connect_blocking().unwrap(),
"a device-addressed autoConnect must not touch the port's flag"
);
assert!(seen.lock().unwrap().contains(&(
"edp_multi".to_string(),
AsynException::AutoConnect,
1
)));
call("asynAutoConnect", "edp_multi", -1, 0);
assert!(!multi.is_auto_connect_blocking().unwrap());
}
#[test]
fn raw_from_escaped_decodes_c_escapes() {
assert_eq!(raw_from_escaped(r"\r\n"), vec![b'\r', b'\n']);
assert_eq!(raw_from_escaped(r"\t"), vec![b'\t']);
assert_eq!(raw_from_escaped(r"\\"), vec![b'\\']);
assert_eq!(raw_from_escaped("AB"), vec![b'A', b'B']);
assert_eq!(raw_from_escaped(r"\x41"), vec![0x41]);
assert_eq!(raw_from_escaped(r"\x4"), vec![0x04]);
assert_eq!(raw_from_escaped(r"\z"), vec![b'z']);
assert_eq!(raw_from_escaped(r"\0"), vec![0]);
assert_eq!(raw_from_escaped(r"a\"), vec![b'a']);
assert_eq!(raw_from_escaped(r"\xg"), vec![b'g']);
assert!(raw_from_escaped("").is_empty());
}
#[test]
fn the_eos_shell_commands_carry_cs_two_second_io_timeout() {
let user = shell_eos_user(0);
assert_eq!(
user.timeout,
Duration::from_secs(2),
"C asynSetEos sets pasynUser->timeout = 2 (asynShellCommands.c:239)"
);
assert_eq!(
user.reason,
crate::user::ASYN_REASON_QUEUE_EVEN_IF_NOT_CONNECTED,
"C asynSetEos stamps ASYN_REASON_QUEUE_EVEN_IF_NOT_CONNECTED \
(asynShellCommands.c:241)"
);
assert_eq!(
user.timeout, SHELL_IO_TIMEOUT,
"the EOS user takes its deadline from the shared shell constant"
);
}
#[test]
fn iocsh_set_input_output_eos_routes_decoded_bytes_to_driver() {
#[derive(Clone, Default)]
struct Recorded {
input: Arc<Mutex<Option<Vec<u8>>>>,
output: Arc<Mutex<Option<Vec<u8>>>>,
}
struct RecordingDriver {
base: PortDriverBase,
rec: Recorded,
}
impl PortDriver for RecordingDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn set_input_eos(&mut self, _user: &AsynUser, eos: &[u8]) -> AsynResult<()> {
*self.rec.input.lock().unwrap() = Some(eos.to_vec());
Ok(())
}
fn set_output_eos(&mut self, _user: &AsynUser, eos: &[u8]) -> AsynResult<()> {
*self.rec.output.lock().unwrap() = Some(eos.to_vec());
Ok(())
}
}
let rec = Recorded::default();
let mgr = Arc::new(PortManager::new());
mgr.register_port(RecordingDriver {
base: PortDriverBase::new("eos_port", 1, PortFlags::default()),
rec: rec.clone(),
})
.unwrap();
let ctx = make_ctx();
let cmds = build_asyn_commands(mgr.clone());
let set_in = cmds
.iter()
.find(|c| c.name == "asynOctetSetInputEos")
.expect("asynOctetSetInputEos must be registered");
let outcome = set_in
.handler
.call(
&[
ArgValue::String("eos_port".to_string()),
ArgValue::Int(0),
ArgValue::String(r"\r\n".to_string()),
],
&ctx,
)
.expect("handler returns Ok");
assert!(matches!(outcome, CommandOutcome::Continue));
assert_eq!(
rec.input.lock().unwrap().as_deref(),
Some(&[b'\r', b'\n'][..]),
"input EOS must reach the driver as decoded CR LF"
);
let set_out = cmds
.iter()
.find(|c| c.name == "asynOctetSetOutputEos")
.expect("asynOctetSetOutputEos must be registered");
set_out
.handler
.call(
&[
ArgValue::String("eos_port".to_string()),
ArgValue::Int(0),
ArgValue::String(r"\n".to_string()),
],
&ctx,
)
.expect("handler returns Ok");
assert_eq!(
rec.output.lock().unwrap().as_deref(),
Some(&[b'\n'][..]),
"output EOS must reach the driver as decoded LF"
);
}
#[test]
fn iocsh_eos_commands_address_the_device_the_addr_names() {
struct MultiPort {
base: PortDriverBase,
}
impl PortDriver for MultiPort {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
}
let mgr = Arc::new(PortManager::new());
mgr.register_port(MultiPort {
base: PortDriverBase::new(
"eos_multi",
4,
PortFlags {
multi_device: true,
..PortFlags::default()
},
),
})
.unwrap();
let ctx = make_ctx();
let cmds = build_asyn_commands(mgr.clone());
let run = |name: &str, args: Vec<ArgValue>| {
cmds.iter()
.find(|c| c.name == name)
.unwrap_or_else(|| panic!("{name} must be registered"))
.handler
.call(&args, &ctx)
.expect("handler returns Ok");
};
run(
"asynOctetSetInputEos",
vec![
ArgValue::String("eos_multi".into()),
ArgValue::Int(1),
ArgValue::String(r"\r\n".into()),
],
);
run(
"asynOctetSetInputEos",
vec![
ArgValue::String("eos_multi".into()),
ArgValue::Int(2),
ArgValue::String(r"\n".into()),
],
);
let handle = mgr.find_port_handle("eos_multi").unwrap();
let eos_at = |addr: i32| {
handle
.get_input_eos_blocking(shell_eos_user(addr))
.expect("readback")
};
assert_eq!(eos_at(1), b"\r\n");
assert_eq!(eos_at(2), b"\n");
assert!(
eos_at(3).is_empty(),
"a device nobody configured has no terminator"
);
run(
"asynOctetGetInputEos",
vec![ArgValue::String("eos_multi".into()), ArgValue::Int(1)],
);
run(
"asynOctetGetOutputEos",
vec![ArgValue::String("eos_multi".into()), ArgValue::Int(2)],
);
}
#[test]
fn escaped_from_raw_is_the_inverse_of_raw_from_escaped() {
let n = SHOW_EOS_BUF_SIZE;
assert_eq!(escaped_from_raw(b"\r\n", n), r"\r\n");
assert_eq!(escaped_from_raw(b"\x01", n), r"\x01");
assert_eq!(escaped_from_raw(b"ab", n), "ab");
assert_eq!(escaped_from_raw(b"", n), "");
for s in [r"\r\n", r"\t", r"\x1b", "ab"] {
assert_eq!(
escaped_from_raw(&raw_from_escaped(s), n),
s,
"escape/unescape must round-trip"
);
}
assert_eq!(escaped_from_raw(b"\0", n), r"\0");
assert_eq!(raw_from_escaped(r"\0"), vec![0u8]); assert_eq!(raw_from_escaped(r"\x00"), vec![0u8]); assert_eq!(escaped_from_raw(b"\x001", n), r"\01"); assert_eq!(raw_from_escaped(r"\01"), vec![0u8, b'1']); assert_eq!(escaped_from_raw(&[0x01; 10], n).len(), 40);
}
#[test]
fn iocsh_interpose_echo_installs_the_layer_on_a_registered_port() {
use crate::interpose::{EomReason, OctetNext, OctetReadResult};
use crate::request::RequestOp;
use std::collections::VecDeque;
struct EchoingLink {
sizes: Arc<Mutex<Vec<usize>>>,
echo: VecDeque<u8>,
}
impl OctetNext for EchoingLink {
fn read(&mut self, _user: &AsynUser, buf: &mut [u8]) -> AsynResult<OctetReadResult> {
match self.echo.pop_front() {
Some(b) => {
buf[0] = b;
Ok(OctetReadResult {
nbytes_transferred: 1,
eom_reason: EomReason::CNT,
})
}
None => Ok(OctetReadResult {
nbytes_transferred: 0,
eom_reason: EomReason::empty(),
}),
}
}
fn write(&mut self, _user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
self.sizes.lock().unwrap().push(data.len());
self.echo.extend(data.iter().copied());
Ok(data.len())
}
fn flush(&mut self, _user: &mut AsynUser) -> AsynResult<()> {
Ok(())
}
}
struct InterposedDriver {
base: PortDriverBase,
link: EchoingLink,
}
impl PortDriver for InterposedDriver {
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]) -> AsynResult<usize> {
self.link.write(user, data)
}
fn io_read_octet_eom(
&mut self,
user: &AsynUser,
buf: &mut [u8],
) -> AsynResult<(usize, EomReason)> {
let r = self.link.read(user, buf)?;
Ok((r.nbytes_transferred, r.eom_reason))
}
}
let sizes = Arc::new(Mutex::new(Vec::new()));
let mgr = Arc::new(PortManager::new());
mgr.register_port(InterposedDriver {
base: PortDriverBase::new("echo_port", 1, PortFlags::default()),
link: EchoingLink {
sizes: sizes.clone(),
echo: VecDeque::new(),
},
})
.unwrap();
let handle = mgr.find_port_handle("echo_port").unwrap();
handle
.submit_blocking(
RequestOp::OctetWrite {
data: b"AB".to_vec(),
},
AsynUser::default(),
)
.expect("plain write succeeds");
assert_eq!(sizes.lock().unwrap().as_slice(), &[2]);
let ctx = make_ctx();
let cmds = build_asyn_commands(mgr.clone());
let echo_cmd = cmds
.iter()
.find(|c| c.name == "asynInterposeEcho")
.expect("asynInterposeEcho must be registered");
echo_cmd
.handler
.call(&[ArgValue::String("echo_port".to_string())], &ctx)
.expect("handler returns Ok");
handle
.submit_blocking(
RequestOp::OctetWrite {
data: b"AB".to_vec(),
},
AsynUser::default(),
)
.expect("echoed write succeeds");
assert_eq!(
sizes.lock().unwrap().as_slice(),
&[2, 1, 1],
"asynInterposeEcho must install the echo layer on the live port"
);
echo_cmd
.handler
.call(&[ArgValue::String("no_such_port".to_string())], &ctx)
.expect("unknown port is reported, not an Err");
}
#[test]
fn iocsh_interpose_delay_installs_the_layer_on_a_registered_port() {
use crate::interpose::{EomReason, OctetNext, OctetReadResult};
use crate::request::RequestOp;
struct CountingLink {
sizes: Arc<Mutex<Vec<usize>>>,
}
impl OctetNext for CountingLink {
fn read(&mut self, _user: &AsynUser, _buf: &mut [u8]) -> AsynResult<OctetReadResult> {
Ok(OctetReadResult {
nbytes_transferred: 0,
eom_reason: EomReason::empty(),
})
}
fn write(&mut self, _user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
self.sizes.lock().unwrap().push(data.len());
Ok(data.len())
}
fn flush(&mut self, _user: &mut AsynUser) -> AsynResult<()> {
Ok(())
}
}
struct InterposedDriver {
base: PortDriverBase,
link: CountingLink,
}
impl PortDriver for InterposedDriver {
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]) -> AsynResult<usize> {
self.link.write(user, data)
}
}
let sizes = Arc::new(Mutex::new(Vec::new()));
let mgr = Arc::new(PortManager::new());
mgr.register_port(InterposedDriver {
base: PortDriverBase::new("delay_port", 1, PortFlags::default()),
link: CountingLink {
sizes: sizes.clone(),
},
})
.unwrap();
let handle = mgr.find_port_handle("delay_port").unwrap();
handle
.submit_blocking(
RequestOp::OctetWrite {
data: b"ABC".to_vec(),
},
AsynUser::default(),
)
.expect("plain write succeeds");
assert_eq!(sizes.lock().unwrap().as_slice(), &[3]);
let ctx = make_ctx();
let cmds = build_asyn_commands(mgr.clone());
let delay_cmd = cmds
.iter()
.find(|c| c.name == "asynInterposeDelay")
.expect("asynInterposeDelay must be registered");
delay_cmd
.handler
.call(
&[
ArgValue::String("delay_port".to_string()),
ArgValue::Int(0),
ArgValue::Double(0.001),
],
&ctx,
)
.expect("handler returns Ok");
handle
.submit_blocking(
RequestOp::OctetWrite {
data: b"ABC".to_vec(),
},
AsynUser::default(),
)
.expect("delayed write succeeds");
assert_eq!(
sizes.lock().unwrap().as_slice(),
&[3, 1, 1, 1],
"asynInterposeDelay must install the delay layer on the live port"
);
delay_cmd
.handler
.call(
&[
ArgValue::String("no_such_port".to_string()),
ArgValue::Int(0),
ArgValue::Double(0.001),
],
&ctx,
)
.expect("unknown port is reported, not an Err");
}
#[test]
fn report_request_op_invokes_driver_report() -> AsynResult<()> {
let mgr = fresh_mgr_with_port("report_port");
let handle = mgr.find_port_handle("report_port")?;
handle.report_blocking(0)?;
handle.report_blocking(2)?;
Ok(())
}
fn make_ctx() -> CommandContext {
use epics_base_rs::server::database::PvDatabase;
let rt = tokio::runtime::Runtime::new().unwrap();
let db = Arc::new(PvDatabase::new());
let bridge = {
let _guard = rt.enter();
epics_base_rs::runtime::task::BlockingBridge::capture()
};
let ctx = CommandContext::new(db, bridge);
std::mem::forget(rt);
ctx
}
#[test]
fn iocsh_trace_setters_route_addr_to_device_announce() {
let mgr = fresh_mgr_with_port("trace_dev_port");
let ctx = make_ctx();
let observed: Arc<Mutex<Vec<(AsynException, i32)>>> = Arc::new(Mutex::new(Vec::new()));
let obs = observed.clone();
mgr.exception_manager().add_callback(move |ev| {
obs.lock().unwrap().push((ev.exception, ev.addr));
});
let cmds = build_asyn_commands(mgr.clone());
let io_cmd = cmds
.iter()
.find(|c| c.name == "asynSetTraceIOMask")
.expect("asynSetTraceIOMask must be registered");
let _ = io_cmd.handler.call(
&[
ArgValue::String("trace_dev_port".to_string()),
ArgValue::Int(5),
ArgValue::String("HEX".to_string()),
],
&ctx,
);
let info_cmd = cmds
.iter()
.find(|c| c.name == "asynSetTraceInfoMask")
.expect("asynSetTraceInfoMask must be registered");
let _ = info_cmd.handler.call(
&[
ArgValue::String("trace_dev_port".to_string()),
ArgValue::Int(7),
ArgValue::String("SOURCE".to_string()),
],
&ctx,
);
let file_cmd = cmds
.iter()
.find(|c| c.name == "asynSetTraceFile")
.expect("asynSetTraceFile must be registered");
let _ = file_cmd.handler.call(
&[
ArgValue::String("trace_dev_port".to_string()),
ArgValue::Int(2),
ArgValue::String("stderr".to_string()),
],
&ctx,
);
let evs = observed.lock().unwrap();
assert!(
evs.iter()
.any(|(e, a)| matches!(e, AsynException::TraceIoMask) && *a == 5),
"asynSetTraceIOMask addr=5 must fire a device-scoped announce; observed: {evs:?}"
);
assert!(
evs.iter()
.any(|(e, a)| matches!(e, AsynException::TraceInfoMask) && *a == 7),
"asynSetTraceInfoMask addr=7 must fire a device-scoped announce; observed: {evs:?}"
);
assert!(
evs.iter()
.any(|(e, a)| matches!(e, AsynException::TraceFile) && *a == 2),
"asynSetTraceFile addr=2 must fire a device-scoped announce; observed: {evs:?}"
);
}
#[test]
fn drv_asyn_ip_port_configure_registers_port() {
let cmd =
drv_asyn_ip_port_configure_command(PortServices::new(Arc::new(TraceManager::new())));
assert_eq!(cmd.name, "drvAsynIPPortConfigure");
assert_eq!(cmd.args.len(), 5);
let ctx = make_ctx();
let result = cmd.handler.call(
&[
ArgValue::String("iocsh_ip_cfg_test".into()),
ArgValue::String("127.0.0.1:9001".into()),
],
&ctx,
);
assert!(result.is_ok(), "command failed: {:?}", result.err());
assert!(
crate::asyn_record::get_port("iocsh_ip_cfg_test").is_some(),
"port must be resolvable via the asyn_record registry"
);
}
#[test]
fn build_configured_ip_port_installs_eos_unless_suppressed() {
let default_port =
build_configured_ip_port("ip_eos_default", "127.0.0.1:9100", false, false).unwrap();
assert_eq!(
default_port.base().interpose_octet.len(),
1,
"default IP port must auto-install the EOS interpose"
);
let suppressed =
build_configured_ip_port("ip_eos_off", "127.0.0.1:9100", false, true).unwrap();
assert_eq!(
suppressed.base().interpose_octet.len(),
0,
"noProcessEos must suppress the EOS interpose"
);
}
#[test]
fn drv_asyn_serial_port_configure_registers_port() {
let cmd = drv_asyn_serial_port_configure_command(PortServices::new(Arc::new(
TraceManager::new(),
)));
assert_eq!(cmd.name, "drvAsynSerialPortConfigure");
assert_eq!(cmd.args.len(), 5);
let ctx = make_ctx();
let result = cmd.handler.call(
&[
ArgValue::String("iocsh_serial_cfg_test".into()),
ArgValue::String("/dev/ttyS0".into()),
],
&ctx,
);
assert!(result.is_ok(), "command failed: {:?}", result.err());
assert!(
crate::asyn_record::get_port("iocsh_serial_cfg_test").is_some(),
"port must be resolvable via the asyn_record registry"
);
}
#[test]
fn drv_asyn_ip_port_configure_rejects_missing_host() {
let cmd =
drv_asyn_ip_port_configure_command(PortServices::new(Arc::new(TraceManager::new())));
let ctx = make_ctx();
let result = cmd
.handler
.call(&[ArgValue::String("iocsh_ip_cfg_nohost".into())], &ctx);
assert!(result.is_err());
assert!(crate::asyn_record::get_port("iocsh_ip_cfg_nohost").is_none());
}
#[test]
fn drv_asyn_prologix_port_configure_registers_port() {
let cmd = drv_asyn_prologix_port_configure_command(PortServices::new(Arc::new(
TraceManager::new(),
)));
assert_eq!(cmd.name, "prologixGPIBConfigure");
assert_eq!(cmd.args.len(), 4);
let ctx = make_ctx();
let result = cmd.handler.call(
&[
ArgValue::String("iocsh_prologix_cfg_test".into()),
ArgValue::String("127.0.0.1:1234".into()),
],
&ctx,
);
assert!(result.is_ok(), "command failed: {:?}", result.err());
assert!(
crate::asyn_record::get_port("iocsh_prologix_cfg_test").is_some(),
"port must be resolvable via the asyn_record registry"
);
}
#[test]
fn drv_asyn_prologix_port_configure_rejects_missing_host() {
let cmd = drv_asyn_prologix_port_configure_command(PortServices::new(Arc::new(
TraceManager::new(),
)));
let ctx = make_ctx();
let result = cmd.handler.call(
&[ArgValue::String("iocsh_prologix_cfg_nohost".into())],
&ctx,
);
assert!(result.is_err());
assert!(crate::asyn_record::get_port("iocsh_prologix_cfg_nohost").is_none());
}
#[test]
fn iocsh_ftdi_port_configure_creates_the_port_an_st_cmd_names() {
let cmd =
drv_asyn_ftdi_port_configure_command(PortServices::new(Arc::new(TraceManager::new())));
assert_eq!(cmd.name, "drvAsynFTDIPortConfigure");
assert_eq!(cmd.args.len(), 9);
let ctx = make_ctx();
let result = cmd.handler.call(
&[
ArgValue::String("iocsh_ftdi_cfg_test".into()),
ArgValue::Int(0x0403),
ArgValue::Int(0x6001),
ArgValue::Int(9600),
ArgValue::Int(1),
ArgValue::Int(0),
ArgValue::Int(1), ],
&ctx,
);
assert!(result.is_ok(), "command failed: {:?}", result.err());
assert!(
crate::asyn_record::get_port("iocsh_ftdi_cfg_test").is_some(),
"port must be resolvable via the asyn_record registry"
);
}
#[test]
fn iocsh_vxi11_configure_creates_the_port_an_st_cmd_names() {
let cmd = vxi11_configure_command(PortServices::new(Arc::new(TraceManager::new())));
assert_eq!(cmd.name, "vxi11Configure");
assert_eq!(cmd.args.len(), 7);
let ctx = make_ctx();
let result = cmd.handler.call(
&[
ArgValue::String("iocsh_vxi11_cfg_test".into()),
ArgValue::String("127.0.0.1".into()),
ArgValue::Int(0),
ArgValue::String("1.0".into()),
ArgValue::String("inst0".into()),
ArgValue::Int(0),
ArgValue::Int(1), ],
&ctx,
);
assert!(result.is_ok(), "command failed: {:?}", result.err());
assert!(
crate::asyn_record::get_port("iocsh_vxi11_cfg_test").is_some(),
"port must be resolvable via the asyn_record registry"
);
}
#[test]
fn iocsh_vxi11_configure_rejects_missing_host() {
let cmd = vxi11_configure_command(PortServices::new(Arc::new(TraceManager::new())));
let ctx = make_ctx();
let result = cmd
.handler
.call(&[ArgValue::String("iocsh_vxi11_cfg_nohost".into())], &ctx);
assert!(result.is_err());
assert!(crate::asyn_record::get_port("iocsh_vxi11_cfg_nohost").is_none());
}
#[test]
fn iocsh_usbtmc_configure_creates_the_port_an_st_cmd_names() {
let cmd = usbtmc_configure_command(PortServices::new(Arc::new(TraceManager::new())));
assert_eq!(cmd.name, "usbtmcConfigure");
assert_eq!(cmd.args.len(), 6);
let ctx = make_ctx();
let result = cmd.handler.call(
&[
ArgValue::String("iocsh_usbtmc_cfg_test".into()),
ArgValue::Int(0x0699),
ArgValue::Int(0x0368),
ArgValue::String(String::new()),
ArgValue::Int(0),
ArgValue::Int(0),
],
&ctx,
);
assert!(result.is_ok(), "command failed: {:?}", result.err());
assert!(
crate::asyn_record::get_port("iocsh_usbtmc_cfg_test").is_some(),
"port must be resolvable via the asyn_record registry"
);
}
}