mirror of
https://github.com/zerotier/sequential-exchange.git
synced 2026-05-22 16:28:28 -07:00
added timestamp updating
This commit is contained in:
+84
-43
@@ -47,9 +47,8 @@ pub const DEFAULT_RESEND_INTERVAL_MS: i64 = 200;
|
||||
pub const DEFAULT_INITIAL_SEQ_NO: SeqNo = 1;
|
||||
|
||||
pub const DEFAULT_SEND_WINDOW_LEN: usize = 64;
|
||||
pub const DEFAULT_RECV_WINDOW_LEN: usize = 32;
|
||||
pub const DEFAULT_RECV_WINDOW_LEN: usize = 64;
|
||||
|
||||
const MAX_CONCURRENCY: usize = 24;
|
||||
pub struct SeqEx<TL: TransportLayer, const SLEN: usize = DEFAULT_SEND_WINDOW_LEN, const RLEN: usize = DEFAULT_RECV_WINDOW_LEN> {
|
||||
/// The interval at which packets will be resent if they have not yet been acknowledged by the
|
||||
/// remote peer.
|
||||
@@ -58,10 +57,13 @@ pub struct SeqEx<TL: TransportLayer, const SLEN: usize = DEFAULT_SEND_WINDOW_LEN
|
||||
pub next_service_timestamp: i64,
|
||||
next_send_seq_no: SeqNo,
|
||||
pre_recv_seq_no: SeqNo,
|
||||
/// This could be made more efficient in SoA format.
|
||||
send_window: [Option<SendEntry<TL>>; SLEN],
|
||||
recv_window: [Option<RecvEntry<TL>>; RLEN],
|
||||
reserved: [SeqNo; MAX_CONCURRENCY],
|
||||
reserved_len: usize,
|
||||
/// The size of this array determines the maximum number of received packets that the
|
||||
/// application may attempt to process concurrently before new received packets start being dropped.
|
||||
concurrent_replies: [SeqNo; RLEN],
|
||||
concurrent_replies_total: usize,
|
||||
}
|
||||
|
||||
struct RecvEntry<TL: TransportLayer> {
|
||||
@@ -122,18 +124,24 @@ impl<TL: TransportLayer> SeqEx<TL> {
|
||||
pre_recv_seq_no: initial_seq_no.wrapping_sub(1),
|
||||
recv_window: core::array::from_fn(|_| None),
|
||||
send_window: core::array::from_fn(|_| None),
|
||||
reserved: core::array::from_fn(|_| 0),
|
||||
reserved_len: 0,
|
||||
concurrent_replies: core::array::from_fn(|_| 0),
|
||||
concurrent_replies_total: 0,
|
||||
}
|
||||
}
|
||||
fn send_window_slot_mut(&mut self, seq_no: SeqNo) -> &mut Option<SendEntry<TL>> {
|
||||
&mut self.send_window[seq_no as usize % self.send_window.len()]
|
||||
}
|
||||
fn send_window_slot(&self, seq_no: SeqNo) -> &Option<SendEntry<TL>> {
|
||||
&self.send_window[seq_no as usize % self.send_window.len()]
|
||||
}
|
||||
|
||||
|
||||
/// Returns whether or not the send window is full.
|
||||
/// If the send window is full calls to `SeqEx::send` will always fail.
|
||||
pub fn is_full(&self) -> bool {
|
||||
// We claim that the window is full one entry before it is actually full for the sake of
|
||||
// making it always possible for both peers to process at least one reply at all times.
|
||||
let next_i = self.next_send_seq_no as usize;
|
||||
self.send_window[next_i % self.send_window.len()].is_some() || self.send_window[(next_i + 1) % self.send_window.len()].is_some()
|
||||
self.is_full_inner(1)
|
||||
}
|
||||
/// Returns the next sequence number to be attached to the next sent packet.
|
||||
/// This should be called before `SeqEx::send`, and the return value should be
|
||||
@@ -172,11 +180,11 @@ impl<TL: TransportLayer> SeqEx<TL> {
|
||||
let next_resend_time = current_time + self.resend_interval;
|
||||
if self.next_service_timestamp > next_resend_time {
|
||||
self.next_service_timestamp = next_resend_time;
|
||||
app.update_service_time(current_time, next_resend_time);
|
||||
app.update_service_time(next_resend_time, current_time);
|
||||
}
|
||||
let i = seq_no as usize % self.send_window.len();
|
||||
debug_assert!(self.send_window[i].is_none());
|
||||
let entry = self.send_window[i].insert(SendEntry {
|
||||
let slot = self.send_window_slot_mut(seq_no);
|
||||
debug_assert!(slot.is_none());
|
||||
let entry = slot.insert(SendEntry {
|
||||
seq_no,
|
||||
reply_no: None,
|
||||
next_resend_time,
|
||||
@@ -194,7 +202,7 @@ impl<TL: TransportLayer> SeqEx<TL> {
|
||||
reply_no: Option<SeqNo>,
|
||||
packet: P,
|
||||
) -> Result<(SeqNo, P, Option<TL::SendData>), Error> {
|
||||
// We only want to accept packets with seq_nos in the range:
|
||||
// We only want to accept packets with sequence numbers in the range:
|
||||
// `self.pre_recv_seq_no < seq_no <= self.pre_recv_seq_no + self.recv_window.len()`.
|
||||
// To check that range we compute `seq_no - (self.pre_recv_seq_no + 1)` and check
|
||||
// if the number wrapped below 0, or if it is above `self.recv_window.len()`.
|
||||
@@ -203,14 +211,22 @@ impl<TL: TransportLayer> SeqEx<TL> {
|
||||
let is_above_range = !is_below_range && normalized_seq_no >= self.recv_window.len() as u32;
|
||||
let is_next = normalized_seq_no == 0;
|
||||
if is_below_range {
|
||||
// Check whether or not we are already replying to this packet.
|
||||
for entry in self.send_window.iter() {
|
||||
// If it is below the range, that means the packet has been received twice.
|
||||
// For every received packet we either send a normal reply or send an empty reply.
|
||||
// If the application sent a normal reply in response to this packet previously, and
|
||||
// that reply has not been acknowledged, then arg `seq_no` will be in the send window,
|
||||
// and if the application is still deciding what to reply with, it will be in the
|
||||
// concurrent replies array.
|
||||
// If either are the case then we know we will eventually send a reply to the packet.
|
||||
// If neither are the case then we must send an empty reply so the remote peer can stop
|
||||
// resending the packet.
|
||||
for entry in &self.send_window {
|
||||
if entry.as_ref().map_or(false, |e| e.reply_no == Some(seq_no)) {
|
||||
return Err(Error::OutOfSequence);
|
||||
}
|
||||
}
|
||||
for i in 0..self.reserved_len {
|
||||
if self.reserved[i] == seq_no {
|
||||
for reply_no in self.concurrent_replies {
|
||||
if reply_no == seq_no {
|
||||
return Err(Error::OutOfSequence);
|
||||
}
|
||||
}
|
||||
@@ -222,8 +238,8 @@ impl<TL: TransportLayer> SeqEx<TL> {
|
||||
// If the send window is full we cannot safely process received packets,
|
||||
// because there would be no way to reply.
|
||||
// We can only process this packet if processing it would make space in the send window.
|
||||
let next_i = self.next_send_seq_no as usize % self.send_window.len();
|
||||
let is_full = self.send_window[next_i].as_ref().map_or(false, |e| Some(e.seq_no) != reply_no) || self.reserved_len >= self.reserved.len();
|
||||
let is_full = self.is_full_inner(0);
|
||||
|
||||
let i = seq_no as usize % self.recv_window.len();
|
||||
if let Some(pre) = self.recv_window[i].as_mut() {
|
||||
if seq_no == pre.seq_no {
|
||||
@@ -238,14 +254,15 @@ impl<TL: TransportLayer> SeqEx<TL> {
|
||||
};
|
||||
}
|
||||
} else {
|
||||
// NOTE: I believe this return is currently unreachable.
|
||||
// This is currently unreachable due to the range check.
|
||||
// `self.pre_recv_seq_no < seq_no <= self.pre_recv_seq_no + self.recv_window.len()`.
|
||||
return Err(Error::OutOfSequence);
|
||||
}
|
||||
}
|
||||
if is_next && !is_full {
|
||||
self.pre_recv_seq_no = seq_no;
|
||||
self.reserved[self.reserved_len] = seq_no;
|
||||
self.reserved_len += 1;
|
||||
self.concurrent_replies[self.concurrent_replies_total] = seq_no;
|
||||
self.concurrent_replies_total += 1;
|
||||
let data = reply_no.and_then(|r| self.take_send(r));
|
||||
Ok((seq_no, packet, data))
|
||||
} else {
|
||||
@@ -262,27 +279,40 @@ impl<TL: TransportLayer> SeqEx<TL> {
|
||||
}
|
||||
}
|
||||
pub fn receive_ack(&mut self, reply_no: SeqNo) {
|
||||
let i = reply_no as usize % self.send_window.len();
|
||||
if let Some(entry) = self.send_window[i].as_mut() {
|
||||
let slot = self.send_window_slot_mut(reply_no);
|
||||
if let Some(entry) = slot.as_mut() {
|
||||
if entry.seq_no == reply_no {
|
||||
entry.next_resend_time = i64::MAX;
|
||||
}
|
||||
}
|
||||
}
|
||||
pub fn receive_empty_reply(&mut self, reply_no: SeqNo) -> Option<TL::SendData> {
|
||||
let i = reply_no as usize % self.send_window.len();
|
||||
if self.send_window[i].as_ref().map_or(false, |e| e.seq_no == reply_no) {
|
||||
let entry = self.send_window[i].take().unwrap();
|
||||
let slot = self.send_window_slot_mut(reply_no);
|
||||
if slot.as_ref().map_or(false, |e| e.seq_no == reply_no) {
|
||||
let entry = slot.take().unwrap();
|
||||
Some(entry.data)
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
fn is_full_inner(&self, bias: u32) -> bool {
|
||||
if self.concurrent_replies_total >= self.concurrent_replies.len() {
|
||||
return true;
|
||||
}
|
||||
for i in 0..self.concurrent_replies_total as u32 + 1 + bias {
|
||||
let slot = self.send_window_slot(self.next_send_seq_no.wrapping_add(i));
|
||||
if slot.is_some() {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
false
|
||||
}
|
||||
|
||||
fn take_send(&mut self, reply_no: SeqNo) -> Option<TL::SendData> {
|
||||
let i = reply_no as usize % self.send_window.len();
|
||||
if self.send_window[i].as_ref().map_or(false, |e| e.seq_no == reply_no) {
|
||||
self.send_window[i].take().map(|e| e.data)
|
||||
let slot = self.send_window_slot_mut(reply_no);
|
||||
if slot.as_ref().map_or(false, |e| e.seq_no == reply_no) {
|
||||
slot.take().map(|e| e.data)
|
||||
} else {
|
||||
None
|
||||
}
|
||||
@@ -292,14 +322,14 @@ impl<TL: TransportLayer> SeqEx<TL> {
|
||||
let i = next_seq_no as usize % self.recv_window.len();
|
||||
|
||||
if self.recv_window[i].as_ref().map_or(false, |pre| pre.seq_no == next_seq_no) {
|
||||
let next_i = self.next_send_seq_no as usize % self.send_window.len();
|
||||
if self.send_window[next_i].is_some() || self.reserved_len >= self.reserved.len() {
|
||||
if self.is_full_inner(0) {
|
||||
return Err(Error::WindowIsFull);
|
||||
}
|
||||
|
||||
let entry = self.recv_window[i].take().unwrap();
|
||||
self.pre_recv_seq_no = next_seq_no;
|
||||
self.reserved[self.reserved_len] = entry.seq_no;
|
||||
self.reserved_len += 1;
|
||||
self.concurrent_replies[self.concurrent_replies_total] = entry.seq_no;
|
||||
self.concurrent_replies_total += 1;
|
||||
let data = entry.reply_no.and_then(|r| self.take_send(r));
|
||||
Ok((entry.seq_no, entry.data, data))
|
||||
} else {
|
||||
@@ -307,20 +337,28 @@ impl<TL: TransportLayer> SeqEx<TL> {
|
||||
}
|
||||
}
|
||||
|
||||
/// This function must be passed a reply number given by `receive_raw` or `pump_raw`, otherwise
|
||||
/// it will do nothing. This reply number can only be used to reply once.
|
||||
///
|
||||
/// 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 fn reply_raw(&mut self, mut app: TL, reply_no: SeqNo, packet_data: TL::SendData) {
|
||||
if self.remove_reservation(reply_no) {
|
||||
let seq_no = self.next_send_seq_no;
|
||||
self.next_send_seq_no = self.next_send_seq_no.wrapping_add(1);
|
||||
|
||||
let i = seq_no as usize % self.send_window.len();
|
||||
let current_time = app.time();
|
||||
let next_resend_time = current_time + self.resend_interval;
|
||||
if self.next_service_timestamp > next_resend_time {
|
||||
self.next_service_timestamp = next_resend_time;
|
||||
app.update_service_time(current_time, next_resend_time);
|
||||
app.update_service_time(next_resend_time, current_time);
|
||||
}
|
||||
debug_assert!(self.send_window[i].is_none());
|
||||
let entry = self.send_window[i].insert(SendEntry {
|
||||
let slot = self.send_window_slot_mut(reply_no);
|
||||
debug_assert!(slot.is_none());
|
||||
let entry = slot.insert(SendEntry {
|
||||
seq_no,
|
||||
reply_no: Some(reply_no),
|
||||
next_resend_time,
|
||||
@@ -332,14 +370,17 @@ impl<TL: TransportLayer> SeqEx<TL> {
|
||||
}
|
||||
pub fn reply_empty_raw(&mut self, mut app: TL, reply_no: SeqNo) {
|
||||
if self.remove_reservation(reply_no) {
|
||||
// Empty replies are only sent once. There is code in `receive_raw` to handle resending
|
||||
// an empty reply in the event that the first one here was dropped by the network.
|
||||
app.send_empty_reply(reply_no);
|
||||
}
|
||||
}
|
||||
fn remove_reservation(&mut self, reply_no: SeqNo) -> bool {
|
||||
for i in 0..self.reserved_len {
|
||||
if self.reserved[i] == reply_no {
|
||||
self.reserved_len -= 1;
|
||||
self.reserved[i] = self.reserved[self.reserved_len];
|
||||
for i in 0..self.concurrent_replies_total {
|
||||
if self.concurrent_replies[i] == reply_no {
|
||||
// swap remove
|
||||
self.concurrent_replies_total -= 1;
|
||||
self.concurrent_replies[i] = self.concurrent_replies[self.concurrent_replies_total];
|
||||
return true;
|
||||
}
|
||||
}
|
||||
@@ -359,7 +400,7 @@ impl<TL: TransportLayer> SeqEx<TL> {
|
||||
}
|
||||
}
|
||||
self.next_service_timestamp = next_activity;
|
||||
app.update_service_time(current_time, next_activity);
|
||||
app.update_service_time(next_activity, current_time);
|
||||
self.resend_interval.min(next_activity - current_time)
|
||||
}
|
||||
|
||||
|
||||
@@ -2,12 +2,11 @@ use crate::{Error, SeqEx, SeqNo, TransportLayer};
|
||||
|
||||
pub struct ReplyGuard<'a, TL: TransportLayer>(&'a mut SeqEx<TL>, TL, SeqNo);
|
||||
impl<'a, TL: TransportLayer> ReplyGuard<'a, TL> {
|
||||
pub fn seq_no(&self) -> SeqNo {
|
||||
self.0.seq_no()
|
||||
}
|
||||
pub fn reply_no(&self) -> SeqNo {
|
||||
self.2
|
||||
}
|
||||
/// 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 fn reply(self, packet_data: TL::SendData) {
|
||||
self.0.reply_raw(self.1.clone(), self.2, packet_data);
|
||||
core::mem::forget(self);
|
||||
|
||||
@@ -22,6 +22,11 @@ pub struct SeqExSync<TL: TransportLayer, const SLEN: usize = DEFAULT_SEND_WINDOW
|
||||
|
||||
pub struct ReplyGuard<'a, TL: TransportLayer>(&'a SeqExSync<TL>, TL, SeqNo);
|
||||
impl<'a, TL: TransportLayer> ReplyGuard<'a, TL> {
|
||||
/// 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 fn reply(self, packet_data: TL::SendData) {
|
||||
let mut seq = self.0.seq_ex.lock().unwrap();
|
||||
seq.reply_raw(self.1.clone(), self.2, packet_data);
|
||||
|
||||
@@ -12,7 +12,7 @@ pub trait TransportLayer: Sized + Clone {
|
||||
|
||||
fn time(&mut self) -> i64;
|
||||
#[allow(unused)]
|
||||
fn update_service_time(&mut self, current_time: i64, timestamp: i64) {}
|
||||
fn update_service_time(&mut self, timestamp: i64, current_time: i64) {}
|
||||
|
||||
fn send(&mut self, seq_no: SeqNo, reply_no: Option<SeqNo>, payload: &Self::SendData);
|
||||
fn send_ack(&mut self, reply_no: SeqNo);
|
||||
|
||||
Reference in New Issue
Block a user