mirror of
https://github.com/zerotier/sequential-exchange.git
synced 2026-05-22 16:28:28 -07:00
cargo fmt
This commit is contained in:
@@ -1,6 +1,6 @@
|
||||
use std::{sync::mpsc::Receiver, thread, time::Duration};
|
||||
|
||||
use seq_ex::sync::{MpscTransport, PacketType, RecvSuccess, MpscGuard, MpscSeqEx};
|
||||
use seq_ex::sync::{MpscGuard, MpscSeqEx, MpscTransport, PacketType, RecvSuccess};
|
||||
|
||||
#[derive(Clone)]
|
||||
enum Packet {
|
||||
@@ -27,19 +27,14 @@ fn process(_: MpscGuard<'_, Packet>, recv_packet: Packet, _: Option<Packet>, val
|
||||
}
|
||||
}
|
||||
|
||||
fn receive<'a>(
|
||||
recv: &Receiver<PacketType<Packet>>,
|
||||
seq: &MpscSeqEx<Packet>,
|
||||
transport: &'a MpscTransport<Packet>,
|
||||
value: &mut f32,
|
||||
) {
|
||||
fn receive(recv: &Receiver<PacketType<Packet>>, seq: &MpscSeqEx<Packet>, transport: &MpscTransport<Packet>, value: &mut f32) {
|
||||
while let Ok(packet) = recv.try_recv() {
|
||||
if !drop_packet() {
|
||||
match packet {
|
||||
PacketType::Ack ( reply_no ) => {
|
||||
PacketType::Ack(reply_no) => {
|
||||
let _ = seq.receive_ack(reply_no);
|
||||
}
|
||||
PacketType::Payload ( seq_no, reply_no, payload ) => {
|
||||
PacketType::Payload(seq_no, reply_no, payload) => {
|
||||
for RecvSuccess { guard, packet, send_data } in seq.receive_all(transport, seq_no, reply_no, payload) {
|
||||
process(guard, packet, send_data, value);
|
||||
}
|
||||
|
||||
+42
-28
@@ -1,22 +1,27 @@
|
||||
use std::{sync::{mpsc::{Receiver, Sender, channel}, Arc, RwLock}, thread, time::{Duration, Instant}, collections::HashMap, ops::Deref};
|
||||
use std::{
|
||||
collections::HashMap,
|
||||
ops::Deref,
|
||||
sync::{
|
||||
mpsc::{channel, Receiver, Sender},
|
||||
Arc, RwLock,
|
||||
},
|
||||
thread,
|
||||
time::{Duration, Instant},
|
||||
};
|
||||
|
||||
use seq_ex::{sync::{PacketType, RecvSuccess, ReplyGuard, SeqExSync}, TransportLayer, SeqNo};
|
||||
use rand_core::{RngCore, OsRng};
|
||||
use serde::{Serialize, Deserialize};
|
||||
use rand_core::{OsRng, RngCore};
|
||||
use seq_ex::{
|
||||
sync::{PacketType, RecvSuccess, ReplyGuard, SeqExSync},
|
||||
SeqNo, TransportLayer,
|
||||
};
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
const FILE_CHUNK_SIZE: usize = 1000;
|
||||
#[derive(Clone, Debug, Serialize, Deserialize)]
|
||||
enum Packet {
|
||||
RequestFile {
|
||||
filename: String,
|
||||
},
|
||||
ConfirmRequestFile {
|
||||
filesize: u64,
|
||||
},
|
||||
FileDownload {
|
||||
filename: String,
|
||||
file_chunk: Vec<u8>,
|
||||
}
|
||||
RequestFile { filename: String },
|
||||
ConfirmRequestFile { filesize: u64 },
|
||||
FileDownload { filename: String, file_chunk: Vec<u8> },
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
@@ -53,12 +58,11 @@ impl TransportLayer for &Transport {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
fn drop_packet() -> bool {
|
||||
OsRng.next_u32() >= (u32::MAX / 4 * 3)
|
||||
}
|
||||
|
||||
fn process<'a>(peer: &Peer, guard: ReplyGuard<'_, &Transport, Packet, Packet>, recv_packet: Packet, sent_packet: Option<Packet>) {
|
||||
fn process(peer: &Peer, guard: ReplyGuard<'_, &Transport, Packet, Packet>, recv_packet: Packet, sent_packet: Option<Packet>) {
|
||||
match (recv_packet, sent_packet) {
|
||||
(Packet::RequestFile { filename }, None) => {
|
||||
let filesystem = peer.filesystem.clone();
|
||||
@@ -75,7 +79,10 @@ fn process<'a>(peer: &Peer, guard: ReplyGuard<'_, &Transport, Packet, Packet>, r
|
||||
let mut i = 0;
|
||||
while i < file.len() {
|
||||
let j = file.len().min(i + FILE_CHUNK_SIZE);
|
||||
seqex.send(&transport, Packet::FileDownload { filename: filename.clone(), file_chunk: file[i..j].to_vec() });
|
||||
seqex.send(
|
||||
&transport,
|
||||
Packet::FileDownload { filename: filename.clone(), file_chunk: file[i..j].to_vec() },
|
||||
);
|
||||
i = j;
|
||||
}
|
||||
}
|
||||
@@ -100,19 +107,17 @@ fn process<'a>(peer: &Peer, guard: ReplyGuard<'_, &Transport, Packet, Packet>, r
|
||||
}
|
||||
}
|
||||
|
||||
fn receive<'a>(
|
||||
peer: &Peer,
|
||||
) {
|
||||
fn receive(peer: &Peer) {
|
||||
while let Ok(packet) = peer.receiver.try_recv() {
|
||||
if drop_packet() {
|
||||
continue;
|
||||
}
|
||||
let parsed_packet = serde_json::from_slice::<PacketType<Packet>>(&packet);
|
||||
match parsed_packet {
|
||||
Ok(PacketType::Ack ( reply_no )) => {
|
||||
Ok(PacketType::Ack(reply_no)) => {
|
||||
let _ = peer.seqex.receive_ack(reply_no);
|
||||
}
|
||||
Ok(PacketType::Payload ( seq_no, reply_no, payload )) => {
|
||||
Ok(PacketType::Payload(seq_no, reply_no, payload)) => {
|
||||
for RecvSuccess { guard, packet, send_data } in peer.seqex.receive_all(&peer.transport, seq_no, reply_no, payload) {
|
||||
process(peer, guard, packet, send_data);
|
||||
}
|
||||
@@ -122,7 +127,6 @@ fn receive<'a>(
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
fn main() {
|
||||
let mut filesystem2 = HashMap::new();
|
||||
let mut file = Vec::from([0u8; 1 << 16]);
|
||||
@@ -138,12 +142,22 @@ fn main() {
|
||||
let (send1, recv2) = channel();
|
||||
let (send2, recv1) = channel();
|
||||
|
||||
let peer1 = Peer { filesystem: Arc::new(RwLock::new(HashMap::new())), seqex: Arc::new(SeqExSync::new(5, 1)), transport: Transport{time: Instant::now(), sender: send1}, receiver: recv1 };
|
||||
let peer2 = Peer { filesystem: Arc::new(RwLock::new(filesystem2)), seqex: Arc::new(SeqExSync::new(5, 1)), transport: Transport{time: Instant::now(), sender: send2}, receiver: recv2 };
|
||||
let peer1 = Peer {
|
||||
filesystem: Arc::new(RwLock::new(HashMap::new())),
|
||||
seqex: Arc::new(SeqExSync::new(5, 1)),
|
||||
transport: Transport { time: Instant::now(), sender: send1 },
|
||||
receiver: recv1,
|
||||
};
|
||||
let peer2 = Peer {
|
||||
filesystem: Arc::new(RwLock::new(filesystem2)),
|
||||
seqex: Arc::new(SeqExSync::new(5, 1)),
|
||||
transport: Transport { time: Instant::now(), sender: send2 },
|
||||
receiver: recv2,
|
||||
};
|
||||
|
||||
peer1.seqex.send(&peer1.transport, Packet::RequestFile{filename: "File1".to_string()});
|
||||
peer1.seqex.send(&peer1.transport, Packet::RequestFile{filename: "File3".to_string()});
|
||||
peer1.seqex.send(&peer1.transport, Packet::RequestFile{filename: "File2".to_string()});
|
||||
peer1.seqex.send(&peer1.transport, Packet::RequestFile { filename: "File1".to_string() });
|
||||
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 {
|
||||
receive(&peer1);
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
use std::sync::mpsc::Receiver;
|
||||
|
||||
use seq_ex::sync::{MpscTransport, PacketType, RecvSuccess, MpscGuard, MpscSeqEx};
|
||||
use seq_ex::sync::{MpscGuard, MpscSeqEx, MpscTransport, PacketType, RecvSuccess};
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
enum Packet {
|
||||
@@ -39,14 +39,14 @@ fn process(guard: MpscGuard<'_, Packet>, recv_packet: Packet, send_packet: Optio
|
||||
|
||||
fn receive(recv: &Receiver<PacketType<Packet>>, seq: &MpscSeqEx<Packet>, transport: &MpscTransport<Packet>) {
|
||||
match recv.recv().unwrap() {
|
||||
PacketType::Ack ( reply_no ) => {
|
||||
PacketType::Ack(reply_no) => {
|
||||
let result = seq.receive_ack(reply_no);
|
||||
if let Ok(Exclamation) = result {
|
||||
// Our Hello World exchange ends right here.
|
||||
print!("\n");
|
||||
}
|
||||
}
|
||||
PacketType::Payload ( seq_no, reply_no, payload ) => {
|
||||
PacketType::Payload(seq_no, reply_no, payload) => {
|
||||
for RecvSuccess { guard, packet, send_data } in seq.receive_all(transport, seq_no, reply_no, payload) {
|
||||
process(guard, packet, send_data)
|
||||
}
|
||||
|
||||
@@ -1,6 +1,10 @@
|
||||
use crate::{Error, SeqEx, SeqNo, TransportLayer, DEFAULT_WINDOW_CAP};
|
||||
|
||||
pub struct ReplyGuard<'a, TL: TransportLayer<SendData = SendData>, SendData, RecvData, const CAP: usize = DEFAULT_WINDOW_CAP>(&'a mut SeqEx<SendData, RecvData, CAP>, TL, SeqNo);
|
||||
pub struct ReplyGuard<'a, TL: TransportLayer<SendData = SendData>, SendData, RecvData, const CAP: usize = DEFAULT_WINDOW_CAP>(
|
||||
&'a mut SeqEx<SendData, RecvData, CAP>,
|
||||
TL,
|
||||
SeqNo,
|
||||
);
|
||||
impl<'a, TL: TransportLayer<SendData = SendData>, SendData, RecvData, const CAP: usize> ReplyGuard<'a, TL, SendData, RecvData, 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.
|
||||
@@ -35,7 +39,10 @@ impl<SendData, RecvData, const CAP: usize> SeqEx<SendData, RecvData, CAP> {
|
||||
self.receive_raw(app.clone(), seq_no, reply_no, packet)
|
||||
.map(|(reply_no, packet, send_data)| RecvSuccess { guard: ReplyGuard(self, app, reply_no), packet, send_data })
|
||||
}
|
||||
pub fn pump<TL: TransportLayer<SendData = SendData>>(&mut self, app: TL) -> Result<RecvSuccess<'_, TL, RecvData, SendData, RecvData, CAP>, Error> {
|
||||
pub fn pump<TL: TransportLayer<SendData = SendData>>(
|
||||
&mut self,
|
||||
app: TL,
|
||||
) -> Result<RecvSuccess<'_, TL, RecvData, SendData, RecvData, CAP>, Error> {
|
||||
self.pump_raw()
|
||||
.map(|(reply_no, packet, send_data)| RecvSuccess { guard: ReplyGuard(self, app, reply_no), packet, send_data })
|
||||
}
|
||||
|
||||
+16
-7
@@ -1,9 +1,8 @@
|
||||
use std::cell::UnsafeCell;
|
||||
use std::sync::Condvar;
|
||||
use std::{
|
||||
cell::UnsafeCell,
|
||||
sync::{
|
||||
mpsc::{channel, Receiver, Sender},
|
||||
Mutex, MutexGuard,
|
||||
Condvar, Mutex, MutexGuard,
|
||||
},
|
||||
time::Instant,
|
||||
};
|
||||
@@ -18,7 +17,11 @@ pub struct SeqExSync<SendData, RecvData, const CAP: usize = DEFAULT_WINDOW_CAP>
|
||||
send_block: Condvar,
|
||||
}
|
||||
|
||||
pub struct ReplyGuard<'a, TL: TransportLayer<SendData = SendData>, SendData, RecvData, const CAP: usize = DEFAULT_WINDOW_CAP>(&'a SeqExSync<SendData, RecvData, CAP>, TL, SeqNo);
|
||||
pub struct ReplyGuard<'a, TL: TransportLayer<SendData = SendData>, SendData, RecvData, const CAP: usize = DEFAULT_WINDOW_CAP>(
|
||||
&'a SeqExSync<SendData, RecvData, CAP>,
|
||||
TL,
|
||||
SeqNo,
|
||||
);
|
||||
impl<'a, TL: TransportLayer<SendData = SendData>, SendData, RecvData, const CAP: usize> ReplyGuard<'a, TL, SendData, RecvData, 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.
|
||||
@@ -78,7 +81,13 @@ impl<SendData, RecvData, const CAP: usize> SeqExSync<SendData, RecvData, CAP> {
|
||||
self.unblock(ret.is_ok());
|
||||
ret.map(|(reply_no, packet, send_data)| RecvSuccess { guard: ReplyGuard(self, app, reply_no), packet, send_data })
|
||||
}
|
||||
pub fn receive_all<TL: TransportLayer<SendData = SendData>>(&self, app: TL, seq_no: SeqNo, reply_no: Option<SeqNo>, packet: RecvData) -> ReplyIter<'_, TL, SendData, RecvData, CAP> {
|
||||
pub fn receive_all<TL: TransportLayer<SendData = SendData>>(
|
||||
&self,
|
||||
app: TL,
|
||||
seq_no: SeqNo,
|
||||
reply_no: Option<SeqNo>,
|
||||
packet: RecvData,
|
||||
) -> ReplyIter<'_, TL, SendData, RecvData, CAP> {
|
||||
if let Ok(g) = self.receive(app.clone(), seq_no, reply_no, packet) {
|
||||
ReplyIter { origin: Some(self), app, first: Some(g) }
|
||||
} else {
|
||||
@@ -172,9 +181,9 @@ impl<Payload: Clone> TransportLayer for &MpscTransport<Payload> {
|
||||
}
|
||||
|
||||
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() ));
|
||||
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 ));
|
||||
let _ = self.channel.send(PacketType::Ack(reply_no));
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user