From cd1c00c64f02694d6eda2fa2ce71214dfffaf0cc Mon Sep 17 00:00:00 2001 From: Monica Moniot Date: Thu, 27 Jul 2023 08:03:05 -0400 Subject: [PATCH] version 1 of seq_queue --- .gitignore | 1 + Cargo.lock | 75 +++++++++++ Cargo.toml | 13 ++ rustfmt.toml | 8 ++ src/lib.rs | 41 ++++++ src/seq_queue.rs | 322 +++++++++++++++++++++++++++++++++++++++++++++++ 6 files changed, 460 insertions(+) create mode 100644 .gitignore create mode 100644 Cargo.lock create mode 100644 Cargo.toml create mode 100644 rustfmt.toml create mode 100644 src/lib.rs create mode 100644 src/seq_queue.rs diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..ea8c4bf --- /dev/null +++ b/.gitignore @@ -0,0 +1 @@ +/target diff --git a/Cargo.lock b/Cargo.lock new file mode 100644 index 0000000..f44f6ae --- /dev/null +++ b/Cargo.lock @@ -0,0 +1,75 @@ +# This file is automatically @generated by Cargo. +# It is not intended for manual editing. +version = 3 + +[[package]] +name = "cfg-if" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "baf1de4339761588bc0619e3cbc0120ee582ebb74b53b4efbf79117bd2da40fd" + +[[package]] +name = "getrandom" +version = "0.2.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "be4136b2a15dd319360be1c07d9933517ccf0be8f16bf62a3bee4f0d618df427" +dependencies = [ + "cfg-if", + "libc", + "wasi", +] + +[[package]] +name = "libc" +version = "0.2.147" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b4668fb0ea861c1df094127ac5f1da3409a82116a4ba74fca2e58ef927159bb3" + +[[package]] +name = "ppv-lite86" +version = "0.2.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5b40af805b3121feab8a3c29f04d8ad262fa8e0561883e7653e024ae4479e6de" + +[[package]] +name = "rand" +version = "0.8.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "34af8d1a0e25924bc5b7c43c079c942339d8f0a8b57c39049bef581b46327404" +dependencies = [ + "libc", + "rand_chacha", + "rand_core", +] + +[[package]] +name = "rand_chacha" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e6c10a63a0fa32252be49d21e7709d4d4baf8d231c2dbce1eaa8141b9b127d88" +dependencies = [ + "ppv-lite86", + "rand_core", +] + +[[package]] +name = "rand_core" +version = "0.6.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ec0be4795e2f6a28069bec0b5ff3e2ac9bafc99e6a9a7dc3547996c5c816922c" +dependencies = [ + "getrandom", +] + +[[package]] +name = "seq_queue" +version = "0.0.1" +dependencies = [ + "rand", +] + +[[package]] +name = "wasi" +version = "0.11.0+wasi-snapshot-preview1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9c8d87e72b64a3b4db28d11ce29237c246188f4f51057d65a7eab63b7987e423" diff --git a/Cargo.toml b/Cargo.toml new file mode 100644 index 0000000..6f35043 --- /dev/null +++ b/Cargo.toml @@ -0,0 +1,13 @@ +[package] +name = "seq_queue" +version = "0.0.1" +authors = ["Monica Moniot"] +edition = "2021" + +[lib] +name = "seq_queue" +path = "src/lib.rs" +doc = true + +[dependencies] +rand = "0.8.5" diff --git a/rustfmt.toml b/rustfmt.toml new file mode 100644 index 0000000..3a3929c --- /dev/null +++ b/rustfmt.toml @@ -0,0 +1,8 @@ +max_width = 150 +edition = "2021" +newline_style = "Unix" +struct_lit_width = 60 +tab_spaces = 4 +use_small_heuristics = "Default" +single_line_if_else_max_width = 0 +use_try_shorthand = true diff --git a/src/lib.rs b/src/lib.rs new file mode 100644 index 0000000..444e247 --- /dev/null +++ b/src/lib.rs @@ -0,0 +1,41 @@ + + +pub type SeqNum = u32; + +pub mod seq_queue; + +//struct App(); + +//use seq_queue::ApplicationLayer; +//impl ApplicationLayer for App { +// type RecvData = (Vec, u64); + +// type RecvDataRef<'a> = (&'a [u8], &'a [u8], u64); + +// type RecvReturn = (); + +// type SendData = Vec; + +// fn send(&self, data: &Self::SendData) { +// todo!() +// } +// fn send_ack(&self, reply_num: SeqNum) { +// todo!() +// } +// fn send_empty_reply(&self, reply_num: SeqNum) { +// todo!() +// } + +// fn deserialize<'a>(data: &'a Self::RecvData) -> Self::RecvDataRef<'a> { +// let s = data.0.split_at(1); +// (s.0, s.1, data.1) +// } +// fn process( +// &self, +// packet: Self::RecvDataRef<'_>, +// reply: seq_queue::ReplyGuard<'_, Self>, +// send_data: Option, +// ) -> Self::RecvReturn { +// todo!() +// } +//} diff --git a/src/seq_queue.rs b/src/seq_queue.rs new file mode 100644 index 0000000..53db24d --- /dev/null +++ b/src/seq_queue.rs @@ -0,0 +1,322 @@ + +use crate::SeqNum; + +const INITIAL_SEQ_NUM: SeqNum = 1; +const SEQ_RECV_WINDOW_LEN: usize = 32; +const SEQ_SEND_WINDOW_LEN: usize = 64; + +pub trait ApplicationLayer: Sized { + type RecvData; + type RecvDataRef<'a>; + type RecvReturn; + + type SendData; + + fn send(&self, data: &Self::SendData); + fn send_ack(&self, reply_num: SeqNum); + fn send_empty_reply(&self, reply_num: SeqNum); + + fn deserialize<'a>(data: &'a Self::RecvData) -> Self::RecvDataRef<'a>; + fn process( + &self, + packet: Self::RecvDataRef<'_>, + reply: ReplyGuard<'_, Self>, + send_data: Option, + ) -> Self::RecvReturn; +} + +pub trait IntoRecvData: Into { + fn as_ref(&self) -> App::RecvDataRef<'_>; +} +impl IntoRecvData for AppInner::RecvData { + fn as_ref(&self) -> AppInner::RecvDataRef<'_> { + AppInner::deserialize(self) + } +} + + +pub struct SeqQueue { + pub retry_interval: i64, + next_send_seq_num: SeqNum, + pre_recv_seq_num: SeqNum, + recv_window: [Option>; SEQ_RECV_WINDOW_LEN], + send_window: [Option>; SEQ_SEND_WINDOW_LEN], +} + +struct RecvEntry { + seq_num: SeqNum, + reply_num: Option, + data: App::RecvData, +} + +struct SendEntry { + seq_num: SeqNum, + reply_num: Option, + next_resent_time: i64, + data: App::SendData, +} + +/// Whenever a packet is received, it must be replied to. +/// This Guard object guarantees that this is the case. +/// If it is dropped without calling `reply` an empty reply will be sent to the remote peer. +pub struct ReplyGuard<'a, App: ApplicationLayer> { + app: Option<&'a App>, + seq_queue: &'a mut SeqQueue, + reply_num: SeqNum, +} + +impl SeqQueue { + pub fn new(retry_interval: i64) -> Self { + Self { + retry_interval, + next_send_seq_num: INITIAL_SEQ_NUM, + pre_recv_seq_num: INITIAL_SEQ_NUM - 1, + recv_window: std::array::from_fn(|_| None), + send_window: std::array::from_fn(|_| None), + } + } + + pub fn is_full(&self) -> bool { + self.send_window[self.next_send_seq_num as usize % self.send_window.len()].is_some() + } + /// If the return value is `false` the queue is full and the packet will not be sent. + /// The caller must either cancel sending, abort the connection, or wait until a call to + /// `receive` returns `Some` and try again. + #[must_use = "The queue might be full causing the packet to not be sent"] + pub fn send_seq( + &mut self, + app: App, + create: impl FnOnce(SeqNum) -> App::SendData, + current_time: i64, + ) -> bool { + let seq_num = self.next_send_seq_num; + self.next_send_seq_num += 1; + + let i = seq_num as usize % self.send_window.len(); + if self.send_window[i].is_some() { + return false; + } + let next_resent_time = current_time + self.retry_interval; + let entry = self.send_window[i].insert(SendEntry { + seq_num, + reply_num: None, + next_resent_time, + data: create(seq_num), + }); + + app.send(&entry.data); + true + } + + pub fn receive( + &mut self, + app: App, + seq_num: SeqNum, + reply_num: Option, + packet: impl IntoRecvData, + ) -> Option { + let normalized_seq_num = seq_num.wrapping_sub(self.pre_recv_seq_num).wrapping_sub(1); + let is_below_range = normalized_seq_num > SeqNum::MAX / 2; + let is_above_range = !is_below_range && normalized_seq_num >= self.recv_window.len() as u32; + let is_next = normalized_seq_num == 0; + if is_below_range { + // Check whether or not we are already replying to this packet. + for entry in self.send_window.iter() { + if entry.as_ref().map_or(false, |e| e.reply_num == Some(seq_num)) { + return None; + } + } + app.send_empty_reply(seq_num); + return None; + } else if is_above_range { + return None; + } + let i = seq_num as usize % self.recv_window.len(); + if let Some(pre) = self.recv_window[i].as_mut() { + if seq_num == pre.seq_num { + if is_next { + self.recv_window[i] = None; + } else { + app.send_ack(seq_num); + return None; + } + } else { + return None; + } + } + if is_next { + // This should reliably handle sequence number overflow. + self.pre_recv_seq_num = seq_num; + let data = reply_num.and_then(|r| self.take_send(r)); + Some(app.process( + packet.as_ref(), + ReplyGuard { app: Some(&app), seq_queue: self, reply_num: seq_num }, + data, + )) + } else { + self.recv_window[i] = Some(RecvEntry { + seq_num, + reply_num, + data: packet.into(), + }); + if let Some(reply_num) = reply_num { + self.receive_ack(reply_num); + } + app.send_ack(seq_num); + None + } + } + + fn take_send(&mut self, reply_num: SeqNum) -> Option { + let i = reply_num as usize % self.send_window.len(); + if self.send_window[i].as_ref().map_or(false, |e| e.seq_num == reply_num) { + self.send_window[i].take().map(|e| e.data) + } else { + None + } + } + pub fn pump(&mut self, app: App) -> Option { + let next_seq_num = self.pre_recv_seq_num.wrapping_add(1); + let i = next_seq_num as usize % self.recv_window.len(); + + if self.recv_window[i].as_ref().map_or(false, |pre| pre.seq_num == next_seq_num) { + self.pre_recv_seq_num = next_seq_num; + let entry = self.recv_window[i].take().unwrap(); + let data = entry.reply_num.and_then(|r| self.take_send(r)); + Some(app.process( + App::deserialize(&entry.data), + ReplyGuard { app: Some(&app), seq_queue: self, reply_num: entry.seq_num }, + data, + )) + } else { + None + } + } + + pub fn receive_ack(&mut self, reply_num: SeqNum) { + let i = reply_num as usize % self.send_window.len(); + if let Some(entry) = self.send_window[i].as_mut() { + if entry.seq_num == reply_num { + entry.next_resent_time = i64::MAX; + } + } + } + pub fn receive_empty_reply(&mut self, reply_num: SeqNum) -> Option { + let i = reply_num as usize % self.send_window.len(); + if self.send_window[i].as_ref().map_or(false, |e| e.seq_num == reply_num) { + let entry = self.send_window[i].take().unwrap(); + Some(entry.data) + } else { + None + } + } + + pub fn service(&mut self, app: App, current_time: i64) -> i64 { + let next_interval = current_time + self.retry_interval; + let mut next_activity = next_interval; + for item in self.send_window.iter_mut() { + if let Some(entry) = item { + if entry.next_resent_time <= current_time { + entry.next_resent_time = next_interval; + app.send(&entry.data); + } else { + next_activity = next_activity.min(entry.next_resent_time); + } + } + } + next_activity - current_time + } + + pub fn window_iter(&self) -> Iter<'_, App> { + Iter(self.send_window.iter()) + } + pub fn window_iter_mut(&mut self) -> IterMut<'_, App> { + IterMut(self.send_window.iter_mut()) + } +} +macro_rules! iterator { + ($iter:ident, {$( $mut:tt )?}) => { + impl<'a, App: ApplicationLayer> Iterator for $iter<'a, App> { + type Item = &'a $($mut)? App::SendData; + fn next(&mut self) -> Option { + while let Some(entry) = self.0.next() { + if let Some(entry) = entry { + return Some(& $($mut)? entry.data) + } + } + None + } + + fn size_hint(&self) -> (usize, Option) { + (0, Some(self.0.len())) + } + } + impl<'a, App: ApplicationLayer> DoubleEndedIterator for $iter<'a, App> { + fn next_back(&mut self) -> Option { + while let Some(entry) = self.0.next_back() { + if let Some(entry) = entry { + return Some(& $($mut)? entry.data) + } + } + None + } + } + } +} + +pub struct Iter<'a, App: ApplicationLayer> (std::slice::Iter<'a, Option>>); +pub struct IterMut<'a, App: ApplicationLayer> (std::slice::IterMut<'a, Option>>); + +iterator!(Iter, { }); +iterator!(IterMut, {mut}); + +impl<'a, App: ApplicationLayer> ReplyGuard<'a, App> { + pub fn is_full(&self) -> bool { + let seq_queue = &self.seq_queue; + seq_queue.send_window[seq_queue.next_send_seq_num as usize % seq_queue.send_window.len()].is_some() + } + /// If the return value is `false` the queue is full and the packet will not be sent. + /// The caller must either cancel the reply or abort the connection. + /// If the reply is cancelled then the remote peer will receive an empty reply instead. + /// + /// A packet can only be replied to once. Once a call to `reply` is successful and returns `true`, + /// all subsequent calls will return `false`. + #[must_use = "The queue might be full causing the packet to not be sent"] + pub fn reply( + &mut self, + create: impl FnOnce(SeqNum, SeqNum) -> App::SendData, + current_time: i64, + ) -> bool { + if let Some(app) = self.app { + let seq_queue = &mut self.seq_queue; + let seq_num = seq_queue.next_send_seq_num; + seq_queue.next_send_seq_num += 1; + + let i = seq_num as usize % seq_queue.send_window.len(); + if seq_queue.send_window[i].is_some() { + return false; + } + let next_resent_time = current_time + seq_queue.retry_interval; + let entry = seq_queue.send_window[i].insert(SendEntry { + seq_num, + reply_num: Some(self.reply_num), + next_resent_time, + data: create(seq_num, self.reply_num), + }); + + app.send(&entry.data); + self.app = None; + true + } else { + false + } + } +} + +impl<'a, App: ApplicationLayer> Drop for ReplyGuard<'a, App> { + fn drop(&mut self) { + if let Some(app) = self.app { + app.send_empty_reply(self.reply_num); + } + } +}