use crate::error::{AsynError, AsynResult, AsynStatus};
use crate::interpose::{OctetInterpose, OctetNext, OctetReadResult};
use crate::user::AsynUser;
const ECHO_STR_SIZE: usize = 16;
fn escaped(bytes: &[u8]) -> String {
crate::escape::escaped_from_raw(bytes, ECHO_STR_SIZE)
}
pub struct EchoInterpose;
impl EchoInterpose {
pub fn new() -> Self {
Self
}
}
impl Default for EchoInterpose {
fn default() -> Self {
Self::new()
}
}
impl OctetInterpose for EchoInterpose {
fn read(
&mut self,
user: &AsynUser,
buf: &mut [u8],
next: &mut dyn OctetNext,
) -> AsynResult<OctetReadResult> {
next.read(user, buf)
}
fn write(
&mut self,
user: &mut AsynUser,
data: &[u8],
next: &mut dyn OctetNext,
) -> AsynResult<usize> {
let mut total = 0;
for byte in data {
let n = match next.write(user, std::slice::from_ref(byte)) {
Ok(n) => n,
Err(e) => return Err(e.with_partial_write(total)),
};
if n != 1 {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: format!("wrote {n} chars instead of 1"),
}
.with_partial_write(total));
}
let mut echo_buf = [0u8; 1];
let echo_user = AsynUser::new(user.reason)
.with_addr(user.addr)
.with_timeout(user.timeout);
let echo_result = match next.read(&echo_user, &mut echo_buf) {
Ok(r) => r,
Err(e) if e.status() == AsynStatus::Timeout => {
return Err(AsynError::Status {
status: AsynStatus::Timeout,
message: format!("timeout reading back char number {total}"),
}
.with_partial_write(total));
}
Err(e) => return Err(e.with_partial_write(total)),
};
let echo = &echo_buf[..echo_result.nbytes_transferred.min(1)];
if echo_result.nbytes_transferred != 1 || echo[0] != *byte {
return Err(AsynError::Status {
status: AsynStatus::Error,
message: format!(
"got back '{}' instead of '{}'",
escaped(echo),
escaped(std::slice::from_ref(byte))
),
}
.with_partial_write(total));
}
total += n;
}
Ok(total)
}
fn flush(&mut self, user: &mut AsynUser, next: &mut dyn OctetNext) -> AsynResult<()> {
next.flush(user)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::interpose::{EomReason, OctetInterposeStack, OctetNext, OctetReadResult};
use crate::user::AsynUser;
use std::collections::VecDeque;
struct EchoBase {
echo_queue: VecDeque<u8>,
written: Vec<u8>,
}
impl EchoBase {
fn new() -> Self {
Self {
echo_queue: VecDeque::new(),
written: Vec::new(),
}
}
}
impl OctetNext for EchoBase {
fn read(&mut self, _user: &AsynUser, buf: &mut [u8]) -> AsynResult<OctetReadResult> {
if let Some(b) = self.echo_queue.pop_front() {
buf[0] = b;
Ok(OctetReadResult {
nbytes_transferred: 1,
eom_reason: EomReason::CNT,
})
} else {
Err(AsynError::Status {
status: AsynStatus::Timeout,
message: "no echo data".into(),
})
}
}
fn write(&mut self, _user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
for &b in data {
self.written.push(b);
self.echo_queue.push_back(b); }
Ok(data.len())
}
fn flush(&mut self, _user: &mut AsynUser) -> AsynResult<()> {
Ok(())
}
}
#[test]
fn test_echo_success() {
let mut stack = OctetInterposeStack::new(false);
stack.install(-1, Box::new(EchoInterpose::new()));
let mut base = EchoBase::new();
let mut user = AsynUser::default();
let n = stack.dispatch_write(&mut user, b"OK", &mut base).unwrap();
assert_eq!(n, 2);
assert_eq!(&base.written, b"OK");
}
#[test]
fn test_echo_mismatch() {
struct BadEchoBase;
impl OctetNext for BadEchoBase {
fn read(&mut self, _user: &AsynUser, buf: &mut [u8]) -> AsynResult<OctetReadResult> {
buf[0] = b'X'; Ok(OctetReadResult {
nbytes_transferred: 1,
eom_reason: EomReason::CNT,
})
}
fn write(&mut self, _user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
Ok(data.len())
}
fn flush(&mut self, _user: &mut AsynUser) -> AsynResult<()> {
Ok(())
}
}
let mut stack = OctetInterposeStack::new(false);
stack.install(-1, Box::new(EchoInterpose::new()));
let mut base = BadEchoBase;
let mut user = AsynUser::default();
let err = stack
.dispatch_write(&mut user, b"A", &mut base)
.unwrap_err();
assert_eq!(err.status(), AsynStatus::Error);
assert_eq!(err.message(), "got back 'X' instead of 'A'");
assert_eq!(err.partial_write(), Some(0));
}
#[test]
fn test_echo_no_response() {
struct NoEchoBase;
impl OctetNext for NoEchoBase {
fn read(&mut self, _user: &AsynUser, _buf: &mut [u8]) -> AsynResult<OctetReadResult> {
Err(AsynError::Status {
status: AsynStatus::Timeout,
message: "timeout".into(),
})
}
fn write(&mut self, _user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
Ok(data.len())
}
fn flush(&mut self, _user: &mut AsynUser) -> AsynResult<()> {
Ok(())
}
}
let mut stack = OctetInterposeStack::new(false);
stack.install(-1, Box::new(EchoInterpose::new()));
let mut base = NoEchoBase;
let mut user = AsynUser::default();
let err = stack
.dispatch_write(&mut user, b"A", &mut base)
.unwrap_err();
assert_eq!(err.status(), AsynStatus::Timeout);
assert_eq!(err.message(), "timeout reading back char number 0");
assert_eq!(err.partial_write(), Some(0));
}
#[test]
fn timeout_message_names_the_char_whose_echo_was_lost() {
struct EchoesOnlyTheFirst {
written: usize,
}
impl OctetNext for EchoesOnlyTheFirst {
fn read(&mut self, _user: &AsynUser, buf: &mut [u8]) -> AsynResult<OctetReadResult> {
if self.written == 1 {
buf[0] = b'A';
return Ok(OctetReadResult {
nbytes_transferred: 1,
eom_reason: EomReason::CNT,
});
}
Err(AsynError::Status {
status: AsynStatus::Timeout,
message: "timeout".into(),
})
}
fn write(&mut self, _user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
self.written += data.len();
Ok(data.len())
}
fn flush(&mut self, _user: &mut AsynUser) -> AsynResult<()> {
Ok(())
}
}
let mut stack = OctetInterposeStack::new(false);
stack.install(-1, Box::new(EchoInterpose::new()));
let mut base = EchoesOnlyTheFirst { written: 0 };
let mut user = AsynUser::default();
let err = stack
.dispatch_write(&mut user, b"AB", &mut base)
.unwrap_err();
assert_eq!(err.status(), AsynStatus::Timeout);
assert_eq!(err.message(), "timeout reading back char number 1");
assert_eq!(err.partial_write(), Some(1));
}
#[test]
fn non_timeout_read_error_propagates_untouched() {
struct BrokenReadBase;
impl OctetNext for BrokenReadBase {
fn read(&mut self, _user: &AsynUser, _buf: &mut [u8]) -> AsynResult<OctetReadResult> {
Err(AsynError::Status {
status: AsynStatus::Disconnected,
message: "port disconnected".into(),
})
}
fn write(&mut self, _user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
Ok(data.len())
}
fn flush(&mut self, _user: &mut AsynUser) -> AsynResult<()> {
Ok(())
}
}
let mut stack = OctetInterposeStack::new(false);
stack.install(-1, Box::new(EchoInterpose::new()));
let mut base = BrokenReadBase;
let mut user = AsynUser::default();
let err = stack
.dispatch_write(&mut user, b"A", &mut base)
.unwrap_err();
assert_eq!(err.status(), AsynStatus::Disconnected);
assert_eq!(err.message(), "port disconnected");
assert_eq!(err.partial_write(), Some(0));
}
#[test]
fn short_echo_and_control_chars_use_the_escaped_mismatch_message() {
struct ShortEchoBase;
impl OctetNext for ShortEchoBase {
fn read(&mut self, _user: &AsynUser, _buf: &mut [u8]) -> AsynResult<OctetReadResult> {
Ok(OctetReadResult {
nbytes_transferred: 0,
eom_reason: EomReason::CNT,
})
}
fn write(&mut self, _user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
Ok(data.len())
}
fn flush(&mut self, _user: &mut AsynUser) -> AsynResult<()> {
Ok(())
}
}
let mut stack = OctetInterposeStack::new(false);
stack.install(-1, Box::new(EchoInterpose::new()));
let mut base = ShortEchoBase;
let mut user = AsynUser::default();
let err = stack
.dispatch_write(&mut user, b"\n", &mut base)
.unwrap_err();
assert_eq!(err.status(), AsynStatus::Error);
assert_eq!(err.message(), "got back '' instead of '\\n'");
}
#[test]
fn short_write_reports_the_char_count_c_reports() {
struct ShortWriteBase;
impl OctetNext for ShortWriteBase {
fn read(&mut self, _user: &AsynUser, _buf: &mut [u8]) -> AsynResult<OctetReadResult> {
unreachable!("write never succeeds, so no echo is read")
}
fn write(&mut self, _user: &mut AsynUser, _data: &[u8]) -> AsynResult<usize> {
Ok(0)
}
fn flush(&mut self, _user: &mut AsynUser) -> AsynResult<()> {
Ok(())
}
}
let mut stack = OctetInterposeStack::new(false);
stack.install(-1, Box::new(EchoInterpose::new()));
let mut base = ShortWriteBase;
let mut user = AsynUser::default();
let err = stack
.dispatch_write(&mut user, b"A", &mut base)
.unwrap_err();
assert_eq!(err.status(), AsynStatus::Error);
assert_eq!(err.message(), "wrote 0 chars instead of 1");
assert_eq!(err.partial_write(), Some(0));
}
}