mirror of
https://github.com/zerotier/sequential-exchange.git
synced 2026-05-22 16:28:28 -07:00
begun adding tokio
This commit is contained in:
Generated
+109
@@ -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"
|
||||
|
||||
+3
-1
@@ -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" }
|
||||
|
||||
@@ -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<u8> {
|
||||
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<u8> {
|
||||
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<u8>), Vec<u8>>;
|
||||
type ReplyGuard<'a> = seq_ex::sync::ReplyGuard<'a, &'a Transport, (SendData, Vec<u8>), Vec<u8>>;
|
||||
|
||||
#[derive(Clone)]
|
||||
struct Transport {
|
||||
sender: Sender<Vec<u8>>,
|
||||
time: Instant,
|
||||
}
|
||||
|
||||
struct Peer {
|
||||
home_dir: PathBuf,
|
||||
transport: Transport,
|
||||
seqex: SeqEx,
|
||||
}
|
||||
|
||||
impl TransportLayer<(SendData, Vec<u8>)> for &Transport {
|
||||
fn time(&mut self) -> i64 {
|
||||
self.time.elapsed().as_millis() as i64
|
||||
}
|
||||
|
||||
fn send(&mut self, _: SeqNo, _: Option<SeqNo>, (_, packet): &(SendData, Vec<u8>)) {
|
||||
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<Peer>, guard: ReplyGuard<'_>, payload: Packet<'_>, data: Option<SendData>) -> 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<Peer>, receiver: &Receiver<Vec<u8>>) -> 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>(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()
|
||||
}
|
||||
@@ -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));
|
||||
|
||||
@@ -12,3 +12,6 @@ pub use single_thread::*;
|
||||
|
||||
#[cfg(feature = "std")]
|
||||
pub mod sync;
|
||||
|
||||
#[cfg(feature = "tokio")]
|
||||
pub mod tokio;
|
||||
|
||||
+19
-13
@@ -122,7 +122,7 @@ impl<SendData, RecvData, const CAP: usize> SeqEx<SendData, RecvData, CAP> {
|
||||
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<SendData, RecvData, const CAP: usize> SeqEx<SendData, RecvData, CAP> {
|
||||
///
|
||||
/// 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<SendData, RecvData, const CAP: usize> SeqEx<SendData, RecvData, CAP> {
|
||||
|
||||
pub fn service(&mut self, mut app: impl TransportLayer<SendData>) -> 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> {
|
||||
|
||||
+273
@@ -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<Packet> = (oneshot::Sender<(Packet, SeqNo)>, Packet);
|
||||
|
||||
pub struct SeqExTokio<Packet, const CAP: usize = DEFAULT_WINDOW_CAP> {
|
||||
seq_ex: Mutex<(SeqEx<SendData<Packet>, Packet, CAP>, usize)>,
|
||||
send_block: Notify,
|
||||
}
|
||||
|
||||
pub struct ReplyGuard<'a, TL: TransportLayer<SendData<Packet>>, Packet, const CAP: usize = DEFAULT_WINDOW_CAP> {
|
||||
origin: &'a SeqExTokio<Packet, CAP>,
|
||||
app: Option<TL>,
|
||||
reply_no: SeqNo,
|
||||
}
|
||||
impl<'a, TL: TransportLayer<SendData<Packet>>, 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<SendData<Packet>>, 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<SendData<Packet>>, Packet, const CAP: usize = DEFAULT_WINDOW_CAP> {
|
||||
origin: Option<&'a SeqExTokio<Packet, CAP>>,
|
||||
app: TL,
|
||||
first: Option<(Packet, ReplyGuard<'a, TL, Packet, CAP>)>,
|
||||
}
|
||||
|
||||
//pub struct SeqExGuard<'a, SendData, RecvData, const CAP: usize>(MutexGuard<'a, (SeqEx<SendData, RecvData, CAP>, usize)>);
|
||||
//impl<'a, SendData, RecvData, const CAP: usize> Deref for SeqExGuard<'a, SendData, RecvData, CAP> {
|
||||
// type Target = SeqEx<SendData, RecvData, CAP>;
|
||||
|
||||
// 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<Packet, const CAP: usize> SeqExTokio<Packet, CAP> {
|
||||
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<TL: TransportLayer<SendData<Packet>>>(
|
||||
&self, app: TL,
|
||||
mut seq: MutexGuard<'_, (SeqEx<SendData<Packet>, Packet, CAP>, usize)>,
|
||||
result: Result<(SeqNo, Packet, Option<SendData<Packet>>), 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<TL: TransportLayer<SendData<Packet>>>(
|
||||
&self,
|
||||
app: TL,
|
||||
seq_no: SeqNo,
|
||||
reply_no: Option<SeqNo>,
|
||||
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<TL: TransportLayer<SendData<Packet>>>(&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<TL: TransportLayer<SendData<Packet>>>(
|
||||
&self,
|
||||
app: TL,
|
||||
seq_no: SeqNo,
|
||||
reply_no: Option<SeqNo>,
|
||||
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<TL: TransportLayer<SendData>>(&self, app: TL, packet_data: SendData) -> Result<(), SendData> {
|
||||
// let mut seq = self.lock();
|
||||
// seq.try_send(app, packet_data)
|
||||
//}
|
||||
async fn send_inner<TL: TransportLayer<SendData<Packet>>>(
|
||||
&self,
|
||||
mut seq: MutexGuard<'_, (SeqEx<SendData<Packet>, 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<TL: TransportLayer<SendData<Packet>>>(&self, app: TL, packet: Packet) -> Option<(Packet, ReplyGuard<'_, TL, Packet, CAP>)> {
|
||||
self.send_with(app, |_| packet).await
|
||||
}
|
||||
//pub fn try_send_with<TL: TransportLayer<SendData>>(&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<TL: TransportLayer<SendData<Packet>>>(&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<TL: TransportLayer<SendData<Packet>>>(&self, app: TL) {
|
||||
let a = task::spawn(async{
|
||||
|
||||
});
|
||||
}
|
||||
pub fn service<TL: TransportLayer<SendData<Packet>>>(&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<SendData, RecvData, const CAP: usize> Default for SeqExTokio<SendData, RecvData, CAP> {
|
||||
// fn default() -> Self {
|
||||
// Self::new(DEFAULT_RESEND_INTERVAL_MS, DEFAULT_INITIAL_SEQ_NO)
|
||||
// }
|
||||
//}
|
||||
|
||||
impl<'a, TL: TransportLayer<SendData<Packet>>, Packet, const CAP: usize> Iterator for ReplyIter<'a, TL, Packet, CAP> {
|
||||
type Item = (Packet, ReplyGuard<'a, TL, Packet, CAP>);
|
||||
fn next(&mut self) -> Option<Self::Item> {
|
||||
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<Packet>: 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<SeqNo>, payload: &Packet);
|
||||
fn send_ack(&mut self, reply_no: SeqNo);
|
||||
}
|
||||
|
||||
impl<Packet, Tl: TokioTransport<Packet>> TransportLayer<SendData<Packet>> 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<SeqNo>, payload: &SendData<Packet>) {
|
||||
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: Clone> {
|
||||
// Payload(SeqNo, Option<SeqNo>, Payload),
|
||||
// Ack(SeqNo),
|
||||
//}
|
||||
|
||||
//#[derive(Clone)]
|
||||
//pub struct MpscTransport<Payload: Clone> {
|
||||
// pub channel: Sender<PacketType<Payload>>,
|
||||
// pub time: Instant,
|
||||
//}
|
||||
//pub type MpscGuard<'a, Packet> = ReplyGuard<'a, &'a MpscTransport<Packet>, Packet, Packet>;
|
||||
//pub type MpscSeqEx<Packet> = SeqExTokio<Packet, Packet>;
|
||||
|
||||
//impl<Payload: Clone> MpscTransport<Payload> {
|
||||
// pub fn new() -> (Self, Receiver<PacketType<Payload>>) {
|
||||
// let (send, recv) = channel();
|
||||
// (Self { channel: send, time: std::time::Instant::now() }, recv)
|
||||
// }
|
||||
// pub fn from_sender(send: Sender<PacketType<Payload>>) -> Self {
|
||||
// Self { channel: send, time: std::time::Instant::now() }
|
||||
// }
|
||||
//}
|
||||
//impl<Payload: Clone> TransportLayer<Payload> for &MpscTransport<Payload> {
|
||||
// fn time(&mut self) -> i64 {
|
||||
// self.time.elapsed().as_millis() as i64
|
||||
// }
|
||||
|
||||
// fn send(&mut self, seq_no: SeqNo, reply_no: Option<SeqNo>, 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));
|
||||
// }
|
||||
//}
|
||||
Reference in New Issue
Block a user