use std::os::unix::io::RawFd;
use std::time::{Duration, Instant};
use crate::error::{AsynError, AsynResult, AsynStatus};
use crate::exception::AsynException;
use crate::interpose::{EomReason, OctetNext, OctetReadResult};
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::trace::TraceMask;
use crate::user::AsynUser;
use crate::{asyn_trace, asyn_trace_io};
use super::option_parse::{bad_number, parse_yn_option, sscanf_int, sscanf_uint};
use super::serial_config::{DataBits, FlowControl, Parity, SerialConfig, StopBits};
mod platform {
#[cfg(not(any(target_os = "vxworks", target_os = "rtems")))]
mod imp {
pub use libc::termios;
pub type Speed = libc::speed_t;
pub use libc::{
B300, B9600, CLOCAL, CREAD, CS5, CS6, CS7, CS8, CSIZE, CSTOPB, IGNBRK, IGNPAR, IXOFF,
IXON, PARENB, PARODD, TCIFLUSH, TCSANOW, VMIN, VTIME,
};
pub use libc::{
cfmakeraw, cfsetispeed, cfsetospeed, tcflush, tcgetattr, tcsendbreak, tcsetattr,
};
pub const CRTSCTS: libc::tcflag_t = libc::CRTSCTS;
pub const O_NOCTTY: libc::c_int = libc::O_NOCTTY;
pub const FLUSH_IO: libc::c_int = libc::TCIOFLUSH;
pub const IXANY: Option<libc::tcflag_t> = Some(libc::IXANY);
pub const SOFT_FLOW_CHARS: Option<(usize, usize)> = Some((libc::VSTART, libc::VSTOP));
pub fn drain(fd: libc::c_int) -> Option<std::io::Result<()>> {
Some(if unsafe { libc::tcdrain(fd) } < 0 {
Err(std::io::Error::last_os_error())
} else {
Ok(())
})
}
}
#[cfg(target_os = "rtems")]
mod imp {
pub type Speed = libc::c_uint;
pub const NCCS: usize = 20;
#[repr(C)]
#[derive(Clone, Copy)]
pub struct termios {
pub c_iflag: libc::c_uint,
pub c_oflag: libc::c_uint,
pub c_cflag: libc::c_uint,
pub c_lflag: libc::c_uint,
pub c_cc: [libc::c_uchar; NCCS],
pub c_ispeed: Speed,
pub c_ospeed: Speed,
}
const _: () = {
use std::mem::{align_of, offset_of, size_of};
assert!(size_of::<libc::c_uint>() == 4, "tcflag_t is `unsigned int`");
assert!(size_of::<libc::c_uchar>() == 1, "cc_t is `unsigned char`");
assert!(size_of::<Speed>() == 4, "speed_t is `unsigned int`");
assert!(offset_of!(termios, c_iflag) == 0, "c_iflag at 0");
assert!(offset_of!(termios, c_oflag) == 4, "c_oflag at 4");
assert!(offset_of!(termios, c_cflag) == 8, "c_cflag at 8");
assert!(offset_of!(termios, c_lflag) == 12, "c_lflag at 12");
assert!(offset_of!(termios, c_cc) == 16, "c_cc at 16, not 17");
assert!(offset_of!(termios, c_ispeed) == 36, "c_ispeed at 36");
assert!(offset_of!(termios, c_ospeed) == 40, "c_ospeed at 40");
assert!(size_of::<termios>() == 44, "sizeof(struct termios) == 44");
assert!(align_of::<termios>() == 4, "alignof(struct termios) == 4");
assert!(
VMIN < NCCS && VTIME < NCCS && VSTART < NCCS && VSTOP < NCCS,
"every c_cc index used here must be inside the array"
);
};
unsafe extern "C" {
pub fn tcgetattr(fd: libc::c_int, t: *mut termios) -> libc::c_int;
pub fn tcsetattr(
fd: libc::c_int,
action: libc::c_int,
t: *const termios,
) -> libc::c_int;
pub fn tcflush(fd: libc::c_int, queue: libc::c_int) -> libc::c_int;
pub fn tcdrain(fd: libc::c_int) -> libc::c_int;
pub fn tcsendbreak(fd: libc::c_int, duration: libc::c_int) -> libc::c_int;
pub fn cfmakeraw(t: *mut termios);
pub fn cfsetispeed(t: *mut termios, speed: Speed) -> libc::c_int;
pub fn cfsetospeed(t: *mut termios, speed: Speed) -> libc::c_int;
}
pub const IGNBRK: libc::c_uint = 0x0000_0001;
pub const IGNPAR: libc::c_uint = 0x0000_0004;
pub const IXON: libc::c_uint = 0x0000_0200;
pub const IXOFF: libc::c_uint = 0x0000_0400;
pub const IXANY: Option<libc::c_uint> = Some(0x0000_0800);
pub const CSIZE: libc::c_uint = 0x0000_0300;
pub const CS5: libc::c_uint = 0x0000_0000;
pub const CS6: libc::c_uint = 0x0000_0100;
pub const CS7: libc::c_uint = 0x0000_0200;
pub const CS8: libc::c_uint = 0x0000_0300;
pub const CSTOPB: libc::c_uint = 0x0000_0400;
pub const CREAD: libc::c_uint = 0x0000_0800;
pub const PARENB: libc::c_uint = 0x0000_1000;
pub const PARODD: libc::c_uint = 0x0000_2000;
pub const CLOCAL: libc::c_uint = 0x0000_8000;
pub const CRTSCTS: libc::c_uint = 0x0001_0000 | 0x0002_0000;
pub const VSTART: usize = 12;
pub const VSTOP: usize = 13;
pub const VMIN: usize = 16;
pub const VTIME: usize = 17;
pub const SOFT_FLOW_CHARS: Option<(usize, usize)> = Some((VSTART, VSTOP));
pub const TCSANOW: libc::c_int = 0;
pub const TCIFLUSH: libc::c_int = 1;
pub const FLUSH_IO: libc::c_int = 3;
pub const O_NOCTTY: libc::c_int = 0x8000;
pub const B300: Speed = 300;
pub const B9600: Speed = 9600;
pub fn drain(fd: libc::c_int) -> Option<std::io::Result<()>> {
Some(if unsafe { tcdrain(fd) } < 0 {
Err(std::io::Error::last_os_error())
} else {
Ok(())
})
}
}
#[cfg(target_os = "vxworks")]
mod imp {
pub use libc::termios;
pub type Speed = libc::speed_t;
pub use libc::{
B300, B9600, CLOCAL, CREAD, CS5, CS6, CS7, CS8, CSIZE, CSTOPB, IGNBRK, IGNPAR, IXOFF,
IXON, PARENB, PARODD, TCIFLUSH, TCSANOW, VMIN, VTIME,
};
pub use libc::{
cfmakeraw, cfsetispeed, cfsetospeed, tcflush, tcgetattr, tcsendbreak, tcsetattr,
};
pub const CRTSCTS: libc::tcflag_t = 0x0001_0000 | 0x0002_0000;
pub const O_NOCTTY: libc::c_int = 0x8000;
pub const FLUSH_IO: libc::c_int = libc::TCIFLUSH;
pub const IXANY: Option<libc::tcflag_t> = None;
pub const SOFT_FLOW_CHARS: Option<(usize, usize)> = None;
pub fn drain(_fd: libc::c_int) -> Option<std::io::Result<()>> {
None
}
}
pub use imp::*;
}
fn option_unsupported_here(key: &str) -> AsynError {
AsynError::Status {
status: AsynStatus::Error,
message: format!("Option {key} not supported on {}", std::env::consts::OS),
}
}
fn set_termios_speed(t: &mut platform::termios, speed: platform::Speed) -> AsynResult<()> {
if unsafe { platform::cfsetispeed(t, speed) } < 0 {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: format!("cfsetispeed returned {}", std::io::Error::last_os_error()),
});
}
if unsafe { platform::cfsetospeed(t, speed) } < 0 {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: format!("cfsetospeed returned {}", std::io::Error::last_os_error()),
});
}
Ok(())
}
impl SerialConfig {
pub fn apply_to_termios(&self, t: &mut platform::termios) -> AsynResult<()> {
let baud = baud_to_speed(self.baud).ok_or_else(|| AsynError::Status {
status: AsynStatus::Error,
message: format!("unsupported baud rate: {}", self.baud),
})?;
set_termios_speed(t, baud)?;
t.c_cflag &= !platform::CSIZE;
t.c_cflag |= match self.data_bits {
DataBits::Five => platform::CS5,
DataBits::Six => platform::CS6,
DataBits::Seven => platform::CS7,
DataBits::Eight => platform::CS8,
};
match self.parity {
Parity::None => {
t.c_cflag &= !platform::PARENB;
}
Parity::Even => {
t.c_cflag |= platform::PARENB;
t.c_cflag &= !platform::PARODD;
}
Parity::Odd => {
t.c_cflag |= platform::PARENB;
t.c_cflag |= platform::PARODD;
}
}
match self.stop_bits {
StopBits::One => t.c_cflag &= !platform::CSTOPB,
StopBits::Two => t.c_cflag |= platform::CSTOPB,
}
let ixany = platform::IXANY.unwrap_or(0);
match self.flow_control {
FlowControl::None => {
t.c_cflag &= !platform::CRTSCTS;
t.c_iflag &= !(platform::IXON | platform::IXOFF | ixany);
}
FlowControl::Hardware => {
t.c_cflag |= platform::CRTSCTS;
t.c_iflag &= !(platform::IXON | platform::IXOFF | ixany);
}
FlowControl::Software => {
t.c_cflag &= !platform::CRTSCTS;
t.c_iflag |= platform::IXON | platform::IXOFF;
}
}
Ok(())
}
}
fn baud_to_speed(baud: u32) -> Option<platform::Speed> {
#[cfg(asyn_baud_code_is_rate)]
{
Some(platform::Speed::from(baud))
}
#[cfg(not(asyn_baud_code_is_rate))]
{
Some(match baud {
50 => libc::B50,
75 => libc::B75,
110 => libc::B110,
134 => libc::B134,
150 => libc::B150,
200 => libc::B200,
300 => libc::B300,
600 => libc::B600,
1200 => libc::B1200,
1800 => libc::B1800,
2400 => libc::B2400,
4800 => libc::B4800,
9600 => libc::B9600,
19200 => libc::B19200,
38400 => libc::B38400,
57600 => libc::B57600,
115200 => libc::B115200,
230400 => libc::B230400,
#[cfg(any(target_os = "linux", target_os = "android"))]
460800 => libc::B460800,
#[cfg(any(target_os = "linux", target_os = "android"))]
500000 => libc::B500000,
#[cfg(any(target_os = "linux", target_os = "android"))]
576000 => libc::B576000,
#[cfg(any(target_os = "linux", target_os = "android"))]
921600 => libc::B921600,
#[cfg(any(target_os = "linux", target_os = "android"))]
1000000 => libc::B1000000,
#[cfg(any(target_os = "linux", target_os = "android"))]
1152000 => libc::B1152000,
#[cfg(any(target_os = "linux", target_os = "android"))]
1500000 => libc::B1500000,
#[cfg(any(target_os = "linux", target_os = "android"))]
2000000 => libc::B2000000,
#[cfg(any(target_os = "linux", target_os = "android"))]
2500000 => libc::B2500000,
#[cfg(any(target_os = "linux", target_os = "android"))]
3000000 => libc::B3000000,
#[cfg(any(target_os = "linux", target_os = "android"))]
3500000 => libc::B3500000,
#[cfg(any(target_os = "linux", target_os = "android"))]
4000000 => libc::B4000000,
_ => return None,
})
}
}
const _: () = assert!(
(platform::B300 == 300 && platform::B9600 == 9600) == cfg!(asyn_baud_code_is_rate),
"asyn_baud_code_is_rate disagrees with this platform's Bxxx codes: add or \
remove this target in build.rs, do not leave the branch mismatched"
);
#[allow(dead_code)]
fn speed_to_baud(speed: platform::Speed) -> u32 {
#[cfg(asyn_baud_code_is_rate)]
{
u32::try_from(speed).unwrap_or(0)
}
#[cfg(not(asyn_baud_code_is_rate))]
{
match speed {
libc::B0 => 0,
libc::B50 => 50,
libc::B75 => 75,
libc::B110 => 110,
libc::B134 => 134,
libc::B150 => 150,
libc::B200 => 200,
libc::B300 => 300,
libc::B600 => 600,
libc::B1200 => 1200,
libc::B1800 => 1800,
libc::B2400 => 2400,
libc::B4800 => 4800,
libc::B9600 => 9600,
libc::B19200 => 19200,
libc::B38400 => 38400,
libc::B57600 => 57600,
libc::B115200 => 115200,
libc::B230400 => 230400,
#[cfg(any(target_os = "linux", target_os = "android"))]
libc::B460800 => 460800,
#[cfg(any(target_os = "linux", target_os = "android"))]
libc::B500000 => 500000,
#[cfg(any(target_os = "linux", target_os = "android"))]
libc::B576000 => 576000,
#[cfg(any(target_os = "linux", target_os = "android"))]
libc::B921600 => 921600,
#[cfg(any(target_os = "linux", target_os = "android"))]
libc::B1000000 => 1000000,
#[cfg(any(target_os = "linux", target_os = "android"))]
libc::B1152000 => 1152000,
#[cfg(any(target_os = "linux", target_os = "android"))]
libc::B1500000 => 1500000,
#[cfg(any(target_os = "linux", target_os = "android"))]
libc::B2000000 => 2000000,
#[cfg(any(target_os = "linux", target_os = "android"))]
libc::B2500000 => 2500000,
#[cfg(any(target_os = "linux", target_os = "android"))]
libc::B3000000 => 3000000,
#[cfg(any(target_os = "linux", target_os = "android"))]
libc::B3500000 => 3500000,
#[cfg(any(target_os = "linux", target_os = "android"))]
libc::B4000000 => 4000000,
_ => 0,
}
}
}
struct SerialIoState {
fd: Option<RawFd>,
n_read: u64,
n_written: u64,
}
impl SerialIoState {
fn new() -> Self {
Self {
fd: None,
n_read: 0,
n_written: 0,
}
}
fn fd_or_err(&self) -> AsynResult<RawFd> {
self.fd.ok_or_else(|| AsynError::Status {
status: AsynStatus::Disconnected,
message: "serial port not open".into(),
})
}
}
fn duration_to_poll_ms(d: Duration) -> i32 {
d.as_millis().min(i32::MAX as u128) as i32
}
impl OctetNext for SerialIoState {
fn read(&mut self, user: &AsynUser, buf: &mut [u8]) -> AsynResult<OctetReadResult> {
let fd = self.fd_or_err()?;
if buf.is_empty() {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: "maxchars 0 Why <=0?".into(),
});
}
let timeout_ms = duration_to_poll_ms(user.timeout);
loop {
let mut pfd = libc::pollfd {
fd,
events: libc::POLLIN,
revents: 0,
};
let ret = unsafe { libc::poll(&mut pfd, 1, timeout_ms) };
if ret < 0 {
let err = std::io::Error::last_os_error();
if err.kind() == std::io::ErrorKind::Interrupted {
continue;
}
return Err(AsynError::Io(err));
}
if ret == 0 {
return Err(AsynError::Status {
status: AsynStatus::Timeout,
message: "serial read timeout".into(),
});
}
let n = unsafe { libc::read(fd, buf.as_mut_ptr() as *mut libc::c_void, buf.len()) };
if n < 0 {
let err = std::io::Error::last_os_error();
if err.kind() == std::io::ErrorKind::Interrupted
|| err.kind() == std::io::ErrorKind::WouldBlock
{
continue;
}
return Err(AsynError::Io(err));
}
if n == 0 {
return Err(AsynError::Status {
status: AsynStatus::Disconnected,
message: "serial port EOF".into(),
});
}
self.n_read += n as u64; return Ok(OctetReadResult {
nbytes_transferred: n as usize,
eom_reason: if n as usize >= buf.len() {
EomReason::CNT
} else {
EomReason::empty()
},
});
}
}
fn write(&mut self, user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
let fd = self.fd_or_err()?;
let deadline = Instant::now() + user.timeout;
let mut total = 0usize;
let result: AsynResult<()> = loop {
if total >= data.len() {
break Ok(());
}
let poll_ms = duration_to_poll_ms(deadline.saturating_duration_since(Instant::now()));
let mut pfd = libc::pollfd {
fd,
events: libc::POLLOUT,
revents: 0,
};
let ret = unsafe { libc::poll(&mut pfd, 1, poll_ms) };
if ret < 0 {
let err = std::io::Error::last_os_error();
if err.kind() == std::io::ErrorKind::Interrupted {
continue;
}
break Err(AsynError::Io(err));
}
if ret == 0 {
break Err(AsynError::Status {
status: AsynStatus::Timeout,
message: "serial write timeout".into(),
});
}
let n = unsafe {
libc::write(
fd,
data[total..].as_ptr() as *const libc::c_void,
data.len() - total,
)
};
if n < 0 {
let err = std::io::Error::last_os_error();
if err.kind() == std::io::ErrorKind::Interrupted
|| err.kind() == std::io::ErrorKind::WouldBlock
{
continue;
}
break Err(AsynError::Io(err));
}
total += n as usize;
self.n_written += n as u64;
if total < data.len() && Instant::now() >= deadline {
break Err(AsynError::Status {
status: AsynStatus::Timeout,
message: "serial write timeout".into(),
});
}
};
match result {
Ok(()) => Ok(total),
Err(e) => Err(e.with_partial_write(total)),
}
}
fn flush(&mut self, _user: &mut AsynUser) -> AsynResult<()> {
if let Some(fd) = self.fd {
let ret = unsafe { platform::tcflush(fd, platform::TCIFLUSH) };
if ret < 0 {
return Err(AsynError::Io(std::io::Error::last_os_error()));
}
}
Ok(())
}
}
pub struct DrvAsynSerialPort {
base: PortDriverBase,
device: String,
baud: u32,
termios: platform::termios,
io: SerialIoState,
saved_termios: Option<platform::termios>,
}
fn seed_termios(config: &SerialConfig) -> AsynResult<platform::termios> {
let mut t: platform::termios = unsafe { std::mem::zeroed() };
unsafe { platform::cfmakeraw(&mut t) };
t.c_cflag |= platform::CREAD | platform::CLOCAL;
t.c_iflag |= platform::IGNBRK | platform::IGNPAR;
t.c_cc[platform::VMIN] = 1;
t.c_cc[platform::VTIME] = 0;
if let Some((vstart, vstop)) = platform::SOFT_FLOW_CHARS {
t.c_cc[vstart] = 0x11; t.c_cc[vstop] = 0x13; }
config.apply_to_termios(&mut t)?;
Ok(t)
}
impl DrvAsynSerialPort {
fn drop_connection(&mut self) {
if let Some(fd) = self.io.fd.take() {
unsafe { libc::close(fd) };
}
self.saved_termios = None;
self.base.set_connected(false);
}
pub fn new(port_name: &str, config_str: &str) -> AsynResult<Self> {
let config = SerialConfig::parse(config_str)?;
let mut base = PortDriverBase::new(
port_name,
1,
PortFlags {
multi_device: false,
can_block: true,
destructible: true,
},
);
base.init_connected(false);
base.auto_connect = true;
base.octet_interrupt_process = true;
Ok(Self {
base,
device: config.device.clone(),
baud: config.baud,
termios: seed_termios(&config)?,
io: SerialIoState::new(),
saved_termios: None,
})
}
pub fn configure(
port_name: &str,
config_str: &str,
no_auto_connect: bool,
no_process_eos: bool,
) -> AsynResult<Self> {
let mut driver = Self::new(port_name, config_str)?;
if no_auto_connect {
driver.base.auto_connect = false;
}
if !no_process_eos {
driver.install_interpose(Box::new(crate::interpose::eos::EosInterpose::default()));
}
Ok(driver)
}
pub fn install_interpose(&mut self, layer: Box<dyn crate::interpose::OctetInterpose>) {
self.base.install_octet_interpose(layer);
}
pub fn send_break(&self, duration_tenths: i32) -> AsynResult<()> {
let fd = self.io.fd_or_err()?;
let ret = unsafe { platform::tcsendbreak(fd, duration_tenths) };
if ret < 0 {
return Err(AsynError::Io(std::io::Error::last_os_error()));
}
Ok(())
}
pub fn drain_output(&self) -> AsynResult<()> {
let fd = self.io.fd_or_err()?;
match platform::drain(fd) {
Some(Ok(())) => Ok(()),
Some(Err(e)) => Err(AsynError::Io(e)),
None => Err(option_unsupported_here("drain")),
}
}
fn get_current_termios(&self) -> AsynResult<platform::termios> {
let fd = self.io.fd_or_err()?;
let mut t: platform::termios = unsafe { std::mem::zeroed() };
let ret = unsafe { platform::tcgetattr(fd, &mut t) };
if ret < 0 {
return Err(AsynError::Io(std::io::Error::last_os_error()));
}
Ok(t)
}
fn apply_termios(&self, t: &platform::termios) -> AsynResult<()> {
let fd = self.io.fd_or_err()?;
let ret = unsafe { platform::tcsetattr(fd, platform::TCSANOW, t) };
if ret < 0 {
return Err(AsynError::Io(std::io::Error::last_os_error()));
}
Ok(())
}
fn apply_options(&mut self) -> AsynResult<()> {
self.termios.c_cflag |= platform::CREAD;
let t = self.termios;
self.apply_termios(&t)
}
}
impl PortDriver for DrvAsynSerialPort {
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::octet_transport_capabilities()
}
fn connect(&mut self, _user: &AsynUser) -> AsynResult<()> {
if self.io.fd.is_some() {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: format!("{}: Link already open!", self.base.port_name),
});
}
let c_path =
std::ffi::CString::new(self.device.as_str()).map_err(|_| AsynError::Status {
status: AsynStatus::Error,
message: "invalid device path (contains NUL)".into(),
})?;
let fd = unsafe {
libc::open(
c_path.as_ptr(),
libc::O_RDWR | platform::O_NOCTTY | libc::O_NONBLOCK,
)
};
if fd < 0 {
return Err(AsynError::Io(std::io::Error::last_os_error()));
}
self.io.fd = Some(fd);
let setup = (|| -> AsynResult<()> {
if unsafe { libc::fcntl(fd, libc::F_SETFD, libc::FD_CLOEXEC) } < 0 {
return Err(AsynError::Io(std::io::Error::last_os_error()));
}
let saved = self.get_current_termios()?;
self.saved_termios = Some(saved);
self.apply_options()?;
unsafe { platform::tcflush(fd, platform::FLUSH_IO) };
Ok(())
})();
if let Err(e) = setup {
if let Some(fd) = self.io.fd.take() {
unsafe { libc::close(fd) };
}
self.saved_termios = None;
return Err(e);
}
self.base.set_connected(true);
asyn_trace!(
Some(self.base.trace),
&self.base.port_name,
TraceMask::FLOW,
"connected to {} at {} baud",
self.device,
self.baud
);
Ok(())
}
fn disconnect(&mut self, _user: &AsynUser) -> AsynResult<()> {
asyn_trace!(
Some(self.base.trace),
&self.base.port_name,
TraceMask::FLOW,
"disconnect"
);
if let (Some(fd), Some(saved)) = (self.io.fd, &self.saved_termios) {
unsafe { platform::tcsetattr(fd, platform::TCSANOW, saved) };
}
if let Some(fd) = self.io.fd.take() {
unsafe { libc::close(fd) };
}
self.saved_termios = None;
self.base.set_connected(false);
Ok(())
}
fn report(&self, out: &mut dyn std::fmt::Write, level: i32) {
use std::fmt::Write as _;
let _ = writeln!(
out,
"Serial line {}: {}",
self.device,
if self.base.is_connected() {
"Connected"
} else {
"Disconnected"
}
);
if level >= 1 {
let _ = writeln!(out, " fd: {}", self.io.fd.unwrap_or(-1));
let _ = writeln!(out, " Characters written: {}", self.io.n_written);
let _ = writeln!(out, " Characters read: {}", self.io.n_read);
self.base.report_params(out, level);
}
}
fn read_octet(&mut self, user: &AsynUser, buf: &mut [u8]) -> AsynResult<usize> {
self.io_read_octet_eom(user, buf).map(|(n, _eom)| n)
}
fn io_read_octet_eom(
&mut self,
user: &AsynUser,
buf: &mut [u8],
) -> AsynResult<(usize, EomReason)> {
self.base.check_ready()?;
let result = match self.io.read(user, buf) {
Ok(r) => r,
Err(e) => {
if e.is_fatal_transport() && self.base.is_connected() {
asyn_trace!(
Some(self.base.trace),
&self.base.port_name,
TraceMask::FLOW,
"read error, disconnecting: {e}"
);
self.drop_connection();
}
return Err(e);
}
};
asyn_trace_io!(
Some(self.base.trace),
&self.base.port_name,
TraceMask::IO_DRIVER,
&buf[..result.nbytes_transferred],
"read"
);
Ok((result.nbytes_transferred, result.eom_reason))
}
fn write_octet(&mut self, user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
self.base.check_ready()?;
asyn_trace_io!(
Some(self.base.trace),
&self.base.port_name,
TraceMask::IO_DRIVER,
data,
"write"
);
match self.io.write(user, data) {
Ok(n) => Ok(n),
Err(e) => {
if e.is_fatal_transport() && self.base.is_connected() {
asyn_trace!(
Some(self.base.trace),
&self.base.port_name,
TraceMask::FLOW,
"write error, disconnecting: {e}"
);
self.drop_connection();
}
Err(e)
}
}
}
fn io_flush(&mut self, user: &mut AsynUser) -> AsynResult<()> {
self.io.flush(user)
}
fn set_option(&mut self, _user: &mut AsynUser, key: &str, value: &str) -> AsynResult<()> {
let key = key.trim().to_ascii_lowercase();
let value = value.trim();
let baud_prev = self.baud;
let termios_prev = self.termios;
match key.as_str() {
"baud" => {
let baud = sscanf_int(value).ok_or_else(bad_number)?;
let speed = u32::try_from(baud)
.ok()
.and_then(baud_to_speed)
.ok_or_else(|| AsynError::Status {
status: AsynStatus::Error,
message: format!("Unsupported data rate ({baud} baud)"),
})?;
set_termios_speed(&mut self.termios, speed)?;
self.baud = baud as u32;
}
"bits" => {
let bits = match value {
"5" => DataBits::Five,
"6" => DataBits::Six,
"7" => DataBits::Seven,
"8" => DataBits::Eight,
_ => {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: "Invalid number of bits.".into(),
});
}
};
self.termios.c_cflag &= !platform::CSIZE;
self.termios.c_cflag |= match bits {
DataBits::Five => platform::CS5,
DataBits::Six => platform::CS6,
DataBits::Seven => platform::CS7,
DataBits::Eight => platform::CS8,
};
}
"parity" => {
let val_lower = value.to_ascii_lowercase();
match val_lower.as_str() {
"none" => self.termios.c_cflag &= !platform::PARENB,
"even" => {
self.termios.c_cflag |= platform::PARENB;
self.termios.c_cflag &= !platform::PARODD;
}
"odd" => {
self.termios.c_cflag |= platform::PARENB;
self.termios.c_cflag |= platform::PARODD;
}
_ => {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: "Invalid parity.".into(),
});
}
}
}
"stop" => match value {
"1" => self.termios.c_cflag &= !platform::CSTOPB,
"2" => self.termios.c_cflag |= platform::CSTOPB,
_ => {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: "Invalid number of stop bits.".into(),
});
}
},
"clocal" => {
if parse_yn_option(&key, value)? {
self.termios.c_cflag |= platform::CLOCAL;
} else {
self.termios.c_cflag &= !platform::CLOCAL;
}
}
"crtscts" => {
if parse_yn_option(&key, value)? {
self.termios.c_cflag |= platform::CRTSCTS;
} else {
self.termios.c_cflag &= !platform::CRTSCTS;
}
}
"ixon" => {
if parse_yn_option(&key, value)? {
self.termios.c_iflag |= platform::IXON;
} else {
self.termios.c_iflag &= !platform::IXON;
}
}
"ixoff" => {
if parse_yn_option(&key, value)? {
self.termios.c_iflag |= platform::IXOFF;
} else {
self.termios.c_iflag &= !platform::IXOFF;
}
}
"ixany" => {
let bit = platform::IXANY.ok_or_else(|| option_unsupported_here("ixany"))?;
if parse_yn_option(&key, value)? {
self.termios.c_iflag |= bit;
} else {
self.termios.c_iflag &= !bit;
}
}
"break" => {
if value != "off" {
let duration = if value.is_empty() || value == "on" {
0 } else {
sscanf_uint(value).ok_or_else(bad_number)? as i32
};
let fd = self.io.fd_or_err()?;
if let Some(res) = platform::drain(fd) {
res.map_err(AsynError::Io)?;
}
let ret = unsafe { platform::tcsendbreak(fd, duration) };
if ret < 0 {
return Err(AsynError::Io(std::io::Error::last_os_error()));
}
}
}
#[cfg(target_os = "linux")]
"rs485_enable"
| "rs485_rts_on_send"
| "rs485_rts_after_send"
| "rs485_delay_rts_before_send"
| "rs485_delay_rts_after_send" => {
self.set_rs485_option(&key, value)?;
}
other => {
if !other.is_empty() {
return Err(AsynError::OptionNotFound(other.to_string()));
}
}
}
if self.io.fd.is_some() {
if let Err(e) = self.apply_options() {
self.baud = baud_prev;
self.termios = termios_prev;
return Err(e);
}
}
Ok(())
}
fn get_option(&self, key: &str) -> AsynResult<String> {
match key {
"baud" => Ok(self.baud.to_string()),
"bits" => Ok(match self.termios.c_cflag & platform::CSIZE {
platform::CS5 => "5",
platform::CS6 => "6",
platform::CS7 => "7",
platform::CS8 => "8",
_ => "?",
}
.to_string()),
"parity" => Ok(if self.termios.c_cflag & platform::PARENB == 0 {
"none"
} else if self.termios.c_cflag & platform::PARODD != 0 {
"odd"
} else {
"even"
}
.to_string()),
"stop" => Ok(if self.termios.c_cflag & platform::CSTOPB != 0 {
"2"
} else {
"1"
}
.to_string()),
"clocal" => Ok(if self.termios.c_cflag & platform::CLOCAL != 0 {
"Y"
} else {
"N"
}
.to_string()),
"crtscts" => Ok(if self.termios.c_cflag & platform::CRTSCTS != 0 {
"Y"
} else {
"N"
}
.to_string()),
"ixon" => Ok(if self.termios.c_iflag & platform::IXON != 0 {
"Y"
} else {
"N"
}
.to_string()),
"ixoff" => Ok(if self.termios.c_iflag & platform::IXOFF != 0 {
"Y"
} else {
"N"
}
.to_string()),
"ixany" => Ok(match platform::IXANY {
Some(bit) if self.termios.c_iflag & bit != 0 => "Y",
_ => "N",
}
.to_string()),
"break" => Ok("off".to_string()),
#[cfg(target_os = "linux")]
"rs485_enable"
| "rs485_rts_on_send"
| "rs485_rts_after_send"
| "rs485_delay_rts_before_send"
| "rs485_delay_rts_after_send" => self.get_rs485_option(key),
_ => self
.base
.options
.get(key)
.cloned()
.ok_or_else(|| AsynError::OptionNotFound(key.to_string())),
}
}
}
#[cfg(target_os = "linux")]
#[repr(C)]
#[derive(Clone, Copy, Default)]
struct SerialRs485 {
flags: u32,
delay_rts_before_send: u32,
delay_rts_after_send: u32,
padding: [u32; 5],
}
#[cfg(target_os = "linux")]
mod rs485_flags {
pub const SER_RS485_ENABLED: u32 = 1 << 0;
pub const SER_RS485_RTS_ON_SEND: u32 = 1 << 1;
pub const SER_RS485_RTS_AFTER_SEND: u32 = 1 << 2;
}
#[cfg(target_os = "linux")]
const TIOCGRS485: libc::c_ulong = 0x542E;
#[cfg(target_os = "linux")]
const TIOCSRS485: libc::c_ulong = 0x542F;
#[cfg(target_os = "linux")]
impl DrvAsynSerialPort {
fn rs485_get(&self, fd: RawFd) -> AsynResult<SerialRs485> {
let mut r: SerialRs485 = SerialRs485::default();
let ret = unsafe { libc::ioctl(fd, TIOCGRS485, &mut r as *mut SerialRs485) };
if ret < 0 {
return Err(AsynError::Io(std::io::Error::last_os_error()));
}
Ok(r)
}
fn rs485_set(&self, fd: RawFd, r: &SerialRs485) -> AsynResult<()> {
let ret = unsafe { libc::ioctl(fd, TIOCSRS485, r as *const SerialRs485) };
if ret < 0 {
return Err(AsynError::Io(std::io::Error::last_os_error()));
}
Ok(())
}
fn set_rs485_option(&mut self, key: &str, value: &str) -> AsynResult<()> {
let fd = self.io.fd.ok_or_else(|| AsynError::Status {
status: AsynStatus::Disconnected,
message: "not connected".into(),
})?;
let mut r = self.rs485_get(fd)?;
let prev = r;
use rs485_flags::*;
match key {
"rs485_enable" => {
if parse_yn_option(key, value)? {
r.flags |= SER_RS485_ENABLED;
} else {
r.flags = 0;
}
}
"rs485_rts_on_send" => {
if parse_yn_option(key, value)? {
r.flags |= SER_RS485_RTS_ON_SEND;
} else {
r.flags &= !SER_RS485_RTS_ON_SEND;
}
}
"rs485_rts_after_send" => {
if parse_yn_option(key, value)? {
r.flags |= SER_RS485_RTS_AFTER_SEND;
} else {
r.flags &= !SER_RS485_RTS_AFTER_SEND;
}
}
"rs485_delay_rts_before_send" => {
r.delay_rts_before_send = sscanf_uint(value).ok_or_else(bad_number)?;
}
"rs485_delay_rts_after_send" => {
r.delay_rts_after_send = sscanf_uint(value).ok_or_else(bad_number)?;
}
_ => {}
}
if let Err(e) = self.rs485_set(fd, &r) {
let _ = self.rs485_set(fd, &prev);
return Err(e);
}
Ok(())
}
fn get_rs485_option(&self, key: &str) -> AsynResult<String> {
let fd = self.io.fd.ok_or_else(|| AsynError::Status {
status: AsynStatus::Disconnected,
message: "not connected".into(),
})?;
let r = self.rs485_get(fd)?;
use rs485_flags::*;
let s = match key {
"rs485_enable" => if r.flags & SER_RS485_ENABLED != 0 {
"Y"
} else {
"N"
}
.to_string(),
"rs485_rts_on_send" => if r.flags & SER_RS485_RTS_ON_SEND != 0 {
"Y"
} else {
"N"
}
.to_string(),
"rs485_rts_after_send" => if r.flags & SER_RS485_RTS_AFTER_SEND != 0 {
"Y"
} else {
"N"
}
.to_string(),
"rs485_delay_rts_before_send" => r.delay_rts_before_send.to_string(),
"rs485_delay_rts_after_send" => r.delay_rts_after_send.to_string(),
_ => {
return Err(AsynError::OptionNotFound(key.to_string()));
}
};
Ok(s)
}
}
impl Drop for DrvAsynSerialPort {
fn drop(&mut self) {
let user = AsynUser::default();
if self.base.is_connected() {
let _ = self.disconnect(&user);
}
}
}
#[cfg(test)]
#[allow(deprecated)]
mod tests {
use super::*;
#[test]
fn test_parse_device() {
let cfg = SerialConfig::parse("/dev/ttyUSB0").unwrap();
assert_eq!(cfg.device, "/dev/ttyUSB0");
assert_eq!(cfg.baud, 9600);
assert_eq!(cfg.data_bits, DataBits::Eight);
assert_eq!(cfg.parity, Parity::None);
assert_eq!(cfg.stop_bits, StopBits::One);
assert_eq!(cfg.flow_control, FlowControl::None);
}
#[test]
fn test_parse_empty_error() {
assert!(SerialConfig::parse("").is_err());
assert!(SerialConfig::parse(" ").is_err());
}
#[test]
fn test_driver_initial_state() {
let drv = DrvAsynSerialPort::new("serial1", "/dev/ttyUSB0").unwrap();
assert!(!drv.base().is_connected());
assert!(drv.base().auto_connect);
assert!(drv.base().flags.can_block);
assert_eq!(drv.base().interpose_octet.len(), 0);
}
#[test]
fn test_configure_installs_eos_unless_suppressed_and_honors_no_auto_connect() {
let default_port =
DrvAsynSerialPort::configure("s_eos_default", "/dev/ttyS0", false, false).unwrap();
assert_eq!(
default_port.base().interpose_octet.len(),
1,
"default serial port must auto-install the EOS interpose"
);
assert!(default_port.base().auto_connect);
let suppressed =
DrvAsynSerialPort::configure("s_eos_off", "/dev/ttyS0", true, true).unwrap();
assert_eq!(
suppressed.base().interpose_octet.len(),
0,
"noProcessEos must suppress the EOS interpose"
);
assert!(!suppressed.base().auto_connect);
}
#[test]
fn test_set_option_baud_disconnected() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
drv.set_option(&mut AsynUser::default(), "baud", "115200")
.unwrap();
assert_eq!(drv.baud, 115200);
assert_eq!(drv.get_option("baud").unwrap(), "115200");
}
#[test]
fn test_set_option_bits() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
drv.set_option(&mut AsynUser::default(), "bits", "7")
.unwrap();
assert_eq!(drv.termios.c_cflag & libc::CSIZE, libc::CS7);
assert_eq!(drv.get_option("bits").unwrap(), "7");
}
#[test]
fn test_set_option_parity() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
drv.set_option(&mut AsynUser::default(), "parity", "even")
.unwrap();
assert_ne!(drv.termios.c_cflag & libc::PARENB, 0);
assert_eq!(drv.termios.c_cflag & libc::PARODD, 0);
assert_eq!(drv.get_option("parity").unwrap(), "even");
drv.set_option(&mut AsynUser::default(), "parity", "odd")
.unwrap();
assert_eq!(drv.get_option("parity").unwrap(), "odd");
}
#[test]
fn test_set_option_stop() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
drv.set_option(&mut AsynUser::default(), "stop", "2")
.unwrap();
assert_ne!(drv.termios.c_cflag & libc::CSTOPB, 0);
assert_eq!(drv.get_option("stop").unwrap(), "2");
}
#[test]
fn test_set_option_invalid_baud() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
assert!(
drv.set_option(&mut AsynUser::default(), "baud", "abc")
.is_err()
);
}
#[test]
fn test_set_option_unsupported_baud() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
#[cfg(not(any(
target_os = "macos",
target_os = "ios",
target_os = "freebsd",
target_os = "netbsd",
target_os = "openbsd",
target_os = "dragonfly"
)))]
{
let err = drv
.set_option(&mut AsynUser::default(), "baud", "12345")
.unwrap_err();
assert_eq!(err.message(), "Unsupported data rate (12345 baud)");
}
#[cfg(any(
target_os = "macos",
target_os = "ios",
target_os = "freebsd",
target_os = "netbsd",
target_os = "openbsd",
target_os = "dragonfly"
))]
{
drv.set_option(&mut AsynUser::default(), "baud", "12345")
.unwrap();
assert_eq!(drv.baud, 12345);
assert_eq!(drv.get_option("baud").unwrap(), "12345");
}
}
#[test]
fn baud_arbitrary_on_bsd_mapped_set_on_linux() {
assert!(baud_to_speed(9600).is_some(), "9600 is standard everywhere");
#[cfg(any(
target_os = "macos",
target_os = "ios",
target_os = "freebsd",
target_os = "netbsd",
target_os = "openbsd",
target_os = "dragonfly"
))]
{
assert_eq!(baud_to_speed(28800), Some(28800 as libc::speed_t));
assert_eq!(baud_to_speed(250000), Some(250000 as libc::speed_t));
assert_eq!(baud_to_speed(115200), Some(115200 as libc::speed_t));
assert_eq!(baud_to_speed(0), Some(0 as libc::speed_t));
}
#[cfg(any(target_os = "linux", target_os = "android"))]
{
assert!(baud_to_speed(460800).is_some(), "Linux maps 460800");
assert!(baud_to_speed(4000000).is_some(), "Linux maps 4000000");
assert!(baud_to_speed(28800).is_none(), "Linux has no B28800");
assert!(
baud_to_speed(250000).is_none(),
"Linux rejects non-standard"
);
assert!(
baud_to_speed(0).is_none(),
"Linux rejects baud 0 (C switch has no case 0)"
);
}
}
#[test]
fn apply_to_termios_errors_on_unmappable_baud() {
let valid = SerialConfig::parse("/dev/ttyS0").unwrap();
let mut t: libc::termios = unsafe { std::mem::zeroed() };
assert!(valid.apply_to_termios(&mut t).is_ok());
#[cfg(any(target_os = "linux", target_os = "android"))]
{
let bad = SerialConfig {
baud: 28800,
..SerialConfig::parse("/dev/ttyS0").unwrap()
};
assert!(bad.apply_to_termios(&mut t).is_err());
}
}
#[test]
fn test_set_option_invalid_bits() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
assert!(
drv.set_option(&mut AsynUser::default(), "bits", "9")
.is_err()
);
}
#[test]
fn test_set_option_key_case_insensitive() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
drv.set_option(&mut AsynUser::default(), "BAUD", "115200")
.unwrap();
assert_eq!(drv.baud, 115200);
drv.set_option(&mut AsynUser::default(), "Parity", "Even")
.unwrap();
assert_eq!(drv.get_option("parity").unwrap(), "even");
}
#[test]
fn test_set_option_value_trimmed() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
drv.set_option(&mut AsynUser::default(), "baud", " 9600 ")
.unwrap();
assert_eq!(drv.baud, 9600);
}
#[test]
fn test_set_option_parity_case_insensitive() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
drv.set_option(&mut AsynUser::default(), "parity", "EVEN")
.unwrap();
assert_eq!(drv.get_option("parity").unwrap(), "even");
drv.set_option(&mut AsynUser::default(), "parity", "None")
.unwrap();
assert_eq!(drv.get_option("parity").unwrap(), "none");
assert!(
drv.set_option(&mut AsynUser::default(), "parity", "n")
.is_err()
);
}
#[test]
fn test_set_option_parity_mark_space_unsupported() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
let err = drv
.set_option(&mut AsynUser::default(), "parity", "mark")
.unwrap_err();
assert_eq!(err.message(), "Invalid parity.");
}
#[test]
fn test_parse_bool_option() {
assert!(parse_yn_option("clocal", "Y").unwrap());
assert!(parse_yn_option("clocal", "y").unwrap());
assert!(!parse_yn_option("clocal", "N").unwrap());
assert!(!parse_yn_option("clocal", "n").unwrap());
for v in &["yes", "1", "true", "no", "0", "false", "maybe", ""] {
assert!(
parse_yn_option("clocal", v).is_err(),
"expected err for '{v}'"
);
}
}
#[test]
fn get_option_break_returns_off() {
let drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
assert_eq!(drv.get_option("break").unwrap(), "off");
}
#[test]
fn set_option_break_on_disconnected_errors_but_off_is_noop() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
drv.set_option(&mut AsynUser::default(), "break", "off")
.unwrap();
let err = drv
.set_option(&mut AsynUser::default(), "break", "on")
.unwrap_err();
assert!(
matches!(
err,
AsynError::Status {
status: AsynStatus::Disconnected,
..
}
),
"break on a disconnected port must error, got {err:?}"
);
let err = drv
.set_option(&mut AsynUser::default(), "break", "notanumber")
.unwrap_err();
assert_eq!(err.message(), "Bad number");
}
#[test]
fn every_serial_option_reports_cs_text() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
let mut err = |key: &str, val: &str| {
drv.set_option(&mut AsynUser::default(), key, val)
.unwrap_err()
.message()
};
assert_eq!(err("baud", "fast"), "Bad number", "C :264");
assert_eq!(err("bits", "9"), "Invalid number of bits.", "C :374");
assert_eq!(err("parity", "mark"), "Invalid parity.", "C :395");
assert_eq!(err("stop", "3"), "Invalid number of stop bits.", "C :406");
assert_eq!(err("clocal", "maybe"), "Invalid clocal value.", "C :419");
assert_eq!(err("crtscts", "maybe"), "Invalid crtscts value.", "C :440");
assert_eq!(err("ixon", "maybe"), "Invalid ixon value.", "C :464");
assert_eq!(err("ixany", "maybe"), "Invalid ixany value.", "C :481");
assert_eq!(err("ixoff", "maybe"), "Invalid ixoff value.", "C :502");
assert_eq!(err("break", "soon"), "Bad number", "C :513");
assert_eq!(err("nosuch", "x"), "Unsupported key \"nosuch\"", "C :595");
}
#[test]
fn a_numeric_option_prefix_parses_like_c_sscanf() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
drv.set_option(&mut AsynUser::default(), "baud", "9600x")
.unwrap();
assert_eq!(drv.baud, 9600, "C's sscanf %d stops at the first non-digit");
let err = drv
.set_option(&mut AsynUser::default(), "baud", "x9600")
.unwrap_err();
assert_eq!(err.message(), "Bad number");
assert_eq!(drv.baud, 9600, "the rejected write left the cache alone");
}
#[test]
fn test_set_option_unknown() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
let err = drv
.set_option(&mut AsynUser::default(), "custom", "value")
.unwrap_err();
assert!(matches!(err, AsynError::OptionNotFound(_)));
assert_eq!(err.message(), "Unsupported key \"custom\"");
assert!(drv.get_option("custom").is_err());
drv.set_option(&mut AsynUser::default(), "", "ignored")
.unwrap();
}
#[test]
fn test_get_option_not_found() {
let drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
assert!(drv.get_option("nonexistent").is_err());
}
#[test]
fn test_read_write_when_disconnected() {
let mut drv = DrvAsynSerialPort::new("s1", "/dev/ttyS0").unwrap();
let user = AsynUser::new(0).with_timeout(Duration::from_secs(1));
let mut buf = [0u8; 32];
assert!(drv.read_octet(&user, &mut buf).is_err());
let mut user = AsynUser::new(0);
assert!(drv.write_octet(&mut user, b"hello").is_err());
}
#[test]
fn test_baud_speed_roundtrip() {
for baud in [
50, 75, 110, 134, 150, 200, 300, 600, 1200, 1800, 2400, 4800, 9600, 19200, 38400,
57600, 115200, 230400,
] {
let speed = baud_to_speed(baud).expect("standard rate must map");
assert_eq!(
speed_to_baud(speed),
baud,
"roundtrip failed for baud={baud}"
);
}
}
fn create_pty_pair() -> Option<(RawFd, RawFd, String)> {
let mut master: RawFd = 0;
let mut slave: RawFd = 0;
let mut name_buf = [0u8; 256];
let ret = unsafe {
libc::openpty(
&mut master,
&mut slave,
name_buf.as_mut_ptr() as *mut libc::c_char,
std::ptr::null_mut(),
std::ptr::null_mut(),
)
};
if ret < 0 {
return None;
}
let name = unsafe {
std::ffi::CStr::from_ptr(name_buf.as_ptr() as *const libc::c_char)
.to_string_lossy()
.into_owned()
};
Some((master, slave, name))
}
struct PtyGuard {
master: RawFd,
slave: RawFd,
}
impl Drop for PtyGuard {
fn drop(&mut self) {
unsafe {
libc::close(self.master);
libc::close(self.slave);
}
}
}
#[test]
fn test_pty_connect_disconnect() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_test", &slave_name).unwrap();
let user = AsynUser::default();
assert!(!drv.base().is_connected());
drv.connect(&user).unwrap();
assert!(drv.base().is_connected());
drv.disconnect(&user).unwrap();
assert!(!drv.base().is_connected());
}
#[test]
fn pty_zero_length_read_rejected_not_eof_teardown() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_maxchars", &slave_name).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
assert!(drv.base().is_connected());
let ruser = AsynUser::new(0).with_timeout(Duration::from_millis(50));
let mut empty: [u8; 0] = [];
let res = drv.read_octet(&ruser, &mut empty);
assert!(
matches!(
res,
Err(AsynError::Status {
status: AsynStatus::Error,
..
})
),
"zero-length serial read must be rejected with asynError, got {res:?}"
);
assert!(
drv.base().is_connected(),
"zero-length serial read must not tear down the connection"
);
}
#[test]
fn pty_empty_key_reapplies_configured_termios() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_reapply", &slave_name).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
let fd = drv.io.fd.expect("connected fd");
let read_cstopb = |fd: RawFd| -> bool {
let mut t: libc::termios = unsafe { std::mem::zeroed() };
assert_eq!(unsafe { libc::tcgetattr(fd, &mut t) }, 0);
(t.c_cflag & libc::CSTOPB) != 0
};
let configured = read_cstopb(fd);
{
let mut t: libc::termios = unsafe { std::mem::zeroed() };
assert_eq!(unsafe { libc::tcgetattr(fd, &mut t) }, 0);
if configured {
t.c_cflag &= !libc::CSTOPB;
} else {
t.c_cflag |= libc::CSTOPB;
}
assert_eq!(unsafe { libc::tcsetattr(fd, libc::TCSANOW, &t) }, 0);
}
assert_eq!(
read_cstopb(fd),
!configured,
"external clobber must take effect before the re-apply"
);
drv.set_option(&mut AsynUser::default(), "", "").unwrap();
assert_eq!(
read_cstopb(fd),
configured,
"empty-key set_option must restore the configured line state"
);
drv.disconnect(&user).unwrap();
}
#[test]
fn test_pty_write_read_roundtrip() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_test", &slave_name).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
let mut user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
drv.write_octet(&mut user, b"hello").unwrap();
let mut buf = [0u8; 32];
let n = unsafe { libc::read(master, buf.as_mut_ptr() as *mut libc::c_void, buf.len()) };
assert!(n > 0);
assert_eq!(&buf[..n as usize], b"hello");
let msg = b"world";
unsafe { libc::write(master, msg.as_ptr() as *const libc::c_void, msg.len()) };
let user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut rbuf = [0u8; 32];
let n = drv.read_octet(&user, &mut rbuf).unwrap();
assert_eq!(&rbuf[..n], b"world");
}
#[test]
fn test_pty_read_timeout() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_test", &slave_name).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
let user = AsynUser::new(0).with_timeout(Duration::from_millis(100));
let mut buf = [0u8; 32];
let err = drv.read_octet(&user, &mut buf).unwrap_err();
match err {
AsynError::Status {
status: AsynStatus::Timeout,
..
} => {}
other => panic!("expected Timeout, got {other:?}"),
}
}
#[test]
fn test_pty_read_error_disconnects() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let mut drv = DrvAsynSerialPort::new("pty_test", &slave_name).unwrap();
drv.connect(&AsynUser::default()).unwrap();
assert!(drv.base().is_connected());
unsafe { libc::close(master) };
let user = AsynUser::new(0).with_timeout(Duration::from_secs(1));
let mut buf = [0u8; 32];
let err = drv.read_octet(&user, &mut buf).unwrap_err();
assert!(
err.is_fatal_transport(),
"expected a fatal transport error, got {err:?}"
);
assert!(
!drv.base().is_connected(),
"DRV-31: fatal read error must set connected=false"
);
}
#[test]
fn test_pty_write_error_disconnects() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let mut drv = DrvAsynSerialPort::new("pty_test", &slave_name).unwrap();
drv.connect(&AsynUser::default()).unwrap();
assert!(drv.base().is_connected());
unsafe { libc::close(master) };
let mut user = AsynUser::new(0).with_timeout(Duration::from_secs(1));
let err = drv.write_octet(&mut user, b"hello world").unwrap_err();
assert!(
err.is_fatal_transport(),
"expected a fatal transport error, got {err:?}"
);
assert!(
!drv.base().is_connected(),
"DRV-31: fatal write error must set connected=false"
);
}
#[test]
fn fatal_transport_error_sees_through_the_eos_interpose_wrapper() {
use crate::interpose::PartialOctetRead;
let wrapped_hangup = AsynError::Status {
status: AsynStatus::Disconnected,
message: "hangup".into(),
}
.with_partial_read(PartialOctetRead {
data: b"AB".to_vec(),
eom_reason: EomReason::empty(),
});
assert!(
wrapped_hangup.is_fatal_transport(),
"a hangup wrapped by the EOS interpose is still fatal"
);
let wrapped_timeout = AsynError::Status {
status: AsynStatus::Timeout,
message: "read timeout".into(),
}
.with_partial_read(PartialOctetRead {
data: b"AB".to_vec(),
eom_reason: EomReason::empty(),
});
assert!(
!wrapped_timeout.is_fatal_transport(),
"a partial-line timeout leaves the fd intact (C returns asynTimeout)"
);
let wrapped_errno =
AsynError::Io(std::io::Error::other("EIO")).with_partial_read(PartialOctetRead {
data: b"AB".to_vec(),
eom_reason: EomReason::empty(),
});
assert!(
wrapped_errno.is_fatal_transport(),
"an errno behind the read carrier is still fatal"
);
let half_written = AsynError::Io(std::io::Error::other("EIO")).with_partial_write(2);
assert!(
half_written.is_fatal_transport(),
"an errno behind the write carrier is still fatal"
);
}
#[test]
fn test_pty_eos_interpose() {
use crate::interpose::eos::{EosConfig, EosInterpose};
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_test", &slave_name).unwrap();
let eos = EosInterpose::new(EosConfig {
input_eos: vec![b'\r', b'\n'],
output_eos: vec![],
});
drv.install_interpose(Box::new(eos));
let user = AsynUser::default();
drv.connect(&user).unwrap();
let msg = b"OK\r\n";
unsafe { libc::write(master, msg.as_ptr() as *const libc::c_void, msg.len()) };
let user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut buf = [0u8; 32];
let (n, _eom) = crate::port::octet_read_chain(&mut drv, &user, &mut buf).unwrap();
assert_eq!(&buf[..n], b"OK");
}
#[test]
fn test_pty_set_option_baud() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_test", &slave_name).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
drv.set_option(&mut AsynUser::default(), "baud", "115200")
.unwrap();
assert_eq!(drv.baud, 115200);
let t = drv.get_current_termios().unwrap();
let actual_speed = unsafe { libc::cfgetospeed(&t) };
assert_eq!(actual_speed, libc::B115200);
}
#[test]
fn clocal_survives_reconnect_and_can_be_set_while_disconnected() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_test", &slave_name).unwrap();
assert_eq!(
drv.get_option("clocal").unwrap(),
"Y",
"C seeds CS8|CLOCAL|CREAD (:1077) and getOption reads the cache"
);
drv.set_option(&mut AsynUser::default(), "clocal", "N")
.unwrap();
assert_eq!(
drv.get_option("clocal").unwrap(),
"N",
"an option set while the line is down is held, not dropped"
);
drv.connect(&AsynUser::default()).unwrap();
let live = drv.get_current_termios().unwrap();
assert_eq!(
live.c_cflag & libc::CLOCAL,
0,
"connect must push the cached clocal=N, not force CLOCAL back on"
);
assert_ne!(
live.c_cflag & libc::CREAD,
0,
"applyOptions still forces CREAD (C :119)"
);
drv.disconnect(&AsynUser::default()).unwrap();
assert_eq!(drv.get_option("clocal").unwrap(), "N");
drv.connect(&AsynUser::default()).unwrap();
let live = drv.get_current_termios().unwrap();
assert_eq!(
live.c_cflag & libc::CLOCAL,
0,
"clocal=N must survive the reconnect (C re-pushes tty->termios)"
);
drv.set_option(&mut AsynUser::default(), "clocal", "Y")
.unwrap();
let live = drv.get_current_termios().unwrap();
assert_ne!(live.c_cflag & libc::CLOCAL, 0);
assert_eq!(drv.get_option("clocal").unwrap(), "Y");
}
#[test]
fn iflag_family_reads_back_per_flag_from_the_cache() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_test", &slave_name).unwrap();
assert_eq!(drv.get_option("ixon").unwrap(), "N");
assert_eq!(drv.get_option("ixoff").unwrap(), "N");
assert_eq!(drv.get_option("ixany").unwrap(), "N");
drv.set_option(&mut AsynUser::default(), "ixon", "Y")
.unwrap();
drv.set_option(&mut AsynUser::default(), "ixany", "Y")
.unwrap();
assert_eq!(
drv.get_option("ixon").unwrap(),
"Y",
"a disconnected readback must report the cache, not a flat N"
);
assert_eq!(drv.get_option("ixoff").unwrap(), "N");
assert_eq!(
drv.get_option("ixany").unwrap(),
"Y",
"ixany is a real cached flag on POSIX (C hard-codes N only on vxWorks)"
);
drv.connect(&AsynUser::default()).unwrap();
let live = drv.get_current_termios().unwrap();
assert_ne!(live.c_iflag & libc::IXON, 0);
assert_eq!(live.c_iflag & libc::IXOFF, 0);
assert_ne!(live.c_iflag & libc::IXANY, 0);
drv.set_option(&mut AsynUser::default(), "ixoff", "Y")
.unwrap();
assert_eq!(drv.get_option("ixoff").unwrap(), "Y");
drv.disconnect(&AsynUser::default()).unwrap();
assert_eq!(drv.get_option("ixoff").unwrap(), "Y");
assert_eq!(drv.get_option("ixon").unwrap(), "Y");
drv.connect(&AsynUser::default()).unwrap();
let live = drv.get_current_termios().unwrap();
assert_ne!(live.c_iflag & libc::IXOFF, 0);
assert_ne!(live.c_iflag & libc::IXON, 0);
drv.set_option(&mut AsynUser::default(), "ixon", "N")
.unwrap();
assert_eq!(drv.get_option("ixon").unwrap(), "N");
assert_eq!(drv.get_option("ixoff").unwrap(), "Y");
assert_eq!(drv.get_option("ixany").unwrap(), "Y");
let live = drv.get_current_termios().unwrap();
assert_eq!(live.c_iflag & libc::IXON, 0);
assert_ne!(live.c_iflag & libc::IXOFF, 0);
}
#[test]
fn cflag_options_set_while_disconnected_apply_at_connect() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_test", &slave_name).unwrap();
drv.set_option(&mut AsynUser::default(), "bits", "7")
.unwrap();
drv.set_option(&mut AsynUser::default(), "parity", "odd")
.unwrap();
drv.set_option(&mut AsynUser::default(), "stop", "2")
.unwrap();
drv.set_option(&mut AsynUser::default(), "crtscts", "Y")
.unwrap();
assert_eq!(drv.get_option("bits").unwrap(), "7");
assert_eq!(drv.get_option("parity").unwrap(), "odd");
assert_eq!(drv.get_option("stop").unwrap(), "2");
assert_eq!(drv.get_option("crtscts").unwrap(), "Y");
drv.connect(&AsynUser::default()).unwrap();
let live = drv.get_current_termios().unwrap();
assert_ne!(
live.c_cflag & libc::CSTOPB,
0,
"stop=2 set while disconnected must land on the device at connect"
);
assert_ne!(
live.c_cflag & libc::CRTSCTS,
0,
"crtscts=Y set while disconnected must land on the device at connect"
);
}
#[test]
fn test_pty_connect_rejects_double_open() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_test", &slave_name).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
let first_fd = drv.io.fd;
assert!(first_fd.is_some());
let err = drv.connect(&user).unwrap_err();
assert!(matches!(err, AsynError::Status { .. }));
assert_eq!(drv.io.fd, first_fd);
assert!(drv.saved_termios.is_some());
}
#[test]
fn test_pty_runtime_integration() {
use crate::runtime::{RuntimeConfig, create_port_runtime};
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let drv = DrvAsynSerialPort::new("pty_rt", &slave_name).unwrap();
let (runtime_handle, _jh) = create_port_runtime(drv, RuntimeConfig::default())
.expect("the port runtime thread must start");
let ph = runtime_handle.port_handle();
let user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
ph.submit_blocking(
crate::request::RequestOp::OctetWrite {
data: b"ping".to_vec(),
},
user,
)
.unwrap();
let mut buf = [0u8; 32];
let n = unsafe { libc::read(master, buf.as_mut_ptr() as *mut libc::c_void, buf.len()) };
assert!(n > 0);
assert_eq!(&buf[..n as usize], b"ping");
let resp = b"pong";
unsafe { libc::write(master, resp.as_ptr() as *const libc::c_void, resp.len()) };
let user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let result = ph
.submit_blocking(crate::request::RequestOp::OctetRead { buf_size: 32 }, user)
.unwrap();
assert_eq!(result.data.as_deref(), Some(b"pong".as_slice()));
runtime_handle.shutdown_and_wait();
}
#[test]
fn test_pty_termios_restored_on_disconnect() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_test", &slave_name).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
assert!(drv.saved_termios.is_some());
let saved = drv.saved_termios.unwrap();
let current = drv.get_current_termios().unwrap();
assert_ne!(
current.c_lflag & libc::ECHO,
saved.c_lflag & libc::ECHO,
"raw mode should have changed ECHO flag"
);
drv.saved_termios = Some(saved);
drv.disconnect(&user).unwrap();
assert!(drv.saved_termios.is_none());
assert!(!drv.base().is_connected());
let c_path = std::ffi::CString::new(slave_name.as_str()).unwrap();
let fd2 = unsafe {
libc::open(
c_path.as_ptr(),
libc::O_RDWR | libc::O_NOCTTY | libc::O_NONBLOCK,
)
};
if fd2 >= 0 {
let mut restored: libc::termios = unsafe { std::mem::zeroed() };
if unsafe { libc::tcgetattr(fd2, &mut restored) } == 0 {
assert_eq!(
restored.c_lflag & libc::ECHO,
saved.c_lflag & libc::ECHO,
"ECHO flag should be restored"
);
assert_eq!(
restored.c_lflag & libc::ICANON,
saved.c_lflag & libc::ICANON,
"ICANON flag should be restored"
);
assert_eq!(
restored.c_cflag & libc::CSIZE,
saved.c_cflag & libc::CSIZE,
"CSIZE should be restored"
);
}
unsafe { libc::close(fd2) };
}
}
#[test]
fn pty_termios_sets_ignbrk_ignpar() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_ignbrk", &slave_name).unwrap();
drv.connect(&AsynUser::default()).unwrap();
let t = drv.get_current_termios().unwrap();
assert_ne!(
t.c_iflag & libc::IGNBRK,
0,
"IGNBRK must be set (C default, drvAsynSerialPort.c:1080)"
);
assert_ne!(
t.c_iflag & libc::IGNPAR,
0,
"IGNPAR must be set (C default, drvAsynSerialPort.c:1080)"
);
}
#[test]
fn pty_termios_sets_xon_xoff_chars() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_xonxoff", &slave_name).unwrap();
drv.connect(&AsynUser::default()).unwrap();
let t = drv.get_current_termios().unwrap();
assert_eq!(t.c_cc[libc::VSTART], 0x11, "VSTART must be ^Q (C default)");
assert_eq!(t.c_cc[libc::VSTOP], 0x13, "VSTOP must be ^S (C default)");
}
#[test]
fn pty_fd_is_non_blocking_at_every_boundary() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_nonblock", &slave_name).unwrap();
drv.connect(&AsynUser::default()).unwrap();
let fd = drv.io.fd.expect("connected");
let nonblocking = |where_: &str| {
let fl = unsafe { libc::fcntl(fd, libc::F_GETFL) };
assert!(fl >= 0, "F_GETFL failed at {where_}");
assert_ne!(
fl & libc::O_NONBLOCK,
0,
"fd must be non-blocking {where_}, flags=0x{fl:x}"
);
};
nonblocking("after connect");
let mut user = AsynUser {
timeout: Duration::from_millis(200),
..Default::default()
};
drv.write_octet(&mut user, b"probe\n").unwrap();
nonblocking("after write");
let mut buf = [0u8; 8];
let _ = drv.read_octet(&user, &mut buf);
nonblocking("after read");
}
#[test]
fn platform_reports_every_facility_present_on_this_host() {
assert_eq!(
platform::IXANY,
Some(libc::IXANY),
"hosted unix has an IXANY bit"
);
assert_eq!(
platform::SOFT_FLOW_CHARS,
Some((libc::VSTART, libc::VSTOP)),
"hosted unix has programmable XON/XOFF characters"
);
assert_eq!(
platform::CRTSCTS,
libc::CRTSCTS,
"the seam must be transparent where libc has the bit"
);
assert_eq!(platform::O_NOCTTY, libc::O_NOCTTY);
assert_eq!(
platform::FLUSH_IO,
libc::TCIOFLUSH,
"hosted unix can discard both queues at once"
);
assert!(
matches!(platform::drain(-1), Some(Err(_))),
"hosted unix has tcdrain; -1 must fail as an fd, not as a facility"
);
}
#[cfg(not(asyn_baud_code_is_rate))]
#[test]
fn a_speed_the_platform_refuses_is_an_error_not_a_silent_no_op() {
let mut t: platform::termios = unsafe { std::mem::zeroed() };
set_termios_speed(&mut t, platform::B9600).expect("9600 is settable everywhere");
let before = unsafe { libc::cfgetospeed(&t) };
let bogus = platform::Speed::MAX / 2;
let r = set_termios_speed(&mut t, bogus);
assert!(
matches!(&r, Err(AsynError::Status { message, .. })
if message.starts_with("cfsetispeed returned")
|| message.starts_with("cfsetospeed returned")),
"must name the call C names, got {r:?}"
);
assert_eq!(
unsafe { libc::cfgetospeed(&t) },
before,
"a refused speed must leave the previous rate in place"
);
}
#[test]
fn every_c_cc_index_the_seam_hands_out_is_inside_the_array() {
let t: platform::termios = unsafe { std::mem::zeroed() };
let n = t.c_cc.len();
assert!(
platform::VMIN < n,
"VMIN {} outside c_cc[{n}]",
platform::VMIN
);
assert!(
platform::VTIME < n,
"VTIME {} outside c_cc[{n}]",
platform::VTIME
);
if let Some((vstart, vstop)) = platform::SOFT_FLOW_CHARS {
assert!(vstart < n, "VSTART {vstart} outside c_cc[{n}]");
assert!(vstop < n, "VSTOP {vstop} outside c_cc[{n}]");
}
}
#[test]
fn unsupported_option_names_the_key_and_the_platform() {
let msg = match option_unsupported_here("ixany") {
AsynError::Status { status, message } => {
assert_eq!(status, AsynStatus::Error);
message
}
other => panic!("expected a Status error, got {other:?}"),
};
assert!(msg.contains("ixany"), "must name the option: {msg}");
assert!(
msg.contains(std::env::consts::OS),
"must name the platform: {msg}"
);
}
#[test]
fn set_option_does_not_commit_cached_config_on_apply_failure() {
let mut drv = DrvAsynSerialPort::new("rollback", "/dev/null").unwrap();
let badfd = unsafe { libc::open(c"/dev/null".as_ptr(), libc::O_RDWR) };
assert!(badfd >= 0, "could not open /dev/null");
drv.io.fd = Some(badfd);
assert_eq!(drv.get_option("baud").unwrap(), "9600");
let r = drv.set_option(&mut AsynUser::default(), "baud", "115200");
assert!(r.is_err(), "set_option must fail when apply fails");
assert_eq!(
drv.get_option("baud").unwrap(),
"9600",
"cached baud must stay 9600 when apply fails (C restores baudPrev)"
);
drv.io.fd = None;
unsafe { libc::close(badfd) };
}
#[test]
fn pty_write_timeout_bounds_total_not_per_poll() {
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_write_total", &slave_name).unwrap();
drv.connect(&AsynUser::default()).unwrap();
let stop = Arc::new(AtomicBool::new(false));
let stop2 = stop.clone();
let reader = std::thread::spawn(move || {
let mut buf = [0u8; 4096];
while !stop2.load(Ordering::Relaxed) {
let n =
unsafe { libc::read(master, buf.as_mut_ptr() as *mut libc::c_void, buf.len()) };
if n <= 0 {
break; }
std::thread::sleep(Duration::from_millis(10));
}
});
let payload = vec![0x5Au8; 512 * 1024];
let mut user = AsynUser::new(0).with_timeout(Duration::from_millis(300));
let start = Instant::now();
let res = drv.write_octet(&mut user, &payload);
let elapsed = start.elapsed();
stop.store(true, Ordering::Relaxed);
drv.disconnect(&AsynUser::default()).ok();
let _ = reader.join();
let err = match res {
Err(e) => e,
other => panic!("expected total-deadline Timeout, got {other:?}"),
};
assert_eq!(err.status(), AsynStatus::Timeout, "got {err:?}");
let sent = err
.partial_write()
.expect("a timed-out serial write must report what it transferred");
assert!(
sent > 0 && sent < payload.len(),
"expected a partial count in 1..{}, got {sent}",
payload.len()
);
assert!(
elapsed < Duration::from_millis(900),
"total write time must be bounded by ~the timeout, took {elapsed:?}"
);
}
#[test]
fn pty_connect_sets_cloexec() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_cloexec", &slave_name).unwrap();
drv.connect(&AsynUser::default()).unwrap();
let fd = drv.io.fd.expect("connected fd");
let flags = unsafe { libc::fcntl(fd, libc::F_GETFD) };
assert!(flags >= 0, "F_GETFD failed");
assert!(
flags & libc::FD_CLOEXEC != 0,
"FD_CLOEXEC must be set after connect (C parity)"
);
}
#[test]
fn pty_report_tracks_byte_counters() {
let (master, slave, slave_name) = match create_pty_pair() {
Some(v) => v,
None => {
eprintln!("openpty not available, skipping test");
return;
}
};
unsafe { libc::close(slave) };
let _guard = PtyGuard { master, slave: -1 };
let mut drv = DrvAsynSerialPort::new("pty_report", &slave_name).unwrap();
drv.connect(&AsynUser::default()).unwrap();
assert_eq!(drv.io.n_written, 0);
assert_eq!(drv.io.n_read, 0);
let mut user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
drv.write_octet(&mut user, b"hello").unwrap();
assert_eq!(drv.io.n_written, 5, "n_written must track bytes written");
let msg = b"world";
unsafe { libc::write(master, msg.as_ptr() as *const libc::c_void, msg.len()) };
let user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut rbuf = [0u8; 32];
let n = drv.read_octet(&user, &mut rbuf).unwrap();
assert!(n > 0);
assert_eq!(drv.io.n_read, n as u64, "n_read must track bytes read");
let mut out = String::new();
drv.report(&mut out, 0);
drv.report(&mut out, 2);
}
}