use super::{Arranging, ArrangingSystem};
use crate::packet::SequenceNumber;
use std::collections::HashMap;
pub struct OrderingSystem<T> {
streams: HashMap<u8, OrderingStream<T>>,
}
impl<T> OrderingSystem<T> {
pub fn new() -> OrderingSystem<T> {
OrderingSystem {
streams: HashMap::with_capacity(32),
}
}
}
impl<'a, T> ArrangingSystem for OrderingSystem<T> {
type Stream = OrderingStream<T>;
fn stream_count(&self) -> usize {
self.streams.len()
}
fn get_or_create_stream(&mut self, stream_id: u8) -> &mut Self::Stream {
self.streams
.entry(stream_id)
.or_insert_with(|| OrderingStream::new(stream_id))
}
}
pub struct OrderingStream<T> {
_stream_id: u8,
storage: HashMap<usize, T>,
expected_index: usize,
unique_item_identifier: u16,
}
impl<T> OrderingStream<T> {
pub fn new(stream_id: u8) -> OrderingStream<T> {
OrderingStream::with_capacity(1024, stream_id)
}
pub fn with_capacity(size: usize, stream_id: u8) -> OrderingStream<T> {
OrderingStream {
storage: HashMap::with_capacity(size),
expected_index: 1,
_stream_id: stream_id,
unique_item_identifier: 0,
}
}
#[cfg(test)]
pub fn stream_id(&self) -> u8 {
self._stream_id
}
#[cfg(test)]
pub fn expected_index(&self) -> usize {
self.expected_index
}
pub fn new_item_identifier(&mut self) -> SequenceNumber {
self.unique_item_identifier = self.unique_item_identifier.wrapping_add(1);
self.unique_item_identifier
}
pub fn iter_mut(&mut self) -> IterMut<T> {
IterMut {
items: &mut self.storage,
expected_index: &mut self.expected_index,
}
}
}
impl<T> Arranging for OrderingStream<T> {
type ArrangingItem = T;
fn arrange(
&mut self,
incoming_offset: usize,
item: Self::ArrangingItem,
) -> Option<Self::ArrangingItem> {
if incoming_offset == self.expected_index {
self.expected_index += 1;
Some(item)
} else if incoming_offset > self.expected_index {
self.storage.insert(incoming_offset, item);
None
} else {
None
}
}
}
pub struct IterMut<'a, T> {
items: &'a mut HashMap<usize, T>,
expected_index: &'a mut usize,
}
impl<'a, T> Iterator for IterMut<'a, T> {
type Item = T;
fn next(&mut self) -> Option<<Self as Iterator>::Item> {
match self.items.remove(&self.expected_index) {
None => None,
Some(e) => {
*self.expected_index += 1;
Some(e)
}
}
}
}
#[cfg(test)]
mod tests {
use super::{Arranging, ArrangingSystem, OrderingSystem};
#[derive(Debug, PartialEq, Clone)]
struct Packet {
pub sequence: usize,
pub ordering_stream: u8,
}
impl Packet {
fn new(sequence: usize, ordering_stream: u8) -> Packet {
Packet {
sequence,
ordering_stream,
}
}
}
#[test]
fn create_stream() {
let mut system: OrderingSystem<Packet> = OrderingSystem::new();
let stream = system.get_or_create_stream(1);
assert_eq!(stream.expected_index(), 1);
assert_eq!(stream.stream_id(), 1);
}
#[test]
fn create_existing_stream() {
let mut system: OrderingSystem<Packet> = OrderingSystem::new();
system.get_or_create_stream(1);
let stream = system.get_or_create_stream(1);
assert_eq!(stream.stream_id(), 1);
}
#[test]
fn can_iterate() {
let mut system: OrderingSystem<Packet> = OrderingSystem::new();
system.get_or_create_stream(1);
let stream = system.get_or_create_stream(1);
let stub_packet1 = Packet::new(1, 1);
let stub_packet2 = Packet::new(2, 1);
let stub_packet3 = Packet::new(3, 1);
let stub_packet4 = Packet::new(4, 1);
let stub_packet5 = Packet::new(5, 1);
{
assert_eq!(
stream.arrange(1, stub_packet1.clone()).unwrap(),
stub_packet1
);
stream.arrange(4, stub_packet4.clone()).is_none();
stream.arrange(5, stub_packet5.clone()).is_none();
stream.arrange(3, stub_packet3.clone()).is_none();
}
{
let mut iterator = stream.iter_mut();
assert_eq!(iterator.next(), None);
}
{
assert_eq!(
stream.arrange(2, stub_packet2.clone()).unwrap(),
stub_packet2
);
}
{
let mut iterator = stream.iter_mut();
assert_eq!(iterator.next().unwrap(), stub_packet3);
assert_eq!(iterator.next().unwrap(), stub_packet4);
assert_eq!(iterator.next().unwrap(), stub_packet5);
}
}
macro_rules! assert_order {
( [$( $x:expr ),*] , [$( $y:expr),*] , $stream_id:expr) => {
{
let mut before: Vec<usize> = Vec::new();
$(
before.push($x);
)*
let mut after: Vec<usize> = Vec::new();
$(
after.push($y);
)*
let mut packets = Vec::new();
for (_, v) in before.iter().enumerate() {
packets.push(Packet::new(*v, $stream_id));
}
let mut ordering_system = OrderingSystem::<Packet>::new();
let stream = ordering_system.get_or_create_stream(1);
let mut ordered_packets = Vec::new();
for packet in packets.into_iter() {
match stream.arrange(packet.sequence, packet.clone()) {
Some(packet) => {
ordered_packets.push(packet.sequence);
let mut iter = stream.iter_mut();
while let Some(packet) = iter.next() {
ordered_packets.push(packet.sequence);
}
}
None => {}
};
}
assert_eq!(after, ordered_packets);
}
};
}
#[test]
fn expect_right_order() {
assert_order!([1, 3, 5, 4, 2], [1, 2, 3, 4, 5], 1);
assert_order!([1, 5, 4, 3, 2], [1, 2, 3, 4, 5], 1);
assert_order!([5, 3, 4, 2, 1], [1, 2, 3, 4, 5], 1);
assert_order!([4, 3, 2, 1, 5], [1, 2, 3, 4, 5], 1);
assert_order!([2, 1, 4, 3, 5], [1, 2, 3, 4, 5], 1);
assert_order!([5, 2, 1, 4, 3], [1, 2, 3, 4, 5], 1);
assert_order!([3, 2, 4, 1, 5], [1, 2, 3, 4, 5], 1);
assert_order!([2, 1, 4, 3, 5], [1, 2, 3, 4, 5], 1);
}
#[test]
fn order_on_multiple_streams() {
assert_order!([1, 3, 5, 4, 2], [1, 2, 3, 4, 5], 1);
assert_order!([1, 5, 4, 3, 2], [1, 2, 3, 4, 5], 2);
assert_order!([5, 3, 4, 2, 1], [1, 2, 3, 4, 5], 3);
assert_order!([4, 3, 2, 1, 5], [1, 2, 3, 4, 5], 4);
assert_order!([2, 1, 4, 3, 5], [1, 2, 3, 4, 5], 5);
assert_order!([5, 2, 1, 4, 3], [1, 2, 3, 4, 5], 6);
assert_order!([3, 2, 4, 1, 5], [1, 2, 3, 4, 5], 7);
assert_order!([2, 1, 4, 3, 5], [1, 2, 3, 4, 5], 8);
}
}