use crate::chunk::chunk_payload_data::ChunkPayloadData; use crate::chunk::chunk_selective_ack::GapAckBlock; use crate::util::*; use alloc::string::String; use alloc::vec::Vec; use std::collections::HashMap; #[derive(Default, Debug)] pub(crate) struct PayloadQueue { // length: usize, chunk_map: HashMap, pub(crate) sorted: Vec, dup_tsn: Vec, n_bytes: usize, } impl PayloadQueue { pub(crate) fn new() -> Self { PayloadQueue::default() } pub(crate) fn update_sorted_keys(&mut self) { self.sorted.sort_by(|a, b| { if sna32lt(*a, *b) { core::cmp::Ordering::Less } else { core::cmp::Ordering::Greater } }); } pub(crate) fn can_push(&self, p: &ChunkPayloadData, cumulative_tsn: u32) -> bool { !(self.chunk_map.contains_key(&p.tsn) || sna32lte(p.tsn, cumulative_tsn)) } pub(crate) fn push_no_check(&mut self, p: ChunkPayloadData) { self.n_bytes += p.user_data.len(); self.sorted.push(p.tsn); self.chunk_map.insert(p.tsn, p); //self.length += 1; self.update_sorted_keys(); } /// push pushes a payload data. If the payload data is already in our queue or /// older than our cumulative_tsn marker, it will be recored as duplications, /// which can later be retrieved using popDuplicates. pub(crate) fn push(&mut self, p: ChunkPayloadData, cumulative_tsn: u32) -> bool { let ok = self.chunk_map.contains_key(&p.tsn); if ok || sna32lte(p.tsn, cumulative_tsn) { // Found the packet, log in dups self.dup_tsn.push(p.tsn); return false; } self.n_bytes += p.user_data.len(); self.sorted.push(p.tsn); self.chunk_map.insert(p.tsn, p); //self.length += 1; self.update_sorted_keys(); true } /// pop pops only if the oldest chunk's TSN matches the given TSN. pub(crate) fn pop(&mut self, tsn: u32) -> Option { if !self.sorted.is_empty() && tsn == self.sorted[0] { self.sorted.remove(0); if let Some(c) = self.chunk_map.remove(&tsn) { //self.length -= 1; self.n_bytes = self.n_bytes.saturating_sub(c.user_data.len()); return Some(c); } } None } /// Removes every queued chunk with a TSN at or before `cumulative_tsn`. /// /// Used when a FORWARD-TSN moves the cumulative TSN point past chunks the /// peer abandoned. The cost is proportional to the number of queued chunks, /// not to the size of the TSN jump: `cumulative_tsn` comes off the wire and /// may be up to 2^31 ahead of the current point. pub(crate) fn pop_up_to(&mut self, cumulative_tsn: u32) { let chunk_map = &mut self.chunk_map; let n_bytes = &mut self.n_bytes; self.sorted.retain(|tsn| { if sna32lte(*tsn, cumulative_tsn) { if let Some(c) = chunk_map.remove(tsn) { *n_bytes = n_bytes.saturating_sub(c.user_data.len()); } false } else { true } }); } /// get returns reference to chunkPayloadData with the given TSN value. pub(crate) fn get(&self, tsn: u32) -> Option<&ChunkPayloadData> { self.chunk_map.get(&tsn) } pub(crate) fn get_mut(&mut self, tsn: u32) -> Option<&mut ChunkPayloadData> { self.chunk_map.get_mut(&tsn) } /// popDuplicates returns an array of TSN values that were found duplicate. pub(crate) fn pop_duplicates(&mut self) -> Vec { core::mem::take(&mut self.dup_tsn) } pub(crate) fn get_gap_ack_blocks(&self, cumulative_tsn: u32) -> Vec { if self.chunk_map.is_empty() { return vec![]; } let mut b = GapAckBlock::default(); let mut gap_ack_blocks = vec![]; for (i, tsn) in self.sorted.iter().enumerate() { let diff = if *tsn >= cumulative_tsn { (*tsn - cumulative_tsn) as u16 } else { 0 }; if i == 0 { b.start = diff; b.end = b.start; } else if b.end + 1 == diff { b.end += 1; } else { gap_ack_blocks.push(b); b.start = diff; b.end = diff; } } gap_ack_blocks.push(b); gap_ack_blocks } pub(crate) fn get_gap_ack_blocks_string(&self, cumulative_tsn: u32) -> String { let mut s = format!("cumTSN={}", cumulative_tsn); for b in self.get_gap_ack_blocks(cumulative_tsn) { s += format!(",{}-{}", b.start, b.end).as_str(); } s } pub(crate) fn mark_as_acked(&mut self, tsn: u32) -> usize { if let Some(c) = self.chunk_map.get_mut(&tsn) { c.acked = true; c.retransmit = false; let n = c.user_data.len(); self.n_bytes -= n; c.user_data.clear(); n } else { 0 } } pub(crate) fn get_last_tsn_received(&self) -> Option<&u32> { self.sorted.last() } pub(crate) fn mark_all_to_retrasmit(&mut self) { for c in self.chunk_map.values_mut() { if c.acked || c.abandoned() { continue; } c.retransmit = true; } } pub(crate) fn get_num_bytes(&self) -> usize { self.n_bytes } pub(crate) fn len(&self) -> usize { //assert_eq!(self.chunk_map.len(), self.length); self.chunk_map.len() } pub(crate) fn is_empty(&self) -> bool { self.len() == 0 } }