use std::time::Duration;
use oms_modbus::*;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
fn rtu_frame(slave: u8, pdu: &[u8]) -> Vec<u8> {
let mut data = vec![slave];
data.extend_from_slice(pdu);
let crc = calculate_crc(&data);
data.push(crc as u8);
data.push((crc >> 8) as u8);
data
}
fn ascii_frame(data: &[u8]) -> String {
let lrc = calculate_lrc(data);
let mut frame = String::with_capacity(1 + data.len() * 2 + 2 + 2);
frame.push(':');
for &b in data {
frame.push_str(&format!("{:02X}", b));
}
frame.push_str(&format!("{:02X}", lrc));
frame.push_str("\r\n");
frame
}
#[tokio::test]
async fn rtu_crc_bitflip_response_recovery() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let n = server_stream.read(&mut buf).await.unwrap();
assert!(n >= 8, "expected RTU request, got {n} bytes");
let slave = buf[0];
let pdu = [0x03, 0x02, 0x00, 0x2A];
let mut frame = rtu_frame(slave, &pdu);
let last = frame.len() - 1;
frame[last] ^= 0x01;
server_stream.write_all(&frame).await.unwrap();
let n = server_stream.read(&mut buf).await.unwrap();
assert!(n >= 8);
let slave = buf[0];
let frame = rtu_frame(slave, &pdu);
server_stream.write_all(&frame).await.unwrap();
});
let client = rtu::RtuClient::with_timeout(client_stream, Duration::from_secs(3));
let err = client.read_holding_registers(1, 0, 1).await.unwrap_err();
assert!(
err.detail().contains("CRC"),
"expected CRC error, got: {}",
err.detail()
);
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![42]);
server.await.unwrap();
}
#[tokio::test]
async fn rtu_crc_error_isolated_per_call() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
let slave = buf[0];
let pdu = [0x03, 0x02, 0x00, 0x2A];
let mut frame = rtu_frame(slave, &pdu);
let idx = frame.len() - 2;
frame[idx] ^= 0xFF; server_stream.write_all(&frame).await.unwrap();
let _n = server_stream.read(&mut buf).await.unwrap();
let slave = buf[0];
server_stream
.write_all(&rtu_frame(slave, &pdu))
.await
.unwrap();
let _n = server_stream.read(&mut buf).await.unwrap();
let slave = buf[0];
let mut frame = rtu_frame(slave, &pdu);
let last = frame.len() - 1;
frame[last] ^= 0x80;
server_stream.write_all(&frame).await.unwrap();
let _n = server_stream.read(&mut buf).await.unwrap();
let slave = buf[0];
server_stream
.write_all(&rtu_frame(slave, &pdu))
.await
.unwrap();
});
let client = rtu::RtuClient::with_timeout(client_stream, Duration::from_secs(3));
let r1 = client.read_holding_registers(1, 0, 1).await;
assert!(r1.unwrap_err().detail().contains("CRC"));
let r2 = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(r2, vec![42]);
let r3 = client.read_holding_registers(1, 0, 1).await;
assert!(r3.unwrap_err().detail().contains("CRC"));
let r4 = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(r4, vec![42]);
server.await.unwrap();
}
#[tokio::test]
async fn rtu_fragmented_response_byte_by_byte() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
let slave = buf[0];
let pdu = [0x03, 0x02, 0x00, 0x2A]; let frame = rtu_frame(slave, &pdu);
for &b in &frame {
server_stream.write_all(&[b]).await.unwrap();
tokio::time::sleep(Duration::from_millis(1)).await;
}
});
let client = rtu::RtuClient::with_timeout(client_stream, Duration::from_secs(3));
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![42]);
server.await.unwrap();
}
#[tokio::test]
async fn rtu_fragmented_frame_three_then_four() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
let slave = buf[0];
let pdu = [0x03, 0x02, 0x00, 0x2A];
let frame = rtu_frame(slave, &pdu);
server_stream.write_all(&frame[..3]).await.unwrap();
tokio::time::sleep(Duration::from_millis(5)).await;
server_stream.write_all(&frame[3..]).await.unwrap();
});
let client = rtu::RtuClient::with_timeout(client_stream, Duration::from_secs(3));
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![42]);
server.await.unwrap();
}
#[tokio::test]
async fn rtu_exception_response_short_frame() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
let slave = buf[0];
let pdu = [0x83, 0x02];
let frame = rtu_frame(slave, &pdu); assert_eq!(frame.len(), 5);
server_stream.write_all(&frame).await.unwrap();
});
let client = rtu::RtuClient::with_timeout(client_stream, Duration::from_secs(3));
let result = client.read_holding_registers(1, 0, 1).await;
let err = result.expect_err("short frame should produce exception or timeout");
let detail = err.detail();
assert!(
detail.contains("Illegal") || detail.contains("timeout") || detail.contains("TIMEOUT"),
"unexpected error: {detail}"
);
server.await.unwrap();
}
#[tokio::test]
async fn rtu_exception_from_diagnostic() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
let slave = buf[0];
let pdu = [0x88, 0x01];
let frame = rtu_frame(slave, &pdu);
assert_eq!(frame.len(), 5);
server_stream.write_all(&frame).await.unwrap();
});
let client = rtu::RtuClient::with_timeout(client_stream, Duration::from_secs(3));
let result = client.diagnostic(1, 0x0001, 0).await;
let err = result.expect_err("diagnostic exception should produce error");
let detail = err.detail();
assert!(
detail.contains("Illegal") || detail.contains("timeout"),
"unexpected error: {detail}"
);
server.await.unwrap();
}
#[tokio::test]
async fn rtu_garbage_then_valid_frame() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let (tx, rx) = tokio::sync::oneshot::channel::<()>();
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
let slave = buf[0];
let pdu = [0x03, 0x02, 0x00, 0x2A];
server_stream
.write_all(&rtu_frame(slave, &pdu))
.await
.unwrap();
let _ = rx.await;
server_stream
.write_all(&[0xFF, 0x00, 0xAA, 0x55])
.await
.unwrap();
let _n = server_stream.read(&mut buf).await.unwrap();
let pdu2 = [0x03, 0x02, 0x00, 0x63]; server_stream
.write_all(&rtu_frame(buf[0], &pdu2))
.await
.unwrap();
});
let client = rtu::RtuClient::with_timeout(client_stream, Duration::from_secs(3));
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![42]);
tx.send(()).unwrap();
tokio::time::sleep(Duration::from_millis(10)).await;
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![99]);
server.await.unwrap();
}
#[tokio::test]
async fn rtu_stale_data_burst_drained_then_request() {
let (client_stream, mut server_stream) = tokio::io::duplex(8192);
let (tx, rx) = tokio::sync::oneshot::channel::<()>();
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
let slave = buf[0];
let pdu = [0x03, 0x02, 0x00, 0x2A];
server_stream
.write_all(&rtu_frame(slave, &pdu))
.await
.unwrap();
let _ = rx.await;
server_stream.write_all(&[0xFF; 128]).await.unwrap();
let _n = server_stream.read(&mut buf).await.unwrap();
let pdu2 = [0x03, 0x02, 0x00, 0x63]; server_stream
.write_all(&rtu_frame(buf[0], &pdu2))
.await
.unwrap();
});
let client = rtu::RtuClient::with_timeout(client_stream, Duration::from_secs(3));
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![42]);
tx.send(()).unwrap();
tokio::time::sleep(Duration::from_millis(10)).await;
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![99]);
server.await.unwrap();
}
#[tokio::test]
async fn ascii_lrc_error_recovery() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
server_stream.write_all(b":010302002AFF\r\n").await.unwrap();
let _n = server_stream.read(&mut buf).await.unwrap();
let data = [0x01, 0x03, 0x02, 0x00, 0x2A];
server_stream
.write_all(ascii_frame(&data).as_bytes())
.await
.unwrap();
});
let client = ascii::AsciiClient::with_timeout(client_stream, Duration::from_millis(200));
let result = client.read_holding_registers(1, 0, 1).await;
assert!(result.is_err());
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![42]);
server.await.unwrap();
}
#[tokio::test]
async fn ascii_hex_decode_error_recovery() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
server_stream.write_all(b":01GG02002AD0\r\n").await.unwrap();
let _n = server_stream.read(&mut buf).await.unwrap();
let data = [0x01, 0x03, 0x02, 0x00, 0x2A];
server_stream
.write_all(ascii_frame(&data).as_bytes())
.await
.unwrap();
});
let client = ascii::AsciiClient::with_timeout(client_stream, Duration::from_millis(200));
let result = client.read_holding_registers(1, 0, 1).await;
assert!(result.is_err());
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![42]);
server.await.unwrap();
}
#[tokio::test]
async fn ascii_garbage_no_colon_prefix() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
server_stream.write_all(b"RANDOM NOISE\r\n").await.unwrap();
let _n = server_stream.read(&mut buf).await.unwrap();
let data = [0x01, 0x03, 0x02, 0x00, 0x2A];
server_stream
.write_all(ascii_frame(&data).as_bytes())
.await
.unwrap();
});
let client = ascii::AsciiClient::with_timeout(client_stream, Duration::from_millis(200));
let result = client.read_holding_registers(1, 0, 1).await;
assert!(result.is_err());
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![42]);
server.await.unwrap();
}
#[tokio::test]
async fn ascii_partial_frame_no_crlf() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
server_stream.write_all(b":010302002AD0").await.unwrap();
let _n = server_stream.read(&mut buf).await.unwrap();
let data = [0x01, 0x03, 0x02, 0x00, 0x63]; server_stream
.write_all(ascii_frame(&data).as_bytes())
.await
.unwrap();
});
let client = ascii::AsciiClient::with_timeout(client_stream, Duration::from_millis(200));
let result = client.read_holding_registers(1, 0, 1).await;
assert!(result.is_err());
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![99]);
server.await.unwrap();
}
#[tokio::test]
async fn ascii_basic_read_holding() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
let data = [0x01, 0x03, 0x02, 0x00, 0x2A]; let frame = oms_modbus::codec::encode_ascii_frame(&data);
server_stream.write_all(&frame).await.unwrap();
tokio::time::sleep(Duration::from_millis(200)).await;
});
let client = ascii::AsciiClient::with_timeout(client_stream, Duration::from_millis(300));
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![42]);
server.await.unwrap();
}
#[tokio::test]
async fn ascii_exception_response_decoded() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
let data = [0x01, 0x83, 0x02];
let frame = oms_modbus::codec::encode_ascii_frame(&data);
server_stream.write_all(&frame).await.unwrap();
tokio::time::sleep(Duration::from_millis(200)).await;
});
let client = ascii::AsciiClient::with_timeout(client_stream, Duration::from_millis(300));
let result = client.read_holding_registers(1, 0, 1).await;
let err = result.expect_err("ASCII exception should produce error");
let detail = err.detail();
assert!(
detail.contains("Illegal") || detail.contains("timeout"),
"unexpected error: {detail}"
);
server.await.unwrap();
}
#[tokio::test]
async fn rtu_slave_id_mismatch_error() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
let pdu = [0x03, 0x02, 0x00, 0x2A];
server_stream.write_all(&rtu_frame(2, &pdu)).await.unwrap();
let _n = server_stream.read(&mut buf).await.unwrap();
let pdu2 = [0x03, 0x02, 0x00, 0x63];
server_stream.write_all(&rtu_frame(1, &pdu2)).await.unwrap();
});
let client = rtu::RtuClient::with_timeout(client_stream, Duration::from_millis(300));
let err = client.read_holding_registers(1, 0, 1).await.unwrap_err();
assert!(
err.detail().contains("slave ID mismatch"),
"expected slave ID mismatch, got: {}",
err.detail()
);
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![99]);
server.await.unwrap();
}
#[tokio::test]
async fn ascii_slave_id_mismatch_error() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
let data = [0x02, 0x03, 0x02, 0x00, 0x2A];
let frame = oms_modbus::codec::encode_ascii_frame(&data);
server_stream.write_all(&frame).await.unwrap();
let _n = server_stream.read(&mut buf).await.unwrap();
let data2 = [0x01, 0x03, 0x02, 0x00, 0x63];
let frame2 = oms_modbus::codec::encode_ascii_frame(&data2);
server_stream.write_all(&frame2).await.unwrap();
});
let client = ascii::AsciiClient::with_timeout(client_stream, Duration::from_millis(300));
let err = client.read_holding_registers(1, 0, 1).await.unwrap_err();
assert!(
err.detail().contains("slave ID mismatch"),
"expected slave ID mismatch, got: {}",
err.detail()
);
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![99]);
server.await.unwrap();
}
#[tokio::test]
async fn ascii_cr_without_lf_not_accepted() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
server_stream.write_all(b":010302002AD0\r").await.unwrap();
let _n = server_stream.read(&mut buf).await.unwrap();
let data = [0x01, 0x03, 0x02, 0x00, 0x63];
server_stream
.write_all(ascii_frame(&data).as_bytes())
.await
.unwrap();
});
let client = ascii::AsciiClient::with_timeout(client_stream, Duration::from_millis(200));
let result = client.read_holding_registers(1, 0, 1).await;
assert!(result.is_err());
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![99]);
server.await.unwrap();
}
#[tokio::test]
async fn ascii_multiple_frames_one_read() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
let data = [0x01, 0x03, 0x02, 0x00, 0x2A];
let frame = oms_modbus::codec::encode_ascii_frame(&data);
server_stream.write_all(&frame).await.unwrap();
server_stream.write_all(&frame).await.unwrap();
tokio::time::sleep(Duration::from_millis(200)).await;
});
let client = ascii::AsciiClient::with_timeout(client_stream, Duration::from_millis(300));
let regs = client.read_holding_registers(1, 0, 1).await.unwrap();
assert_eq!(regs, vec![42]);
server.await.unwrap();
}
#[tokio::test]
async fn rtu_frame_shorter_than_5_bytes_rejected() {
let (client_stream, mut server_stream) = tokio::io::duplex(64);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
server_stream.write_all(&[0x01, 0x03, 0x02]).await.unwrap();
tokio::time::sleep(Duration::from_millis(300)).await;
});
let client = rtu::RtuClient::with_timeout(client_stream, Duration::from_millis(100));
let err = client.read_holding_registers(1, 0, 1).await.unwrap_err();
let detail = err.detail();
assert!(
detail.contains("timed out") || detail.contains("TIMEOUT"),
"expected timeout for <5 byte frame, got: {detail}"
);
server.await.unwrap();
}
#[tokio::test]
async fn rtu_exactly_4_bytes_then_close() {
let (client_stream, mut server_stream) = tokio::io::duplex(64);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
server_stream
.write_all(&[0x01, 0x03, 0x02, 0x00])
.await
.unwrap();
drop(server_stream);
});
let client = rtu::RtuClient::with_timeout(client_stream, Duration::from_millis(200));
let err = client.read_holding_registers(1, 0, 1).await.unwrap_err();
let detail = err.detail();
assert!(
detail.contains("timed out") || detail.contains("TIMEOUT") || detail.contains("closed"),
"expected timeout or connection-closed for 4-byte + close, got: {detail}"
);
server.await.unwrap();
}
#[tokio::test]
async fn rtu_minimum_valid_frame_5_bytes() {
let (client_stream, mut server_stream) = tokio::io::duplex(1024);
let server = tokio::spawn(async move {
let mut buf = [0u8; 256];
let _n = server_stream.read(&mut buf).await.unwrap();
let slave = buf[0];
let pdu = [0x83, 0x02];
server_stream
.write_all(&rtu_frame(slave, &pdu))
.await
.unwrap();
assert_eq!(rtu_frame(slave, &pdu).len(), 5);
});
let client = rtu::RtuClient::with_timeout(client_stream, Duration::from_secs(3));
let result = client.read_holding_registers(1, 0, 1).await;
let err = result.expect_err("5-byte frame should produce exception or timeout");
let detail = err.detail();
assert!(
detail.contains("Illegal") || detail.contains("timeout"),
"expected exception or timeout from 5-byte frame, got: {detail}"
);
server.await.unwrap();
}