use std::collections::BTreeMap;
use bytes::Bytes;
use crate::codec::OrderedMap;
use crate::types::messaging::DeliveryState;
#[derive(Debug, Clone)]
pub struct UnsettledEntry {
pub delivery_tag: Bytes,
pub state: Option<DeliveryState>,
pub settled: bool,
}
#[derive(Debug, Default)]
pub struct UnsettledMap {
entries: BTreeMap<u32, UnsettledEntry>,
}
impl UnsettledMap {
pub fn new() -> Self {
UnsettledMap {
entries: BTreeMap::new(),
}
}
pub fn len(&self) -> usize {
self.entries.len()
}
pub fn is_empty(&self) -> bool {
self.entries.is_empty()
}
pub fn insert(&mut self, delivery_id: u32, delivery_tag: Bytes, state: Option<DeliveryState>) {
self.entries.insert(
delivery_id,
UnsettledEntry {
delivery_tag,
state,
settled: false,
},
);
}
pub fn get(&self, delivery_id: u32) -> Option<&UnsettledEntry> {
self.entries.get(&delivery_id)
}
pub fn contains(&self, delivery_id: u32) -> bool {
self.entries.contains_key(&delivery_id)
}
pub fn settle(&mut self, delivery_id: u32) -> Option<UnsettledEntry> {
self.entries.remove(&delivery_id)
}
pub fn apply_disposition(
&mut self,
first: u32,
last: u32,
state: Option<&DeliveryState>,
settled: bool,
) -> Vec<u32> {
let ids: Vec<u32> = self
.entries
.range(first..=last)
.map(|(id, _)| *id)
.collect();
for id in &ids {
if let Some(entry) = self.entries.get_mut(id) {
if state.is_some() {
entry.state = state.cloned();
}
entry.settled = settled;
}
}
if settled {
for id in &ids {
self.entries.remove(id);
}
}
ids
}
pub fn snapshot(&self) -> OrderedMap<Bytes, DeliveryState> {
let mut map = OrderedMap::with_capacity(self.entries.len());
for entry in self.entries.values() {
if let Some(state) = &entry.state {
map.push(entry.delivery_tag.clone(), state.clone());
}
}
map
}
pub fn iter(&self) -> impl Iterator<Item = (u32, &UnsettledEntry)> {
self.entries.iter().map(|(id, e)| (*id, e))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::types::messaging::{Accepted, Rejected};
use crate::types::definitions::AmqpError;
#[test]
fn insert_settle_and_range() {
let mut m = UnsettledMap::new();
for id in 1..=5 {
m.insert(id, Bytes::from(vec![id as u8]), None);
}
assert_eq!(m.len(), 5);
let affected = m.apply_disposition(
2,
4,
Some(&DeliveryState::Accepted(Accepted::default())),
true,
);
assert_eq!(affected, vec![2, 3, 4]);
assert_eq!(m.len(), 2); assert!(m.contains(1) && m.contains(5));
m.apply_disposition(
1,
1,
Some(&DeliveryState::Rejected(Rejected {
error: Some(crate::types::definitions::Error::new(AmqpError::NotFound, None)),
})),
false,
);
let e = m.get(1).unwrap();
assert!(!e.settled);
assert!(matches!(e.state, Some(DeliveryState::Rejected(_))));
}
#[test]
fn snapshot_preserves_states_in_order() {
let mut m = UnsettledMap::new();
m.insert(2, Bytes::from_static(b"b"), Some(DeliveryState::Accepted(Accepted::default())));
m.insert(1, Bytes::from_static(b"a"), Some(DeliveryState::Accepted(Accepted::default())));
let snap = m.snapshot();
let tags: Vec<_> = snap.iter().map(|(k, _)| k.clone()).collect();
assert_eq!(tags, vec![Bytes::from_static(b"a"), Bytes::from_static(b"b")]);
}
}