fix(multiplexer): Remove disconnected protocols from muxer loop (#16)
This commit is contained in:
parent
a5a027690a
commit
f53921c160
1 changed files with 39 additions and 28 deletions
|
|
@ -27,51 +27,62 @@ const MAX_SEGMENT_PAYLOAD_LENGTH: usize = 65535;
|
||||||
|
|
||||||
pub type Payload = Vec<u8>;
|
pub type Payload = Vec<u8>;
|
||||||
|
|
||||||
fn tx_round<TBearer>(
|
enum TxStepError {
|
||||||
|
BearerError(std::io::Error),
|
||||||
|
IngressDisconnected,
|
||||||
|
IngressEmpty,
|
||||||
|
}
|
||||||
|
|
||||||
|
fn tx_step<TBearer>(
|
||||||
bearer: &mut TBearer,
|
bearer: &mut TBearer,
|
||||||
ingress: &MuxIngress,
|
ingress_id: u16,
|
||||||
|
ingress_rx: &mut Receiver<Payload>,
|
||||||
clock: Instant,
|
clock: Instant,
|
||||||
) -> Result<u16, std::io::Error>
|
) -> Result<(), TxStepError>
|
||||||
where
|
where
|
||||||
TBearer: Bearer,
|
TBearer: Bearer,
|
||||||
{
|
{
|
||||||
let mut writes = 0u16;
|
match ingress_rx.try_recv() {
|
||||||
|
Ok(payload) => {
|
||||||
|
let chunks = payload.chunks(MAX_SEGMENT_PAYLOAD_LENGTH);
|
||||||
|
|
||||||
for (id, rx) in ingress.iter() {
|
for chunk in chunks {
|
||||||
match rx.try_recv() {
|
bearer
|
||||||
Ok(payload) => {
|
.write_segment(clock, ingress_id, chunk)
|
||||||
let chunks = payload.chunks(MAX_SEGMENT_PAYLOAD_LENGTH);
|
.map_err(TxStepError::BearerError)?;
|
||||||
|
}
|
||||||
|
|
||||||
for chunk in chunks {
|
Ok(())
|
||||||
bearer.write_segment(clock, *id, chunk)?;
|
}
|
||||||
writes += 1;
|
Err(TryRecvError::Disconnected) => Err(TxStepError::IngressDisconnected),
|
||||||
}
|
Err(TryRecvError::Empty) => Err(TxStepError::IngressEmpty),
|
||||||
}
|
|
||||||
Err(TryRecvError::Disconnected) => {
|
|
||||||
//TODO: remove handle from list
|
|
||||||
warn!("protocol handle {} disconnected", id);
|
|
||||||
}
|
|
||||||
Err(TryRecvError::Empty) => (),
|
|
||||||
};
|
|
||||||
}
|
}
|
||||||
|
|
||||||
Ok(writes)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn tx_loop<TBearer>(bearer: &mut TBearer, ingress: MuxIngress)
|
fn tx_loop<TBearer>(bearer: &mut TBearer, ingress: MuxIngress)
|
||||||
where
|
where
|
||||||
TBearer: Bearer,
|
TBearer: Bearer,
|
||||||
{
|
{
|
||||||
|
let mut rx_map: HashMap<_, _> = ingress.into_iter().collect();
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
let clock = Instant::now();
|
let clock = Instant::now();
|
||||||
match tx_round(bearer, &ingress, clock) {
|
|
||||||
Err(err) => {
|
rx_map.retain(|id, rx| match tx_step(bearer, *id, rx, clock) {
|
||||||
|
Err(TxStepError::BearerError(err)) => {
|
||||||
error!("{:?}", err);
|
error!("{:?}", err);
|
||||||
panic!();
|
panic!();
|
||||||
}
|
}
|
||||||
Ok(0) => thread::sleep(Duration::from_millis(10)),
|
Err(TxStepError::IngressDisconnected) => {
|
||||||
Ok(_) => (),
|
warn!("protocol handle {} disconnected", id);
|
||||||
};
|
false
|
||||||
|
}
|
||||||
|
Err(TxStepError::IngressEmpty) => {
|
||||||
|
thread::sleep(Duration::from_millis(10));
|
||||||
|
true
|
||||||
|
}
|
||||||
|
Ok(_) => true,
|
||||||
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -111,7 +122,7 @@ pub struct Channel(pub Sender<Payload>, pub Receiver<Payload>);
|
||||||
type ChannelProtocolHandle = (u16, Channel);
|
type ChannelProtocolHandle = (u16, Channel);
|
||||||
type ChannelIngressHandle = (u16, Receiver<Payload>);
|
type ChannelIngressHandle = (u16, Receiver<Payload>);
|
||||||
type ChannelEgressHandle = (u16, Sender<Payload>);
|
type ChannelEgressHandle = (u16, Sender<Payload>);
|
||||||
type MuxIngress<'a> = &'a [ChannelIngressHandle];
|
type MuxIngress = Vec<ChannelIngressHandle>;
|
||||||
type DemuxerEgress = Vec<ChannelEgressHandle>;
|
type DemuxerEgress = Vec<ChannelEgressHandle>;
|
||||||
|
|
||||||
pub struct Multiplexer {
|
pub struct Multiplexer {
|
||||||
|
|
@ -146,7 +157,7 @@ impl Multiplexer {
|
||||||
let (ingress, egress): (Vec<_>, Vec<_>) = multiplex_handles.into_iter().unzip();
|
let (ingress, egress): (Vec<_>, Vec<_>) = multiplex_handles.into_iter().unzip();
|
||||||
|
|
||||||
let mut tx_bearer = bearer.clone();
|
let mut tx_bearer = bearer.clone();
|
||||||
let tx_thread = thread::spawn(move || tx_loop(&mut tx_bearer, ingress.as_slice()));
|
let tx_thread = thread::spawn(move || tx_loop(&mut tx_bearer, ingress));
|
||||||
|
|
||||||
let mut rx_bearer = bearer.clone();
|
let mut rx_bearer = bearer.clone();
|
||||||
let rx_thread = thread::spawn(move || rx_loop(&mut rx_bearer, egress));
|
let rx_thread = thread::spawn(move || rx_loop(&mut rx_bearer, egress));
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue