cargo fmt

This commit is contained in:
Monica Moniot
2023-08-23 20:58:33 -04:00
parent 85523b3c3f
commit 9eb2bdb5e1
5 changed files with 57 additions and 37 deletions
+1 -2
View File
@@ -75,8 +75,7 @@ fn peer_main(transport: MpscTransport<Payload>, mut recv: mpsc::Receiver<Packet<
if let Some(peer) = peer_weak.upgrade() {
let tl = transport.clone();
task::spawn(async move {
let result = peer.receive(&tl, packet).await;
if let Ok((g, payload)) = result {
if let Ok((g, payload)) = peer.receive(&tl, packet).await {
receive(g, payload).await;
}
});
+1 -1
View File
@@ -294,7 +294,7 @@ impl<SendData, RecvData, const CAP: usize> SeqEx<SendData, RecvData, CAP> {
Err(TryError::WaitingForRecv)
} else {
Err(TryError::WaitingForReply)
}
};
}
}
Ok(())
+3 -2
View File
@@ -1,4 +1,4 @@
use crate::{TryRecvError, Packet, TryError, RecvOkRaw, SeqEx, SeqNo, TransportLayer, DEFAULT_WINDOW_CAP};
use crate::{Packet, RecvOkRaw, SeqEx, SeqNo, TransportLayer, TryError, TryRecvError, DEFAULT_WINDOW_CAP};
pub struct ReplyGuard<'a, TL: TransportLayer<SendData>, SendData, RecvData, const CAP: usize = DEFAULT_WINDOW_CAP> {
seq: &'a mut SeqEx<SendData, RecvData, CAP>,
@@ -225,7 +225,8 @@ impl<SendData, RecvData, const CAP: usize> SeqEx<SendData, RecvData, CAP> {
app: TL,
packet: Packet<RecvData>,
) -> Result<(RecvOk<'_, TL, SendData, RecvData, CAP>, bool), RecvError> {
self.receive_raw(app.clone(), packet).map(|(r, do_pump)| (RecvOk::from_raw(self, app, r), do_pump))
self.receive_raw(app.clone(), packet)
.map(|(r, do_pump)| (RecvOk::from_raw(self, app, r), do_pump))
}
pub fn try_pump<TL: TransportLayer<SendData>>(&mut self, app: TL) -> Result<(RecvOk<'_, TL, SendData, RecvData, CAP>, bool), TryError> {
self.try_pump_raw().map(|(r, do_pump)| (RecvOk::from_raw(self, app, r), do_pump))
+30 -19
View File
@@ -7,7 +7,7 @@ use std::{
};
use crate::{
Packet, TryError, RecvError, RecvOkRaw, SeqEx, SeqNo, TransportLayer, DEFAULT_INITIAL_SEQ_NO, DEFAULT_RESEND_INTERVAL_MS, DEFAULT_WINDOW_CAP,
Packet, RecvError, RecvOkRaw, SeqEx, SeqNo, TransportLayer, TryError, DEFAULT_INITIAL_SEQ_NO, DEFAULT_RESEND_INTERVAL_MS, DEFAULT_WINDOW_CAP,
};
pub struct SeqExSync<SendData, RecvData, const CAP: usize = DEFAULT_WINDOW_CAP> {
@@ -34,7 +34,8 @@ impl<'a, TL: TransportLayer<SendData>, SendData, RecvData, const CAP: usize> Rep
let app = self.app.take().unwrap();
let mut inner = self.seq.inner.lock().unwrap();
let seq_no = inner.seq.seq_no();
inner.seq
inner
.seq
.reply_raw(app, self.reply_no, self.is_holding_lock, seq_cst, packet_data(seq_no, self.reply_no));
self.seq.notify_reply(inner);
}
@@ -98,7 +99,12 @@ pub struct RecvIter<'a, TL: TransportLayer<SendData>, SendData, RecvData, const
impl<SendData, RecvData, const CAP: usize> SeqExSync<SendData, RecvData, CAP> {
pub fn new(retry_interval: i64, initial_seq_no: SeqNo) -> Self {
Self {
inner: Mutex::new(SeqExInner { seq: SeqEx::new(retry_interval, initial_seq_no), recv_waiters: 0, reply_sender_waiters: false, reply_receiver_waiters: false }),
inner: Mutex::new(SeqExInner {
seq: SeqEx::new(retry_interval, initial_seq_no),
recv_waiters: 0,
reply_sender_waiters: false,
reply_receiver_waiters: false,
}),
wait_on_recv: Condvar::default(),
wait_on_reply_receiver: Condvar::default(),
wait_on_reply_sender: Condvar::default(),
@@ -180,32 +186,39 @@ impl<SendData, RecvData, const CAP: usize> SeqExSync<SendData, RecvData, CAP> {
}
}
pub fn receive_all<TL: TransportLayer<SendData>>(
&self,
app: TL,
packet: Packet<RecvData>,
) -> RecvIter<'_, TL, SendData, RecvData, CAP> {
pub fn receive_all<TL: TransportLayer<SendData>>(&self, app: TL, packet: Packet<RecvData>) -> RecvIter<'_, TL, SendData, RecvData, CAP> {
let ret = self.receive(app.clone(), packet);
if let Some((first, do_pump)) = ret {
RecvIter { seq: do_pump.then_some(self), app, first: Some(first), blocking: true }
RecvIter {
seq: do_pump.then_some(self),
app,
first: Some(first),
blocking: true,
}
} else {
RecvIter { seq: None, app, first: None, blocking: true }
}
}
pub fn try_receive_all<TL: TransportLayer<SendData>>(
&self,
app: TL,
packet: Packet<RecvData>,
) -> RecvIter<'_, TL, SendData, RecvData, CAP> {
pub fn try_receive_all<TL: TransportLayer<SendData>>(&self, app: TL, packet: Packet<RecvData>) -> RecvIter<'_, TL, SendData, RecvData, CAP> {
let ret = self.try_receive(app.clone(), packet);
if let Ok((first, do_pump)) = ret {
RecvIter { seq: do_pump.then_some(self), app, first: Some(first), blocking: false }
RecvIter {
seq: do_pump.then_some(self),
app,
first: Some(first),
blocking: false,
}
} else {
RecvIter { seq: None, app, first: None, blocking: false }
}
}
pub fn try_send_with<TL: TransportLayer<SendData>, F: FnOnce(SeqNo) -> SendData>(&self, app: TL, seq_cst: bool, packet_data: F) -> Result<(), (TryError, F)> {
pub fn try_send_with<TL: TransportLayer<SendData>, F: FnOnce(SeqNo) -> SendData>(
&self,
app: TL,
seq_cst: bool,
packet_data: F,
) -> Result<(), (TryError, F)> {
let mut inner = self.inner.lock().unwrap();
inner.seq.try_send_with(app, seq_cst, packet_data)
}
@@ -257,9 +270,7 @@ impl<SendData, RecvData, const CAP: usize> Default for SeqExSync<SendData, RecvD
}
}
impl<'a, TL: TransportLayer<SendData>, SendData, RecvData, const CAP: usize> Iterator
for RecvIter<'a, TL, SendData, RecvData, CAP>
{
impl<'a, TL: TransportLayer<SendData>, SendData, RecvData, const CAP: usize> Iterator for RecvIter<'a, TL, SendData, RecvData, CAP> {
type Item = RecvOk<'a, TL, SendData, RecvData, CAP>;
fn next(&mut self) -> Option<Self::Item> {
if let Some(item) = self.first.take() {
+22 -13
View File
@@ -4,7 +4,9 @@ use tokio::{
time,
};
use crate::{Packet, TryError, RecvOkRaw, SeqEx, SeqNo, TransportLayer, DEFAULT_INITIAL_SEQ_NO, DEFAULT_RESEND_INTERVAL_MS, DEFAULT_WINDOW_CAP, RecvError};
use crate::{
Packet, RecvError, RecvOkRaw, SeqEx, SeqNo, TransportLayer, TryError, DEFAULT_INITIAL_SEQ_NO, DEFAULT_RESEND_INTERVAL_MS, DEFAULT_WINDOW_CAP,
};
type Sender<SendData, RecvData> = (oneshot::Sender<Option<(SeqNo, bool, RecvData)>>, SendData);
type Receiver<RecvData> = (oneshot::Sender<(SeqNo, bool, RecvData)>, RecvData);
@@ -37,7 +39,8 @@ impl<'a, TL: TokioLayer<SendData = SendData>, SendData, RecvData, const CAP: usi
let seq_no = inner.seq.seq_no();
let pre_ts = inner.seq.next_service_timestamp;
inner.seq
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 != inner.seq.next_service_timestamp).then_some(inner.seq.next_service_timestamp);
@@ -134,11 +137,11 @@ pub struct ServiceState {
}
struct IntoOneshot<'a, RecvData>(RecvData, &'a mut Option<oneshot::Receiver<(SeqNo, bool, RecvData)>>);
impl<'a, RecvData> Into<Receiver<RecvData>> for IntoOneshot<'a, RecvData> {
fn into(self) -> Receiver<RecvData> {
impl<'a, RecvData> From<IntoOneshot<'a, RecvData>> for Receiver<RecvData> {
fn from(value: IntoOneshot<'a, RecvData>) -> Self {
let (rx, tx) = oneshot::channel();
*self.1 = Some(tx);
(rx, self.0)
*value.1 = Some(tx);
(rx, value.0)
}
}
@@ -147,7 +150,11 @@ impl<SendData, RecvData, const CAP: usize> SeqExTokio<SendData, RecvData, CAP> {
let (update_queue, recv_service_update) = mpsc::channel(8);
(
Self {
seq_ex: Mutex::new(SeqExInner { seq: SeqEx::new(retry_interval, initial_seq_no), recv_waiters: 0, reply_waiters: false }),
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,
@@ -243,11 +250,13 @@ impl<SendData, RecvData, const CAP: usize> SeqExTokio<SendData, RecvData, CAP> {
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)
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)
}
}
}
}
@@ -270,7 +279,7 @@ impl<SendData, RecvData, const CAP: usize> SeqExTokio<SendData, RecvData, CAP> {
inner.reply_waiters = true;
Err((TryError::WaitingForReply, p))
}
Ok(()) => Ok((pre_ts != inner.seq.next_service_timestamp).then_some(inner.seq.next_service_timestamp))
Ok(()) => Ok((pre_ts != inner.seq.next_service_timestamp).then_some(inner.seq.next_service_timestamp)),
}
}