use std::collections::{BTreeMap, BTreeSet, HashSet};
use super::Id;
use crate::frozen::{self, Broken};
#[derive(Debug, Clone, Copy)]
pub struct Filter {
pub start: Id,
pub end: Id,
pub count: Option<usize>,
pub owner: Option<u32>,
pub min_idle: u64,
}
impl Default for Filter {
fn default() -> Filter {
Filter {
start: Id::MIN,
end: Id::MAX,
count: None,
owner: None,
min_idle: 0,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Nack {
time: u64,
count: u64,
owner: u32,
}
impl Nack {
const NOBODY: u32 = u32::MAX;
#[must_use]
#[inline]
pub fn time(&self) -> u64 {
self.time
}
#[must_use]
#[inline]
pub fn count(&self) -> u64 {
self.count
}
#[must_use]
#[inline]
pub fn owner(&self) -> Option<u32> {
(self.owner != Nack::NOBODY).then_some(self.owner)
}
#[must_use]
#[inline]
pub fn idle(&self, now: u64) -> u64 {
if self.owner == Nack::NOBODY {
return u64::MAX;
}
now.saturating_sub(self.time)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Retry {
Down,
Keep,
Max,
At(u64),
}
impl Retry {
#[must_use]
pub fn applied(self, had: u64) -> u64 {
match self {
Retry::Down => had.saturating_sub(1),
Retry::Keep => had,
Retry::Max => i64::MAX as u64,
Retry::At(n) => n,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Consumer {
name: Vec<u8>,
seen: u64,
active: Option<u64>,
pending: BTreeSet<Id>,
}
impl Consumer {
#[must_use]
#[inline]
pub fn name(&self) -> &[u8] {
&self.name
}
#[must_use]
#[inline]
pub fn seen(&self) -> u64 {
self.seen
}
#[must_use]
#[inline]
pub fn active(&self) -> Option<u64> {
self.active
}
#[must_use]
#[inline]
pub fn len(&self) -> usize {
self.pending.len()
}
#[must_use]
#[inline]
pub fn is_empty(&self) -> bool {
self.pending.is_empty()
}
pub fn pending(&self) -> impl Iterator<Item = Id> + '_ {
self.pending.iter().copied()
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct Group {
last: Id,
read: Option<u64>,
pending: BTreeMap<Id, Nack>,
consumers: Vec<Option<Consumer>>,
nacked: usize,
}
impl Group {
#[must_use]
pub fn new(last: Id, read: Option<u64>) -> Group {
Group {
last,
read,
pending: BTreeMap::new(),
consumers: Vec::new(),
nacked: 0,
}
}
#[must_use]
#[inline]
pub fn last_id(&self) -> Id {
self.last
}
#[must_use]
#[inline]
pub fn entries_read(&self) -> Option<u64> {
self.read
}
pub fn set_id(&mut self, last: Id, read: Option<u64>) {
self.last = last;
self.read = read;
}
#[must_use]
#[inline]
pub fn pending_len(&self) -> usize {
self.pending.len()
}
#[must_use]
#[inline]
pub fn nacked_len(&self) -> usize {
self.nacked
}
#[must_use]
pub fn pending_bounds(&self) -> Option<(Id, Id)> {
let low = *self.pending.keys().next()?;
let high = *self.pending.keys().next_back()?;
Some((low, high))
}
#[must_use]
#[inline]
pub fn nack(&self, id: Id) -> Option<&Nack> {
self.pending.get(&id)
}
#[must_use]
pub fn slot(&self, name: &[u8]) -> Option<u32> {
self.consumers
.iter()
.position(|c| c.as_ref().is_some_and(|c| c.name == name))
.map(|at| at as u32)
}
#[must_use]
#[inline]
pub fn consumer(&self, slot: u32) -> Option<&Consumer> {
self.consumers.get(slot as usize)?.as_ref()
}
#[must_use]
pub fn consumer_named(&self, name: &[u8]) -> Option<&Consumer> {
self.consumers
.iter()
.flatten()
.find(|c| c.name.as_slice() == name)
}
pub fn consumers(&self) -> impl Iterator<Item = &Consumer> + '_ {
self.consumers.iter().flatten()
}
pub fn consumer_or_create(&mut self, name: &[u8], now: u64) -> u32 {
if let Some(at) = self.slot(name) {
let c = self.consumers[at as usize]
.as_mut()
.expect("the slot the search just found");
c.seen = now;
return at;
}
self.consumers.push(Some(Consumer {
name: name.to_vec(),
seen: now,
active: None,
pending: BTreeSet::new(),
}));
(self.consumers.len() - 1) as u32
}
pub fn create_consumer(&mut self, name: &[u8], now: u64) -> bool {
if self.slot(name).is_some() {
return false;
}
self.consumer_or_create(name, now);
true
}
pub fn delete_consumer(&mut self, name: &[u8]) -> u64 {
let Some(at) = self.slot(name) else {
return 0;
};
let gone = self.consumers[at as usize]
.take()
.expect("the slot the search just found");
for id in &gone.pending {
self.pending.remove(id);
}
gone.pending.len() as u64
}
pub fn touch(&mut self, slot: u32, now: u64, read: bool) {
if let Some(Some(c)) = self.consumers.get_mut(slot as usize) {
c.seen = now;
if read {
c.active = Some(now);
}
}
}
pub fn deliver(&mut self, slot: u32, id: Id, now: u64) -> bool {
let Some(Some(c)) = self.consumers.get_mut(slot as usize) else {
return false;
};
c.pending.insert(id);
self.pending.insert(
id,
Nack {
time: now,
count: 1,
owner: slot,
},
);
if id > self.last {
self.last = id;
}
true
}
pub fn skip(&mut self, id: Id) {
if id > self.last {
self.last = id;
}
}
pub fn redeliver(&mut self, id: Id, now: u64) -> bool {
let Some(nack) = self.pending.get_mut(&id) else {
return false;
};
nack.time = now;
nack.count += 1;
true
}
pub fn set_read(&mut self, read: Option<u64>) {
self.read = read;
}
pub fn ack(&mut self, id: Id) -> bool {
let Some(nack) = self.pending.remove(&id) else {
return false;
};
match self.consumers.get_mut(nack.owner as usize) {
Some(Some(c)) => {
c.pending.remove(&id);
}
_ => self.nacked -= usize::from(nack.owner == Nack::NOBODY),
}
true
}
pub fn release(&mut self, id: Id, retry: Retry) -> bool {
let Some(nack) = self.pending.get_mut(&id) else {
return false;
};
let was = std::mem::replace(&mut nack.owner, Nack::NOBODY);
nack.count = retry.applied(nack.count);
nack.time = 0;
if was == Nack::NOBODY {
return true;
}
self.nacked += 1;
if let Some(Some(c)) = self.consumers.get_mut(was as usize) {
c.pending.remove(&id);
}
true
}
pub fn force_release(&mut self, id: Id, retry: Retry) {
if self.release(id, retry) {
return;
}
self.pending.insert(
id,
Nack {
time: 0,
count: retry.applied(0),
owner: Nack::NOBODY,
},
);
self.nacked += 1;
}
pub fn forget(&mut self, id: Id) -> bool {
self.ack(id)
}
pub fn claim(&mut self, id: Id, slot: u32, time: u64, count: Option<u64>, bump: bool) -> bool {
if !matches!(self.consumers.get(slot as usize), Some(Some(_))) {
return false;
}
let Some(nack) = self.pending.get_mut(&id) else {
return false;
};
let was = nack.owner;
nack.owner = slot;
nack.time = time;
if let Some(n) = count {
nack.count = n;
} else if bump {
nack.count += 1;
}
if was != slot {
if was == Nack::NOBODY {
self.nacked -= 1;
} else if let Some(Some(c)) = self.consumers.get_mut(was as usize) {
c.pending.remove(&id);
}
if let Some(Some(c)) = self.consumers.get_mut(slot as usize) {
c.pending.insert(id);
}
}
true
}
pub fn force(&mut self, id: Id, slot: u32, time: u64, count: u64) -> bool {
let Some(Some(c)) = self.consumers.get_mut(slot as usize) else {
return false;
};
c.pending.insert(id);
self.pending.insert(
id,
Nack {
time,
count,
owner: slot,
},
);
true
}
pub fn pending_range<F>(&self, want: Filter, now: u64, mut f: F) -> usize
where
F: FnMut(Id, &Nack, Option<&Consumer>) -> bool,
{
let mut seen = 0;
for (&id, nack) in self.pending.range(want.start..=want.end) {
if want.count.is_some_and(|n| seen >= n) {
break;
}
if want.owner.is_some_and(|c| Some(c) != nack.owner()) {
continue;
}
if nack.idle(now) < want.min_idle {
continue;
}
let who = match nack.owner() {
Some(slot) => match self.consumers.get(slot as usize) {
Some(Some(c)) => Some(c),
_ => continue,
},
None => None,
};
seen += 1;
if !f(id, nack, who) {
break;
}
}
seen
}
pub fn pending_counts(&self) -> impl Iterator<Item = (&[u8], usize)> + '_ {
self.consumers
.iter()
.flatten()
.filter(|c| !c.pending.is_empty())
.map(|c| (c.name.as_slice(), c.pending.len()))
}
#[must_use]
pub fn claimable(
&self,
start: Id,
min_idle: u64,
now: u64,
limit: usize,
out: &mut Vec<Id>,
) -> Option<Id> {
for (tried, (&id, nack)) in self.pending.range(start..).enumerate() {
if out.len() >= limit || tried >= limit * 10 {
return Some(id);
}
if nack.idle(now) >= min_idle {
out.push(id);
}
}
None
}
pub(super) fn freeze(&self, out: &mut Vec<u8>) {
frozen::put_uint(out, self.last.ms);
frozen::put_uint(out, self.last.seq);
put_opt(out, self.read);
frozen::put_uint(out, self.consumers.len() as u64);
for slot in &self.consumers {
match slot {
None => out.push(0),
Some(c) => {
out.push(1);
frozen::put_bytes(out, &c.name);
frozen::put_uint(out, c.seen);
put_opt(out, c.active);
}
}
}
frozen::put_uint(out, self.pending.len() as u64);
for (id, nack) in &self.pending {
frozen::put_uint(out, id.ms);
frozen::put_uint(out, id.seq);
frozen::put_uint(out, nack.time);
frozen::put_uint(out, nack.count);
frozen::put_uint(out, u64::from(nack.owner));
}
}
pub(super) fn thaw(cut: &mut frozen::Cut<'_>) -> Result<Group, Broken> {
let last = Id::new(cut.uint()?, cut.uint()?);
let read = take_opt(cut)?;
let n = usize::try_from(cut.uint()?).map_err(|_| Broken::Short)?;
if n > cut.rest().len() {
return Err(Broken::Short);
}
let mut consumers: Vec<Option<Consumer>> = Vec::with_capacity(n);
let mut names = HashSet::with_capacity(n);
for _ in 0..n {
match cut.byte()? {
0 => consumers.push(None),
1 => {
let name = cut.bytes()?;
if !names.insert(name) {
return Err(Broken::Body);
}
consumers.push(Some(Consumer {
name: name.to_vec(),
seen: cut.uint()?,
active: take_opt(cut)?,
pending: BTreeSet::new(),
}));
}
_ => return Err(Broken::Body),
}
}
let n = usize::try_from(cut.uint()?).map_err(|_| Broken::Short)?;
if n > cut.rest().len() {
return Err(Broken::Short);
}
let mut group = Group {
last,
read,
pending: BTreeMap::new(),
consumers,
nacked: 0,
};
let mut prev = None;
for _ in 0..n {
let id = Id::new(cut.uint()?, cut.uint()?);
if prev.is_some_and(|p| p >= id) {
return Err(Broken::Body);
}
prev = Some(id);
let nack = Nack {
time: cut.uint()?,
count: cut.uint()?,
owner: u32::try_from(cut.uint()?).map_err(|_| Broken::Body)?,
};
if nack.owner == Nack::NOBODY {
group.nacked += 1;
} else {
match group.consumers.get_mut(nack.owner as usize) {
Some(Some(c)) => {
c.pending.insert(id);
}
_ => return Err(Broken::Body),
}
}
group.pending.insert(id, nack);
}
Ok(group)
}
#[must_use]
pub fn memory_bytes(&self) -> usize {
let each = std::mem::size_of::<(Id, Nack)>();
let pending = self.pending.len() * (each + each / 16);
let consumers: usize = self
.consumers
.iter()
.map(|slot| {
std::mem::size_of::<Option<Consumer>>()
+ slot.as_ref().map_or(0, |c| {
let each = std::mem::size_of::<Id>();
c.name.capacity() + c.pending.len() * (each + each / 16)
})
})
.sum();
pending + consumers
}
}
fn put_opt(out: &mut Vec<u8>, v: Option<u64>) {
match v {
None => out.push(0),
Some(n) => {
out.push(1);
frozen::put_uint(out, n);
}
}
}
fn take_opt(cut: &mut frozen::Cut<'_>) -> Result<Option<u64>, Broken> {
match cut.byte()? {
0 => Ok(None),
1 => Ok(Some(cut.uint()?)),
_ => Err(Broken::Body),
}
}
#[cfg(test)]
mod tests {
use super::*;
fn group() -> Group {
Group::new(Id::MIN, Some(0))
}
#[test]
fn a_consumer_appears_by_turning_up() {
let mut g = group();
assert_eq!(g.slot(b"alice"), None);
let at = g.consumer_or_create(b"alice", 100);
assert_eq!(g.slot(b"alice"), Some(at));
assert_eq!(g.consumer_or_create(b"alice", 200), at);
assert_eq!(g.consumers().count(), 1);
assert_eq!(g.consumer(at).expect("alice").seen(), 200);
}
#[test]
fn creating_a_consumer_twice_says_so() {
let mut g = group();
assert!(g.create_consumer(b"alice", 1));
assert!(!g.create_consumer(b"alice", 2));
}
#[test]
fn delivering_moves_the_bookmark_and_fills_both_indexes() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 10);
assert!(g.deliver(a, Id::new(5, 0), 10));
assert!(g.deliver(a, Id::new(7, 0), 12));
assert_eq!(g.last_id(), Id::new(7, 0));
assert_eq!(g.pending_len(), 2);
assert_eq!(
g.consumer(a).expect("alice").pending().collect::<Vec<_>>(),
vec![Id::new(5, 0), Id::new(7, 0)]
);
let nack = g.nack(Id::new(5, 0)).expect("a nack");
assert_eq!((nack.count(), nack.time()), (1, 10));
}
#[test]
fn acking_takes_it_out_of_both_indexes() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
g.deliver(a, Id::new(5, 0), 1);
assert!(g.ack(Id::new(5, 0)));
assert_eq!(g.pending_len(), 0);
assert!(g.consumer(a).expect("alice").is_empty());
assert!(!g.ack(Id::new(5, 0)));
}
#[test]
fn acking_does_not_move_the_bookmark() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
g.deliver(a, Id::new(5, 0), 1);
g.ack(Id::new(5, 0));
assert_eq!(g.last_id(), Id::new(5, 0));
}
#[test]
fn a_claim_moves_it_between_consumers() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
let b = g.consumer_or_create(b"bob", 1);
g.deliver(a, Id::new(5, 0), 100);
assert!(g.claim(Id::new(5, 0), b, 500, None, true));
assert!(g.consumer(a).expect("alice").is_empty());
assert_eq!(
g.consumer(b).expect("bob").pending().collect::<Vec<_>>(),
vec![Id::new(5, 0)]
);
let nack = g.nack(Id::new(5, 0)).expect("a nack");
assert_eq!((nack.count(), nack.time()), (2, 500));
}
#[test]
fn a_claim_that_does_not_bump_is_justid() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
let b = g.consumer_or_create(b"bob", 1);
g.deliver(a, Id::new(5, 0), 100);
g.claim(Id::new(5, 0), b, 500, None, false);
assert_eq!(g.nack(Id::new(5, 0)).expect("a nack").count(), 1);
}
#[test]
fn a_retry_count_replaces_rather_than_adds() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
g.deliver(a, Id::new(5, 0), 100);
g.claim(Id::new(5, 0), a, 500, Some(9), true);
assert_eq!(g.nack(Id::new(5, 0)).expect("a nack").count(), 9);
}
#[test]
fn claiming_back_to_the_same_consumer_keeps_it_there() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
g.deliver(a, Id::new(5, 0), 100);
assert!(g.claim(Id::new(5, 0), a, 500, None, true));
assert_eq!(
g.consumer(a).expect("alice").pending().collect::<Vec<_>>(),
vec![Id::new(5, 0)]
);
}
#[test]
fn nothing_pending_cannot_be_claimed_without_force() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
assert!(!g.claim(Id::new(5, 0), a, 500, None, true));
assert!(g.force(Id::new(5, 0), a, 500, 1));
assert_eq!(g.pending_len(), 1);
}
#[test]
fn deleting_a_consumer_gives_up_its_work() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
let b = g.consumer_or_create(b"bob", 1);
g.deliver(a, Id::new(5, 0), 1);
g.deliver(a, Id::new(6, 0), 1);
g.deliver(b, Id::new(7, 0), 1);
assert_eq!(g.delete_consumer(b"alice"), 2);
assert_eq!(g.pending_len(), 1);
assert!(g.nack(Id::new(5, 0)).is_none());
assert!(g.nack(Id::new(7, 0)).is_some());
assert_eq!(g.delete_consumer(b"alice"), 0);
assert_eq!(g.last_id(), Id::new(7, 0));
}
#[test]
fn a_deleted_slot_is_not_reused() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
g.delete_consumer(b"alice");
let b = g.consumer_or_create(b"bob", 1);
assert_ne!(a, b);
assert_eq!(g.consumers().count(), 1);
}
#[test]
fn idle_is_measured_from_the_last_hand_out() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
g.deliver(a, Id::new(5, 0), 1_000);
assert_eq!(g.nack(Id::new(5, 0)).expect("a nack").idle(4_000), 3_000);
assert_eq!(g.nack(Id::new(5, 0)).expect("a nack").idle(500), 0);
}
#[test]
fn the_summary_is_the_two_ends_and_the_counts() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
let b = g.consumer_or_create(b"bob", 1);
for ms in [3u64, 5, 9] {
g.deliver(a, Id::new(ms, 0), 1);
}
g.deliver(b, Id::new(11, 0), 1);
assert_eq!(g.pending_bounds(), Some((Id::new(3, 0), Id::new(11, 0))));
let counts: Vec<_> = g
.pending_counts()
.map(|(n, c)| (String::from_utf8_lossy(n).into_owned(), c))
.collect();
assert_eq!(counts, vec![("alice".into(), 3), ("bob".into(), 1)]);
}
#[test]
fn a_consumer_with_nothing_is_left_out_of_the_summary() {
let mut g = group();
g.consumer_or_create(b"alice", 1);
assert_eq!(g.pending_counts().count(), 0);
assert_eq!(g.pending_bounds(), None);
}
#[test]
fn the_pending_range_takes_both_ends_and_the_filters() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
let b = g.consumer_or_create(b"bob", 1);
g.deliver(a, Id::new(3, 0), 100);
g.deliver(b, Id::new(5, 0), 100);
g.deliver(a, Id::new(9, 0), 900);
let seen = |g: &Group, start, end, owner, idle| {
let mut out = Vec::new();
let want = Filter {
start,
end,
owner,
min_idle: idle,
..Filter::default()
};
g.pending_range(want, 1_000, |id, _, c| {
out.push((
id,
String::from_utf8_lossy(c.expect("an owner").name()).into_owned(),
));
true
});
out
};
assert_eq!(seen(&g, Id::MIN, Id::MAX, None, 0).len(), 3);
assert_eq!(seen(&g, Id::new(4, 0), Id::new(9, 0), None, 0).len(), 2);
assert_eq!(
seen(&g, Id::MIN, Id::MAX, Some(a), 0),
vec![
(Id::new(3, 0), "alice".into()),
(Id::new(9, 0), "alice".into())
]
);
assert_eq!(seen(&g, Id::MIN, Id::MAX, None, 500).len(), 2);
}
#[test]
fn a_count_stops_the_pending_range() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
for ms in 1..=10u64 {
g.deliver(a, Id::new(ms, 0), 1);
}
let mut out = Vec::new();
let want = Filter {
count: Some(4),
..Filter::default()
};
let seen = g.pending_range(want, 1, |id, _, _| {
out.push(id);
true
});
assert_eq!((seen, out.len()), (4, 4));
}
#[test]
fn the_callback_can_stop_the_pending_range() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
for ms in 1..=10u64 {
g.deliver(a, Id::new(ms, 0), 1);
}
let mut out = Vec::new();
g.pending_range(Filter::default(), 1, |id, _, _| {
out.push(id);
out.len() < 3
});
assert_eq!(out.len(), 3);
}
#[test]
fn claimable_takes_the_idle_ones_and_says_where_to_carry_on() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
for ms in 1..=10u64 {
g.deliver(a, Id::new(ms, 0), if ms <= 5 { 100 } else { 900 });
}
let mut out = Vec::new();
let cursor = g.claimable(Id::MIN, 500, 1_000, 100, &mut out);
assert_eq!(cursor, None, "the scan reached the end");
assert_eq!(out, (1..=5).map(|ms| Id::new(ms, 0)).collect::<Vec<_>>());
out.clear();
let cursor = g.claimable(Id::MIN, 500, 1_000, 3, &mut out);
assert_eq!(out.len(), 3);
assert_eq!(cursor, Some(Id::new(4, 0)));
}
#[test]
fn a_delivery_leaves_the_read_counter_to_the_stream() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
g.deliver(a, Id::new(1, 0), 1);
assert_eq!(g.entries_read(), Some(0), "the group did not count it");
g.set_read(Some(1));
assert_eq!(g.entries_read(), Some(1));
g.set_read(None);
assert_eq!(g.entries_read(), None, "and it can be given up on");
}
#[test]
fn setting_the_id_leaves_the_pending_list_alone() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
g.deliver(a, Id::new(5, 0), 1);
g.set_id(Id::MIN, Some(0));
assert_eq!(g.last_id(), Id::MIN);
assert_eq!(g.pending_len(), 1, "somebody is still holding it");
}
#[test]
fn releasing_takes_the_entry_out_of_the_consumers_hands() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
g.deliver(a, Id::new(5, 0), 100);
assert!(g.release(Id::new(5, 0), Retry::Keep));
assert_eq!(g.pending_len(), 1, "it is still the group's problem");
assert_eq!(
g.consumer(a).expect("alice").pending().count(),
0,
"and no longer alice's"
);
let nack = g.nack(Id::new(5, 0)).expect("a nack");
assert_eq!(nack.owner(), None);
assert_eq!(nack.count(), 1, "Keep left the count where it was");
assert_eq!(nack.idle(100), u64::MAX);
let mut out = Vec::new();
assert_eq!(g.claimable(Id::MIN, u64::MAX, 100, 10, &mut out), None);
assert_eq!(out, vec![Id::new(5, 0)]);
assert!(g.release(Id::new(5, 0), Retry::Keep));
assert_eq!(g.nacked_len(), 1);
assert!(!g.release(Id::new(9, 0), Retry::Keep));
}
#[test]
fn the_three_words_differ_only_in_the_delivery_count() {
let count = |retry| {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
g.deliver(a, Id::new(5, 0), 1);
g.claim(Id::new(5, 0), a, 2, None, true);
g.release(Id::new(5, 0), retry);
g.nack(Id::new(5, 0)).expect("a nack").count()
};
assert_eq!(count(Retry::Down), 1, "one off, not back to nothing");
assert_eq!(count(Retry::Keep), 2);
assert_eq!(count(Retry::Max), i64::MAX as u64);
assert_eq!(count(Retry::At(7)), 7);
}
#[test]
fn forcing_a_release_is_the_same_call_twice() {
let mut g = group();
g.force_release(Id::new(5, 0), Retry::Keep);
assert_eq!(g.pending_len(), 1);
assert_eq!(
g.nack(Id::new(5, 0)).expect("a nack").count(),
0,
"there was no earlier count to keep"
);
g.force_release(Id::new(5, 0), Retry::At(4));
assert_eq!(g.pending_len(), 1);
assert_eq!(g.nacked_len(), 1);
assert_eq!(g.nack(Id::new(5, 0)).expect("a nack").count(), 4);
}
#[test]
fn the_nacked_count_matches_a_full_scan() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
let b = g.consumer_or_create(b"bob", 1);
for ms in 1..=6u64 {
g.deliver(a, Id::new(ms, 0), 100);
}
let scan = |g: &Group| {
(1..=9u64)
.filter(|&ms| g.nack(Id::new(ms, 0)).is_some_and(|n| n.owner().is_none()))
.count()
};
let agrees = |g: &Group| assert_eq!(g.nacked_len(), scan(g), "the field drifted");
agrees(&g);
g.release(Id::new(1, 0), Retry::Keep);
g.release(Id::new(2, 0), Retry::Keep);
agrees(&g);
g.claim(Id::new(1, 0), b, 200, None, true);
agrees(&g);
assert!(g.ack(Id::new(2, 0)));
assert!(g.ack(Id::new(3, 0)));
agrees(&g);
g.force_release(Id::new(9, 0), Retry::Down);
agrees(&g);
assert_eq!(g.nacked_len(), 1);
}
#[test]
fn a_released_entry_is_not_anybodys_pending_work() {
let mut g = group();
let a = g.consumer_or_create(b"alice", 1);
g.deliver(a, Id::new(3, 0), 100);
g.deliver(a, Id::new(5, 0), 100);
g.release(Id::new(3, 0), Retry::Keep);
let seen = |g: &Group, owner| {
let mut out = Vec::new();
let want = Filter {
owner,
..Filter::default()
};
g.pending_range(want, 1_000, |id, _, c| {
out.push((id, c.map(|c| c.name().to_vec())));
true
});
out
};
assert_eq!(
seen(&g, None),
vec![
(Id::new(3, 0), None),
(Id::new(5, 0), Some(b"alice".to_vec()))
]
);
assert_eq!(
seen(&g, Some(a)),
vec![(Id::new(5, 0), Some(b"alice".to_vec()))]
);
assert_eq!(
g.pending_counts()
.map(|(name, n)| (name.to_vec(), n))
.collect::<Vec<_>>(),
vec![(b"alice".to_vec(), 1)]
);
}
#[test]
fn a_frozen_group_with_a_pending_entry_nobody_could_hold_is_refused() {
let mut g = group();
let slot = g.consumer_or_create(b"alice", 1_000);
g.deliver(slot, Id::new(1, 0), 1_000);
let mut bytes = Vec::new();
g.freeze(&mut bytes);
assert_eq!(Group::thaw(&mut frozen::Cut::new(&bytes)), Ok(g));
let mut bad = bytes.clone();
*bad.last_mut().expect("a body") = 7;
assert_eq!(Group::thaw(&mut frozen::Cut::new(&bad)), Err(Broken::Body));
for cut in 0..bytes.len() {
assert!(
Group::thaw(&mut frozen::Cut::new(&bytes[..cut])).is_err(),
"cut at {cut}"
);
}
}
#[test]
fn a_frozen_group_that_names_one_consumer_twice_is_refused() {
let mut g = group();
g.consumer_or_create(b"alice", 1_000);
g.consumer_or_create(b"carol", 1_000);
let mut bytes = Vec::new();
g.freeze(&mut bytes);
let at = bytes
.windows(5)
.position(|w| w == b"carol")
.expect("the second name");
bytes[at..at + 5].copy_from_slice(b"alice");
assert_eq!(
Group::thaw(&mut frozen::Cut::new(&bytes)),
Err(Broken::Body)
);
}
}