From 85523b3c3f8dd7eb1824b4988f06cb6bbe611292 Mon Sep 17 00:00:00 2001 From: Monica Moniot Date: Wed, 23 Aug 2023 20:54:40 -0400 Subject: [PATCH] fixed tokio --- examples/hello_world_tokio.rs | 27 +-- src/lib.rs | 4 +- src/tokio.rs | 383 +++++++++++++++------------------- 3 files changed, 177 insertions(+), 237 deletions(-) diff --git a/examples/hello_world_tokio.rs b/examples/hello_world_tokio.rs index e720bc0..b9191e2 100644 --- a/examples/hello_world_tokio.rs +++ b/examples/hello_world_tokio.rs @@ -3,7 +3,7 @@ use std::sync::Arc; use tokio::{sync::mpsc, task}; use seq_ex::{ - tokio::{AsyncRecvError, MpscTransport, ReplyGuard, SeqExTokio}, + tokio::{MpscTransport, ReplyGuard, SeqExTokio}, Packet, }; @@ -17,7 +17,7 @@ enum Payload { } use Payload::*; -async fn receive(reply_guard: ReplyGuard<'_, &MpscTransport, Payload>, payload: Payload) -> Option<()> { +async fn receive(reply_guard: ReplyGuard<'_, &MpscTransport, Payload, Payload>, payload: Payload) -> Option<()> { match payload { Hello => { print!("Hello"); @@ -42,7 +42,7 @@ async fn receive(reply_guard: ReplyGuard<'_, &MpscTransport, Payload>, Some(()) } -async fn say_hello(seq: &SeqExTokio, transport: &MpscTransport) -> Option<()> { +async fn say_hello(seq: &SeqExTokio, transport: &MpscTransport) -> Option<()> { let (reply_guard, payload) = seq.send(transport, false, Hello).await.ok()?; if payload != Space { return None; @@ -59,20 +59,8 @@ async fn say_hello(seq: &SeqExTokio, transport: &MpscTransport Some(()) } -async fn pump_all(peer: Arc>, transport: MpscTransport) { - if let Some((g, payload)) = peer.pump(&transport).await { - spawn_pump(peer.clone(), transport.clone()); - receive(g, payload).await; - } -} -fn spawn_pump(peer: Arc>, transport: MpscTransport) { - task::spawn(async move { - pump_all(peer, transport).await; - }); -} - -fn peer_main(transport: MpscTransport, mut recv: mpsc::Receiver>) -> Arc> { - let (seq, mut service) = SeqExTokio::::new_default(); +fn peer_main(transport: MpscTransport, mut recv: mpsc::Receiver>) -> Arc> { + let (seq, mut service) = SeqExTokio::new_default(); let peer = Arc::new(seq); let peer_weak = Arc::downgrade(&peer); let tl = transport.clone(); @@ -87,10 +75,7 @@ fn peer_main(transport: MpscTransport, mut recv: mpsc::Receiver = (oneshot::Sender>, Payload); +type Sender = (oneshot::Sender>, SendData); +type Receiver = (oneshot::Sender<(SeqNo, bool, RecvData)>, RecvData); -pub struct SeqExTokio { - seq_ex: Mutex<(SeqEx, Payload, CAP>, usize, bool)>, - send_block: Notify, - reply_block: Notify, +pub struct SeqExTokio { + seq_ex: Mutex>, + wait_on_recv: Notify, + wait_on_reply: Notify, update_queue: mpsc::Sender, } -pub struct ReplyGuard<'a, TL: TokioLayer, Payload, const CAP: usize = DEFAULT_WINDOW_CAP> { - seq: &'a SeqExTokio, +struct SeqExInner { + seq: SeqEx, Receiver, CAP>, + recv_waiters: usize, + reply_waiters: bool, +} + +pub struct ReplyGuard<'a, TL: TokioLayer, SendData, RecvData, const CAP: usize = DEFAULT_WINDOW_CAP> { + seq: &'a SeqExTokio, app: Option, reply_no: SeqNo, is_holding_lock: bool, } -impl<'a, TL: TokioLayer, Payload, const CAP: usize> ReplyGuard<'a, TL, Payload, CAP> { - fn new(seq: &'a SeqExTokio, app: TL, reply_no: SeqNo, is_holding_lock: bool) -> Self { +impl<'a, TL: TokioLayer, SendData, RecvData, const CAP: usize> ReplyGuard<'a, TL, SendData, RecvData, CAP> { + fn new(seq: &'a SeqExTokio, app: TL, reply_no: SeqNo, is_holding_lock: bool) -> Self { ReplyGuard { seq, app: Some(app), reply_no, is_holding_lock } } - fn try_reply_with_inner(&mut self, app: TL, seq_cst: bool, packet_data: impl FnOnce(SeqNo, SeqNo) -> SendData) -> Option { - let mut seq = self.seq.seq_ex.lock().unwrap(); - let seq_no = seq.0.seq_no(); + fn try_reply_with_inner(&mut self, app: TL, seq_cst: bool, packet_data: impl FnOnce(SeqNo, SeqNo) -> Sender) -> Option { + let mut inner = self.seq.seq_ex.lock().unwrap(); + let seq_no = inner.seq.seq_no(); - let pre_ts = seq.0.next_service_timestamp; - seq.0 + let pre_ts = inner.seq.next_service_timestamp; + inner.seq .reply_raw(app, self.reply_no, self.is_holding_lock, seq_cst, packet_data(seq_no, self.reply_no)); - let ret = (pre_ts != seq.0.next_service_timestamp).then_some(seq.0.next_service_timestamp); - if seq.2 { - seq.2 = false; - drop(seq); - self.seq.reply_block.notify_waiters(); - } + let ret = (pre_ts != inner.seq.next_service_timestamp).then_some(inner.seq.next_service_timestamp); + + self.seq.notify_reply(inner); ret } ///// If you need to reply more than once, say to fragment a large file, then include in your @@ -45,14 +49,14 @@ impl<'a, TL: TokioLayer, Payload, const CAP: usize> ReplyGuar ///// The identifier will tell the remote peer which packets contain fragments of the file, ///// and since each fragment will be received in order it will be trivial for them to reconstruct ///// the original file. - pub async fn reply(self, seq_cst: bool, packet_data: Payload) -> Result<(ReplyGuard<'a, TL, Payload, CAP>, Payload), AsyncError> { + pub async fn reply(self, seq_cst: bool, packet_data: SendData) -> Result<(ReplyGuard<'a, TL, SendData, RecvData, CAP>, RecvData), AsyncError> { self.reply_with(seq_cst, |_, _| packet_data).await } pub async fn reply_with( mut self, seq_cst: bool, - packet_data: impl FnOnce(SeqNo, SeqNo) -> Payload, - ) -> Result<(ReplyGuard<'a, TL, Payload, CAP>, Payload), AsyncError> { + packet_data: impl FnOnce(SeqNo, SeqNo) -> SendData, + ) -> Result<(ReplyGuard<'a, TL, SendData, RecvData, CAP>, RecvData), AsyncError> { let app = self.app.take().unwrap(); let (tx, rx) = oneshot::channel(); @@ -61,31 +65,27 @@ impl<'a, TL: TokioLayer, Payload, const CAP: usize> ReplyGuar if let Some(update_ts) = update_ts { let _ = self.seq.update_queue.send(update_ts).await; } - let (reply_no, seq_cst, recv_data) = rx.await.map_err(|_| AsyncError::SeqExClosed)?.ok_or(AsyncError::ReceivedAck)?; + let (reply_no, seq_cst, recv_data) = rx.await.map_err(|_| AsyncError::SeqExClosed)?.ok_or(AsyncError::EndOfExchange)?; Ok((Self::new(self.seq, app, reply_no, seq_cst), recv_data)) } pub fn to_components(mut self) -> (TL, SeqNo, bool) { (self.app.take().unwrap(), self.reply_no, self.is_holding_lock) } - pub unsafe fn from_components(seq: &'a SeqExTokio, app: TL, reply_no: SeqNo, is_holding_lock: bool) -> Self { + pub unsafe fn from_components(seq: &'a SeqExTokio, app: TL, reply_no: SeqNo, is_holding_lock: bool) -> Self { Self::new(seq, app, reply_no, is_holding_lock) } } -impl<'a, TL: TokioLayer, Payload, const CAP: usize> Drop for ReplyGuard<'a, TL, Payload, CAP> { +impl<'a, TL: TokioLayer, SendData, RecvData, const CAP: usize> Drop for ReplyGuard<'a, TL, SendData, RecvData, CAP> { fn drop(&mut self) { if let Some(app) = self.app.take() { - let mut seq = self.seq.seq_ex.lock().unwrap(); - seq.0.ack_raw(app, self.reply_no, self.is_holding_lock); - if seq.2 { - seq.2 = false; - drop(seq); - self.seq.reply_block.notify_waiters(); - } + let mut inner = self.seq.seq_ex.lock().unwrap(); + inner.seq.ack_raw(app, self.reply_no, self.is_holding_lock); + self.seq.notify_reply(inner); } } } -impl<'a, TL: TokioLayer, Payload, const CAP: usize> std::fmt::Debug for ReplyGuard<'a, TL, Payload, CAP> { +impl<'a, TL: TokioLayer, SendData, RecvData, const CAP: usize> std::fmt::Debug for ReplyGuard<'a, TL, SendData, RecvData, CAP> { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.debug_struct("ReplyGuard") .field("reply_no", &self.reply_no) @@ -98,18 +98,16 @@ impl<'a, TL: TokioLayer, Payload, const CAP: usize> std::fmt: pub enum AsyncRecvError { DroppedTooEarly, DroppedDuplicate, - WaitingForRecv, - WaitingForReply, AsyncReply, + SeqExClosed, } impl std::fmt::Display for AsyncRecvError { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { AsyncRecvError::DroppedTooEarly => write!(f, "packet arrived too early"), AsyncRecvError::DroppedDuplicate => write!(f, "packet was a duplicate"), - AsyncRecvError::WaitingForRecv => write!(f, "can't process until another packet is received"), - AsyncRecvError::WaitingForReply => write!(f, "can't process until a reply is finished"), AsyncRecvError::AsyncReply => write!(f, "packet was an async reply"), + AsyncRecvError::SeqExClosed => write!(f, "the instance of SeqEx was dropped"), } } } @@ -117,43 +115,41 @@ impl std::error::Error for AsyncRecvError {} #[derive(Clone, Debug, PartialEq, Eq)] pub enum AsyncError { - ReceivedAck, + EndOfExchange, SeqExClosed, } impl std::fmt::Display for AsyncError { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { - AsyncError::ReceivedAck => write!(f, "peer replied with an ack"), + AsyncError::EndOfExchange => write!(f, "peer replied with an ack"), AsyncError::SeqExClosed => write!(f, "the window was closed before a reply could be received"), } } } impl std::error::Error for AsyncError {} -pub struct AsyncRecvIter<'a, TL: TokioLayer, P: Into, Payload, const CAP: usize = DEFAULT_WINDOW_CAP> { - seq: Option<&'a SeqExTokio>, - app: TL, - first: Option<(ReplyGuard<'a, TL, Payload, CAP>, P)>, -} -pub struct RecvIter<'a, TL: TokioLayer, P: Into, Payload, const CAP: usize = DEFAULT_WINDOW_CAP> { - seq: Option<&'a SeqExTokio>, - app: TL, - first: Option<(ReplyGuard<'a, TL, Payload, CAP>, P)>, -} - pub struct ServiceState { next_service_timestamp: i64, recv_service_update: mpsc::Receiver, } -impl SeqExTokio { +struct IntoOneshot<'a, RecvData>(RecvData, &'a mut Option>); +impl<'a, RecvData> Into> for IntoOneshot<'a, RecvData> { + fn into(self) -> Receiver { + let (rx, tx) = oneshot::channel(); + *self.1 = Some(tx); + (rx, self.0) + } +} + +impl SeqExTokio { pub fn new(retry_interval: i64, initial_seq_no: SeqNo) -> (Self, ServiceState) { let (update_queue, recv_service_update) = mpsc::channel(8); ( Self { - seq_ex: Mutex::new((SeqEx::new(retry_interval, initial_seq_no), 0, false)), - send_block: Notify::new(), - reply_block: Notify::new(), + seq_ex: Mutex::new(SeqExInner { seq: SeqEx::new(retry_interval, initial_seq_no), recv_waiters: 0, reply_waiters: false }), + wait_on_recv: Notify::new(), + wait_on_reply: Notify::new(), update_queue, }, ServiceState { next_service_timestamp: i64::MAX, recv_service_update }, @@ -163,160 +159,153 @@ impl SeqExTokio { Self::new(DEFAULT_RESEND_INTERVAL_MS, DEFAULT_INITIAL_SEQ_NO) } - fn map_raw_or_reply<'a, TL: TokioLayer, P: Into>( - &'a self, - app: TL, - seq: &mut SeqEx, Payload, CAP>, - ret: RecvOkRaw, P>, - ) -> Option<(ReplyGuard<'a, TL, Payload, CAP>, P)> { - match ret { - RecvOkRaw::Payload { reply_no, seq_cst, recv_data } => Some((ReplyGuard::new(self, app, reply_no, seq_cst), recv_data)), - RecvOkRaw::Reply { reply_no, seq_cst, recv_data, send_data: (tx, _) } => { - if tx.send(Some((reply_no, seq_cst, recv_data.into()))).is_err() { - // Send an ack if no one is receiving the reply on the other end. - // Could occur if the future holding the receiver is dropped. - seq.ack_raw(app, reply_no, seq_cst); - } - None - } - RecvOkRaw::Ack { send_data: (tx, _) } => { - drop(tx); - None - } + fn notify_reply(&self, mut inner: MutexGuard<'_, SeqExInner>) { + if inner.reply_waiters { + inner.reply_waiters = false; + drop(inner); + self.wait_on_reply.notify_waiters(); } } - - pub fn receive, P: Into>( + fn receive_inner>( &self, app: TL, - packet: Packet

, - ) -> Result<(ReplyGuard<'_, TL, Payload, CAP>, P), AsyncRecvError> { - let mut seq = self.seq_ex.lock().unwrap(); - return match seq.0.receive_raw(app.clone(), packet) { - Ok(ret) => { - if seq.1 > 0 { - seq.1 -= 1; - self.send_block.notify_one(); + packet: Packet>, + ) -> Result<(ReplyGuard<'_, TL, SendData, RecvData, CAP>, RecvData), Option> { + let mut inner = self.seq_ex.lock().unwrap(); + return match inner.seq.receive_raw(app.clone(), packet) { + Ok((recv_data, do_pump)) => { + // pump first, handle return value second. + let mut total_recv = 1; + if do_pump { + while let Ok((data, do_pump)) = inner.seq.try_pump_raw() { + total_recv += 1; + match data { + RecvOkRaw::Payload { reply_no, seq_cst, recv_data } => { + if recv_data.0.send((reply_no, seq_cst, recv_data.1)).is_err() { + // Send an ack if no one is receiving the reply on the other end. + // Could occur if the future holding the receiver is dropped. + inner.seq.ack_raw(app.clone(), reply_no, seq_cst); + } + } + RecvOkRaw::Reply { reply_no, seq_cst, recv_data, send_data } => { + if send_data.0.send(Some((reply_no, seq_cst, recv_data.1))).is_err() { + inner.seq.ack_raw(app.clone(), reply_no, seq_cst); + } + } + RecvOkRaw::Ack { send_data: (tx, _) } => { + let _ = tx.send(None); + } + } + if !do_pump { + break; + } + } } - self.map_raw_or_reply(app, &mut seq.0, ret).ok_or(AsyncRecvError::AsyncReply) + + let ret = match recv_data { + RecvOkRaw::Payload { reply_no, seq_cst, recv_data } => Ok((ReplyGuard::new(self, app, reply_no, seq_cst), recv_data.0)), + RecvOkRaw::Reply { reply_no, seq_cst, recv_data, send_data: (tx, _) } => { + if tx.send(Some((reply_no, seq_cst, recv_data.0))).is_err() { + // Send an ack if no one is receiving the reply on the other end. + // Could occur if the future holding the receiver is dropped. + inner.seq.ack_raw(app, reply_no, seq_cst); + } + Err(Some(AsyncRecvError::AsyncReply)) + } + RecvOkRaw::Ack { send_data: (tx, _) } => { + let _ = tx.send(None); + Err(Some(AsyncRecvError::AsyncReply)) + } + }; + if inner.recv_waiters > 0 { + let to_wake = inner.recv_waiters.min(total_recv); + inner.recv_waiters -= to_wake; + drop(inner); + for _ in 0..to_wake { + self.wait_on_recv.notify_one(); + } + } + ret } - Err(crate::RecvError::DroppedTooEarly) => Err(AsyncRecvError::DroppedTooEarly), - Err(crate::RecvError::DroppedDuplicate) => Err(AsyncRecvError::DroppedDuplicate), - Err(crate::RecvError::WaitingForRecv) => Err(AsyncRecvError::WaitingForRecv), - Err(crate::RecvError::WaitingForReply) => Err(AsyncRecvError::WaitingForReply), + Err(RecvError::DroppedTooEarly) => Err(Some(AsyncRecvError::DroppedTooEarly)), + Err(RecvError::DroppedDuplicate) => Err(Some(AsyncRecvError::DroppedDuplicate)), + Err(RecvError::WaitingForRecv) | Err(RecvError::WaitingForReply) => Err(None), }; } - fn try_pump_inner>( + pub async fn receive>( &self, app: TL, - blocking: bool, - ) -> Result<(ReplyGuard<'_, TL, Payload, CAP>, Payload), TryError> { - let mut seq = self.seq_ex.lock().unwrap(); - // Enforce that only one thread may pump at a time. - loop { - match seq.0.try_pump_raw() { - Ok(ret) => { - if seq.1 > 0 { - seq.1 -= 1; - self.send_block.notify_one(); - } - if let Some(ret) = self.map_raw_or_reply(app.clone(), &mut seq.0, ret) { - return Ok(ret); - } - } - Err(TryError::WaitingForRecv) => return Err(TryError::WaitingForRecv), - Err(TryError::WaitingForReply) => { - seq.2 |= blocking; - return Err(TryError::WaitingForReply); - } + packet: Packet, + ) -> Result<(ReplyGuard<'_, TL, SendData, RecvData, CAP>, RecvData), AsyncRecvError> { + let mut tx = None; + let packet = packet.map(|r| IntoOneshot(r, &mut tx)); + match self.receive_inner(app.clone(), packet) { + Ok(ret) => Ok(ret), + Err(Some(e)) => Err(e), + Err(None) => if let Some(tx) = tx { + let (reply_no, is_holding_lock, data) = tx.await.map_err(|_| AsyncRecvError::SeqExClosed)?; + Ok((ReplyGuard::new(self, app, reply_no, is_holding_lock), data)) + } else { + Err(AsyncRecvError::DroppedDuplicate) } } } - pub fn try_pump>(&self, app: TL) -> Result<(ReplyGuard<'_, TL, Payload, CAP>, Payload), TryError> { - self.try_pump_inner(app, false) - } - pub async fn pump>(&self, app: TL) -> Option<(ReplyGuard<'_, TL, Payload, CAP>, Payload)> { - loop { - match self.try_pump_inner(app.clone(), true) { - Ok(ret) => return Some(ret), - Err(TryError::WaitingForRecv) => return None, - Err(TryError::WaitingForReply) => { - self.reply_block.notified().await; - } - } - } - } - pub fn try_receive_all, P: Into>( - &self, - app: TL, - packet: Packet

, - ) -> RecvIter<'_, TL, P, Payload, CAP> { - match self.receive(app.clone(), packet) { - Ok(ret) => RecvIter { seq: Some(self), app, first: Some(ret) }, - Err(AsyncRecvError::AsyncReply) => RecvIter { seq: Some(self), app, first: None }, - Err(_) => RecvIter { seq: None, app, first: None }, - } - } - pub fn receive_all, P: Into>( - &self, - app: TL, - packet: Packet

, - ) -> AsyncRecvIter<'_, TL, P, Payload, CAP> { - match self.receive(app.clone(), packet) { - Ok(ret) => AsyncRecvIter { seq: Some(self), app, first: Some(ret) }, - Err(AsyncRecvError::AsyncReply) => AsyncRecvIter { seq: Some(self), app, first: None }, - Err(_) => AsyncRecvIter { seq: None, app, first: None }, - } - } - fn try_send_with_inner, F: FnOnce(SeqNo) -> SendData>( + fn send_with_inner, F: FnOnce(SeqNo) -> Sender>( &self, app: TL, - blocking: bool, seq_cst: bool, packet_data: F, - ) -> Result, F> { - let mut seq = self.seq_ex.lock().unwrap(); - let pre_ts = seq.0.next_service_timestamp; - let result = seq.0.try_send_with(app, seq_cst, packet_data); - if let Err(e) = result { - seq.1 += blocking as usize; - Err(e) - } else { - Ok((pre_ts != seq.0.next_service_timestamp).then_some(seq.0.next_service_timestamp)) + ) -> Result, (TryError, F)> { + let mut inner = self.seq_ex.lock().unwrap(); + let pre_ts = inner.seq.next_service_timestamp; + let result = inner.seq.try_send_with(app, seq_cst, packet_data); + match result { + Err((TryError::WaitingForRecv, p)) => { + inner.recv_waiters += 1; + Err((TryError::WaitingForRecv, p)) + } + Err((TryError::WaitingForReply, p)) => { + inner.reply_waiters = true; + Err((TryError::WaitingForReply, p)) + } + Ok(()) => Ok((pre_ts != inner.seq.next_service_timestamp).then_some(inner.seq.next_service_timestamp)) } } - pub async fn send>( + pub async fn send>( &self, app: TL, seq_cst: bool, - packet_data: Payload, - ) -> Result<(ReplyGuard<'_, TL, Payload, CAP>, Payload), AsyncError> { + packet_data: SendData, + ) -> Result<(ReplyGuard<'_, TL, SendData, RecvData, CAP>, RecvData), AsyncError> { self.send_with(app, seq_cst, |_| packet_data).await } - pub async fn send_with>( + pub async fn send_with>( &self, app: TL, seq_cst: bool, - packet_data: impl FnOnce(SeqNo) -> Payload, - ) -> Result<(ReplyGuard<'_, TL, Payload, CAP>, Payload), AsyncError> { + packet_data: impl FnOnce(SeqNo) -> SendData, + ) -> Result<(ReplyGuard<'_, TL, SendData, RecvData, CAP>, RecvData), AsyncError> { let (rx, tx) = oneshot::channel(); let mut pf = |s| (rx, packet_data(s)); loop { - let ret = self.try_send_with_inner(app.clone(), true, seq_cst, pf); - match ret { + match self.send_with_inner(app.clone(), seq_cst, pf) { Ok(update) => { if let Some(update) = update { let _ = self.update_queue.send(update).await; } - let (reply_no, locked, recv_data) = tx.await.map_err(|_| AsyncError::SeqExClosed)?.ok_or(AsyncError::ReceivedAck)?; + let (reply_no, locked, recv_data) = tx.await.map_err(|_| AsyncError::SeqExClosed)?.ok_or(AsyncError::EndOfExchange)?; return Ok((ReplyGuard::new(self, app, reply_no, locked), recv_data)); } - Err(p) => { + Err((TryError::WaitingForRecv, p)) => { pf = p; - self.send_block.notified().await; + self.wait_on_recv.notified().await; + } + Err((TryError::WaitingForReply, p)) => { + pf = p; + self.wait_on_reply.notified().await; } } } @@ -324,7 +313,7 @@ impl SeqExTokio { /// This function must be called with the same ServiceState instance returned upon creation of /// the given SeqExTokio instance. - pub async fn service_task>(&self, mut app: TL, state: &mut ServiceState) { + pub async fn service_task>(&self, mut app: TL, state: &mut ServiceState) { let mut result = None; if state.next_service_timestamp < i64::MAX { let diff = state.next_service_timestamp - app.time(); @@ -340,60 +329,26 @@ impl SeqExTokio { if let Some(up) = result { state.next_service_timestamp = state.next_service_timestamp.min(up); } else { - let mut seq = self.seq_ex.lock().unwrap(); - seq.0.service(app.clone()); - state.next_service_timestamp = seq.0.next_service_timestamp; - } - } -} - -impl<'a, TL: TokioLayer, P: Into, Payload, const CAP: usize> RecvIter<'a, TL, P, Payload, CAP> { - pub fn take_first(&mut self) -> Option<(ReplyGuard<'a, TL, Payload, CAP>, P)> { - self.first.take() - } -} -impl<'a, TL: TokioLayer, P: Into, Payload, const CAP: usize> Iterator for RecvIter<'a, TL, P, Payload, CAP> { - type Item = (ReplyGuard<'a, TL, Payload, CAP>, Payload); - - fn next(&mut self) -> Option { - if let Some(g) = self.first.take() { - Some((g.0, g.1.into())) - } else if let Some(seq) = self.seq { - seq.try_pump(self.app.clone()).ok() - } else { - None - } - } -} -impl<'a, TL: TokioLayer, P: Into, Payload, const CAP: usize> AsyncRecvIter<'a, TL, P, Payload, CAP> { - pub fn take_first(&mut self) -> Option<(ReplyGuard<'a, TL, Payload, CAP>, P)> { - self.first.take() - } - - pub async fn next(&mut self) -> Option<(ReplyGuard<'a, TL, Payload, CAP>, Payload)> { - if let Some(g) = self.first.take() { - Some((g.0, g.1.into())) - } else if let Some(seq) = self.seq { - seq.pump(self.app.clone()).await - } else { - None + let mut inner = self.seq_ex.lock().unwrap(); + inner.seq.service(app.clone()); + state.next_service_timestamp = inner.seq.next_service_timestamp; } } } pub trait TokioLayer: Clone { - type Payload; + type SendData; fn time(&mut self) -> i64; - fn send(&mut self, packet: Packet<&Self::Payload>); + fn send(&mut self, packet: Packet<&Self::SendData>); } -impl TransportLayer> for TL { +impl TransportLayer> for TL { fn time(&mut self) -> i64 { self.time() } - fn send(&mut self, packet: Packet<&SendData>) { + fn send(&mut self, packet: Packet<&Sender>) { self.send(packet.map(|p| &p.1)) } } @@ -414,7 +369,7 @@ impl MpscTransport { } } impl TokioLayer for &MpscTransport { - type Payload = Payload; + type SendData = Payload; fn time(&mut self) -> i64 { self.time.elapsed().as_millis() as i64