use std::collections::{BTreeMap, HashMap, btree_map};
use std::sync::Mutex;
use weida_protocol::header::OrderingMode;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct Gap {
pub expected: u64,
pub seen: u64,
}
impl Gap {
pub fn missed(&self) -> u64 {
self.seen.saturating_sub(self.expected)
}
}
pub(crate) struct Sequencer {
enabled: bool,
counters: Mutex<HashMap<Box<str>, u64>>,
}
impl Sequencer {
pub(crate) fn new(ordering: OrderingMode) -> Sequencer {
Sequencer {
enabled: ordering != OrderingMode::None,
counters: Mutex::new(HashMap::new()),
}
}
pub(crate) fn next(&self, scope: &str) -> Option<u64> {
if !self.enabled {
return None;
}
let mut counters = self.counters.lock().expect("sequencer poisoned");
let next = counters.entry_ref_or_insert(scope);
Some(next)
}
#[cfg(test)]
fn tracked_scopes(&self) -> usize {
self.counters.lock().expect("sequencer poisoned").len()
}
#[cfg(test)]
fn allocated(&self) -> bool {
self.counters.lock().expect("sequencer poisoned").capacity() > 0
}
}
trait CounterMap {
fn entry_ref_or_insert(&mut self, scope: &str) -> u64;
}
impl CounterMap for HashMap<Box<str>, u64> {
fn entry_ref_or_insert(&mut self, scope: &str) -> u64 {
match self.get_mut(scope) {
Some(counter) => {
let issued = *counter;
*counter = counter.wrapping_add(1);
issued
}
None => {
self.insert(scope.into(), 1);
0
}
}
}
}
pub(crate) struct GapDetector {
enabled: bool,
max_scopes: usize,
expected: Mutex<HashMap<Box<str>, u64>>,
}
impl GapDetector {
pub(crate) fn new(ordering: OrderingMode, max_scopes: usize) -> GapDetector {
GapDetector {
enabled: ordering == OrderingMode::PerProducerDetect,
max_scopes,
expected: Mutex::new(HashMap::new()),
}
}
pub(crate) fn observe(&self, scope: &str, sequence: u64) -> Option<Gap> {
if !self.enabled {
return None;
}
let mut expected = self.expected.lock().expect("gap detector poisoned");
let next = expected.get_mut(scope);
match next {
Some(next) => {
let gap = (sequence > *next).then_some(Gap {
expected: *next,
seen: sequence,
});
if sequence >= *next {
*next = sequence.wrapping_add(1);
}
gap
}
None => {
if expected.len() >= self.max_scopes {
return None;
}
expected.insert(scope.into(), sequence.wrapping_add(1));
None
}
}
}
#[cfg(test)]
fn allocated(&self) -> bool {
self.expected
.lock()
.expect("gap detector poisoned")
.capacity()
> 0
}
}
pub(crate) struct Reassembler<T> {
enabled: bool,
max_hold: usize,
max_scopes: usize,
state: Mutex<Reorder<T>>,
}
struct Reorder<T> {
scopes: HashMap<Box<str>, Scope<T>>,
held: usize,
}
struct Scope<T> {
next: u64,
pending: BTreeMap<u64, T>,
}
impl<T> Reassembler<T> {
pub(crate) fn new(
ordering: OrderingMode,
max_hold: usize,
max_scopes: usize,
) -> Reassembler<T> {
Reassembler {
enabled: ordering == OrderingMode::PerProducerReassemble,
max_hold,
max_scopes,
state: Mutex::new(Reorder {
scopes: HashMap::new(),
held: 0,
}),
}
}
pub(crate) fn enabled(&self) -> bool {
self.enabled
}
pub(crate) fn admit(
&self,
scope: &str,
sequence: Option<u64>,
item: T,
) -> Vec<(T, Option<Gap>)> {
if !self.enabled {
return vec![(item, None)];
}
let Some(sequence) = sequence else {
return vec![(item, None)];
};
let mut state = self.state.lock().expect("reassembler poisoned");
let Reorder { scopes, held } = &mut *state;
if !scopes.contains_key(scope) {
if scopes.len() >= self.max_scopes {
return vec![(item, None)];
}
scopes.insert(
scope.into(),
Scope {
next: sequence,
pending: BTreeMap::new(),
},
);
}
let entry = scopes.get_mut(scope).expect("present");
if sequence < entry.next {
return vec![(item, None)];
}
if sequence == entry.next && entry.pending.is_empty() {
entry.next = sequence.wrapping_add(1);
return vec![(item, None)];
}
match entry.pending.entry(sequence) {
btree_map::Entry::Occupied(_) => return vec![(item, None)],
btree_map::Entry::Vacant(slot) => slot.insert(item),
};
*held += 1;
let mut out = Vec::new();
*held -= drain_in_order(entry, &mut out);
while *held > self.max_hold {
let Some(victim) = widest_scope(scopes) else {
break;
};
let Some((seq, item)) = victim.pending.pop_first() else {
break;
};
let gap = Gap {
expected: victim.next,
seen: seq,
};
victim.next = seq.wrapping_add(1);
out.push((item, Some(gap)));
*held -= 1;
*held -= drain_in_order(victim, &mut out);
}
out
}
#[cfg(test)]
fn held(&self) -> usize {
self.state.lock().expect("reassembler poisoned").held
}
#[cfg(test)]
fn allocated(&self) -> bool {
self.state
.lock()
.expect("reassembler poisoned")
.scopes
.capacity()
> 0
}
}
fn drain_in_order<T>(scope: &mut Scope<T>, out: &mut Vec<(T, Option<Gap>)>) -> usize {
let mut released = 0;
while let Some(entry) = scope.pending.first_entry() {
if *entry.key() != scope.next {
break;
}
out.push((entry.remove(), None));
scope.next = scope.next.wrapping_add(1);
released += 1;
}
released
}
fn widest_scope<T>(scopes: &mut HashMap<Box<str>, Scope<T>>) -> Option<&mut Scope<T>> {
scopes
.values_mut()
.max_by_key(|scope| scope.pending.len())
.filter(|scope| !scope.pending.is_empty())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_disabled_sequencer_numbers_nothing_and_allocates_nothing() {
let sequencer = Sequencer::new(OrderingMode::None);
for _ in 0..1000 {
assert_eq!(sequencer.next("/md"), None);
}
assert_eq!(sequencer.tracked_scopes(), 0);
assert!(
!sequencer.allocated(),
"an off sequencer must not allocate a table"
);
}
#[test]
fn a_disabled_detector_reports_nothing_and_allocates_nothing() {
let detector = GapDetector::new(OrderingMode::None, 1024);
for seq in [0, 5, 1, 900] {
assert_eq!(detector.observe("/md", seq), None);
}
assert!(
!detector.allocated(),
"an off detector must not allocate a table"
);
}
#[test]
fn numbers_are_dense_and_per_scope() {
let sequencer = Sequencer::new(OrderingMode::PerProducerDetect);
assert_eq!(sequencer.next("a"), Some(0));
assert_eq!(sequencer.next("a"), Some(1));
assert_eq!(sequencer.next("b"), Some(0));
assert_eq!(sequencer.next("a"), Some(2));
assert_eq!(sequencer.tracked_scopes(), 2);
}
#[test]
fn a_missing_number_is_a_gap_of_that_size() {
let detector = GapDetector::new(OrderingMode::PerProducerDetect, 1024);
assert_eq!(detector.observe("a", 0), None);
assert_eq!(detector.observe("a", 1), None);
let gap = detector.observe("a", 4).expect("a gap");
assert_eq!(
gap,
Gap {
expected: 2,
seen: 4
}
);
assert_eq!(gap.missed(), 2);
assert_eq!(detector.observe("a", 5), None);
}
#[test]
fn scopes_do_not_interfere() {
let detector = GapDetector::new(OrderingMode::PerProducerDetect, 1024);
assert_eq!(detector.observe("a", 0), None);
assert_eq!(detector.observe("b", 0), None);
assert_eq!(detector.observe("a", 1), None);
assert_eq!(detector.observe("b", 7).expect("gap").missed(), 6);
}
#[test]
fn an_out_of_order_arrival_is_not_a_gap() {
let detector = GapDetector::new(OrderingMode::PerProducerDetect, 1024);
assert_eq!(detector.observe("a", 0), None);
assert_eq!(detector.observe("a", 2).expect("gap").missed(), 1);
assert_eq!(detector.observe("a", 1), None);
assert_eq!(detector.observe("a", 4).expect("gap").missed(), 1);
}
#[test]
fn the_scope_table_stops_growing_at_its_cap() {
let detector = GapDetector::new(OrderingMode::PerProducerDetect, 4);
for i in 0..100 {
detector.observe(&format!("scope-{i}"), 0);
}
assert_eq!(
detector.expected.lock().expect("poisoned").len(),
4,
"a peer must not be able to size this table"
);
}
fn items(released: Vec<(u64, Option<Gap>)>) -> Vec<u64> {
released.into_iter().map(|(item, _)| item).collect()
}
#[test]
fn a_disabled_reassembler_holds_nothing_and_allocates_nothing() {
for mode in [OrderingMode::None, OrderingMode::PerProducerDetect] {
let reassembler: Reassembler<u64> = Reassembler::new(mode, 256, 1024);
assert!(!reassembler.enabled());
for seq in 0..1000 {
assert_eq!(items(reassembler.admit("/md", Some(999 - seq), seq)), [seq]);
}
assert_eq!(reassembler.held(), 0);
assert!(
!reassembler.allocated(),
"an off reassembler must not allocate a table"
);
}
}
#[test]
fn out_of_order_arrivals_are_released_in_sequence_order() {
let r: Reassembler<u64> = Reassembler::new(OrderingMode::PerProducerReassemble, 256, 1024);
assert_eq!(items(r.admit("a", Some(0), 0)), [0]);
assert!(items(r.admit("a", Some(2), 2)).is_empty());
assert!(items(r.admit("a", Some(3), 3)).is_empty());
assert_eq!(r.held(), 2);
assert_eq!(items(r.admit("a", Some(1), 1)), [1, 2, 3]);
assert_eq!(r.held(), 0, "the hold must drain");
}
#[test]
fn a_run_released_in_order_carries_no_gap() {
let r: Reassembler<u64> = Reassembler::new(OrderingMode::PerProducerReassemble, 256, 1024);
r.admit("a", Some(0), 0);
r.admit("a", Some(2), 2);
let released = r.admit("a", Some(1), 1);
assert!(
released.iter().all(|(_, gap)| gap.is_none()),
"nothing was missed: the hole was filled"
);
}
#[test]
fn the_hold_releases_out_of_order_at_its_cap_and_reports_the_gap() {
let r: Reassembler<u64> = Reassembler::new(OrderingMode::PerProducerReassemble, 2, 1024);
assert_eq!(items(r.admit("a", Some(0), 0)), [0]);
assert!(r.admit("a", Some(2), 2).is_empty());
assert!(r.admit("a", Some(3), 3).is_empty());
let released = r.admit("a", Some(4), 4);
assert_eq!(items(released.clone()), [2, 3, 4]);
assert_eq!(
released[0].1.expect("a gap"),
Gap {
expected: 1,
seen: 2
}
);
assert!(
released[1..].iter().all(|(_, gap)| gap.is_none()),
"3 and 4 followed 2 in order"
);
assert_eq!(r.held(), 0, "the hold drained behind the forced release");
}
#[test]
fn a_repeated_number_is_delivered_and_does_not_move_the_hold() {
let r: Reassembler<u64> = Reassembler::new(OrderingMode::PerProducerReassemble, 4, 1024);
let mut released = items(r.admit("a", Some(3), 30));
assert_eq!(released, [30], "the position starts at the first number");
assert!(r.admit("a", Some(5), 50).is_empty());
released.extend(items(r.admit("a", Some(5), 51)));
assert_eq!(r.held(), 1, "a repeat must not inflate the hold");
released.extend(items(r.admit("a", Some(4), 40)));
released.sort_unstable();
assert_eq!(released, [30, 40, 50, 51]);
assert_eq!(r.held(), 0, "the hold must drain");
}
#[test]
fn the_hold_never_exceeds_its_cap() {
let r: Reassembler<u64> = Reassembler::new(OrderingMode::PerProducerReassemble, 8, 1024);
r.admit("a", Some(0), 0);
for seq in 2..500u64 {
r.admit("a", Some(seq), seq);
assert!(r.held() <= 8, "the hold must be bounded by its cap");
}
}
#[test]
fn one_stalled_scope_does_not_spend_another_scope_s_budget() {
let r: Reassembler<u64> = Reassembler::new(OrderingMode::PerProducerReassemble, 4, 1024);
r.admit("stalled", Some(0), 0);
r.admit("busy", Some(0), 100);
r.admit("stalled", Some(2), 2);
r.admit("stalled", Some(3), 3);
r.admit("busy", Some(102), 102);
r.admit("busy", Some(103), 103);
let released = r.admit("busy", Some(104), 104);
assert!(!released.is_empty(), "the cap must release something");
assert!(r.held() <= 4);
}
#[test]
fn an_unnumbered_or_late_arrival_passes_straight_through() {
let r: Reassembler<u64> = Reassembler::new(OrderingMode::PerProducerReassemble, 256, 1024);
assert_eq!(items(r.admit("a", None, 7)), [7]);
assert_eq!(items(r.admit("a", Some(5), 5)), [5]);
assert_eq!(items(r.admit("a", Some(4), 4)), [4]);
assert_eq!(r.held(), 0);
}
#[test]
fn the_reassembler_scope_table_stops_growing_at_its_cap() {
let r: Reassembler<u64> = Reassembler::new(OrderingMode::PerProducerReassemble, 256, 4);
for i in 0..100u64 {
r.admit(&format!("scope-{i}"), Some(0), i);
r.admit(&format!("scope-{i}"), Some(2), i);
}
assert_eq!(
r.state.lock().expect("poisoned").scopes.len(),
4,
"a peer must not be able to size this table"
);
assert_eq!(r.held(), 4, "only the tracked scopes hold anything");
}
}