From 69a7872e0d76a328081bc861a35077d4eb40a340 Mon Sep 17 00:00:00 2001 From: Monica Moniot Date: Wed, 9 Aug 2023 12:07:35 -0400 Subject: [PATCH] added timestamp updating --- src/seq_queue.rs | 127 +++++++++++++++++++++++++++-------------- src/single_thread.rs | 11 ++-- src/sync.rs | 5 ++ src/transport_layer.rs | 2 +- 4 files changed, 95 insertions(+), 50 deletions(-) diff --git a/src/seq_queue.rs b/src/seq_queue.rs index a0df572..b99e947 100644 --- a/src/seq_queue.rs +++ b/src/seq_queue.rs @@ -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 { /// 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>; SLEN], recv_window: [Option>; 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 { @@ -122,18 +124,24 @@ impl SeqEx { 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> { + &mut self.send_window[seq_no as usize % self.send_window.len()] + } + fn send_window_slot(&self, seq_no: SeqNo) -> &Option> { + &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 SeqEx { 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 SeqEx { reply_no: Option, packet: P, ) -> Result<(SeqNo, P, Option), 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 SeqEx { 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 SeqEx { // 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 SeqEx { }; } } 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 SeqEx { } } 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 { - 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 { - 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 SeqEx { 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 SeqEx { } } + /// 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 SeqEx { } 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 SeqEx { } } 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) } diff --git a/src/single_thread.rs b/src/single_thread.rs index d8cce40..2b17dbc 100644 --- a/src/single_thread.rs +++ b/src/single_thread.rs @@ -2,12 +2,11 @@ use crate::{Error, SeqEx, SeqNo, TransportLayer}; pub struct ReplyGuard<'a, TL: TransportLayer>(&'a mut SeqEx, 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); diff --git a/src/sync.rs b/src/sync.rs index d3d3585..d5b9c90 100644 --- a/src/sync.rs +++ b/src/sync.rs @@ -22,6 +22,11 @@ pub struct SeqExSync(&'a SeqExSync, 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); diff --git a/src/transport_layer.rs b/src/transport_layer.rs index 482aaa2..d2c0bd8 100644 --- a/src/transport_layer.rs +++ b/src/transport_layer.rs @@ -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, payload: &Self::SendData); fn send_ack(&mut self, reply_no: SeqNo);