diff --git a/Cargo.lock b/Cargo.lock index b49f1c2..ac844df 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2,6 +2,45 @@ # It is not intended for manual editing. version = 3 +[[package]] +name = "addr2line" +version = "0.20.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f4fa78e18c64fce05e902adecd7a5eed15a5e0a3439f7b0e169f0252214865e3" +dependencies = [ + "gimli", +] + +[[package]] +name = "adler" +version = "1.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f26201604c87b1e01bd3d98f8d5d9a8fcbb815e8cedb41ffccbeb4bf593a35fe" + +[[package]] +name = "backtrace" +version = "0.3.68" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4319208da049c43661739c5fade2ba182f09d1dc2299b32298d3a31692b17e12" +dependencies = [ + "addr2line", + "cc", + "cfg-if", + "libc", + "miniz_oxide", + "object", + "rustc-demangle", +] + +[[package]] +name = "cc" +version = "1.0.82" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "305fe645edc1442a0fa8b6726ba61d422798d37a52e12eaecf4b022ebbb88f01" +dependencies = [ + "libc", +] + [[package]] name = "cfg-if" version = "1.0.0" @@ -19,6 +58,18 @@ dependencies = [ "wasi", ] +[[package]] +name = "gimli" +version = "0.27.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b6c80984affa11d98d1b88b66ac8853f143217b399d3c74116778ff8fdb4ed2e" + +[[package]] +name = "half" +version = "1.8.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "eabb4a44450da02c90444cf74558da904edde8fb4e9035a9a6a4e15445af0bd7" + [[package]] name = "itoa" version = "1.0.9" @@ -31,6 +82,36 @@ version = "0.2.147" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b4668fb0ea861c1df094127ac5f1da3409a82116a4ba74fca2e58ef927159bb3" +[[package]] +name = "memchr" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2dffe52ecf27772e601905b7522cb4ef790d2cc203488bbd0e2fe85fcb74566d" + +[[package]] +name = "miniz_oxide" +version = "0.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e7810e0be55b428ada41041c41f32c9f1a42817901b4ccf45fa3d4b6561e74c7" +dependencies = [ + "adler", +] + +[[package]] +name = "object" +version = "0.31.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8bda667d9f2b5051b8833f59f3bf748b28ef54f850f4fcb389a252aa383866d1" +dependencies = [ + "memchr", +] + +[[package]] +name = "pin-project-lite" +version = "0.2.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "12cc1b0bf1727a77a54b6654e7b5f1af8604923edc8b81885f8ec92f9e3f0a05" + [[package]] name = "proc-macro2" version = "1.0.66" @@ -58,6 +139,12 @@ dependencies = [ "getrandom", ] +[[package]] +name = "rustc-demangle" +version = "0.1.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d626bb9dae77e28219937af045c257c28bfd3f69333c512553507f5f9798cb76" + [[package]] name = "ryu" version = "1.0.15" @@ -70,7 +157,9 @@ version = "0.1.0" dependencies = [ "rand_core", "serde", + "serde_cbor", "serde_json", + "tokio", ] [[package]] @@ -82,6 +171,16 @@ dependencies = [ "serde_derive", ] +[[package]] +name = "serde_cbor" +version = "0.11.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2bef2ebfde456fb76bbcf9f59315333decc4fda0b2b44b420243c11e0f5ec1f5" +dependencies = [ + "half", + "serde", +] + [[package]] name = "serde_derive" version = "1.0.183" @@ -115,6 +214,16 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "tokio" +version = "1.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "17ed6077ed6cd6c74735e21f37eb16dc3935f96878b1fe961074089cc80893f9" +dependencies = [ + "backtrace", + "pin-project-lite", +] + [[package]] name = "unicode-ident" version = "1.0.11" diff --git a/Cargo.toml b/Cargo.toml index c97516e..89688df 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -10,12 +10,14 @@ path = "src/lib.rs" doc = true [features] -default = ["std", "serde"] +default = ["std", "serde", "tokio"] std = [] [dependencies] serde = { version = "1.0.183", default-features = false, features = ["derive"], optional = true } +tokio = { version = "1.31.0", default-features = false, features = ["sync", "rt"], optional = true } [dev-dependencies] rand_core = { version = "0.6.4", features = ["getrandom"]} serde_json = { version = "1.0.104" } +serde_cbor = { version = "0.11.2" } diff --git a/async_file_share.rs b/async_file_share.rs new file mode 100644 index 0000000..1efd22c --- /dev/null +++ b/async_file_share.rs @@ -0,0 +1,204 @@ +use std::{ + collections::HashMap, + io::Read, + ops::Deref, + path::PathBuf, + sync::{ + mpsc::{channel, Receiver, Sender}, + Arc, RwLock, + }, + thread, + time::{Duration, Instant}, +}; + +use rand_core::{OsRng, RngCore}; +use seq_ex::{sync::RecvSuccess, SeqNo, TransportLayer}; +use serde::{Deserialize, Serialize}; +use smol::fs::File; +use smol::prelude::*; + +/// serde_cbor minimal format is both smaller and faster than default format. +/// The overhead is ~33% faster. +fn to_writer_minimal(value: &impl serde::Serialize, w: &mut impl serde_cbor::ser::Write) -> serde_cbor::Result<()> { + value.serialize(&mut serde_cbor::Serializer::new(w).packed_format().legacy_enums())?; + Ok(()) +} + +fn drop_packet() -> bool { + OsRng.next_u32() >= (u32::MAX / 4 * 3) +} + +const FILE_CHUNK_SIZE: usize = 1000; +const DOWNLOAD_LIMIT: u64 = 1000000; + +#[derive(Debug)] +enum SendData { + RequestFile { filename: String }, + ConfirmFileSize { filename: String, file: File }, + ConfirmDownload, + FileDownload, +} + +#[derive(Serialize, Deserialize)] +enum Packet<'a> { + RequestFile { filename: &'a str }, + ConfirmFileSize { filesize: u64 }, + ConfirmDownload, + FileDownload { filename: &'a str, file_chunk: &'a [u8] }, + FileDownloadComplete { filename: &'a str }, +} + +const PACKET_TYPE_PAYLOAD: u8 = 0; +const PACKET_TYPE_REPLY: u8 = 1; +const PACKET_TYPE_ACK: u8 = 2; + +fn create_payload(seq_no: SeqNo, packet: &Packet<'_>) -> Vec { + let mut p = vec![PACKET_TYPE_PAYLOAD]; + p.extend(&seq_no.to_be_bytes()); + to_writer_minimal(packet, &mut p); + p +} +fn create_reply(seq_no: SeqNo, reply_no: SeqNo, packet: &Packet<'_>) -> Vec { + let mut p = vec![PACKET_TYPE_REPLY]; + p.extend(&seq_no.to_be_bytes()); + p.extend(&reply_no.to_be_bytes()); + to_writer_minimal(packet, &mut p); + p +} +macro_rules! reply { + ($guard:expr, $send_data:expr, $packet: expr) => { + $guard.reply_with(|seq_no, reply_no| ($send_data, create_reply(seq_no, reply_no, &$packet))) + }; +} +macro_rules! send { + ($peer:expr, $send_data:expr, $packet: expr) => { + $peer + .seqex + .send_with(&$peer.transport, |seq_no| ($send_data, create_payload(seq_no, &$packet))) + }; +} + +type SeqEx = seq_ex::sync::SeqExSync<(SendData, Vec), Vec>; +type ReplyGuard<'a> = seq_ex::sync::ReplyGuard<'a, &'a Transport, (SendData, Vec), Vec>; + +#[derive(Clone)] +struct Transport { + sender: Sender>, + time: Instant, +} + +struct Peer { + home_dir: PathBuf, + transport: Transport, + seqex: SeqEx, +} + +impl TransportLayer<(SendData, Vec)> for &Transport { + fn time(&mut self) -> i64 { + self.time.elapsed().as_millis() as i64 + } + + fn send(&mut self, _: SeqNo, _: Option, (_, packet): &(SendData, Vec)) { + let _ = self.sender.send(packet.clone()); + } + fn send_ack(&mut self, reply_no: SeqNo) { + let mut p = vec![PACKET_TYPE_ACK]; + p.extend(&reply_no.to_be_bytes()); + self.sender.send(p); + } +} + +async fn process(peer: Arc, guard: ReplyGuard<'_>, payload: Packet<'_>, data: Option) -> Option<()> { + match (payload, data) { + (Packet::RequestFile { filename }, None) => { + let file = File::open(peer.home_dir.join(filename)).await.ok()?; + let metadata = file.metadata().await.ok()?; + let filesize = metadata.len(); + reply!( + guard, + SendData::ConfirmFileSize { filename: filename.to_string(), file }, + Packet::ConfirmFileSize { filesize } + ); + } + (Packet::ConfirmFileSize { filesize }, Some(SendData::RequestFile { filename })) => { + let path = peer.home_dir.join(filename); + if filesize > DOWNLOAD_LIMIT { + return None; + } + let file = File::open(&path).await; + if file.is_ok() { + return None; + } + reply!(guard, SendData::ConfirmDownload, Packet::ConfirmDownload); + } + (Packet::ConfirmDownload, Some(SendData::ConfirmFileSize { filename, mut file })) => { + drop(guard); + const BUFFERED_CHUNKS: usize = 10; + let mut buffer = [0u8; BUFFERED_CHUNKS * FILE_CHUNK_SIZE]; + while let Ok(n) = file.read(&mut buffer).await { + if n == 0 { + break; + } + let mut i = 0; + while i < n { + let j = n.min(i + FILE_CHUNK_SIZE); + send!( + peer, + SendData::FileDownload, + Packet::FileDownload { filename: &filename, file_chunk: &buffer[i..j] } + ); + i = j; + } + } + send!(peer, SendData::FileDownload, Packet::FileDownloadComplete { filename: &filename }); + } + _ => {} + } + Some(()) +} + +async fn receive(peer: &Arc, receiver: &Receiver>) -> Option<()> { + while let Ok(packet) = receiver.try_recv() { + if drop_packet() { + continue; + } + let iter = match *packet.get(0)? { + PACKET_TYPE_ACK => { + let reply_no = SeqNo::from_be_bytes(packet.get(1..5)?.try_into().ok()?); + let _ = peer.seqex.receive_ack(reply_no); + return None; + } + PACKET_TYPE_PAYLOAD => { + let seq_no = SeqNo::from_be_bytes(packet.get(1..5)?.try_into().ok()?); + peer.seqex.receive_all(&peer.transport, seq_no, None, packet) + } + PACKET_TYPE_REPLY => { + let seq_no = SeqNo::from_be_bytes(packet.get(1..5)?.try_into().ok()?); + let reply_no = SeqNo::from_be_bytes(packet.get(5..9)?.try_into().ok()?); + peer.seqex.receive_all(&peer.transport, seq_no, Some(reply_no), packet) + } + _ => return None, + }; + for RecvSuccess { guard, packet, send_data } in iter { + let offset = match packet[0] { + PACKET_TYPE_PAYLOAD => 5, + PACKET_TYPE_REPLY => 9, + _ => return None, + }; + if let Ok(parsed_packet) = serde_cbor::from_slice::(packet.get(offset..)?) { + process(peer.clone(), guard, parsed_packet, send_data.map(|d| d.0)).await; + } + } + } + Some(()) +} + +fn main() { + let alice_root = std::path::Path::new("examples").join("alice_home"); + let bob_root = std::path::Path::new("examples").join("bob_home"); +} + +#[test] +fn test() { + main() +} diff --git a/examples/file_download.rs b/examples/file_download.rs index 8d2be99..2c65666 100644 --- a/examples/file_download.rs +++ b/examples/file_download.rs @@ -127,13 +127,13 @@ fn receive(peer: &Peer) { fn main() { let mut filesystem2 = HashMap::new(); - let mut file = vec![0; 1 << 16]; + let mut file = vec![0; 1 << 14]; OsRng.fill_bytes(&mut file); filesystem2.insert("File1".to_string(), file); - let mut file = vec![0; 1 << 18]; + let mut file = vec![0; 1 << 16]; OsRng.fill_bytes(&mut file); filesystem2.insert("File2".to_string(), file); - let mut file = vec![0; 1 << 20]; + let mut file = vec![0; 1 << 18]; OsRng.fill_bytes(&mut file); filesystem2.insert("File3".to_string(), file); @@ -157,7 +157,7 @@ fn main() { peer1.seqex.send(&peer1.transport, Packet::RequestFile { filename: "File3".to_string() }); peer1.seqex.send(&peer1.transport, Packet::RequestFile { filename: "File2".to_string() }); - for _ in 0..300 { + for _ in 0..3000 { receive(&peer1); receive(&peer2); thread::sleep(Duration::from_millis(1)); diff --git a/src/lib.rs b/src/lib.rs index 08a29e1..6a331a3 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -12,3 +12,6 @@ pub use single_thread::*; #[cfg(feature = "std")] pub mod sync; + +#[cfg(feature = "tokio")] +pub mod tokio; diff --git a/src/seq_queue.rs b/src/seq_queue.rs index 4543d77..1824c43 100644 --- a/src/seq_queue.rs +++ b/src/seq_queue.rs @@ -122,7 +122,7 @@ impl SeqEx { pub fn new(retry_interval: i64, initial_seq_no: SeqNo) -> Self { Self { resend_interval: retry_interval, - next_service_timestamp: i64::MIN, + next_service_timestamp: i64::MAX, next_send_seq_no: initial_seq_no, pre_recv_seq_no: initial_seq_no.wrapping_sub(1), recv_window: core::array::from_fn(|_| None), @@ -151,6 +151,7 @@ impl SeqEx { /// /// When `packet_data` is sent to the remote peer, the receiver should be able to quickly read /// the sequence number off of it. + #[inline] pub fn seq_no(&self) -> SeqNo { self.next_send_seq_no } @@ -370,19 +371,24 @@ impl SeqEx { pub fn service(&mut self, mut app: impl TransportLayer) -> i64 { let current_time = app.time(); - let next_interval = current_time + self.resend_interval; - let mut next_activity = i64::MAX; - for entry in self.send_window.iter_mut().flatten() { - if entry.next_resend_time <= current_time { - entry.next_resend_time = next_interval; - app.send(entry.seq_no, entry.reply_no, &entry.data); - } else { - next_activity = next_activity.min(entry.next_resend_time); + let real_interval = if self.next_service_timestamp <= current_time { + let next_resend_time = current_time + self.resend_interval; + let mut next_activity = i64::MAX; + for entry in self.send_window.iter_mut().flatten() { + if entry.next_resend_time <= current_time { + entry.next_resend_time = next_resend_time; + app.send(entry.seq_no, entry.reply_no, &entry.data); + } else { + next_activity = next_activity.min(entry.next_resend_time); + } } - } - self.next_service_timestamp = next_activity; - app.update_service_time(next_activity, current_time); - self.resend_interval.min(next_activity - current_time) + self.next_service_timestamp = next_activity; + app.update_service_time(next_activity, current_time); + next_activity - current_time + } else { + self.next_service_timestamp - current_time + }; + self.resend_interval.min(real_interval) } pub fn iter(&self) -> Iter<'_, SendData> { diff --git a/src/tokio.rs b/src/tokio.rs new file mode 100644 index 0000000..6303bfe --- /dev/null +++ b/src/tokio.rs @@ -0,0 +1,273 @@ +use std::{ + ops::{Deref, DerefMut}, + sync::{ + Mutex, MutexGuard, + }, + time::Instant, +}; +use tokio::{task, sync::{Notify, oneshot}}; + +use crate::{Error, SeqEx, SeqNo, TransportLayer, DEFAULT_INITIAL_SEQ_NO, DEFAULT_RESEND_INTERVAL_MS, DEFAULT_WINDOW_CAP}; + +type SendData = (oneshot::Sender<(Packet, SeqNo)>, Packet); + +pub struct SeqExTokio { + seq_ex: Mutex<(SeqEx, Packet, CAP>, usize)>, + send_block: Notify, +} + +pub struct ReplyGuard<'a, TL: TransportLayer>, Packet, const CAP: usize = DEFAULT_WINDOW_CAP> { + origin: &'a SeqExTokio, + app: Option, + reply_no: SeqNo, +} +impl<'a, TL: TransportLayer>, Packet, const CAP: usize> ReplyGuard<'a, TL, Packet, CAP> { + /// If you need to reply more than once, say to fragment a large file, then include in your + /// first reply some identifier, and then `send` all fragments with the same included identifier. + /// 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, packet: Packet) -> Option<(Packet, ReplyGuard<'a, TL, Packet, CAP>)> { + self.reply_with(|_, _| packet).await + } + pub async fn reply_with(mut self, packet: impl FnOnce(SeqNo, SeqNo) -> Packet) -> Option<(Packet, ReplyGuard<'a, TL, Packet, CAP>)> { + let (tx, rx) = oneshot::channel(); + let mut seq = self.origin.seq_ex.lock().unwrap(); + let seq_no = seq.0.seq_no(); + let app = self.app.take().unwrap(); + seq.0.reply_raw(app.clone(), self.reply_no, (tx, packet(seq_no, self.reply_no))); + let (packet, reply_no) = rx.await.ok()?; + let g = ReplyGuard { origin: self.origin, app: Some(app), reply_no }; + Some((packet, g)) + } +} +impl<'a, TL: TransportLayer>, Packet, const CAP: usize> Drop for ReplyGuard<'a, TL, Packet, CAP> { + fn drop(&mut self) { + if let Some(app) = self.app.take() { + let mut seq = self.origin.seq_ex.lock().unwrap(); + seq.0.ack_raw(app, self.reply_no); + } + } +} + +pub struct ReplyIter<'a, TL: TransportLayer>, Packet, const CAP: usize = DEFAULT_WINDOW_CAP> { + origin: Option<&'a SeqExTokio>, + app: TL, + first: Option<(Packet, ReplyGuard<'a, TL, Packet, CAP>)>, +} + +//pub struct SeqExGuard<'a, SendData, RecvData, const CAP: usize>(MutexGuard<'a, (SeqEx, usize)>); +//impl<'a, SendData, RecvData, const CAP: usize> Deref for SeqExGuard<'a, SendData, RecvData, CAP> { +// type Target = SeqEx; + +// fn deref(&self) -> &Self::Target { +// &self.0 .0 +// } +//} +//impl<'a, SendData, RecvData, const CAP: usize> DerefMut for SeqExGuard<'a, SendData, RecvData, CAP> { +// fn deref_mut(&mut self) -> &mut Self::Target { +// &mut self.0 .0 +// } +//} + +impl SeqExTokio { + pub fn new(retry_interval: i64, initial_seq_no: SeqNo) -> Self { + Self { + seq_ex: Mutex::new((SeqEx::new(retry_interval, initial_seq_no), 0)), + send_block: Notify::new(), + } + } + + fn process>>( + &self, app: TL, + mut seq: MutexGuard<'_, (SeqEx, Packet, CAP>, usize)>, + result: Result<(SeqNo, Packet, Option>), Error> + ) -> Option<(Packet, ReplyGuard<'_, TL, Packet, CAP>)> { + if let Ok((reply_no, packet, send_data)) = result { + if seq.1 > 0 { + self.send_block.notify_one(); + } + if let Some((tx, _)) = send_data { + if let Err(e) = tx.send((packet, reply_no)) { + // Allow the drop code to be run + seq.0.ack_raw(app, reply_no); + } + None + } else { + Some((packet, ReplyGuard { origin: self, app: Some(app), reply_no })) + } + } else { + None + } + } + pub fn receive>>( + &self, + app: TL, + seq_no: SeqNo, + reply_no: Option, + packet: Packet, + ) -> Option<(Packet, ReplyGuard<'_, TL, Packet, CAP>)> { + let mut seq = self.seq_ex.lock().unwrap(); + let result = seq.0.receive_raw(app.clone(), seq_no, reply_no, packet); + self.process(app, seq, result) + } + pub fn pump>>(&self, app: TL) -> Option<(Packet, ReplyGuard<'_, TL, Packet, CAP>)> { + let mut seq = self.seq_ex.lock().unwrap(); + let result = seq.0.pump_raw(); + self.process(app, seq, result) + } + pub fn receive_all>>( + &self, + app: TL, + seq_no: SeqNo, + reply_no: Option, + packet: Packet, + ) -> ReplyIter<'_, TL, Packet, CAP> { + if let Some(g) = self.receive(app.clone(), seq_no, reply_no, packet) { + ReplyIter { origin: Some(self), app, first: Some(g) } + } else { + ReplyIter { origin: None, app, first: None } + } + } + pub fn receive_ack(&self, reply_no: SeqNo) { + let mut seq = self.seq_ex.lock().unwrap(); + // We drop the sender to notify the receiver that no packet was received. + if let Ok(_) = seq.0.receive_ack(reply_no) { + if seq.1 > 0 { + self.send_block.notify_one(); + } + } + } + //pub fn try_send>(&self, app: TL, packet_data: SendData) -> Result<(), SendData> { + // let mut seq = self.lock(); + // seq.try_send(app, packet_data) + //} + async fn send_inner>>( + &self, + mut seq: MutexGuard<'_, (SeqEx, Packet, CAP>, usize)>, + app: TL, + mut tx: oneshot::Sender<(Packet, SeqNo)>, + mut packet: Packet + ) { + while let Err(e) = seq.0.try_send(app.clone(), (tx, packet)) { + (tx, packet) = e; + seq.1 += 1; + drop(seq); + self.send_block.notified().await; + seq = self.seq_ex.lock().unwrap(); + seq.1 -= 1; + } + } + /// If this future is dropped then the remote peer's reply to this packet will also be dropped. + pub async fn send>>(&self, app: TL, packet: Packet) -> Option<(Packet, ReplyGuard<'_, TL, Packet, CAP>)> { + self.send_with(app, |_| packet).await + } + //pub fn try_send_with>(&self, app: TL, packet_data: impl FnOnce(SeqNo) -> SendData) -> Result<(), SendData> { + // let mut seq = self.lock(); + // let seq_no = seq.seq_no(); + // seq.try_send(app, packet_data(seq_no)) + //} + pub async fn send_with>>(&self, app: TL, packet: impl FnOnce(SeqNo) -> Packet) -> Option<(Packet, ReplyGuard<'_, TL, Packet, CAP>)> { + let (tx, rx) = oneshot::channel(); + let seq = self.seq_ex.lock().unwrap(); + let seq_no = seq.0.seq_no(); + self.send_inner(seq, app.clone(), tx, packet(seq_no)).await; + // This can only return an error if the sender was dropped. + let (packet, reply_no) = rx.await.ok()?; + Some((packet, ReplyGuard { origin: self, app: Some(app), reply_no })) + } + + pub async fn main>>(&self, app: TL) { + let a = task::spawn(async{ + + }); + } + pub fn service>>(&self, app: TL) -> i64 { + self.seq_ex.lock().unwrap().0.service(app) + } + + //pub fn lock(&self) -> SeqExGuard<'_, SendData, RecvData, CAP> { + // SeqExGuard(self.seq_ex.lock().unwrap()) + //} +} +//impl Default for SeqExTokio { +// fn default() -> Self { +// Self::new(DEFAULT_RESEND_INTERVAL_MS, DEFAULT_INITIAL_SEQ_NO) +// } +//} + +impl<'a, TL: TransportLayer>, Packet, const CAP: usize> Iterator for ReplyIter<'a, TL, Packet, CAP> { + type Item = (Packet, ReplyGuard<'a, TL, Packet, CAP>); + fn next(&mut self) -> Option { + if let Some(g) = self.first.take() { + Some(g) + } else if let Some(origin) = self.origin { + origin.pump(self.app.clone()) + } else { + None + } + } +} + +pub trait TokioTransport: Clone { + fn time(&mut self) -> i64; + #[allow(unused)] + fn update_service_time(&mut self, timestamp: i64, current_time: i64) {} + + fn send(&mut self, seq_no: SeqNo, reply_no: Option, payload: &Packet); + fn send_ack(&mut self, reply_no: SeqNo); +} + +impl> TransportLayer> for (Tl, ) { + fn time(&mut self) -> i64 { + todo!() + } + fn update_service_time(&mut self, timestamp: i64, current_time: i64) { + + } + + fn send(&mut self, seq_no: SeqNo, reply_no: Option, payload: &SendData) { + todo!() + } + + fn send_ack(&mut self, reply_no: SeqNo) { + todo!() + } +} + +//#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +//#[derive(Clone)] +//pub enum PacketType { +// Payload(SeqNo, Option, Payload), +// Ack(SeqNo), +//} + +//#[derive(Clone)] +//pub struct MpscTransport { +// pub channel: Sender>, +// pub time: Instant, +//} +//pub type MpscGuard<'a, Packet> = ReplyGuard<'a, &'a MpscTransport, Packet, Packet>; +//pub type MpscSeqEx = SeqExTokio; + +//impl MpscTransport { +// pub fn new() -> (Self, Receiver>) { +// let (send, recv) = channel(); +// (Self { channel: send, time: std::time::Instant::now() }, recv) +// } +// pub fn from_sender(send: Sender>) -> Self { +// Self { channel: send, time: std::time::Instant::now() } +// } +//} +//impl TransportLayer for &MpscTransport { +// fn time(&mut self) -> i64 { +// self.time.elapsed().as_millis() as i64 +// } + +// fn send(&mut self, seq_no: SeqNo, reply_no: Option, payload: &Payload) { +// let _ = self.channel.send(PacketType::Payload(seq_no, reply_no, payload.clone())); +// } +// fn send_ack(&mut self, reply_no: SeqNo) { +// let _ = self.channel.send(PacketType::Ack(reply_no)); +// } +//}