use std::io::{Read, Write};
use std::net::{TcpStream, UdpSocket};
use std::time::Duration;
use super::option_parse::parse_yn_option;
use crate::error::{AsynError, AsynResult, AsynStatus};
use crate::exception::AsynException;
use crate::interpose::com::{ComInterpose, ComPortOptions};
use crate::interpose::{EomReason, OctetInterpose, OctetNext, OctetReadResult};
use crate::port::{PortDriver, PortDriverBase, PortFlags};
use crate::trace::TraceMask;
use crate::user::AsynUser;
use crate::{asyn_trace, asyn_trace_io};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum IpProtocol {
#[default]
Tcp,
TcpReusePort,
Udp,
UdpReusePort,
UdpBroadcast,
UdpBroadcastReusePort,
Unix,
Http,
Com,
}
impl IpProtocol {
fn broadcast(self) -> bool {
matches!(self, Self::UdpBroadcast | Self::UdpBroadcastReusePort)
}
fn reuse_port(self) -> bool {
matches!(
self,
Self::TcpReusePort | Self::UdpReusePort | Self::UdpBroadcastReusePort
)
}
}
#[derive(Debug, Clone)]
pub struct IpPortConfig {
pub host: String,
pub port: u16,
pub local_port: Option<u16>,
pub protocol: IpProtocol,
pub connect_timeout: Duration,
pub no_delay: bool,
}
const DEFAULT_CONNECT_TIMEOUT: Duration = Duration::from_secs(5);
impl IpPortConfig {
pub fn parse(spec: &str) -> AsynResult<Self> {
let spec = spec.trim();
if let Some(path) = spec
.strip_prefix("unix://")
.or_else(|| spec.strip_prefix("UNIX://"))
{
if path.is_empty() {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: "empty unix socket path".into(),
});
}
return Ok(Self {
host: path.to_string(),
port: 0,
local_port: None,
protocol: IpProtocol::Unix,
connect_timeout: DEFAULT_CONNECT_TIMEOUT,
no_delay: false,
});
}
let (addr_part, proto) = split_protocol(spec)?;
let addr_part = addr_part.trim();
let (host, port, local_port) = parse_host_port(addr_part, spec)?;
Ok(Self {
host,
port,
local_port,
protocol: proto,
connect_timeout: DEFAULT_CONNECT_TIMEOUT,
no_delay: true,
})
}
}
fn split_protocol(spec: &str) -> AsynResult<(&str, IpProtocol)> {
let Some(blank) = spec.find(' ') else {
return Ok((spec, IpProtocol::Tcp));
};
let (addr_part, rest) = spec.split_at(blank);
let token: String = rest
.split_whitespace()
.next()
.unwrap_or("")
.chars()
.take(5)
.collect();
Ok((addr_part, protocol_from_token(&token)?))
}
fn protocol_from_token(token: &str) -> AsynResult<IpProtocol> {
match token.to_ascii_lowercase().as_str() {
"" | "tcp" => Ok(IpProtocol::Tcp),
"tcp&" => Ok(IpProtocol::TcpReusePort),
"http" => Ok(IpProtocol::Http),
"udp" => Ok(IpProtocol::Udp),
"udp&" => Ok(IpProtocol::UdpReusePort),
"udp*" => Ok(IpProtocol::UdpBroadcast),
"udp*&" => Ok(IpProtocol::UdpBroadcastReusePort),
"com" => Ok(IpProtocol::Com),
_ => Err(AsynError::Status {
status: AsynStatus::Error,
message: format!("Unknown protocol \"{token}\"."),
}),
}
}
fn parse_host_port(addr_part: &str, orig_spec: &str) -> AsynResult<(String, u16, Option<u16>)> {
if addr_part.starts_with('[') {
let bracket_end = addr_part.find(']').ok_or_else(|| AsynError::Status {
status: AsynStatus::Error,
message: format!("missing closing bracket in IPv6 address: '{orig_spec}'"),
})?;
let host = addr_part[1..bracket_end].to_string();
if host.is_empty() {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: "empty IPv6 address".into(),
});
}
let rest = &addr_part[bracket_end + 1..];
let rest = rest.strip_prefix(':').ok_or_else(|| AsynError::Status {
status: AsynStatus::Error,
message: format!("expected ':port' after IPv6 bracket: '{orig_spec}'"),
})?;
let parts: Vec<&str> = rest.splitn(2, ':').collect();
let port: u16 = parts[0].parse().map_err(|_| AsynError::Status {
status: AsynStatus::Error,
message: format!("invalid port number: '{}'", parts[0]),
})?;
let local_port = if parts.len() > 1 {
Some(parts[1].parse::<u16>().map_err(|_| AsynError::Status {
status: AsynStatus::Error,
message: format!("invalid local port: '{}'", parts[1]),
})?)
} else {
None
};
return Ok((host, port, local_port));
}
let parts: Vec<&str> = addr_part.splitn(3, ':').collect();
if parts.len() < 2 {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: format!("invalid IP port spec: expected host:port, got '{orig_spec}'"),
});
}
let host = parts[0].to_string();
if host.is_empty() {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: "empty hostname".into(),
});
}
let port: u16 = parts[1].parse().map_err(|_| AsynError::Status {
status: AsynStatus::Error,
message: format!("invalid port number: '{}'", parts[1]),
})?;
let local_port = if parts.len() > 2 {
Some(parts[2].parse::<u16>().map_err(|_| AsynError::Status {
status: AsynStatus::Error,
message: format!("invalid local port: '{}'", parts[2]),
})?)
} else {
None
};
Ok((host, port, local_port))
}
enum IpIoInner {
Tcp(TcpStream),
Udp(UdpSocket, std::net::SocketAddr),
#[cfg(unix)]
Unix(std::os::unix::net::UnixStream),
}
fn write_with_retry(
stream: &mut impl Write,
data: &[u8],
deadline: std::time::Instant,
) -> AsynResult<usize> {
let mut offset = 0;
while offset < data.len() {
if std::time::Instant::now() > deadline {
return Err(AsynError::Status {
status: AsynStatus::Timeout,
message: "write timeout".into(),
}
.with_partial_write(offset));
}
match stream.write(&data[offset..]) {
Ok(0) => {
return Err(AsynError::Status {
status: AsynStatus::Timeout,
message: "write returned 0 bytes".into(),
}
.with_partial_write(offset));
}
Ok(n) => offset += n,
Err(ref e)
if e.kind() == std::io::ErrorKind::WouldBlock
|| e.kind() == std::io::ErrorKind::Interrupted =>
{
std::thread::sleep(std::time::Duration::from_millis(10));
}
Err(e) => return Err(AsynError::Io(e).with_partial_write(offset)),
}
}
Ok(offset)
}
struct IpIoState {
inner: Option<IpIoInner>,
}
impl OctetNext for IpIoState {
fn read(&mut self, user: &AsynUser, buf: &mut [u8]) -> AsynResult<OctetReadResult> {
let inner = self.inner.as_mut().ok_or_else(|| AsynError::Status {
status: AsynStatus::Disconnected,
message: "not connected".into(),
})?;
if buf.is_empty() {
return Err(maxchars_zero_error());
}
match inner {
IpIoInner::Tcp(stream) => {
let _ = stream.set_read_timeout(Some(socket_poll_timeout(user.timeout)));
match stream.read(buf) {
Ok(0) => Ok(OctetReadResult {
nbytes_transferred: 0,
eom_reason: EomReason::END,
}),
Ok(n) => Ok(OctetReadResult {
nbytes_transferred: n,
eom_reason: if n >= buf.len() {
EomReason::CNT
} else {
EomReason::empty()
},
}),
Err(e) => Err(classify_read_error(e)),
}
}
IpIoInner::Udp(socket, _peer) => {
let _ = socket.set_read_timeout(Some(socket_poll_timeout(user.timeout)));
match socket.recv_from(buf) {
Ok((n, _src)) => Ok(OctetReadResult {
nbytes_transferred: n,
eom_reason: if n >= buf.len() {
EomReason::CNT
} else {
EomReason::empty()
},
}),
Err(e) => Err(classify_read_error(e)),
}
}
#[cfg(unix)]
IpIoInner::Unix(stream) => {
let _ = stream.set_read_timeout(Some(socket_poll_timeout(user.timeout)));
match stream.read(buf) {
Ok(0) => Ok(OctetReadResult {
nbytes_transferred: 0,
eom_reason: EomReason::END,
}),
Ok(n) => Ok(OctetReadResult {
nbytes_transferred: n,
eom_reason: if n >= buf.len() {
EomReason::CNT
} else {
EomReason::empty()
},
}),
Err(e) => Err(classify_read_error(e)),
}
}
}
}
fn write(&mut self, user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
let inner = self.inner.as_mut().ok_or_else(|| AsynError::Status {
status: AsynStatus::Disconnected,
message: "not connected".into(),
})?;
if data.is_empty() {
return Ok(0);
}
let deadline = std::time::Instant::now() + socket_poll_timeout(user.timeout);
match inner {
IpIoInner::Tcp(stream) => {
stream.set_write_timeout(Some(socket_poll_timeout(user.timeout)))?;
write_with_retry(stream, data, deadline)
}
IpIoInner::Udp(socket, peer) => {
socket.set_write_timeout(Some(socket_poll_timeout(user.timeout)))?;
Ok(socket.send_to(data, *peer)?)
}
#[cfg(unix)]
IpIoInner::Unix(stream) => {
stream.set_write_timeout(Some(socket_poll_timeout(user.timeout)))?;
write_with_retry(stream, data, deadline)
}
}
}
fn flush(&mut self, _user: &mut AsynUser) -> AsynResult<()> {
let mut scratch = [0u8; 4096];
match self.inner.as_mut() {
Some(IpIoInner::Tcp(stream)) => {
let restore = stream.set_nonblocking(true);
loop {
match stream.read(&mut scratch) {
Ok(0) => break, Ok(_) => continue,
Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => break,
Err(ref e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
Err(_) => break, }
}
if restore.is_ok() {
let _ = stream.set_nonblocking(false);
}
}
Some(IpIoInner::Udp(socket, _peer)) => {
let restore = socket.set_nonblocking(true);
loop {
match socket.recv_from(&mut scratch) {
Ok(_) => continue,
Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => break,
Err(ref e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
Err(_) => break,
}
}
if restore.is_ok() {
let _ = socket.set_nonblocking(false);
}
}
#[cfg(unix)]
Some(IpIoInner::Unix(stream)) => {
let restore = stream.set_nonblocking(true);
loop {
match stream.read(&mut scratch) {
Ok(0) => break,
Ok(_) => continue,
Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => break,
Err(ref e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
Err(_) => break,
}
}
if restore.is_ok() {
let _ = stream.set_nonblocking(false);
}
}
None => {}
}
Ok(())
}
}
pub struct DrvAsynIPPort {
base: PortDriverBase,
config: IpPortConfig,
io: IpIoState,
disconnect_on_read_timeout: bool,
host_info: String,
com: Option<ComState>,
}
struct ComState {
octet: ComInterpose,
options: ComPortOptions,
}
struct ComLink<'a> {
io: &'a mut IpIoState,
com: &'a mut ComInterpose,
}
impl OctetNext for ComLink<'_> {
fn read(&mut self, user: &AsynUser, buf: &mut [u8]) -> AsynResult<OctetReadResult> {
self.com.read(user, buf, self.io)
}
fn write(&mut self, user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
self.com.write(user, data, self.io)
}
fn flush(&mut self, user: &mut AsynUser) -> AsynResult<()> {
self.com.flush(user, self.io)
}
}
pub(crate) fn is_nonfatal_read_timeout(kind: std::io::ErrorKind) -> bool {
matches!(
kind,
std::io::ErrorKind::TimedOut
| std::io::ErrorKind::WouldBlock
| std::io::ErrorKind::Interrupted
)
}
pub(crate) fn maxchars_zero_error() -> AsynError {
AsynError::Status {
status: AsynStatus::Error,
message: "maxchars 0. Why <=0?".into(),
}
}
struct ReadFailure {
disconnect: bool,
error: AsynError,
}
fn classify_read_failure(
e: AsynError,
disconnect_on_read_timeout: bool,
timeout: Duration,
device_name: &str,
) -> ReadFailure {
let is_timeout = e.status() == AsynStatus::Timeout;
let disconnect = (disconnect_on_read_timeout && is_timeout && timeout > Duration::ZERO)
|| e.is_fatal_transport();
if !disconnect {
return ReadFailure {
disconnect: false,
error: e,
};
}
let partial = e.partial_read().cloned();
let restamped = AsynError::Status {
status: AsynStatus::Error,
message: format!("{device_name} read error: {}", e.message()),
};
ReadFailure {
disconnect: true,
error: match partial {
Some(p) => restamped.with_partial_read(p),
None => restamped,
},
}
}
struct NegotiationLink<'a> {
io: &'a mut IpIoState,
disconnect_on_read_timeout: bool,
device_name: &'a str,
teardown: bool,
}
impl OctetNext for NegotiationLink<'_> {
fn read(&mut self, user: &AsynUser, buf: &mut [u8]) -> AsynResult<OctetReadResult> {
match self.io.read(user, buf) {
Ok(r) => Ok(r),
Err(e) => {
let failure = classify_read_failure(
e,
self.disconnect_on_read_timeout,
user.timeout,
self.device_name,
);
if failure.disconnect {
self.teardown = true;
}
Err(failure.error)
}
}
}
fn write(&mut self, user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
let r = self.io.write(user, data);
if let Err(ref e) = r {
if e.is_fatal_transport() {
self.teardown = true;
}
}
r
}
fn flush(&mut self, user: &mut AsynUser) -> AsynResult<()> {
self.io.flush(user)
}
}
fn classify_read_error(e: std::io::Error) -> AsynError {
if is_nonfatal_read_timeout(e.kind()) {
AsynError::Status {
status: AsynStatus::Timeout,
message: "read timeout".into(),
}
} else {
AsynError::Io(e)
}
}
pub(crate) fn socket_poll_timeout(timeout: Duration) -> Duration {
if timeout.is_zero() {
Duration::from_millis(1)
} else {
timeout
}
}
impl DrvAsynIPPort {
fn drop_connection(&mut self) {
self.io.inner = None;
if self.config.protocol != IpProtocol::Http {
self.base.set_connected(false);
}
}
fn read_octet_core(
&mut self,
user: &AsynUser,
buf: &mut [u8],
) -> AsynResult<(usize, EomReason)> {
if self.config.protocol == IpProtocol::Http && self.io.inner.is_none() {
self.connect(&AsynUser::default())?;
}
self.base.check_ready()?;
let result = self.with_base_link(|link| link.read(user, buf));
match result {
Ok(r) => {
asyn_trace_io!(
Some(self.base.trace),
&self.base.port_name,
TraceMask::IO_DRIVER,
&buf[..r.nbytes_transferred],
"read"
);
let eof = r.eom_reason.contains(EomReason::END);
if eof && self.base.is_connected() {
self.drop_connection();
}
Ok((r.nbytes_transferred, r.eom_reason))
}
Err(e) => {
let failure = classify_read_failure(
e,
self.disconnect_on_read_timeout,
user.timeout,
&self.host_info,
);
if failure.disconnect && self.base.is_connected() {
asyn_trace!(
Some(self.base.trace),
&self.base.port_name,
TraceMask::FLOW,
"read error, disconnecting: {}",
failure.error
);
self.drop_connection();
}
Err(failure.error)
}
}
}
pub fn new(port_name: &str, config_str: &str) -> AsynResult<Self> {
let config = IpPortConfig::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;
let com = (config.protocol == IpProtocol::Com).then(|| ComState {
octet: ComInterpose::new(),
options: ComPortOptions::new(),
});
Ok(Self {
base,
config,
io: IpIoState { inner: None },
disconnect_on_read_timeout: false,
host_info: config_str.to_string(),
com,
})
}
pub fn new_configured(
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)?;
Self::apply_ip_port_configure(&mut driver.base, no_auto_connect, no_process_eos);
Ok(driver)
}
fn with_base_link<T>(&mut self, f: impl FnOnce(&mut dyn OctetNext) -> T) -> T {
match self.com.as_mut() {
Some(com) => {
let mut link = ComLink {
io: &mut self.io,
com: &mut com.octet,
};
f(&mut link)
}
None => f(&mut self.io),
}
}
fn with_negotiation<T>(
&mut self,
f: impl FnOnce(&mut ComPortOptions, &mut NegotiationLink) -> T,
) -> Option<T> {
let com = self.com.as_mut()?;
let mut link = NegotiationLink {
io: &mut self.io,
disconnect_on_read_timeout: self.disconnect_on_read_timeout,
device_name: &self.host_info,
teardown: false,
};
let out = f(&mut com.options, &mut link);
let teardown = link.teardown;
if teardown && self.base.is_connected() {
self.drop_connection();
}
Some(out)
}
fn restore_com_settings(&mut self) {
let Some(Err(e)) = self.with_negotiation(|com, link| com.restore_settings(link)) else {
return;
};
asyn_trace!(
Some(self.base.trace),
&self.base.port_name,
TraceMask::ERROR,
"Unable to restore parameters for port {}: {}",
self.base.port_name,
e.message()
);
}
pub fn install_interpose(&mut self, layer: Box<dyn crate::interpose::OctetInterpose>) {
self.base.install_octet_interpose(layer);
}
pub(crate) fn apply_ip_port_configure(
base: &mut crate::port::PortDriverBase,
no_auto_connect: bool,
no_process_eos: bool,
) {
base.octet_interrupt_process = true;
base.auto_connect = !no_auto_connect;
if !no_process_eos {
base.install_octet_interpose(Box::new(crate::interpose::eos::EosInterpose::default()));
}
}
fn new_socket(
&self,
domain: socket2::Domain,
ty: socket2::Type,
protocol: socket2::Protocol,
) -> AsynResult<socket2::Socket> {
let socket = socket2::Socket::new(domain, ty, Some(protocol))?;
if self.config.protocol.broadcast() {
socket.set_broadcast(true).map_err(|e| AsynError::Status {
status: AsynStatus::Error,
message: format!("Can't set {} socket BROADCAST option: {e}", self.host_info),
})?;
}
if self.config.protocol.reuse_port() {
#[cfg(unix)]
socket.set_reuse_port(true).map_err(|e| AsynError::Status {
status: AsynStatus::Error,
message: format!(
"Can't set {} socket SO_REUSEPORT option: {e}",
self.host_info
),
})?;
#[cfg(not(unix))]
socket
.set_reuse_address(true)
.map_err(|e| AsynError::Status {
status: AsynStatus::Error,
message: format!(
"Can't set {} socket SO_REUSEPORT option: {e}",
self.host_info
),
})?;
}
Ok(socket)
}
fn connect_tcp(&mut self) -> AsynResult<TcpStream> {
use std::net::ToSocketAddrs;
let addr_str = format!("{}:{}", self.config.host, self.config.port);
let addrs: Vec<std::net::SocketAddr> = addr_str
.to_socket_addrs()
.map_err(|e| AsynError::Status {
status: AsynStatus::Error,
message: format!("failed to resolve '{addr_str}': {e}"),
})?
.collect();
let mut last_err: Option<AsynError> = None;
for remote_addr in &addrs {
let domain = if remote_addr.is_ipv6() {
socket2::Domain::IPV6
} else {
socket2::Domain::IPV4
};
let socket =
match self.new_socket(domain, socket2::Type::STREAM, socket2::Protocol::TCP) {
Ok(s) => s,
Err(e) => {
last_err = Some(e);
continue;
}
};
if let Some(local_port) = self.config.local_port {
let local_addr: std::net::SocketAddr = if remote_addr.is_ipv6() {
(std::net::Ipv6Addr::UNSPECIFIED, local_port).into()
} else {
(std::net::Ipv4Addr::UNSPECIFIED, local_port).into()
};
if let Err(e) = socket.set_reuse_address(true) {
last_err = Some(AsynError::Io(e));
continue;
}
if let Err(e) = socket.bind(&local_addr.into()) {
last_err = Some(AsynError::Io(e));
continue;
}
}
match socket.connect_timeout(&(*remote_addr).into(), self.config.connect_timeout) {
Ok(()) => return Ok(TcpStream::from(socket)),
Err(e) => match socket.peer_addr() {
Ok(_) => return Ok(TcpStream::from(socket)),
Err(_) => last_err = Some(AsynError::Io(e)),
},
}
}
Err(last_err.unwrap_or_else(|| AsynError::Status {
status: AsynStatus::Error,
message: format!("no addresses found for '{addr_str}'"),
}))
}
fn connect_udp(&mut self) -> AsynResult<(UdpSocket, std::net::SocketAddr)> {
use std::net::ToSocketAddrs;
let remote = format!("{}:{}", self.config.host, self.config.port);
let peer = remote
.to_socket_addrs()
.map_err(|e| AsynError::Status {
status: AsynStatus::Error,
message: format!("UDP resolve '{remote}': {e}"),
})?
.next()
.ok_or_else(|| AsynError::Status {
status: AsynStatus::Error,
message: format!("UDP resolve '{remote}': no addresses"),
})?;
let local_port = self.config.local_port.unwrap_or(0);
let (domain, local_addr): (socket2::Domain, std::net::SocketAddr) = if peer.is_ipv6() {
(
socket2::Domain::IPV6,
(std::net::Ipv6Addr::UNSPECIFIED, local_port).into(),
)
} else {
(
socket2::Domain::IPV4,
(std::net::Ipv4Addr::UNSPECIFIED, local_port).into(),
)
};
let socket = self.new_socket(domain, socket2::Type::DGRAM, socket2::Protocol::UDP)?;
socket
.bind(&local_addr.into())
.map_err(|e| AsynError::Status {
status: AsynStatus::Error,
message: format!("UDP bind '{local_addr}' failed: {e}"),
})?;
Ok((UdpSocket::from(socket), peer))
}
#[cfg(unix)]
fn connect_unix(&mut self) -> AsynResult<std::os::unix::net::UnixStream> {
let stream = std::os::unix::net::UnixStream::connect(&self.config.host).map_err(|e| {
AsynError::Status {
status: AsynStatus::Error,
message: format!("unix connect to '{}': {e}", self.config.host),
}
})?;
Ok(stream)
}
}
impl PortDriver for DrvAsynIPPort {
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.inner.is_some() {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: format!("{}: Link already open!", self.base.port_name),
});
}
match self.config.protocol {
IpProtocol::Tcp | IpProtocol::TcpReusePort | IpProtocol::Com => {
let stream = self.connect_tcp()?;
if self.config.no_delay {
stream.set_nodelay(true)?;
}
self.io.inner = Some(IpIoInner::Tcp(stream));
}
IpProtocol::Udp
| IpProtocol::UdpReusePort
| IpProtocol::UdpBroadcast
| IpProtocol::UdpBroadcastReusePort => {
let (socket, peer) = self.connect_udp()?;
self.io.inner = Some(IpIoInner::Udp(socket, peer));
}
#[cfg(unix)]
IpProtocol::Unix => {
let stream = self.connect_unix()?;
self.io.inner = Some(IpIoInner::Unix(stream));
}
#[cfg(not(unix))]
IpProtocol::Unix => {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: "Unix domain sockets not supported on this platform".into(),
});
}
IpProtocol::Http => {
let stream = self.connect_tcp()?;
stream.set_nodelay(true)?;
self.io.inner = Some(IpIoInner::Tcp(stream));
}
}
self.base.set_connected(true);
asyn_trace!(
Some(self.base.trace),
&self.base.port_name,
TraceMask::FLOW,
"connected to {}:{} ({:?})",
self.config.host,
self.config.port,
self.config.protocol
);
self.restore_com_settings();
Ok(())
}
fn disconnect(&mut self, _user: &AsynUser) -> AsynResult<()> {
asyn_trace!(
Some(self.base.trace),
&self.base.port_name,
TraceMask::FLOW,
"disconnect"
);
self.io.inner = None;
self.base.set_connected(false);
Ok(())
}
fn read_octet(&mut self, user: &AsynUser, buf: &mut [u8]) -> AsynResult<usize> {
self.read_octet_core(user, buf).map(|(n, _eom)| n)
}
fn io_read_octet_eom(
&mut self,
user: &AsynUser,
buf: &mut [u8],
) -> AsynResult<(usize, EomReason)> {
self.read_octet_core(user, buf)
}
fn write_octet(&mut self, user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
if self.config.protocol == IpProtocol::Http && self.io.inner.is_none() {
self.connect(&AsynUser::default())?;
}
self.base.check_ready()?;
asyn_trace_io!(
Some(self.base.trace),
&self.base.port_name,
TraceMask::IO_DRIVER,
data,
"write"
);
match self.with_base_link(|link| link.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.with_base_link(|link| link.flush(user))
}
fn set_option(&mut self, user: &mut AsynUser, key: &str, value: &str) -> AsynResult<()> {
if self.com.is_some() && ComPortOptions::owns_key(key) {
return self
.with_negotiation(|com, link| com.set_option(user, link, key, value))
.expect("com is Some");
}
if key.eq_ignore_ascii_case("noDelay") {
let enabled = parse_yn_option("noDelay", value)?;
self.config.no_delay = enabled;
if let Some(IpIoInner::Tcp(ref stream)) = self.io.inner {
stream.set_nodelay(enabled)?;
}
} else if key.eq_ignore_ascii_case("disconnectOnReadTimeout") {
self.disconnect_on_read_timeout = parse_yn_option("disconnectOnReadTimeout", value)?;
} else if key.eq_ignore_ascii_case("hostInfo") {
let new_config = IpPortConfig::parse(value)?;
if self.io.inner.is_some() {
self.io.inner = None;
self.base.set_connected(false);
std::thread::sleep(Duration::from_millis(20));
}
self.config.host = new_config.host;
self.config.port = new_config.port;
self.config.local_port = new_config.local_port;
self.config.protocol = new_config.protocol;
self.host_info = value.to_string();
} else if !key.is_empty() {
return Err(AsynError::OptionNotFound(key.to_string()));
}
Ok(())
}
fn get_option(&self, key: &str) -> AsynResult<String> {
if let Some(com) = self.com.as_ref() {
if ComPortOptions::owns_key(key) {
return com.options.get_option(key);
}
}
if key.eq_ignore_ascii_case("disconnectOnReadTimeout") {
Ok(if self.disconnect_on_read_timeout {
"Y"
} else {
"N"
}
.to_string())
} else if key.eq_ignore_ascii_case("hostInfo") {
Ok(self.host_info.clone())
} else {
self.base
.options
.get(key)
.cloned()
.ok_or_else(|| AsynError::OptionNotFound(key.to_string()))
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::net::TcpListener;
use std::thread;
fn chain_read(drv: &mut DrvAsynIPPort, user: &AsynUser, buf: &mut [u8]) -> AsynResult<usize> {
crate::port::octet_read_chain(drv, user, buf).map(|(n, _eom)| n)
}
fn chain_read_eom(
drv: &mut DrvAsynIPPort,
user: &AsynUser,
buf: &mut [u8],
) -> AsynResult<(usize, EomReason)> {
crate::port::octet_read_chain(drv, user, buf)
}
fn chain_flush(drv: &mut DrvAsynIPPort, user: &mut AsynUser) -> AsynResult<()> {
crate::port::octet_flush_chain(drv, user)
}
#[test]
fn test_parse_tcp_default() {
let cfg = IpPortConfig::parse("localhost:5025").unwrap();
assert_eq!(cfg.host, "localhost");
assert_eq!(cfg.port, 5025);
assert_eq!(cfg.protocol, IpProtocol::Tcp);
assert_eq!(cfg.local_port, None);
}
#[test]
fn test_parse_tcp_explicit() {
let cfg = IpPortConfig::parse("192.168.1.1:8080 tcp").unwrap();
assert_eq!(cfg.host, "192.168.1.1");
assert_eq!(cfg.port, 8080);
assert_eq!(cfg.protocol, IpProtocol::Tcp);
}
#[test]
fn test_parse_udp() {
let cfg = IpPortConfig::parse("device:9000 udp").unwrap();
assert_eq!(cfg.protocol, IpProtocol::Udp);
}
#[test]
fn test_parse_local_port() {
let cfg = IpPortConfig::parse("host:5025:4000").unwrap();
assert_eq!(cfg.local_port, Some(4000));
}
#[test]
fn test_parse_invalid_no_port() {
assert!(IpPortConfig::parse("hostname_only").is_err());
}
#[test]
fn test_parse_invalid_port_number() {
assert!(IpPortConfig::parse("host:abc").is_err());
}
#[test]
fn test_parse_empty_host() {
assert!(IpPortConfig::parse(":5025").is_err());
}
#[test]
fn test_driver_initial_state() {
let drv = DrvAsynIPPort::new("iptest", "localhost:5025").unwrap();
assert!(!drv.base().is_connected());
assert!(drv.base().auto_connect);
assert!(drv.base().flags.can_block);
}
#[test]
fn new_configured_no_auto_connect_matches_the_oracle_port_model() {
let drv = DrvAsynIPPort::new_configured("ORACLEASYN", "localhost:1", true, false).unwrap();
assert_eq!(
drv.capabilities(),
crate::interfaces::octet_transport_capabilities()
);
assert!(!drv.base().is_connected());
assert!(!drv.base().auto_connect);
assert_eq!(drv.get_option("hostinfo").unwrap(), "localhost:1");
assert_eq!(drv.get_option("disconnectOnReadTimeout").unwrap(), "N");
}
fn start_echo_server() -> (TcpListener, u16) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
(listener, port)
}
#[test]
fn test_connect_disconnect() {
let (listener, port) = start_echo_server();
let _handle = thread::spawn(move || {
let _ = listener.accept();
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).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 test_read_write_octet_roundtrip() {
let (listener, port) = start_echo_server();
let handle = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
let mut buf = [0u8; 256];
let n = stream.read(&mut buf).unwrap();
stream.write_all(&buf[..n]).unwrap();
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).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 user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut buf = [0u8; 32];
let n = drv.read_octet(&user, &mut buf).unwrap();
assert_eq!(&buf[..n], b"hello");
handle.join().unwrap();
}
#[test]
fn http_multi_segment_response_not_truncated() {
let (listener, port) = start_echo_server();
let handle = thread::spawn(move || {
let (mut s1, _) = listener.accept().unwrap();
let mut req = [0u8; 64];
let _ = s1.read(&mut req).unwrap();
s1.write_all(b"AAAABBBB").unwrap();
drop(s1); let (mut s2, _) = listener.accept().unwrap();
let _ = s2.read(&mut req).unwrap();
s2.write_all(b"CCCC").unwrap();
drop(s2);
});
let mut drv = DrvAsynIPPort::new("httptest", &format!("127.0.0.1:{port} HTTP")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
assert!(drv.base().is_connected());
let mut wuser = AsynUser::new(0).with_timeout(Duration::from_secs(2));
drv.write_octet(&mut wuser, b"GET /1\r\n").unwrap();
let ruser = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut buf = [0u8; 4];
let (n1, eom1) = drv.io_read_octet_eom(&ruser, &mut buf).unwrap();
assert_eq!(&buf[..n1], b"AAAA");
assert!(!eom1.contains(EomReason::END), "first chunk is not EOF");
assert!(
drv.base().is_connected(),
"HTTP must stay connected mid-response (no truncation)"
);
assert!(
drv.io.inner.is_some(),
"socket must remain open mid-response"
);
let (n2, _) = drv.io_read_octet_eom(&ruser, &mut buf).unwrap();
assert_eq!(
&buf[..n2],
b"BBBB",
"second chunk must be the rest of the response"
);
let (n3, eom3) = drv.io_read_octet_eom(&ruser, &mut buf).unwrap();
assert_eq!(n3, 0);
assert!(eom3.contains(EomReason::END));
assert!(
drv.base().is_connected(),
"HTTP logical connection must not flap on per-transaction EOF"
);
assert!(drv.io.inner.is_none(), "socket released after EOF");
drv.write_octet(&mut wuser, b"GET /2\r\n").unwrap();
assert!(
drv.io.inner.is_some(),
"socket must reopen lazily for the next transaction"
);
let mut buf2 = [0u8; 16];
let (n4, _) = drv.io_read_octet_eom(&ruser, &mut buf2).unwrap();
assert_eq!(&buf2[..n4], b"CCCC");
handle.join().unwrap();
}
#[test]
fn test_read_timeout() {
let (listener, port) = start_echo_server();
let _handle = thread::spawn(move || {
let (_stream, _) = listener.accept().unwrap();
thread::sleep(Duration::from_secs(5));
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).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_server_disconnect_eof() {
let (listener, port) = start_echo_server();
let handle = thread::spawn(move || {
let (stream, _) = listener.accept().unwrap();
drop(stream);
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
thread::sleep(Duration::from_millis(50));
let user = AsynUser::new(0).with_timeout(Duration::from_secs(1));
let mut buf = [0u8; 32];
let (n, eom) = drv.io_read_octet_eom(&user, &mut buf).unwrap();
assert_eq!(n, 0);
assert!(eom.contains(EomReason::END));
assert!(!drv.base().is_connected());
handle.join().unwrap();
}
#[test]
fn test_is_fatal_transport_error_classification() {
assert!(
AsynError::Status {
status: AsynStatus::Disconnected,
message: "EOF".into(),
}
.is_fatal_transport()
);
assert!(
AsynError::Io(std::io::Error::new(
std::io::ErrorKind::ConnectionReset,
"rst"
))
.is_fatal_transport()
);
assert!(
!AsynError::Status {
status: AsynStatus::Timeout,
message: "read timeout".into(),
}
.is_fatal_transport()
);
}
#[test]
fn test_write_error_disconnects() {
let (listener, port) = start_echo_server();
let handle = thread::spawn(move || {
let (stream, _) = listener.accept().unwrap();
drop(stream); });
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).unwrap();
drv.connect(&AsynUser::default()).unwrap();
assert!(drv.base().is_connected());
handle.join().unwrap();
thread::sleep(Duration::from_millis(50));
let mut user = AsynUser::new(0).with_timeout(Duration::from_secs(1));
let mut last: AsynResult<usize> = Ok(0);
for _ in 0..200 {
last = drv.write_octet(&mut user, b"ping\n");
if last.is_err() {
break;
}
thread::sleep(Duration::from_millis(5));
}
assert!(last.is_err(), "expected a write to the dead peer to fail");
assert!(
!drv.base().is_connected(),
"DRV-5: fatal write error must set connected=false"
);
}
#[test]
fn test_partial_read() {
let (listener, port) = start_echo_server();
let handle = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream.write_all(b"he").unwrap();
stream.flush().unwrap();
thread::sleep(Duration::from_millis(50));
stream.write_all(b"llo").unwrap();
stream.flush().unwrap();
thread::sleep(Duration::from_millis(200));
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
let user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut buf = [0u8; 32];
let n1 = drv.read_octet(&user, &mut buf).unwrap();
assert!(n1 > 0);
assert!(n1 <= 5);
handle.join().unwrap();
}
#[test]
fn com_port_negotiates_rfc2217_on_connect_and_escapes_the_data_path() {
use crate::interpose::com::{IAC, SB, SE};
const COM_PORT_OPTION: u8 = 44;
const WILL_B: u8 = 251;
const DO_B: u8 = 253;
fn ack(subcmd: u8, values: &[u8]) -> Vec<u8> {
let mut v = vec![IAC, SB, COM_PORT_OPTION, subcmd + 100];
v.extend_from_slice(values);
v.extend_from_slice(&[IAC, SE]);
v
}
fn sb(payload: &[u8]) -> Vec<u8> {
let mut v = vec![IAC, SB, COM_PORT_OPTION];
v.extend_from_slice(payload);
v.extend_from_slice(&[IAC, SE]);
v
}
let (listener, port) = start_echo_server();
let server = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(2)))
.unwrap();
let mut reply = vec![IAC, WILL_B, 0, IAC, DO_B, 0, IAC, DO_B, COM_PORT_OPTION];
reply.extend_from_slice(&ack(11, &[0])); reply.extend_from_slice(&ack(1, &[0x00, 0x00, 0x25, 0x80])); reply.extend_from_slice(&ack(2, &[8])); reply.extend_from_slice(&ack(3, &[1])); reply.extend_from_slice(&ack(4, &[1])); reply.extend_from_slice(&ack(5, &[1])); reply.extend_from_slice(&ack(5, &[1])); stream.write_all(&reply).unwrap();
stream.flush().unwrap();
let mut want = vec![
IAC,
DO_B,
0, IAC,
WILL_B,
0, IAC,
WILL_B,
COM_PORT_OPTION, ];
want.extend_from_slice(&sb(&[11, 0])); want.extend_from_slice(&sb(&[1, 0x00, 0x00, 0x25, 0x80])); want.extend_from_slice(&sb(&[2, 8]));
want.extend_from_slice(&sb(&[3, 1]));
want.extend_from_slice(&sb(&[4, 1]));
want.extend_from_slice(&sb(&[5, 1])); want.extend_from_slice(&sb(&[5, 1])); let mut got = vec![0u8; want.len()];
stream.read_exact(&mut got).unwrap();
assert_eq!(got, want, "connect must emit C's RFC 2217 handshake");
let mut got = [0u8; 4];
stream.read_exact(&mut got).unwrap();
assert_eq!(got, [b'A', IAC, IAC, b'B'], "device write must be stuffed");
stream.write_all(&[b'X', IAC, IAC, b'Y']).unwrap();
stream.flush().unwrap();
thread::sleep(Duration::from_millis(200));
});
let mut drv = DrvAsynIPPort::new("comtest", &format!("127.0.0.1:{port} COM")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
assert_eq!(drv.get_option("baud").unwrap(), "9600");
assert_eq!(drv.get_option("bits").unwrap(), "8");
assert_eq!(drv.get_option("parity").unwrap(), "none");
assert_eq!(drv.get_option("stop").unwrap(), "1");
let mut wuser = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let n = drv.write_octet(&mut wuser, &[b'A', IAC, b'B']).unwrap();
assert_eq!(n, 3);
let ruser = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut buf = [0u8; 16];
let n = drv.read_octet(&ruser, &mut buf).unwrap();
assert_eq!(&buf[..n], &[b'X', IAC, b'Y']);
server.join().unwrap();
}
#[test]
fn set_option_on_a_live_com_port_negotiates_with_the_server() {
use crate::interpose::com::{IAC, SB, SE};
let (listener, port) = start_echo_server();
let server = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(2)))
.unwrap();
stream.write_all(&[IAC, 252, 0]).unwrap(); stream.flush().unwrap();
let mut hs = [0u8; 3];
stream.read_exact(&mut hs).unwrap();
assert_eq!(hs, [IAC, 253, 0]);
let mut req = [0u8; 10];
stream.read_exact(&mut req).unwrap();
assert_eq!(
req,
[IAC, SB, 44, 1, 0x00, 0x01, 0xC2, 0x00, IAC, SE],
"115200 baud must go out as a big-endian SET-BAUDRATE"
);
let mut ack = vec![IAC, SB, 44, 101, 0x00, 0x01, 0xC2, 0x00];
ack.extend_from_slice(&[IAC, SE]);
stream.write_all(&ack).unwrap();
stream.flush().unwrap();
thread::sleep(Duration::from_millis(200));
});
let mut drv = DrvAsynIPPort::new("comtest2", &format!("127.0.0.1:{port} COM")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
assert!(drv.base.is_connected());
drv.set_option(&mut AsynUser::default(), "baud", "115200")
.unwrap();
assert_eq!(drv.get_option("baud").unwrap(), "115200");
server.join().unwrap();
}
#[test]
fn a_timed_out_negotiation_read_honours_disconnect_on_read_timeout() {
fn silent_server() -> (u16, thread::JoinHandle<()>) {
let (listener, port) = start_echo_server();
let handle = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_millis(100)))
.unwrap();
let deadline = std::time::Instant::now() + Duration::from_millis(2600);
let mut scratch = [0u8; 64];
while std::time::Instant::now() < deadline {
match stream.read(&mut scratch) {
Ok(0) => break,
Ok(_) => continue,
Err(_) => continue, }
}
});
(port, handle)
}
let (port, server) = silent_server();
let mut drv = DrvAsynIPPort::new("comto1", &format!("127.0.0.1:{port} COM")).unwrap();
drv.connect(&AsynUser::default()).unwrap();
assert!(
drv.base.is_connected(),
"handshake failure alone must not close the socket (C only asynPrints)"
);
assert!(drv.io.inner.is_some());
server.join().unwrap();
let (port, server) = silent_server();
let mut drv = DrvAsynIPPort::new("comto2", &format!("127.0.0.1:{port} COM")).unwrap();
drv.set_option(&mut AsynUser::default(), "disconnectOnReadTimeout", "Y")
.unwrap();
drv.connect(&AsynUser::default()).unwrap();
assert!(
!drv.base.is_connected(),
"a timed-out negotiation read must close the socket, as C readRaw does"
);
assert!(drv.io.inner.is_none());
server.join().unwrap();
}
#[test]
fn test_eos_interpose_with_tcp() {
use crate::interpose::eos::{EosConfig, EosInterpose};
let (listener, port) = start_echo_server();
let handle = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream.write_all(b"OK\r\n").unwrap();
stream.flush().unwrap();
thread::sleep(Duration::from_millis(200));
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).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 user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut buf = [0u8; 32];
let n = chain_read(&mut drv, &user, &mut buf).unwrap();
assert_eq!(&buf[..n], b"OK");
handle.join().unwrap();
}
#[test]
fn disconnect_on_read_timeout_fires_through_the_eos_interpose() {
use crate::interpose::eos::{EosConfig, EosInterpose};
let (listener, port) = start_echo_server();
let handle = thread::spawn(move || {
let (stream, _) = listener.accept().unwrap();
thread::sleep(Duration::from_millis(300));
drop(stream);
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).unwrap();
drv.install_interpose(Box::new(EosInterpose::new(EosConfig {
input_eos: vec![b'\n'],
output_eos: vec![],
})));
drv.set_option(&mut AsynUser::default(), "disconnectOnReadTimeout", "Y")
.unwrap();
drv.connect(&AsynUser::default()).unwrap();
assert!(drv.base().is_connected());
let user = AsynUser::new(0).with_timeout(Duration::from_millis(50));
let mut buf = [0u8; 32];
let err = drv
.read_octet(&user, &mut buf)
.expect_err("a silent peer must time the read out");
assert_eq!(
err.status(),
AsynStatus::Error,
"C readRaw overwrites the status on the disconnect branch"
);
assert!(
err.message().contains("read error"),
"C's teardown text (drvAsynIPPort.c:801-803), got {:?}",
err.message()
);
assert!(
!drv.base().is_connected(),
"disconnectOnReadTimeout must drop the connection (C drvAsynIPPort.c:798-806)"
);
handle.join().unwrap();
}
#[test]
fn read_timeout_without_the_option_keeps_the_connection() {
use crate::interpose::eos::{EosConfig, EosInterpose};
let (listener, port) = start_echo_server();
let handle = thread::spawn(move || {
let (stream, _) = listener.accept().unwrap();
thread::sleep(Duration::from_millis(300));
drop(stream);
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).unwrap();
drv.install_interpose(Box::new(EosInterpose::new(EosConfig {
input_eos: vec![b'\n'],
output_eos: vec![],
})));
drv.connect(&AsynUser::default()).unwrap();
let user = AsynUser::new(0).with_timeout(Duration::from_millis(50));
let mut buf = [0u8; 32];
let err = drv.read_octet(&user, &mut buf).unwrap_err();
assert_eq!(err.status(), AsynStatus::Timeout);
assert!(
drv.base().is_connected(),
"a plain read timeout leaves the socket intact (C returns asynTimeout, no close)"
);
handle.join().unwrap();
}
#[test]
fn a_torn_down_read_reports_asyn_error_and_still_carries_the_partial_bytes() {
use crate::interpose::PartialOctetRead;
let timeout = Duration::from_millis(50);
let partial = PartialOctetRead {
data: b"OK\n".to_vec(),
eom_reason: EomReason::empty(),
};
let timed_out = AsynError::Status {
status: AsynStatus::Timeout,
message: "read timeout".into(),
}
.with_partial_read(partial);
let f = classify_read_failure(timed_out, true, timeout, "127.0.0.1:5001");
assert!(f.disconnect, "disconnectOnReadTimeout + timeout > 0");
assert_eq!(f.error.status(), AsynStatus::Error, "C :805");
assert_eq!(
f.error.message(),
"127.0.0.1:5001 read error: read timeout",
"C :801-803"
);
assert_eq!(
f.error.partial_read().map(|p| p.nbytes_transferred()),
Some(3),
"the transfer C writes out at :824 survives the restamp"
);
}
#[test]
fn the_read_failure_rule_matches_c_should_disconnect_at_every_boundary() {
let timeout = || AsynError::Status {
status: AsynStatus::Timeout,
message: "read timeout".into(),
};
let fatal = || AsynError::Io(std::io::Error::other("cable yanked"));
let dev = "host:1";
let f = classify_read_failure(timeout(), false, Duration::from_secs(1), dev);
assert!(!f.disconnect);
assert_eq!(f.error.status(), AsynStatus::Timeout);
let f = classify_read_failure(timeout(), true, Duration::ZERO, dev);
assert!(!f.disconnect);
assert_eq!(f.error.status(), AsynStatus::Timeout);
let f = classify_read_failure(timeout(), true, Duration::from_millis(1), dev);
assert!(f.disconnect);
assert_eq!(f.error.status(), AsynStatus::Error);
let f = classify_read_failure(fatal(), false, Duration::ZERO, dev);
assert!(f.disconnect);
assert_eq!(f.error.status(), AsynStatus::Error);
assert_eq!(f.error.message(), "host:1 read error: IO: cable yanked");
}
#[test]
fn test_read_write_when_disconnected() {
let mut drv = DrvAsynIPPort::new("iptest", "127.0.0.1:9999").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_set_option_nodelay() {
let mut drv = DrvAsynIPPort::new("iptest", "127.0.0.1:5025").unwrap();
drv.set_option(&mut AsynUser::default(), "noDelay", "Y")
.unwrap();
assert!(drv.config.no_delay);
drv.set_option(&mut AsynUser::default(), "noDelay", "n")
.unwrap();
assert!(!drv.config.no_delay);
assert!(
drv.set_option(&mut AsynUser::default(), "noDelay", "1")
.is_err()
);
assert!(
drv.set_option(&mut AsynUser::default(), "noDelay", "maybe")
.is_err()
);
}
#[test]
fn test_udp_connect_and_roundtrip() {
let server = UdpSocket::bind("127.0.0.1:0").unwrap();
let server_port = server.local_addr().unwrap().port();
let handle = thread::spawn(move || {
let mut buf = [0u8; 256];
let (n, src) = server.recv_from(&mut buf).unwrap();
server.send_to(&buf[..n], src).unwrap();
});
let mut drv =
DrvAsynIPPort::new("udptest", &format!("127.0.0.1:{server_port} udp")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
assert!(drv.base().is_connected());
let mut user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
drv.write_octet(&mut user, b"ping").unwrap();
let user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut buf = [0u8; 32];
let n = drv.read_octet(&user, &mut buf).unwrap();
assert_eq!(&buf[..n], b"ping");
handle.join().unwrap();
}
#[test]
fn test_udp_accepts_reply_from_any_peer() {
let server = UdpSocket::bind("127.0.0.1:0").unwrap();
let server_port = server.local_addr().unwrap().port();
let local_port = UdpSocket::bind("127.0.0.1:0")
.unwrap()
.local_addr()
.unwrap()
.port();
let handle = thread::spawn(move || {
let mut buf = [0u8; 256];
let (n, _src) = server.recv_from(&mut buf).unwrap();
assert_eq!(&buf[..n], b"ping");
let other = UdpSocket::bind("127.0.0.1:0").unwrap();
other
.send_to(b"pong", format!("127.0.0.1:{local_port}"))
.unwrap();
});
let mut drv = DrvAsynIPPort::new(
"udptest",
&format!("127.0.0.1:{server_port}:{local_port} udp"),
)
.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"ping").unwrap();
let user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut buf = [0u8; 32];
let n = drv.read_octet(&user, &mut buf).unwrap();
assert_eq!(&buf[..n], b"pong");
handle.join().unwrap();
}
#[test]
fn test_udp_empty_datagram_is_not_eof() {
let server = UdpSocket::bind("127.0.0.1:0").unwrap();
let server_port = server.local_addr().unwrap().port();
let mut drv =
DrvAsynIPPort::new("udptest", &format!("127.0.0.1:{server_port} udp")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
let mut wuser = AsynUser::new(0).with_timeout(Duration::from_secs(2));
drv.write_octet(&mut wuser, b"hello").unwrap();
let mut sbuf = [0u8; 16];
let (_n, src) = server.recv_from(&mut sbuf).unwrap();
server.send_to(&[], src).unwrap();
let ruser = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut buf = [0u8; 32];
let n = drv.read_octet(&ruser, &mut buf).unwrap();
assert_eq!(n, 0);
assert!(drv.base().is_connected());
server.send_to(b"world", src).unwrap();
let n = drv.read_octet(&ruser, &mut buf).unwrap();
assert_eq!(&buf[..n], b"world");
}
#[test]
fn test_udp_zero_length_write_sends_nothing() {
let server = UdpSocket::bind("127.0.0.1:0").unwrap();
server
.set_read_timeout(Some(Duration::from_millis(200)))
.unwrap();
let server_port = server.local_addr().unwrap().port();
let mut drv =
DrvAsynIPPort::new("udptest", &format!("127.0.0.1:{server_port} udp")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
let mut wuser = AsynUser::new(0).with_timeout(Duration::from_secs(1));
drv.write_octet(&mut wuser, b"").unwrap();
let mut sbuf = [0u8; 16];
assert!(
server.recv_from(&mut sbuf).is_err(),
"zero-length write must not emit a datagram"
);
drv.write_octet(&mut wuser, b"data").unwrap();
let (n, _src) = server.recv_from(&mut sbuf).unwrap();
assert_eq!(&sbuf[..n], b"data");
}
#[test]
fn test_disconnect_on_read_timeout() {
let (listener, port) = start_echo_server();
let _handle = thread::spawn(move || {
let (_stream, _) = listener.accept().unwrap();
thread::sleep(Duration::from_secs(5));
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).unwrap();
drv.set_option(&mut AsynUser::default(), "disconnectOnReadTimeout", "Y")
.unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
assert!(drv.base().is_connected());
let user = AsynUser::new(0).with_timeout(Duration::from_millis(50));
let mut buf = [0u8; 32];
let _ = drv.read_octet(&user, &mut buf);
assert!(!drv.base().is_connected());
}
#[test]
fn socket_poll_timeout_floors_zero_to_one_ms() {
assert_eq!(
socket_poll_timeout(Duration::ZERO),
Duration::from_millis(1)
);
assert_eq!(
socket_poll_timeout(Duration::from_micros(500)),
Duration::from_micros(500)
);
assert_eq!(
socket_poll_timeout(Duration::from_secs(2)),
Duration::from_secs(2)
);
}
#[test]
fn zero_timeout_read_polls_without_disconnect() {
let (listener, port) = start_echo_server();
let _handle = thread::spawn(move || {
let (_stream, _) = listener.accept().unwrap();
thread::sleep(Duration::from_secs(5));
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).unwrap();
drv.set_option(&mut AsynUser::default(), "disconnectOnReadTimeout", "Y")
.unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
assert!(drv.base().is_connected());
let user = AsynUser::new(0).with_timeout(Duration::ZERO);
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 (not a Duration::ZERO error), got {other:?}"),
}
assert!(
drv.base().is_connected(),
"zero-timeout read must not disconnect even with disconnectOnReadTimeout=Y"
);
}
#[test]
fn classify_read_error_eintr_and_wouldblock_are_nonfatal_timeout() {
use std::io::{Error, ErrorKind};
for kind in [
ErrorKind::TimedOut,
ErrorKind::WouldBlock,
ErrorKind::Interrupted,
] {
let err = classify_read_error(Error::from(kind));
match err {
AsynError::Status {
status: AsynStatus::Timeout,
..
} => {}
other => panic!("{kind:?} must map to a non-fatal Timeout, got {other:?}"),
}
assert!(
!classify_read_error(Error::from(kind)).is_fatal_transport(),
"{kind:?} must not be a fatal transport error"
);
}
let reset = classify_read_error(Error::from(ErrorKind::ConnectionReset));
assert!(matches!(reset, AsynError::Io(_)));
assert!(reset.is_fatal_transport());
}
#[test]
fn zero_timeout_write_attempts_send_not_instant_timeout() {
let (listener, port) = start_echo_server();
let handle = thread::spawn(move || {
let (mut s, _) = listener.accept().unwrap();
let mut buf = [0u8; 16];
let _ = s.read(&mut buf);
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
let mut wuser = AsynUser::new(0).with_timeout(Duration::ZERO);
let n = drv.write_octet(&mut wuser, b"x").unwrap();
assert_eq!(
n, 1,
"zero-timeout write of a writable socket must attempt the send, not instant-timeout"
);
handle.join().unwrap();
}
#[test]
fn zero_length_read_request_rejected_not_eof_teardown() {
let (listener, port) = start_echo_server();
let handle = thread::spawn(move || {
let (_s, _) = listener.accept().unwrap();
thread::sleep(Duration::from_millis(50));
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
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 read must be rejected with asynError, got {res:?}"
);
assert!(
drv.base().is_connected(),
"zero-length read must not tear down the connection"
);
handle.join().unwrap();
}
#[test]
fn test_disconnect_on_read_timeout_value_validation() {
let mut drv = DrvAsynIPPort::new("iptest", "127.0.0.1:5025 tcp").unwrap();
drv.set_option(&mut AsynUser::default(), "disconnectOnReadTimeout", "Y")
.unwrap();
assert_eq!(drv.get_option("disconnectOnReadTimeout").unwrap(), "Y");
drv.set_option(&mut AsynUser::default(), "disconnectOnReadTimeout", "n")
.unwrap();
assert_eq!(drv.get_option("disconnectOnReadTimeout").unwrap(), "N");
for bad in ["1", "yes", "true", "", "maybe"] {
assert!(
drv.set_option(&mut AsynUser::default(), "disconnectOnReadTimeout", bad)
.is_err(),
"value {bad:?} should be rejected"
);
}
assert_eq!(drv.get_option("disconnectOnReadTimeout").unwrap(), "N");
}
#[test]
fn connect_rejects_double_open() {
let (listener, port) = start_echo_server();
let _handle = thread::spawn(move || {
let _ = listener.accept();
thread::sleep(Duration::from_secs(1));
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
assert!(drv.base().is_connected());
let err = drv.connect(&user).unwrap_err();
assert!(matches!(err, AsynError::Status { .. }));
assert!(drv.base().is_connected());
}
#[test]
fn test_set_option_host_info() {
let mut drv = DrvAsynIPPort::new("iptest", "127.0.0.1:5025").unwrap();
drv.set_option(&mut AsynUser::default(), "hostInfo", "192.168.1.1:8080")
.unwrap();
assert_eq!(drv.config.host, "192.168.1.1");
assert_eq!(drv.config.port, 8080);
}
#[test]
fn test_set_option_host_info_disconnects() {
let (listener, port) = start_echo_server();
let _handle = thread::spawn(move || {
let _ = listener.accept();
thread::sleep(Duration::from_secs(1));
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
assert!(drv.base().is_connected());
drv.set_option(&mut AsynUser::default(), "hostInfo", "127.0.0.1:9999")
.unwrap();
assert!(!drv.base().is_connected());
assert_eq!(drv.config.port, 9999);
}
#[test]
fn host_info_reparse_updates_protocol_and_flags() {
let mut drv = DrvAsynIPPort::new("iptest", "127.0.0.1:5025 tcp").unwrap();
assert_eq!(drv.config.protocol, IpProtocol::Tcp);
drv.set_option(&mut AsynUser::default(), "hostInfo", "127.0.0.1:5026 udp")
.unwrap();
assert_eq!(drv.config.protocol, IpProtocol::Udp);
assert_eq!(drv.config.port, 5026);
drv.set_option(&mut AsynUser::default(), "hostInfo", "127.0.0.1:5027 udp*")
.unwrap();
assert_eq!(drv.config.protocol, IpProtocol::UdpBroadcast);
drv.set_option(&mut AsynUser::default(), "hostInfo", "127.0.0.1:5028 udp&")
.unwrap();
assert_eq!(drv.config.protocol, IpProtocol::UdpReusePort);
drv.set_option(&mut AsynUser::default(), "hostInfo", "127.0.0.1:5029 udp*&")
.unwrap();
assert_eq!(drv.config.protocol, IpProtocol::UdpBroadcastReusePort);
drv.set_option(&mut AsynUser::default(), "hostInfo", "127.0.0.1:5030 tcp&")
.unwrap();
assert_eq!(drv.config.protocol, IpProtocol::TcpReusePort);
}
#[test]
fn host_info_reparse_clears_omitted_local_port() {
let mut drv = DrvAsynIPPort::new("iptest", "127.0.0.1:5025:12345 tcp").unwrap();
assert_eq!(drv.config.local_port, Some(12345));
drv.set_option(&mut AsynUser::default(), "hostInfo", "127.0.0.1:5026 tcp")
.unwrap();
assert_eq!(
drv.config.local_port, None,
"local_port must reset on hostInfo reparse"
);
}
#[test]
fn host_info_option_key_is_case_insensitive_for_get_and_set() {
let mut drv = DrvAsynIPPort::new("iptest", "127.0.0.1:5025 tcp").unwrap();
drv.set_option(&mut AsynUser::default(), "hostinfo", "10.0.0.5:1234 udp")
.unwrap();
assert_eq!(drv.config.host, "10.0.0.5");
assert_eq!(drv.config.port, 1234);
assert_eq!(drv.config.protocol, IpProtocol::Udp);
assert_eq!(drv.get_option("hostinfo").unwrap(), "10.0.0.5:1234 udp");
assert_eq!(drv.get_option("hostInfo").unwrap(), "10.0.0.5:1234 udp");
drv.set_option(&mut AsynUser::default(), "DISCONNECTONREADTIMEOUT", "Y")
.unwrap();
assert_eq!(drv.get_option("disconnectonreadtimeout").unwrap(), "Y");
}
#[test]
fn unsupported_option_key_is_rejected() {
let mut drv = DrvAsynIPPort::new("iptest", "127.0.0.1:5025 tcp").unwrap();
let err = drv
.set_option(&mut AsynUser::default(), "bogusKey", "value")
.unwrap_err();
assert!(matches!(err, AsynError::OptionNotFound(_)));
assert_eq!(err.message(), "Unsupported key \"bogusKey\"");
assert!(drv.get_option("bogusKey").is_err());
drv.set_option(&mut AsynUser::default(), "", "ignored")
.unwrap();
}
#[test]
fn every_ip_option_reports_cs_text() {
let mut drv = DrvAsynIPPort::new("iptest", "127.0.0.1:5025 tcp").unwrap();
assert_eq!(
drv.set_option(&mut AsynUser::default(), "disconnectOnReadTimeout", "maybe")
.unwrap_err()
.message(),
"Invalid disconnectOnReadTimeout value."
);
assert_eq!(
drv.set_option(&mut AsynUser::default(), "noDelay", "1")
.unwrap_err()
.message(),
"Invalid noDelay value."
);
drv.set_option(&mut AsynUser::default(), "disconnectOnReadTimeout", "y")
.unwrap();
assert!(drv.disconnect_on_read_timeout);
}
#[test]
fn test_parse_tcp_reuse_port() {
let cfg = IpPortConfig::parse("host:5025 TCP&").unwrap();
assert_eq!(cfg.protocol, IpProtocol::TcpReusePort);
assert_eq!(cfg.host, "host");
assert_eq!(cfg.port, 5025);
}
#[test]
fn test_parse_tcp_reuse_port_lowercase() {
let cfg = IpPortConfig::parse("host:5025 tcp&").unwrap();
assert_eq!(cfg.protocol, IpProtocol::TcpReusePort);
}
#[test]
fn test_parse_udp_reuse_port() {
let cfg = IpPortConfig::parse("192.168.1.10:9000 UDP&").unwrap();
assert_eq!(cfg.protocol, IpProtocol::UdpReusePort);
assert_eq!(cfg.host, "192.168.1.10");
}
#[test]
fn test_parse_udp_broadcast() {
let cfg = IpPortConfig::parse("192.168.1.255:9000 UDP*").unwrap();
assert_eq!(cfg.protocol, IpProtocol::UdpBroadcast);
assert_eq!(cfg.host, "192.168.1.255");
}
#[test]
fn test_parse_udp_broadcast_reuse_port() {
let cfg = IpPortConfig::parse("192.168.1.255:9000 UDP*&").unwrap();
assert_eq!(cfg.protocol, IpProtocol::UdpBroadcastReusePort);
}
#[test]
fn test_parse_unix_socket() {
let cfg = IpPortConfig::parse("unix:///tmp/asyn.sock").unwrap();
assert_eq!(cfg.protocol, IpProtocol::Unix);
assert_eq!(cfg.host, "/tmp/asyn.sock");
assert_eq!(cfg.port, 0);
}
#[test]
fn test_parse_unix_empty_path() {
assert!(IpPortConfig::parse("unix://").is_err());
}
#[test]
fn test_parse_ipv6_brackets() {
let cfg = IpPortConfig::parse("[::1]:5025").unwrap();
assert_eq!(cfg.host, "::1");
assert_eq!(cfg.port, 5025);
assert_eq!(cfg.protocol, IpProtocol::Tcp);
}
#[test]
fn test_parse_ipv6_with_local_port() {
let cfg = IpPortConfig::parse("[::1]:5025:4000").unwrap();
assert_eq!(cfg.host, "::1");
assert_eq!(cfg.port, 5025);
assert_eq!(cfg.local_port, Some(4000));
}
#[test]
fn test_parse_ipv6_with_proto() {
let cfg = IpPortConfig::parse("[fe80::1]:9000 UDP").unwrap();
assert_eq!(cfg.host, "fe80::1");
assert_eq!(cfg.port, 9000);
assert_eq!(cfg.protocol, IpProtocol::Udp);
}
#[test]
fn test_parse_case_insensitive() {
assert_eq!(
IpPortConfig::parse("h:1 Tcp").unwrap().protocol,
IpProtocol::Tcp
);
assert_eq!(
IpPortConfig::parse("h:1 Udp").unwrap().protocol,
IpProtocol::Udp
);
assert_eq!(
IpPortConfig::parse("h:1 Tcp&").unwrap().protocol,
IpProtocol::TcpReusePort
);
assert_eq!(
IpPortConfig::parse("h:1 Udp&").unwrap().protocol,
IpProtocol::UdpReusePort
);
assert_eq!(
IpPortConfig::parse("h:1 Udp*").unwrap().protocol,
IpProtocol::UdpBroadcast
);
assert_eq!(
IpPortConfig::parse("h:1 Udp*&").unwrap().protocol,
IpProtocol::UdpBroadcastReusePort
);
}
#[test]
fn com_protocol_is_accepted_case_insensitively() {
for spec in ["1.2.3.4:5000 COM", "1.2.3.4:5000 com", "host:23 Com"] {
let cfg = IpPortConfig::parse(spec).unwrap();
assert_eq!(cfg.protocol, IpProtocol::Com, "{spec}");
}
let cfg = IpPortConfig::parse("1.2.3.4:5000 COM").unwrap();
assert_eq!(cfg.host, "1.2.3.4");
assert_eq!(cfg.port, 5000);
}
#[test]
fn com_port_installs_the_iac_layer_at_configure_time() {
let com = DrvAsynIPPort::new("comport", "1.2.3.4:5000 COM").unwrap();
assert!(com.com.is_some());
assert_eq!(
com.base.interpose_octet.len(),
0,
"the IAC layer is the base link, not a reorderable stack entry"
);
let tcp = DrvAsynIPPort::new("tcpport", "1.2.3.4:5000 TCP").unwrap();
assert!(tcp.com.is_none());
assert_eq!(tcp.base.interpose_octet.len(), 0);
}
#[test]
fn eos_pushed_after_com_sits_above_it_and_sees_unstuffed_bytes() {
use crate::interpose::com::IAC;
let (listener, port) = start_echo_server();
let server = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(2)))
.unwrap();
stream.write_all(&[IAC, 252, 0]).unwrap(); stream.flush().unwrap();
let mut hs = [0u8; 3];
stream.read_exact(&mut hs).unwrap();
stream.write_all(&[b'A', IAC, IAC, b'B', b'\n']).unwrap();
stream.flush().unwrap();
thread::sleep(Duration::from_millis(300));
});
let mut drv = crate::iocsh::build_configured_ip_port(
"com_eos",
&format!("127.0.0.1:{port} COM"),
false,
false,
)
.unwrap();
assert_eq!(
drv.base.interpose_octet.len(),
1,
"EOS is the only stack layer; COM is the base"
);
drv.connect(&AsynUser::default()).unwrap();
drv.set_input_eos(&AsynUser::default(), b"\n").unwrap();
let user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut buf = [0u8; 3];
let (n, _eom) = chain_read_eom(&mut drv, &user, &mut buf).unwrap();
assert_eq!(
&buf[..n],
&[b'A', IAC, b'B'],
"the 3-byte read must carry 3 bytes of device data — an escape byte \
must not eat one of them, which is what happens when COM sits above EOS"
);
server.join().unwrap();
}
#[test]
fn com_option_interface_is_layered_over_the_ip_drivers() {
let com = DrvAsynIPPort::new("comport", "1.2.3.4:5000 COM").unwrap();
assert_eq!(com.get_option("baud").unwrap(), "9600");
assert_eq!(com.get_option("parity").unwrap(), "none");
assert_eq!(com.get_option("crtscts").unwrap(), "N");
assert_eq!(com.get_option("hostInfo").unwrap(), "1.2.3.4:5000 COM");
assert_eq!(com.get_option("disconnectOnReadTimeout").unwrap(), "N");
let tcp = DrvAsynIPPort::new("tcpport", "1.2.3.4:5000 TCP").unwrap();
assert!(tcp.get_option("baud").is_err());
}
#[test]
fn unknown_protocol_token_is_rejected_by_name() {
let err = IpPortConfig::parse("1.2.3.4:5000 SCTP").unwrap_err();
assert_eq!(err.message(), "Unknown protocol \"SCTP\".");
let err = IpPortConfig::parse("1.2.3.4:5000 tcpsocket").unwrap_err();
assert_eq!(err.message(), "Unknown protocol \"tcpso\".");
}
#[test]
fn only_the_first_token_after_the_blank_is_the_protocol() {
let cfg = IpPortConfig::parse("1.2.3.4:5000 UDP extra junk").unwrap();
assert_eq!(cfg.host, "1.2.3.4");
assert_eq!(cfg.port, 5000);
assert_eq!(cfg.protocol, IpProtocol::Udp);
}
#[cfg(unix)]
#[test]
fn test_unix_socket_connect_roundtrip() {
use std::os::unix::net::UnixListener;
let sock_path = format!("/tmp/asyn_test_{}.sock", std::process::id());
let _ = std::fs::remove_file(&sock_path);
let listener = UnixListener::bind(&sock_path).unwrap();
let sock_path2 = sock_path.clone();
let handle = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
let mut buf = [0u8; 256];
let n = stream.read(&mut buf).unwrap();
stream.write_all(&buf[..n]).unwrap();
});
let mut drv = DrvAsynIPPort::new("unixtest", &format!("unix://{sock_path}")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
assert!(drv.base().is_connected());
let mut user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
drv.write_octet(&mut user, b"unix_hello").unwrap();
let user = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut buf = [0u8; 32];
let n = drv.read_octet(&user, &mut buf).unwrap();
assert_eq!(&buf[..n], b"unix_hello");
handle.join().unwrap();
let _ = std::fs::remove_file(&sock_path2);
}
#[test]
fn io_flush_drains_stale_tcp_input() {
let (listener, port) = start_echo_server();
let handle = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream.write_all(b"STALE_PROMPT>").unwrap();
stream.flush().unwrap();
thread::sleep(Duration::from_millis(100));
let mut buf = [0u8; 64];
let n = stream.read(&mut buf).unwrap();
assert_eq!(&buf[..n], b"CMD");
stream.write_all(b"RESPONSE").unwrap();
stream.flush().unwrap();
thread::sleep(Duration::from_millis(100));
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).unwrap();
let user = AsynUser::default();
drv.connect(&user).unwrap();
thread::sleep(Duration::from_millis(50));
let mut fuser = AsynUser::new(0).with_timeout(Duration::from_secs(2));
drv.io_flush(&mut fuser).unwrap();
let mut wuser = AsynUser::new(0).with_timeout(Duration::from_secs(2));
drv.write_octet(&mut wuser, b"CMD").unwrap();
let ruser = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut buf = [0u8; 64];
let n = drv.read_octet(&ruser, &mut buf).unwrap();
assert_eq!(
&buf[..n],
b"RESPONSE",
"io_flush must drain stale input; got {:?}",
String::from_utf8_lossy(&buf[..n])
);
handle.join().unwrap();
}
#[test]
fn io_flush_resets_eos_interpose_buffer() {
use crate::interpose::eos::{EosConfig, EosInterpose};
let (listener, port) = start_echo_server();
let handle = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream.write_all(b"OLD_LINE_DATA\n").unwrap();
stream.flush().unwrap();
thread::sleep(Duration::from_millis(150));
let mut buf = [0u8; 64];
let n = stream.read(&mut buf).unwrap();
assert_eq!(&buf[..n], b"CMD");
stream.write_all(b"NEW\n").unwrap();
stream.flush().unwrap();
thread::sleep(Duration::from_millis(150));
});
let mut drv = DrvAsynIPPort::new("iptest", &format!("127.0.0.1:{port}")).unwrap();
drv.install_interpose(Box::new(EosInterpose::new(EosConfig {
input_eos: vec![b'\n'],
output_eos: vec![],
})));
let user = AsynUser::default();
drv.connect(&user).unwrap();
thread::sleep(Duration::from_millis(50));
let ruser = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut small = [0u8; 4];
let n = chain_read(&mut drv, &ruser, &mut small).unwrap();
assert_eq!(&small[..n], b"OLD_");
let mut fuser = AsynUser::new(0).with_timeout(Duration::from_secs(2));
chain_flush(&mut drv, &mut fuser).unwrap();
let mut wuser = AsynUser::new(0).with_timeout(Duration::from_secs(2));
drv.write_octet(&mut wuser, b"CMD").unwrap();
let ruser2 = AsynUser::new(0).with_timeout(Duration::from_secs(2));
let mut buf = [0u8; 64];
let n = chain_read(&mut drv, &ruser2, &mut buf).unwrap();
assert_eq!(
&buf[..n],
b"NEW",
"io_flush must reset the EOS interpose buffer; got {:?}",
String::from_utf8_lossy(&buf[..n])
);
handle.join().unwrap();
}
#[test]
fn io_flush_noop_when_disconnected() {
let mut drv = DrvAsynIPPort::new("iptest", "127.0.0.1:9999").unwrap();
let mut user = AsynUser::new(0);
drv.io_flush(&mut user).unwrap();
}
#[test]
fn test_udp_broadcast_flag() {
let cfg = IpPortConfig::parse("255.255.255.255:9000 UDP*").unwrap();
let drv = DrvAsynIPPort::new("bcast_test", "255.255.255.255:9000 UDP*").unwrap();
assert_eq!(cfg.protocol, IpProtocol::UdpBroadcast);
assert_eq!(drv.config.protocol, IpProtocol::UdpBroadcast);
}
}