use std::collections::BTreeMap;
use std::sync::Mutex;
#[derive(Debug, Default)]
struct TrackState {
pending: BTreeMap<u64, (String, bool)>,
next_index: u64,
}
#[derive(Debug, Default)]
pub(crate) struct Watermark {
state: Mutex<TrackState>,
}
impl Watermark {
pub(crate) fn deliver(&self, sequence: &str) -> u64 {
let mut state = self.state.lock().expect("watermark mutex poisoned");
let index = state.next_index;
state.next_index += 1;
state.pending.insert(index, (sequence.to_owned(), false));
index
}
#[allow(clippy::significant_drop_tightening)]
pub(crate) fn settle(&self, index: u64) -> Option<String> {
let mut state = self.state.lock().expect("watermark mutex poisoned");
if let Some(entry) = state.pending.get_mut(&index) {
entry.1 = true;
}
let mut advanced = None;
while let Some((&front, (_, done))) = state.pending.first_key_value() {
if !done {
break;
}
let (sequence, _) = state.pending.remove(&front).expect("front exists");
advanced = Some(sequence);
}
advanced
}
pub(crate) fn drained(&self) -> bool {
self.state
.lock()
.expect("watermark mutex poisoned")
.pending
.is_empty()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_watermark_advances_only_contiguously() {
let track = Watermark::default();
let a = track.deliver("100");
let b = track.deliver("101");
let c = track.deliver("102");
assert_eq!(track.settle(b), None);
assert_eq!(track.settle(c), None);
assert_eq!(track.settle(a), Some("102".to_owned()));
assert!(track.drained());
}
#[test]
fn an_unsettled_delivery_wedges_the_watermark() {
let track = Watermark::default();
let _skipped = track.deliver("100");
let b = track.deliver("101");
assert_eq!(track.settle(b), None);
assert!(!track.drained());
}
}