use std::collections::VecDeque;
use std::io::{self, Read, Write};
use std::time::Duration;
use serialport::{ClearBuffer, SerialPort};
use crate::error::{Error, Result};
use crate::protocol::BAUD;
pub trait Transport {
fn send(&mut self, data: &[u8]) -> Result<()>;
fn send_recv(&mut self, data: &[u8], wait: Duration) -> Result<Vec<u8>>;
fn pace(&self, d: Duration) -> Duration {
d
}
}
pub struct SerialTransport {
name: String,
port: Box<dyn SerialPort>,
low_latency: bool,
}
impl SerialTransport {
pub fn open(path: &str, timeout: Duration) -> Result<Self> {
let builder = serialport::new(path, BAUD)
.data_bits(serialport::DataBits::Eight)
.parity(serialport::Parity::None)
.stop_bits(serialport::StopBits::One)
.timeout(timeout);
let serial_err = |source| Error::Serial {
port: path.to_owned(),
source,
};
#[cfg(target_os = "linux")]
let (port, low_latency): (Box<dyn SerialPort>, bool) = {
use std::os::fd::AsRawFd;
let native = builder.open_native().map_err(serial_err)?;
let low_latency = crate::low_latency::enable(native.as_raw_fd()).is_ok();
(Box::new(native), low_latency)
};
#[cfg(not(target_os = "linux"))]
let (port, low_latency): (Box<dyn SerialPort>, bool) =
(builder.open().map_err(serial_err)?, false);
Ok(Self {
name: path.to_owned(),
port,
low_latency,
})
}
pub fn name(&self) -> &str {
&self.name
}
pub fn low_latency(&self) -> bool {
self.low_latency
}
fn serial_err(&self, source: serialport::Error) -> Error {
Error::Serial {
port: self.name.clone(),
source,
}
}
}
impl std::fmt::Debug for SerialTransport {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SerialTransport")
.field("name", &self.name)
.field("low_latency", &self.low_latency)
.finish_non_exhaustive()
}
}
impl Transport for SerialTransport {
fn send(&mut self, data: &[u8]) -> Result<()> {
self.port.write_all(data)?;
Ok(())
}
fn send_recv(&mut self, data: &[u8], wait: Duration) -> Result<Vec<u8>> {
self.port
.clear(ClearBuffer::Input)
.map_err(|e| self.serial_err(e))?;
self.port.write_all(data)?;
let wire = Duration::from_secs_f64(data.len() as f64 * 10.0 / f64::from(BAUD));
std::thread::sleep(wire + wait);
let n = self.port.bytes_to_read().map_err(|e| self.serial_err(e))? as usize;
let mut buf = vec![0u8; n];
let mut filled = 0;
while filled < n {
match self.port.read(&mut buf[filled..]) {
Ok(0) => break,
Ok(k) => filled += k,
Err(e) if e.kind() == io::ErrorKind::Interrupted => continue,
Err(e) if e.kind() == io::ErrorKind::TimedOut => break,
Err(e) => return Err(e.into()),
}
}
buf.truncate(filled);
Ok(buf)
}
}
#[derive(Debug, Default)]
pub struct MockTransport {
pub sent: Vec<Vec<u8>>,
pub replies: VecDeque<Vec<u8>>,
pub echo_tx: bool,
pub echo_truncate: Option<usize>,
pub fail_io: bool,
}
impl MockTransport {
pub fn with_replies<I>(replies: I) -> Self
where
I: IntoIterator<Item = Vec<u8>>,
{
Self {
replies: replies.into_iter().collect(),
..Self::default()
}
}
fn check_fail(&self) -> Result<()> {
if self.fail_io {
Err(Error::Io(io::Error::new(
io::ErrorKind::BrokenPipe,
"mock transport failure",
)))
} else {
Ok(())
}
}
}
impl Transport for MockTransport {
fn send(&mut self, data: &[u8]) -> Result<()> {
self.sent.push(data.to_vec());
self.check_fail()
}
fn send_recv(&mut self, data: &[u8], _wait: Duration) -> Result<Vec<u8>> {
self.sent.push(data.to_vec());
self.check_fail()?;
let reply = self.replies.pop_front().unwrap_or_default();
if self.echo_tx {
let keep = self.echo_truncate.unwrap_or(data.len()).min(data.len());
let mut echoed = data[..keep].to_vec();
echoed.extend_from_slice(&reply);
Ok(echoed)
} else {
Ok(reply)
}
}
fn pace(&self, _d: Duration) -> Duration {
Duration::ZERO
}
}