#![allow(dead_code)]
pub mod com;
pub mod delay;
pub mod echo;
pub mod eos;
pub mod flush;
use bitflags::bitflags;
use crate::error::AsynResult;
use crate::user::AsynUser;
bitflags! {
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct EomReason: u32 {
const CNT = 0x01;
const EOS = 0x02;
const END = 0x04;
}
}
#[derive(Debug, Clone)]
pub struct OctetReadResult {
pub nbytes_transferred: usize,
pub eom_reason: EomReason,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PartialOctetRead {
pub data: Vec<u8>,
pub eom_reason: EomReason,
}
impl PartialOctetRead {
pub fn nbytes_transferred(&self) -> usize {
self.data.len()
}
}
pub trait OctetNext: Send + Sync {
fn read(&mut self, user: &AsynUser, buf: &mut [u8]) -> AsynResult<OctetReadResult>;
fn write(&mut self, user: &mut AsynUser, data: &[u8]) -> AsynResult<usize>;
fn flush(&mut self, user: &mut AsynUser) -> AsynResult<()>;
}
pub trait OctetInterpose: Send + Sync {
fn read(
&mut self,
user: &AsynUser,
buf: &mut [u8],
next: &mut dyn OctetNext,
) -> AsynResult<OctetReadResult>;
fn write(
&mut self,
user: &mut AsynUser,
data: &[u8],
next: &mut dyn OctetNext,
) -> AsynResult<usize>;
fn flush(&mut self, user: &mut AsynUser, next: &mut dyn OctetNext) -> AsynResult<()>;
fn attach_port(&mut self, _multi_device: bool) {}
fn set_input_eos(&mut self, _addr: i32, _eos: &[u8]) {}
fn set_output_eos(&mut self, _addr: i32, _eos: &[u8]) {}
fn connection_changed(&mut self) {}
}
pub const PORT_CHAIN: i32 = -1;
pub struct OctetInterposeStack {
chains: std::collections::BTreeMap<i32, Vec<Box<dyn OctetInterpose>>>,
multi_device: bool,
}
impl OctetInterposeStack {
pub fn new(multi_device: bool) -> Self {
Self {
chains: std::collections::BTreeMap::new(),
multi_device,
}
}
pub fn install(&mut self, addr: i32, mut layer: Box<dyn OctetInterpose>) {
layer.attach_port(self.multi_device);
self.chains
.entry(crate::port::eos_device_key(self.multi_device, addr))
.or_default()
.insert(0, layer);
}
fn chain_key(&self, addr: i32) -> i32 {
let key = crate::port::eos_device_key(self.multi_device, addr);
if key != PORT_CHAIN && self.chains.get(&key).is_some_and(|c| !c.is_empty()) {
key
} else {
PORT_CHAIN
}
}
fn chain_mut(&mut self, addr: i32) -> Option<&mut Vec<Box<dyn OctetInterpose>>> {
let key = self.chain_key(addr);
self.chains.get_mut(&key).filter(|c| !c.is_empty())
}
pub fn len(&self) -> usize {
self.chains.values().map(Vec::len).sum()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn len_for(&self, addr: i32) -> usize {
self.chains
.get(&self.chain_key(addr))
.map_or(0, |c| c.len())
}
pub fn set_input_eos(&mut self, addr: i32, eos: &[u8]) {
if let Some(chain) = self.chain_mut(addr) {
for layer in chain {
layer.set_input_eos(addr, eos);
}
}
}
pub fn set_output_eos(&mut self, addr: i32, eos: &[u8]) {
if let Some(chain) = self.chain_mut(addr) {
for layer in chain {
layer.set_output_eos(addr, eos);
}
}
}
pub fn connection_changed(&mut self) {
for chain in self.chains.values_mut() {
for layer in chain {
layer.connection_changed();
}
}
}
pub fn dispatch_read(
&mut self,
user: &AsynUser,
buf: &mut [u8],
base: &mut dyn OctetNext,
) -> AsynResult<OctetReadResult> {
let addr = user.addr;
let Some(layers) = self.chain_mut(addr) else {
return base.read(user, buf);
};
InterposeChain {
layers: layers.as_mut_slice(),
base,
}
.read(user, buf)
}
pub fn dispatch_write(
&mut self,
user: &mut AsynUser,
data: &[u8],
base: &mut dyn OctetNext,
) -> AsynResult<usize> {
let addr = user.addr;
let Some(layers) = self.chain_mut(addr) else {
return base.write(user, data);
};
InterposeChain {
layers: layers.as_mut_slice(),
base,
}
.write(user, data)
}
pub fn dispatch_flush(
&mut self,
user: &mut AsynUser,
base: &mut dyn OctetNext,
) -> AsynResult<()> {
let addr = user.addr;
let Some(layers) = self.chain_mut(addr) else {
return base.flush(user);
};
InterposeChain {
layers: layers.as_mut_slice(),
base,
}
.flush(user)
}
}
impl Default for OctetInterposeStack {
fn default() -> Self {
Self::new(false)
}
}
struct InterposeChain<'a> {
layers: &'a mut [Box<dyn OctetInterpose>],
base: &'a mut dyn OctetNext,
}
impl OctetNext for InterposeChain<'_> {
fn read(&mut self, user: &AsynUser, buf: &mut [u8]) -> AsynResult<OctetReadResult> {
if let Some((first, rest)) = self.layers.split_first_mut() {
let mut next = InterposeChain {
layers: rest,
base: self.base,
};
first.read(user, buf, &mut next)
} else {
self.base.read(user, buf)
}
}
fn write(&mut self, user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
if let Some((first, rest)) = self.layers.split_first_mut() {
let mut next = InterposeChain {
layers: rest,
base: self.base,
};
first.write(user, data, &mut next)
} else {
self.base.write(user, data)
}
}
fn flush(&mut self, user: &mut AsynUser) -> AsynResult<()> {
if let Some((first, rest)) = self.layers.split_first_mut() {
let mut next = InterposeChain {
layers: rest,
base: self.base,
};
first.flush(user, &mut next)
} else {
self.base.flush(user)
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::user::AsynUser;
struct MockBase {
read_data: Vec<u8>,
written: Vec<u8>,
flushed: bool,
}
impl MockBase {
fn new(data: &[u8]) -> Self {
Self {
read_data: data.to_vec(),
written: Vec::new(),
flushed: false,
}
}
}
impl OctetNext for MockBase {
fn read(&mut self, _user: &AsynUser, buf: &mut [u8]) -> AsynResult<OctetReadResult> {
let n = self.read_data.len().min(buf.len());
buf[..n].copy_from_slice(&self.read_data[..n]);
Ok(OctetReadResult {
nbytes_transferred: n,
eom_reason: EomReason::CNT,
})
}
fn write(&mut self, _user: &mut AsynUser, data: &[u8]) -> AsynResult<usize> {
self.written.extend_from_slice(data);
Ok(data.len())
}
fn flush(&mut self, _user: &mut AsynUser) -> AsynResult<()> {
self.flushed = true;
Ok(())
}
}
struct PassthroughInterpose;
impl OctetInterpose for PassthroughInterpose {
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> {
next.write(user, data)
}
fn flush(&mut self, user: &mut AsynUser, next: &mut dyn OctetNext) -> AsynResult<()> {
next.flush(user)
}
}
struct UppercaseInterpose;
impl OctetInterpose for UppercaseInterpose {
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 upper: Vec<u8> = data.iter().map(|b| b.to_ascii_uppercase()).collect();
next.write(user, &upper)
}
fn flush(&mut self, user: &mut AsynUser, next: &mut dyn OctetNext) -> AsynResult<()> {
next.flush(user)
}
}
#[test]
fn test_empty_stack_passthrough() {
let mut stack = OctetInterposeStack::new(false);
let mut base = MockBase::new(b"hello");
let user = AsynUser::default();
let mut buf = [0u8; 32];
let result = stack.dispatch_read(&user, &mut buf, &mut base).unwrap();
assert_eq!(result.nbytes_transferred, 5);
assert_eq!(&buf[..5], b"hello");
}
#[test]
fn test_single_passthrough_layer() {
let mut stack = OctetInterposeStack::new(false);
stack.install(-1, Box::new(PassthroughInterpose));
let mut base = MockBase::new(b"world");
let user = AsynUser::default();
let mut buf = [0u8; 32];
let result = stack.dispatch_read(&user, &mut buf, &mut base).unwrap();
assert_eq!(result.nbytes_transferred, 5);
assert_eq!(&buf[..5], b"world");
}
#[test]
fn test_uppercase_interpose_write() {
let mut stack = OctetInterposeStack::new(false);
stack.install(-1, Box::new(UppercaseInterpose));
let mut base = MockBase::new(b"");
let mut user = AsynUser::default();
let n = stack
.dispatch_write(&mut user, b"hello", &mut base)
.unwrap();
assert_eq!(n, 5);
assert_eq!(&base.written, b"HELLO");
}
#[test]
fn test_multi_layer_chain() {
let mut stack = OctetInterposeStack::new(false);
stack.install(-1, Box::new(PassthroughInterpose));
stack.install(-1, Box::new(UppercaseInterpose));
assert_eq!(stack.len(), 2);
let mut base = MockBase::new(b"");
let mut user = AsynUser::default();
stack.dispatch_write(&mut user, b"test", &mut base).unwrap();
assert_eq!(&base.written, b"TEST");
}
#[test]
fn test_flush_dispatch() {
let mut stack = OctetInterposeStack::new(false);
stack.install(-1, Box::new(PassthroughInterpose));
let mut base = MockBase::new(b"");
let mut user = AsynUser::default();
stack.dispatch_flush(&mut user, &mut base).unwrap();
assert!(base.flushed);
}
#[test]
fn the_last_layer_installed_is_the_outermost() {
struct Tag(u8);
impl OctetInterpose for Tag {
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 tagged = vec![self.0];
tagged.extend_from_slice(data);
next.write(user, &tagged)
}
fn flush(&mut self, user: &mut AsynUser, next: &mut dyn OctetNext) -> AsynResult<()> {
next.flush(user)
}
}
let mut stack = OctetInterposeStack::new(false);
stack.install(-1, Box::new(Tag(b'A')));
stack.install(-1, Box::new(Tag(b'B')));
stack.install(-1, Box::new(Tag(b'C')));
let mut base = MockBase::new(b"");
let mut user = AsynUser::default();
stack.dispatch_write(&mut user, b"x", &mut base).unwrap();
assert_eq!(&base.written, b"ABCx");
}
#[test]
fn a_device_addressed_interpose_serves_only_that_device() {
struct Tag(u8);
impl OctetInterpose for Tag {
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 tagged = vec![self.0];
tagged.extend_from_slice(data);
next.write(user, &tagged)
}
fn flush(&mut self, user: &mut AsynUser, next: &mut dyn OctetNext) -> AsynResult<()> {
next.flush(user)
}
}
let write_from = |stack: &mut OctetInterposeStack, addr: i32| {
let mut base = MockBase::new(b"");
let mut user = AsynUser::new(0).with_addr(addr);
stack.dispatch_write(&mut user, b"x", &mut base).unwrap();
base.written
};
let mut stack = OctetInterposeStack::new(true);
stack.install(4, Box::new(Tag(b'D')));
stack.install(2, Box::new(Tag(b'E')));
assert_eq!(
write_from(&mut stack, 4),
b"Dx",
"device 4 runs its own layer"
);
assert_eq!(
write_from(&mut stack, 2),
b"Ex",
"device 2 runs its own layer"
);
assert_eq!(
write_from(&mut stack, 7),
b"x",
"a device with no interpose of its own runs none — C's findInterface \
falls back to the port's list, which is empty here"
);
assert_eq!(stack.len(), 2, "two layers on the port, one per device");
assert_eq!(stack.len_for(4), 1);
assert_eq!(stack.len_for(7), 0);
stack.install(PORT_CHAIN, Box::new(Tag(b'P')));
assert_eq!(
write_from(&mut stack, 7),
b"Px",
"a device with no chain of its own falls back to the port's"
);
assert_eq!(
write_from(&mut stack, 4),
b"Dx",
"...and a device WITH one keeps running only its own: C gives it the \
driver's interface as pPrev, not the port's interposes (:2211-2215)"
);
let mut single = OctetInterposeStack::new(false);
single.install(0, Box::new(Tag(b'S')));
assert_eq!(write_from(&mut single, 0), b"Sx");
assert_eq!(write_from(&mut single, 4), b"Sx");
assert_eq!(write_from(&mut single, -1), b"Sx");
}
}