use std::cmp::Ordering;
use std::collections::VecDeque;
use yo_common::num::{DIGITS_MAX, i64_digits, push_u64, u64_digits};
use crate::frozen::{self, Broken};
use crate::listpack::{self, Entry, Listpack};
pub mod groups;
pub use groups::{Consumer, Filter, Group, Nack, Retry};
pub const NODE_BYTES: usize = 4096;
pub const NODE_ENTRIES: usize = 100;
const LIVE: i64 = 0;
const DELETED: i64 = 1;
const SAME_FIELDS: i64 = 2;
const MASTER_FIELDS: usize = 3;
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Default)]
pub struct Id {
pub ms: u64,
pub seq: u64,
}
impl Id {
pub const MIN: Id = Id { ms: 0, seq: 0 };
pub const MAX: Id = Id {
ms: u64::MAX,
seq: u64::MAX,
};
#[must_use]
#[inline]
pub const fn new(ms: u64, seq: u64) -> Id {
Id { ms, seq }
}
#[must_use]
pub const fn next(self) -> Option<Id> {
if self.seq != u64::MAX {
Some(Id {
ms: self.ms,
seq: self.seq + 1,
})
} else if self.ms != u64::MAX {
Some(Id {
ms: self.ms + 1,
seq: 0,
})
} else {
None
}
}
#[must_use]
pub const fn prev(self) -> Option<Id> {
if self.seq != 0 {
Some(Id {
ms: self.ms,
seq: self.seq - 1,
})
} else if self.ms != 0 {
Some(Id {
ms: self.ms - 1,
seq: u64::MAX,
})
} else {
None
}
}
#[must_use]
pub fn to_bytes(self) -> [u8; 16] {
let mut out = [0u8; 16];
out[..8].copy_from_slice(&self.ms.to_be_bytes());
out[8..].copy_from_slice(&self.seq.to_be_bytes());
out
}
#[must_use]
pub fn from_bytes(bytes: [u8; 16]) -> Id {
let mut ms = [0u8; 8];
let mut seq = [0u8; 8];
ms.copy_from_slice(&bytes[..8]);
seq.copy_from_slice(&bytes[8..]);
Id {
ms: u64::from_be_bytes(ms),
seq: u64::from_be_bytes(seq),
}
}
pub fn write_to(self, out: &mut Vec<u8>) {
push_u64(out, self.ms);
out.push(b'-');
push_u64(out, self.seq);
}
#[must_use]
pub fn to_vec(self) -> Vec<u8> {
let mut out = Vec::with_capacity(41);
self.write_to(&mut out);
out
}
#[must_use]
pub fn parse(s: &[u8], default: u64) -> Option<Id> {
let (ms, seq) = match s.iter().position(|c| *c == b'-') {
Some(at) => (&s[..at], Some(&s[at + 1..])),
None => (s, None),
};
Some(Id {
ms: digits(ms)?,
seq: match seq {
Some(seq) => digits(seq)?,
None => default,
},
})
}
}
fn digits(s: &[u8]) -> Option<u64> {
if s.is_empty() || s.len() > 20 {
return None;
}
let mut n = 0u64;
for c in s {
let d = c.wrapping_sub(b'0');
if d > 9 {
return None;
}
n = n.checked_mul(10)?.checked_add(u64::from(d))?;
}
Some(n)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Refused {
NotGreater,
Zero,
Full,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct Limits {
pub max_node_bytes: usize,
pub max_node_entries: usize,
}
impl Default for Limits {
fn default() -> Limits {
Limits {
max_node_bytes: NODE_BYTES,
max_node_entries: NODE_ENTRIES,
}
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub enum Refs {
#[default]
Keep,
Drop,
Acked,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Fate {
Missing,
Gone,
Held,
}
impl Fate {
#[must_use]
#[inline]
pub fn code(self) -> i64 {
match self {
Fate::Missing => -1,
Fate::Gone => 1,
Fate::Held => 2,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct Node {
master: Id,
lp: Listpack,
}
#[derive(Debug, Clone, Copy)]
struct Edges {
added: u64,
length: u64,
first: Option<Id>,
last: Id,
max_deleted: Id,
}
impl Edges {
fn estimate(&self, id: Id) -> Option<u64> {
if self.added == 0 {
return Some(0);
}
if id >= self.last {
return Some(self.added);
}
let first = self.first?;
if self.max_deleted != Id::MIN && self.max_deleted >= first {
return None;
}
let behind = self.added - self.length;
match id.cmp(&first) {
Ordering::Less => Some(behind),
Ordering::Equal => Some(behind + 1),
Ordering::Greater => None,
}
}
fn holed_from(&self, id: Id) -> bool {
if self.length == 0 || self.max_deleted == Id::MIN {
return false;
}
if self.first.is_some_and(|first| first > self.max_deleted) {
return false;
}
id <= self.max_deleted
}
fn lag(&self, group: &Group) -> Option<u64> {
if let Some(read) = self.estimate(group.last_id()) {
return Some(self.added.saturating_sub(read));
}
let read = group.entries_read()?;
if self.holed_from(group.last_id()) {
return None;
}
Some(self.added.saturating_sub(read))
}
fn on_deliver(&self, group: &Group, id: Id) -> Option<u64> {
match group.entries_read() {
Some(read) if !self.holed_from(id) => Some(read + 1),
_ => self.estimate(id),
}
}
}
const FORM_NODES: u8 = 1;
#[derive(Debug, Clone, Copy)]
pub(crate) struct Cursor {
epoch: u64,
next: Id,
master: Id,
byte: usize,
}
#[derive(Debug, Clone, Default)]
pub struct Stream {
nodes: VecDeque<Node>,
length: u64,
last: Id,
max_deleted: Id,
added: u64,
groups: Vec<(Vec<u8>, Group)>,
epoch: u64,
}
impl PartialEq for Stream {
fn eq(&self, other: &Stream) -> bool {
self.nodes == other.nodes
&& self.length == other.length
&& self.last == other.last
&& self.max_deleted == other.max_deleted
&& self.added == other.added
&& self.groups == other.groups
}
}
impl Eq for Stream {}
impl Stream {
#[must_use]
pub fn new() -> Stream {
Stream::default()
}
pub(crate) fn raw_nodes(&self) -> impl Iterator<Item = (Id, &[u8])> + '_ {
self.nodes.iter().map(|n| (n.master, n.lp.as_bytes()))
}
pub(crate) fn push_raw_node(&mut self, master: Id, lp: Listpack) -> bool {
if self.nodes.back().is_some_and(|n| n.master >= master) {
return false;
}
self.nodes.push_back(Node { master, lp });
true
}
pub(crate) fn set_counters(
&mut self,
length: u64,
last: Id,
max_deleted: Id,
added: u64,
) -> bool {
if length > added || max_deleted > last || (self.nodes.is_empty() && length != 0) {
return false;
}
self.length = length;
self.last = last;
self.max_deleted = max_deleted;
self.added = added;
true
}
pub(crate) fn push_group(&mut self, name: &[u8], group: Group) -> bool {
if self.groups.iter().any(|(had, _)| had == name) {
return false;
}
self.groups.push((name.to_vec(), group));
true
}
pub fn freeze(&self, out: &mut Vec<u8>) {
out.push(FORM_NODES);
frozen::put_uint(out, self.length);
frozen::put_uint(out, self.last.ms);
frozen::put_uint(out, self.last.seq);
frozen::put_uint(out, self.max_deleted.ms);
frozen::put_uint(out, self.max_deleted.seq);
frozen::put_uint(out, self.added);
frozen::put_uint(out, self.nodes.len() as u64);
for node in &self.nodes {
frozen::put_uint(out, node.master.ms);
frozen::put_uint(out, node.master.seq);
frozen::put_bytes(out, node.lp.as_bytes());
}
frozen::put_uint(out, self.groups.len() as u64);
for (name, group) in &self.groups {
frozen::put_bytes(out, name);
group.freeze(out);
}
}
pub fn thaw(bytes: &[u8]) -> Result<Stream, Broken> {
let mut cut = frozen::Cut::new(bytes);
if cut.byte()? != FORM_NODES {
return Err(Broken::Form);
}
let length = cut.uint()?;
let last = Id::new(cut.uint()?, cut.uint()?);
let max_deleted = Id::new(cut.uint()?, cut.uint()?);
let added = cut.uint()?;
if length > added || max_deleted > last {
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 nodes = VecDeque::with_capacity(n);
let mut prev: Option<Id> = None;
for _ in 0..n {
let master = Id::new(cut.uint()?, cut.uint()?);
if prev.is_some_and(|p| p >= master) {
return Err(Broken::Body);
}
prev = Some(master);
let lp = Listpack::from_bytes(cut.bytes()?).map_err(|_| Broken::Body)?;
nodes.push_back(Node { master, lp });
}
if nodes.is_empty() && length != 0 {
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 groups: Vec<(Vec<u8>, Group)> = Vec::with_capacity(n);
for _ in 0..n {
let name = cut.bytes()?;
if groups.iter().any(|(had, _)| had == name) {
return Err(Broken::Body);
}
groups.push((name.to_vec(), Group::thaw(&mut cut)?));
}
Ok(Stream {
nodes,
length,
last,
max_deleted,
added,
groups,
epoch: 0,
})
}
#[must_use]
#[inline]
pub fn len(&self) -> u64 {
self.length
}
#[must_use]
#[inline]
pub fn is_empty(&self) -> bool {
self.length == 0
}
#[must_use]
#[inline]
pub fn last_id(&self) -> Id {
self.last
}
#[must_use]
#[inline]
pub fn max_deleted_id(&self) -> Id {
self.max_deleted
}
#[must_use]
#[inline]
pub fn added(&self) -> u64 {
self.added
}
#[must_use]
pub fn first_id(&self) -> Option<Id> {
let mut found = None;
self.walk(Id::MIN, Id::MAX, Some(1), &mut |id, _| {
found = Some(id);
false
});
found
}
#[must_use]
pub fn top_id(&self) -> Option<Id> {
let mut found = None;
self.rev_range(Id::MIN, Id::MAX, Some(1), |id, _| {
found = Some(id);
false
});
found
}
pub fn set_id(
&mut self,
last: Id,
added: Option<u64>,
max_deleted: Option<Id>,
) -> Result<(), Refused> {
if self.top_id().is_some_and(|top| last < top) {
return Err(Refused::NotGreater);
}
self.last = last;
if let Some(added) = added {
self.added = added;
}
if let Some(id) = max_deleted {
self.max_deleted = id;
}
Ok(())
}
#[must_use]
pub fn auto_id(&self, now: u64) -> Option<Id> {
if now > self.last.ms {
Some(Id { ms: now, seq: 0 })
} else {
self.last.next()
}
}
#[must_use]
pub fn auto_seq(&self, ms: u64) -> Option<Id> {
if ms > self.last.ms {
Some(Id { ms, seq: 0 })
} else if ms == self.last.ms {
self.last.next().filter(|id| id.ms == ms)
} else {
None
}
}
pub fn append(
&mut self,
id: Id,
fields: &[(&[u8], &[u8])],
limits: Limits,
) -> Result<(), Refused> {
if id == Id::MIN {
return Err(Refused::Zero);
}
if id <= self.last {
return Err(Refused::NotGreater);
}
let size: usize = fields
.iter()
.map(|(f, v)| f.len() + v.len() + 11)
.sum::<usize>()
+ 32;
let fits = match self.nodes.back() {
Some(node) => {
let (count, deleted) = counts(&node.lp);
node.lp.byte_len() + size < limits.max_node_bytes
&& (count + deleted) < limits.max_node_entries as u64
}
None => false,
};
if !fits {
self.nodes.push_back(Node {
master: id,
lp: master_of(fields),
});
}
let node = self.nodes.back_mut().expect("a node was just made sure of");
let same = same_fields(&node.lp, fields);
write_entry(&mut node.lp, node.master, id, fields, same);
if bump(&mut node.lp, 1, 0) {
self.epoch = self.epoch.wrapping_add(1);
}
self.length += 1;
self.added += 1;
self.last = id;
Ok(())
}
pub fn delete(&mut self, id: Id) -> bool {
if !self.remove(id) {
return false;
}
self.max_deleted = self.max_deleted.max(id);
true
}
pub fn delete_ref(&mut self, id: Id, refs: Refs) -> Fate {
if refs == Refs::Acked && self.still_wanted(id) {
return Fate::Held;
}
if !self.delete(id) {
return Fate::Missing;
}
if refs == Refs::Drop {
self.drop_refs(id);
}
Fate::Gone
}
pub fn ack_delete(&mut self, group: &[u8], id: Id, refs: Refs) -> (Fate, bool) {
let Some(g) = self.group_mut(group) else {
return (Fate::Missing, false);
};
if !g.ack(id) {
return (Fate::Missing, false);
}
if refs == Refs::Acked && self.still_wanted(id) {
return (Fate::Held, false);
}
let gone = self.delete(id);
if refs == Refs::Drop {
self.drop_refs(id);
}
(Fate::Gone, gone)
}
pub fn nack(&mut self, group: &[u8], id: Id, retry: Retry, force: bool) -> Option<bool> {
let here = self.contains(id);
let g = self.group_mut(group)?;
if g.release(id, retry) {
return Some(true);
}
if force && here {
g.force_release(id, retry);
return Some(true);
}
Some(false)
}
fn still_wanted(&self, id: Id) -> bool {
self.groups
.iter()
.any(|(_, g)| g.nack(id).is_some() || id > g.last_id())
}
fn drop_refs(&mut self, id: Id) {
for (_, g) in &mut self.groups {
g.forget(id);
}
}
fn remove(&mut self, id: Id) -> bool {
let Some(at) = self.node_of(id) else {
return false;
};
let node = &self.nodes[at];
let Some((offset, flags)) = find(&node.lp, node.master, id) else {
return false;
};
if flags & DELETED != 0 {
return false;
}
let (count, _) = counts(&node.lp);
if count == 1 {
self.nodes.remove(at);
} else {
let node = &mut self.nodes[at];
set_int(&mut node.lp, offset, flags | DELETED);
bump(&mut node.lp, -1, 1);
}
self.epoch = self.epoch.wrapping_add(1);
self.length -= 1;
true
}
fn edges(&self) -> Edges {
Edges {
added: self.added,
length: self.length,
first: self.first_id(),
last: self.last,
max_deleted: self.max_deleted,
}
}
#[must_use]
pub fn lag(&self, group: &Group) -> Option<u64> {
self.edges().lag(group)
}
pub fn trim_maxlen(&mut self, len: u64, exact: bool, limit: Option<u64>) -> u64 {
let mut gone = 0;
while self.length > len && !limit.is_some_and(|cap| gone >= cap) {
let Some(node) = self.nodes.front() else {
break;
};
let (count, _) = counts(&node.lp);
if self.length - count >= len {
self.length -= count;
gone += count;
self.nodes.pop_front();
continue;
}
if !exact {
break;
}
let Some(id) = self.first_id() else { break };
self.remove(id);
gone += 1;
}
gone
}
pub fn trim_minid(&mut self, id: Id, exact: bool, limit: Option<u64>) -> u64 {
let mut gone = 0;
while let Some(node) = self.nodes.front() {
if limit.is_some_and(|cap| gone >= cap) {
break;
}
let (count, _) = counts(&node.lp);
if last_of(node) < id {
self.length -= count;
gone += count;
self.nodes.pop_front();
continue;
}
if !exact {
break;
}
let Some(first) = self.first_id() else { break };
if first >= id {
break;
}
self.remove(first);
gone += 1;
}
gone
}
pub fn range<F>(&self, start: Id, end: Id, count: Option<usize>, mut f: F) -> usize
where
F: FnMut(Id, Fields<'_>) -> bool,
{
self.walk(start, end, count, &mut f)
}
pub fn rev_range<'s, F>(&'s self, start: Id, end: Id, count: Option<usize>, mut f: F) -> usize
where
F: FnMut(Id, Fields<'_>) -> bool,
{
let mut seen = 0;
let mut buf: Vec<(Id, Fields<'s>)> = Vec::new();
let last = self.node_from(end);
for node in self.nodes.iter().take(last + 1).rev() {
if node.master > end {
continue;
}
if last_of(node) < start {
break;
}
buf.clear();
each(&node.lp, node.master, None, &mut |id, _, fields| {
if id >= start && id <= end {
buf.push((id, fields));
}
id <= end
});
for (id, fields) in buf.drain(..).rev() {
if count.is_some_and(|want| seen >= want) {
return seen;
}
seen += 1;
if !f(id, fields) {
return seen;
}
}
}
seen
}
#[must_use]
pub fn contains(&self, id: Id) -> bool {
let Some(at) = self.node_of(id) else {
return false;
};
let node = &self.nodes[at];
find(&node.lp, node.master, id).is_some_and(|(_, flags)| flags & DELETED == 0)
}
pub fn create_group(&mut self, name: &[u8], last: Id, read: Option<u64>) -> bool {
if self.group(name).is_some() {
return false;
}
self.groups.push((name.to_vec(), Group::new(last, read)));
true
}
pub fn destroy_group(&mut self, name: &[u8]) -> bool {
let Some(at) = self.groups.iter().position(|(n, _)| n == name) else {
return false;
};
self.groups.remove(at);
true
}
#[must_use]
pub fn group(&self, name: &[u8]) -> Option<&Group> {
self.groups
.iter()
.find(|(n, _)| n.as_slice() == name)
.map(|(_, g)| g)
}
pub fn group_mut(&mut self, name: &[u8]) -> Option<&mut Group> {
self.groups
.iter_mut()
.find(|(n, _)| n.as_slice() == name)
.map(|(_, g)| g)
}
pub fn groups(&self) -> impl Iterator<Item = (&[u8], &Group)> + '_ {
self.groups.iter().map(|(n, g)| (n.as_slice(), g))
}
pub fn read_group<F>(
&mut self,
group: &[u8],
consumer: &[u8],
count: Option<usize>,
noack: bool,
now: u64,
mut f: F,
) -> Option<usize>
where
F: FnMut(Id, Fields<'_>) -> bool,
{
let edges = self.edges();
let Stream {
nodes,
groups,
epoch,
..
} = self;
let epoch = *epoch;
let (_, g) = groups.iter_mut().find(|(n, _)| n.as_slice() == group)?;
let slot = g.consumer_or_create(consumer, now);
let Some(from) = g.last_id().next() else {
g.touch(slot, now, false);
return Some(0);
};
let resume = g.resume(epoch, from);
let (seen, mark) = walk_nodes(nodes, from, Id::MAX, count, resume, &mut |id, fields| {
let read = edges.on_deliver(g, id);
if noack {
g.skip(id);
} else {
g.deliver(slot, id, now);
}
g.set_read(read);
f(id, fields)
});
let next = g.last_id().next();
g.set_resume(match (mark, next) {
(Some((master, byte)), Some(next)) => Some(Cursor {
epoch,
next,
master,
byte,
}),
_ => None,
});
g.touch(slot, now, seen > 0);
Some(seen)
}
pub fn read_group_pending<F>(
&mut self,
group: &[u8],
consumer: &[u8],
after: Id,
count: Option<usize>,
now: u64,
mut f: F,
) -> Option<usize>
where
F: FnMut(Id, Option<Fields<'_>>) -> bool,
{
let g = self.group_mut(group)?;
let slot = g.consumer_or_create(consumer, now);
let ids: Vec<Id> = g
.consumer(slot)
.expect("the slot that was just made")
.pending()
.filter(|&id| id > after)
.take(count.unwrap_or(usize::MAX))
.collect();
for &id in &ids {
g.redeliver(id, now);
}
g.touch(slot, now, !ids.is_empty());
let mut seen = 0;
for &id in &ids {
seen += 1;
let mut go = true;
let mut found = false;
self.walk(id, id, Some(1), &mut |got, fields| {
found = true;
go = f(got, Some(fields));
false
});
if !found {
go = f(id, None);
}
if !go {
break;
}
}
Some(seen)
}
#[allow(clippy::too_many_arguments)]
pub fn claim(
&mut self,
group: &[u8],
consumer: &[u8],
ids: &[Id],
min_idle: u64,
time: u64,
retry: Option<u64>,
bump: bool,
force: bool,
now: u64,
gone: &mut Vec<Id>,
) -> Option<Vec<Id>> {
self.group_mut(group)?.consumer_or_create(consumer, now);
let mut took = Vec::new();
for &id in ids {
let here = self.contains(id);
let g = self.group_mut(group)?;
let slot = g.consumer_or_create(consumer, now);
match g.nack(id) {
Some(nack) => {
if !here {
g.forget(id);
gone.push(id);
continue;
}
if nack.idle(now) < min_idle {
continue;
}
if g.claim(id, slot, time, retry, bump) {
took.push(id);
}
}
None => {
if force && here && g.force(id, slot, time, retry.unwrap_or(1)) {
took.push(id);
}
}
}
}
if !took.is_empty() {
let g = self.group_mut(group).expect("the group found a moment ago");
let slot = g.consumer_or_create(consumer, now);
g.touch(slot, now, true);
}
Some(took)
}
#[allow(clippy::too_many_arguments)]
pub fn autoclaim(
&mut self,
group: &[u8],
consumer: &[u8],
start: Id,
min_idle: u64,
count: usize,
bump: bool,
now: u64,
gone: &mut Vec<Id>,
) -> Option<(Option<Id>, Vec<Id>)> {
let mut ids = Vec::new();
let cursor = self
.group(group)?
.claimable(start, min_idle, now, count, &mut ids);
let took = self.claim(
group, consumer, &ids, min_idle, now, None, bump, false, now, gone,
)?;
Some((cursor, took))
}
#[must_use]
pub fn memory_bytes(&self) -> usize {
let nodes: usize = self
.nodes
.iter()
.map(|node| node.lp.byte_len() + std::mem::size_of::<Node>())
.sum();
let groups: usize = self
.groups
.iter()
.map(|(name, g)| {
name.capacity() + std::mem::size_of::<(Vec<u8>, Group)>() + g.memory_bytes()
})
.sum();
nodes + groups
}
#[must_use]
pub fn nodes(&self) -> usize {
self.nodes.len()
}
fn walk<F>(&self, start: Id, end: Id, count: Option<usize>, f: &mut F) -> usize
where
F: FnMut(Id, Fields<'_>) -> bool,
{
walk_nodes(&self.nodes, start, end, count, None, f).0
}
fn node_from(&self, id: Id) -> usize {
node_from(&self.nodes, id)
}
fn node_of(&self, id: Id) -> Option<usize> {
let at = self.node_from(id);
let node = self.nodes.get(at)?;
(node.master <= id && id <= last_of(node)).then_some(at)
}
}
fn node_from(nodes: &VecDeque<Node>, id: Id) -> usize {
let after = nodes.partition_point(|node| node.master <= id);
after.saturating_sub(1)
}
fn walk_nodes<F>(
nodes: &VecDeque<Node>,
start: Id,
end: Id,
count: Option<usize>,
resume: Option<(Id, usize)>,
f: &mut F,
) -> (usize, Option<(Id, usize)>)
where
F: FnMut(Id, Fields<'_>) -> bool,
{
let mut seen = 0;
let mut stop = false;
let first = node_from(nodes, start);
let mut from =
resume.and_then(|(master, byte)| (nodes.get(first)?.master == master).then_some(byte));
let mut mark = None;
for node in nodes.iter().skip(first) {
if node.master > end {
break;
}
let at = each(&node.lp, node.master, from.take(), &mut |id, _, fields| {
if id > end {
stop = true;
return false;
}
if id < start {
return true;
}
if count.is_some_and(|want| seen >= want) {
stop = true;
return false;
}
seen += 1;
if !f(id, fields) {
stop = true;
return false;
}
true
});
mark = Some((node.master, at));
if stop {
break;
}
}
(seen, mark)
}
#[derive(Debug, Clone)]
pub struct Fields<'a> {
names: Option<listpack::Iter<'a>>,
body: listpack::Iter<'a>,
left: usize,
}
impl<'a> Iterator for Fields<'a> {
type Item = (Entry<'a>, Entry<'a>);
fn next(&mut self) -> Option<(Entry<'a>, Entry<'a>)> {
if self.left == 0 {
return None;
}
self.left -= 1;
let name = match &mut self.names {
Some(names) => names.next()?,
None => self.body.next()?,
};
Some((name, self.body.next()?))
}
fn size_hint(&self) -> (usize, Option<usize>) {
(self.left, Some(self.left))
}
}
impl ExactSizeIterator for Fields<'_> {}
impl Fields<'_> {
#[must_use]
#[inline]
pub fn is_empty(&self) -> bool {
self.left == 0
}
}
fn master_of(fields: &[(&[u8], &[u8])]) -> Listpack {
let mut lp = Listpack::new();
push_int(&mut lp, 0);
push_int(&mut lp, 0);
push_int(&mut lp, fields.len() as i64);
for (name, _) in fields {
lp.push(name);
}
push_int(&mut lp, 0);
lp
}
fn same_fields(lp: &Listpack, fields: &[(&[u8], &[u8])]) -> bool {
let mut it = lp.iter();
let (_, _, want) = match (it.next(), it.next(), it.next()) {
(Some(_), Some(_), Some(Entry::Int(n))) => ((), (), n),
_ => return false,
};
if want != fields.len() as i64 {
return false;
}
fields.iter().all(|(name, _)| match it.next() {
Some(Entry::Str(s)) => s == *name,
Some(Entry::Int(n)) => {
let mut buf = [0u8; DIGITS_MAX];
i64_digits(&mut buf, n) == *name
}
None => false,
})
}
fn counts(lp: &Listpack) -> (u64, u64) {
let mut it = lp.iter();
let count = int_or_zero(it.next());
let deleted = int_or_zero(it.next());
(count.max(0) as u64, deleted.max(0) as u64)
}
fn bump(lp: &mut Listpack, live: i64, dead: i64) -> bool {
let was = lp.byte_len();
let (count, deleted) = counts(lp);
let mut buf = [0u8; DIGITS_MAX];
let at = count as i64 + live;
lp.replace(0, u64_digits(&mut buf, at.max(0) as u64));
let at = deleted as i64 + dead;
lp.replace(1, u64_digits(&mut buf, at.max(0) as u64));
lp.byte_len() != was
}
fn int_or_zero(entry: Option<Entry<'_>>) -> i64 {
match entry {
Some(Entry::Int(n)) => n,
_ => 0,
}
}
fn push_int(lp: &mut Listpack, n: i64) {
let mut buf = [0u8; DIGITS_MAX];
lp.push(i64_digits(&mut buf, n));
}
fn set_int(lp: &mut Listpack, index: usize, n: i64) {
let mut buf = [0u8; DIGITS_MAX];
lp.replace(index, i64_digits(&mut buf, n));
}
fn write_entry(lp: &mut Listpack, master: Id, id: Id, fields: &[(&[u8], &[u8])], same: bool) {
let flags = if same { LIVE | SAME_FIELDS } else { LIVE };
push_int(lp, flags);
push_int(lp, id.ms.wrapping_sub(master.ms) as i64);
push_int(lp, id.seq.wrapping_sub(master.seq) as i64);
if same {
for (_, value) in fields {
lp.push(value);
}
push_int(lp, fields.len() as i64 + 3);
} else {
push_int(lp, fields.len() as i64);
for (name, value) in fields {
lp.push(name);
lp.push(value);
}
push_int(lp, fields.len() as i64 * 2 + 4);
}
}
fn last_of(node: &Node) -> Id {
let mut last = node.master;
each(&node.lp, node.master, None, &mut |id, _, _| {
last = id;
true
});
last
}
fn find(lp: &Listpack, master: Id, id: Id) -> Option<(usize, i64)> {
let mut at = None;
walk_node(lp, master, None, &mut |mark, flags, _| {
if mark.id == id {
if let Some(index) = mark.index {
at = Some((index, flags));
}
return false;
}
mark.id < id
});
at
}
fn each<'a, F>(lp: &'a Listpack, master: Id, from: Option<usize>, f: &mut F) -> usize
where
F: FnMut(Id, usize, Fields<'a>) -> bool,
{
walk_node(lp, master, from, &mut |mark, flags, fields| {
if flags & DELETED != 0 {
return true;
}
f(mark.id, mark.byte, fields)
})
}
fn walk_node<'a, F>(lp: &'a Listpack, master: Id, from: Option<usize>, f: &mut F) -> usize
where
F: FnMut(Mark, i64, Fields<'a>) -> bool,
{
let mut it = lp.iter();
let (Some(_), Some(_), Some(Entry::Int(masters))) = (it.next(), it.next(), it.next()) else {
return it.offset();
};
let masters = masters.max(0) as usize;
let names = it.clone();
let mut index = MASTER_FIELDS;
for _ in 0..=masters {
if it.next().is_none() {
return it.offset();
}
index += 1;
}
let counting = from.is_none();
if let Some(byte) = from.filter(|byte| *byte > it.offset()) {
it = lp.iter_at(byte);
}
loop {
let at = it.offset();
let element = index;
let (Some(Entry::Int(flags)), Some(Entry::Int(ms)), Some(Entry::Int(seq))) =
(it.next(), it.next(), it.next())
else {
return at;
};
index += 3;
let id = Id {
ms: master.ms.wrapping_add(ms as u64),
seq: master.seq.wrapping_add(seq as u64),
};
let same = flags & SAME_FIELDS != 0;
let (fields, skip) = if same {
let fields = Fields {
names: Some(names.clone()),
body: it.clone(),
left: masters,
};
(fields, masters + 1)
} else {
let Some(Entry::Int(n)) = it.next() else {
return at;
};
index += 1;
let fields = Fields {
names: None,
body: it.clone(),
left: n.max(0) as usize,
};
(fields, n.max(0) as usize * 2 + 1)
};
let mark = Mark {
id,
index: counting.then_some(element),
byte: at,
};
if !f(mark, flags, fields) {
return at;
}
for _ in 0..skip {
if it.next().is_none() {
return it.offset();
}
index += 1;
}
}
}
#[derive(Debug, Clone, Copy)]
struct Mark {
id: Id,
index: Option<usize>,
byte: usize,
}
#[cfg(test)]
mod tests {
use super::*;
use crate::many;
fn pairs<'a>(of: &'a [(&'a str, &'a str)]) -> Vec<(&'a [u8], &'a [u8])> {
of.iter()
.map(|(f, v)| (f.as_bytes(), v.as_bytes()))
.collect()
}
type Flat = (Id, Vec<(Vec<u8>, Vec<u8>)>);
fn dump(s: &Stream) -> Vec<Flat> {
let mut out = Vec::new();
s.range(Id::MIN, Id::MAX, None, |id, fields| {
out.push((id, fields.map(|(f, v)| (f.to_vec(), v.to_vec())).collect()));
true
});
out
}
fn add(s: &mut Stream, ms: u64, seq: u64, fields: &[(&str, &str)]) {
s.append(Id::new(ms, seq), &pairs(fields), Limits::default())
.expect("an append");
}
#[test]
fn an_entry_comes_back_as_it_went_in() {
let mut s = Stream::new();
add(&mut s, 5, 0, &[("sensor", "1"), ("reading", "23.4")]);
let got = dump(&s);
assert_eq!(got.len(), 1);
assert_eq!(got[0].0, Id::new(5, 0));
assert_eq!(
got[0].1,
vec![
(b"sensor".to_vec(), b"1".to_vec()),
(b"reading".to_vec(), b"23.4".to_vec())
]
);
}
#[test]
fn entries_come_back_in_order() {
let mut s = Stream::new();
for ms in 1..200u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
let got = dump(&s);
assert_eq!(got.len(), 199);
for (at, (id, _)) in got.iter().enumerate() {
assert_eq!(*id, Id::new(at as u64 + 1, 0));
}
assert_eq!(s.len(), 199);
assert_eq!(s.added(), 199);
assert_eq!(s.last_id(), Id::new(199, 0));
}
#[test]
fn a_long_stream_is_many_nodes() {
let mut s = Stream::new();
for ms in 1..=1000u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
assert_eq!(s.nodes(), 10, "a hundred entries a node");
assert_eq!(dump(&s).len(), 1000);
}
#[test]
fn the_field_names_are_not_stored_twice() {
let mut shared = Stream::new();
let mut apart = Stream::new();
for ms in 1..=100u64 {
add(
&mut shared,
ms,
0,
&[("temperature_celsius", "21"), ("relative_humidity", "44")],
);
let a = format!("temperature_celsius{ms}");
let b = format!("relative_humidity{ms}");
apart
.append(
Id::new(ms, 0),
&[(a.as_bytes(), b"21"), (b.as_bytes(), b"44")],
Limits::default(),
)
.expect("an append");
}
assert!(
shared.memory_bytes() * 3 < apart.memory_bytes(),
"{} against {}",
shared.memory_bytes(),
apart.memory_bytes()
);
}
#[test]
fn an_entry_costs_about_two_dozen_bytes() {
let mut s = Stream::new();
for ms in 1..=10_000u64 {
let reading = format!("{:.3}", ms as f64 / 7.0);
s.append(
Id::new(ms, 0),
&[(b"sensor", b"a4"), (b"reading", reading.as_bytes())],
Limits::default(),
)
.expect("an append");
}
let each = s.memory_bytes() as f64 / 10_000.0;
assert!(each < 32.0, "{each:.2} bytes an entry");
}
#[test]
fn an_entry_with_its_own_fields_still_reads_back() {
let mut s = Stream::new();
add(&mut s, 1, 0, &[("a", "1"), ("b", "2")]);
add(&mut s, 2, 0, &[("c", "3")]);
add(&mut s, 3, 0, &[("a", "4"), ("b", "5")]);
let got = dump(&s);
assert_eq!(got[1].1, vec![(b"c".to_vec(), b"3".to_vec())]);
assert_eq!(
got[2].1,
vec![
(b"a".to_vec(), b"4".to_vec()),
(b"b".to_vec(), b"5".to_vec())
]
);
}
#[test]
fn the_order_of_the_names_matters() {
let mut s = Stream::new();
add(&mut s, 1, 0, &[("a", "1"), ("b", "2")]);
add(&mut s, 2, 0, &[("b", "3"), ("a", "4")]);
let got = dump(&s);
assert_eq!(
got[1].1,
vec![
(b"b".to_vec(), b"3".to_vec()),
(b"a".to_vec(), b"4".to_vec())
]
);
}
#[test]
fn an_id_must_beat_the_last_one() {
let mut s = Stream::new();
add(&mut s, 5, 5, &[("n", "x")]);
let f = pairs(&[("n", "x")]);
for id in [Id::new(5, 5), Id::new(5, 4), Id::new(1, 0)] {
assert_eq!(
s.append(id, &f, Limits::default()),
Err(Refused::NotGreater),
"{id:?}"
);
}
assert_eq!(s.append(Id::new(5, 6), &f, Limits::default()), Ok(()));
}
#[test]
fn nothing_can_be_added_at_zero() {
let mut s = Stream::new();
assert_eq!(
s.append(Id::MIN, &pairs(&[("n", "x")]), Limits::default()),
Err(Refused::Zero)
);
}
#[test]
fn a_range_takes_both_ends() {
let mut s = Stream::new();
for ms in 1..=10u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
let mut seen = Vec::new();
s.range(Id::new(3, 0), Id::new(6, 0), None, |id, _| {
seen.push(id.ms);
true
});
assert_eq!(seen, vec![3, 4, 5, 6]);
}
#[test]
fn a_range_that_lands_between_entries() {
let mut s = Stream::new();
for ms in [10u64, 20, 30] {
add(&mut s, ms, 0, &[("n", "x")]);
}
let mut seen = Vec::new();
s.range(Id::new(11, 0), Id::new(29, 0), None, |id, _| {
seen.push(id.ms);
true
});
assert_eq!(seen, vec![20]);
let mut none = 0;
s.range(Id::new(31, 0), Id::MAX, None, |_, _| {
none += 1;
true
});
assert_eq!(none, 0);
}
#[test]
fn a_count_stops_the_walk() {
let mut s = Stream::new();
for ms in 1..=500u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
let mut seen = 0;
let answered = s.range(Id::MIN, Id::MAX, Some(7), |_, _| {
seen += 1;
true
});
assert_eq!((seen, answered), (7, 7));
}
#[test]
fn the_callback_can_stop_the_walk() {
let mut s = Stream::new();
for ms in 1..=500u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
let mut seen = 0;
s.range(Id::MIN, Id::MAX, None, |_, _| {
seen += 1;
seen < 3
});
assert_eq!(seen, 3);
}
#[test]
fn a_reverse_range_is_the_forward_one_backwards() {
let mut s = Stream::new();
for ms in 1..=350u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
let mut forward = Vec::new();
s.range(Id::new(50, 0), Id::new(300, 0), None, |id, _| {
forward.push(id);
true
});
let mut back = Vec::new();
s.rev_range(Id::new(50, 0), Id::new(300, 0), None, |id, _| {
back.push(id);
true
});
back.reverse();
assert_eq!(forward, back);
assert_eq!(forward.len(), 251);
}
#[test]
fn a_reverse_range_takes_a_count_from_the_new_end() {
let mut s = Stream::new();
for ms in 1..=350u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
let mut seen = Vec::new();
s.rev_range(Id::MIN, Id::MAX, Some(3), |id, _| {
seen.push(id.ms);
true
});
assert_eq!(seen, vec![350, 349, 348]);
}
#[test]
fn deleting_leaves_the_rest_readable() {
let mut s = Stream::new();
for ms in 1..=10u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
assert!(s.delete(Id::new(4, 0)));
assert!(!s.delete(Id::new(4, 0)), "twice is not twice");
assert!(!s.delete(Id::new(99, 0)));
assert_eq!(s.len(), 9);
assert_eq!(s.max_deleted_id(), Id::new(4, 0));
let seen: Vec<u64> = dump(&s).iter().map(|(id, _)| id.ms).collect();
assert_eq!(seen, vec![1, 2, 3, 5, 6, 7, 8, 9, 10]);
}
#[test]
fn deleting_the_first_entry_moves_the_first_id() {
let mut s = Stream::new();
for ms in 1..=5u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
assert_eq!(s.first_id(), Some(Id::new(1, 0)));
s.delete(Id::new(1, 0));
assert_eq!(s.first_id(), Some(Id::new(2, 0)));
}
#[test]
fn emptying_a_node_drops_it() {
let mut s = Stream::new();
for ms in 1..=250u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
assert_eq!(s.nodes(), 3);
for ms in 1..=100u64 {
assert!(s.delete(Id::new(ms, 0)), "{ms}");
}
assert_eq!(s.nodes(), 2);
assert_eq!(s.len(), 150);
assert_eq!(dump(&s).len(), 150);
}
#[test]
fn an_emptied_stream_still_remembers_its_last_id() {
let mut s = Stream::new();
add(&mut s, 7, 0, &[("n", "x")]);
s.delete(Id::new(7, 0));
assert!(s.is_empty());
assert_eq!(s.last_id(), Id::new(7, 0));
assert_eq!(s.added(), 1);
assert_eq!(s.first_id(), None);
assert_eq!(
s.append(Id::new(7, 0), &pairs(&[("n", "x")]), Limits::default()),
Err(Refused::NotGreater),
"a deleted id is still used up"
);
}
#[test]
fn trimming_to_a_length_takes_the_oldest() {
let mut s = Stream::new();
for ms in 1..=1000u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
assert_eq!(s.trim_maxlen(150, true, None), 850);
assert_eq!(s.len(), 150);
assert_eq!(s.first_id(), Some(Id::new(851, 0)));
assert_eq!(s.last_id(), Id::new(1000, 0));
}
#[test]
fn an_approximate_trim_stops_at_a_node() {
let mut s = Stream::new();
for ms in 1..=1000u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
assert_eq!(s.trim_maxlen(150, false, None), 800);
assert_eq!(s.len(), 200, "left at the node boundary above 150");
assert_eq!(s.nodes(), 2);
}
#[test]
fn trimming_to_a_length_that_is_already_met_does_nothing() {
let mut s = Stream::new();
for ms in 1..=10u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
assert_eq!(s.trim_maxlen(50, true, None), 0);
assert_eq!(s.len(), 10);
}
#[test]
fn trimming_to_zero_empties_it() {
let mut s = Stream::new();
for ms in 1..=250u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
assert_eq!(s.trim_maxlen(0, true, None), 250);
assert!(s.is_empty());
assert_eq!(s.nodes(), 0);
assert_eq!(s.last_id(), Id::new(250, 0));
}
#[test]
fn trimming_below_an_id_takes_everything_under_it() {
let mut s = Stream::new();
for ms in 1..=1000u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
assert_eq!(s.trim_minid(Id::new(400, 0), true, None), 399);
assert_eq!(s.first_id(), Some(Id::new(400, 0)));
assert_eq!(s.len(), 601);
}
#[test]
fn an_approximate_minid_trim_stops_at_a_node() {
let mut s = Stream::new();
for ms in 1..=1000u64 {
add(&mut s, ms, 0, &[("n", "x")]);
}
assert_eq!(s.trim_minid(Id::new(450, 0), false, None), 400);
assert_eq!(s.first_id(), Some(Id::new(401, 0)));
}
#[test]
fn several_entries_share_a_millisecond() {
let mut s = Stream::new();
for seq in 0..250u64 {
add(&mut s, 5, seq, &[("n", "x")]);
}
let got = dump(&s);
assert_eq!(got.len(), 250);
for (at, (id, _)) in got.iter().enumerate() {
assert_eq!(*id, Id::new(5, at as u64), "at {at}");
}
let mut seen = Vec::new();
s.range(Id::new(5, 100), Id::new(5, 102), None, |id, _| {
seen.push(id.seq);
true
});
assert_eq!(seen, vec![100, 101, 102]);
}
#[test]
fn a_sequence_that_crosses_a_node() {
let mut s = Stream::new();
for seq in 0..300u64 {
add(&mut s, 1, seq, &[("n", "x")]);
}
assert!(s.nodes() > 1);
let got = dump(&s);
assert_eq!(got.len(), 300);
assert_eq!(got[299].0, Id::new(1, 299));
assert!(s.delete(Id::new(1, 250)));
assert_eq!(dump(&s).len(), 299);
}
#[test]
fn an_entry_with_no_fields_is_still_an_entry() {
let mut s = Stream::new();
s.append(Id::new(1, 0), &[], Limits::default())
.expect("an append");
add(&mut s, 2, 0, &[("n", "x")]);
let got = dump(&s);
assert_eq!(got.len(), 2);
assert!(got[0].1.is_empty());
}
#[test]
fn a_value_that_looks_like_a_number_comes_back_as_it_went_in() {
let mut s = Stream::new();
add(&mut s, 1, 0, &[("n", "007"), ("m", "7")]);
let got = dump(&s);
assert_eq!(got[0].1[0].1, b"007".to_vec());
assert_eq!(got[0].1[1].1, b"7".to_vec());
}
#[test]
fn the_auto_id_follows_the_clock_and_never_goes_back() {
let mut s = Stream::new();
assert_eq!(s.auto_id(1000), Some(Id::new(1000, 0)));
add(&mut s, 1000, 0, &[("n", "x")]);
assert_eq!(s.auto_id(1000), Some(Id::new(1000, 1)), "same millisecond");
assert_eq!(s.auto_id(900), Some(Id::new(1000, 1)), "clock went back");
assert_eq!(s.auto_id(1001), Some(Id::new(1001, 0)));
}
#[test]
fn an_explicit_millisecond_takes_the_next_sequence() {
let mut s = Stream::new();
add(&mut s, 5, 0, &[("n", "x")]);
assert_eq!(s.auto_seq(5), Some(Id::new(5, 1)));
assert_eq!(s.auto_seq(6), Some(Id::new(6, 0)));
assert_eq!(s.auto_seq(4), None, "below the last one");
}
#[test]
fn an_id_reads_and_writes() {
for (text, default, want) in [
(&b"5"[..], 0, Some(Id::new(5, 0))),
(b"5", u64::MAX, Some(Id::new(5, u64::MAX))),
(b"5-3", 0, Some(Id::new(5, 3))),
(b"0-0", 0, Some(Id::MIN)),
(b"", 0, None),
(b"-1", 0, None),
(b"5-", 0, None),
(b"a", 0, None),
(b"5-a", 0, None),
(b"18446744073709551616", 0, None),
] {
assert_eq!(
Id::parse(text, default),
want,
"{:?}",
String::from_utf8_lossy(text)
);
}
assert_eq!(Id::new(5, 3).to_vec(), b"5-3".to_vec());
}
#[test]
fn an_id_round_trips_through_its_bytes() {
for id in [Id::MIN, Id::MAX, Id::new(1, 2), Id::new(u64::MAX, 0)] {
assert_eq!(Id::from_bytes(id.to_bytes()), id);
}
assert!(Id::new(1, 2).to_bytes() < Id::new(1, 3).to_bytes());
assert!(Id::new(1, u64::MAX).to_bytes() < Id::new(2, 0).to_bytes());
}
#[test]
fn stepping_an_id_carries_and_stops() {
assert_eq!(Id::new(1, 2).next(), Some(Id::new(1, 3)));
assert_eq!(Id::new(1, u64::MAX).next(), Some(Id::new(2, 0)));
assert_eq!(Id::MAX.next(), None);
assert_eq!(Id::new(1, 3).prev(), Some(Id::new(1, 2)));
assert_eq!(Id::new(2, 0).prev(), Some(Id::new(1, u64::MAX)));
assert_eq!(Id::MIN.prev(), None);
}
#[test]
fn the_node_size_changes_nothing_but_the_node_count() {
let last = many(400u64);
let mut want = None;
for entries in [1usize, 2, 7, 100, 4096] {
let mut s = Stream::new();
let limits = Limits {
max_node_bytes: NODE_BYTES,
max_node_entries: entries,
};
for ms in 1..=last {
let value = format!("v{ms}");
s.append(Id::new(ms, 0), &[(b"n", value.as_bytes())], limits)
.expect("an append");
}
for ms in (1..=last).step_by(7) {
s.delete(Id::new(ms, 0));
}
let got = dump(&s);
match &want {
None => want = Some(got),
Some(want) => assert_eq!(&got, want, "at {entries} entries a node"),
}
}
}
#[test]
fn a_tiny_byte_limit_still_works() {
let limits = Limits {
max_node_bytes: 1,
max_node_entries: NODE_ENTRIES,
};
let mut s = Stream::new();
for ms in 1..=20u64 {
s.append(Id::new(ms, 0), &[(b"n", b"x")], limits)
.expect("an append");
}
assert_eq!(s.nodes(), 20);
assert_eq!(dump(&s).len(), 20);
}
const REDIS_DUMP: &str = "1b0110000000000000000100000000000000014070\
700000001f000301010102019374656d70657261747572655f63656c736975731491\
72656c61746976655f68756d69646974791200010201000100011501370105010301\
0001010116013801050102010101000117013901050100010201000101018673656e\
736f72078161020601ff030301010101020400406440640000000f00239a5c2c7208\
ea0a";
fn unhex(s: &str) -> Vec<u8> {
let digits: Vec<u8> = s.bytes().filter(|b| !b.is_ascii_whitespace()).collect();
digits
.chunks(2)
.map(|pair| {
let of = |b: u8| (b as char).to_digit(16).expect("a hex digit") as u8;
of(pair[0]) << 4 | of(pair[1])
})
.collect()
}
fn rdb_len(bytes: &[u8], at: usize) -> (usize, usize) {
match bytes[at] >> 6 {
0 => (usize::from(bytes[at] & 0x3F), 1),
1 => (
usize::from(bytes[at] & 0x3F) << 8 | usize::from(bytes[at + 1]),
2,
),
other => panic!("the fixture used length form {other}"),
}
}
fn redis_node() -> (Id, Vec<u8>, Vec<u8>) {
let dump = unhex(REDIS_DUMP);
assert_eq!(dump[0], 0x1B, "RDB_TYPE_STREAM_LISTPACKS_3");
let (nodes, n) = rdb_len(&dump, 1);
assert_eq!(nodes, 1, "the fixture is one node");
let mut at = 1 + n;
let (key, n) = rdb_len(&dump, at);
assert_eq!(key, 16, "a node key is an id in sixteen bytes");
at += n;
let master = Id::from_bytes(dump[at..at + 16].try_into().expect("sixteen bytes"));
at += 16;
let (len, n) = rdb_len(&dump, at);
at += n;
let lp = dump[at..at + len].to_vec();
let rest = dump[at + len..dump.len() - 10].to_vec();
(master, lp, rest)
}
#[test]
fn a_node_is_written_the_way_redis_writes_one() {
let mut s = Stream::new();
add(
&mut s,
1,
1,
&[("temperature_celsius", "21"), ("relative_humidity", "55")],
);
add(
&mut s,
1,
2,
&[("temperature_celsius", "22"), ("relative_humidity", "56")],
);
add(
&mut s,
2,
1,
&[("temperature_celsius", "23"), ("relative_humidity", "57")],
);
add(&mut s, 3, 1, &[("sensor", "a")]);
assert!(s.delete(Id::new(1, 2)));
let (master, lp, _) = redis_node();
assert_eq!(s.nodes(), 1, "all four fit in one node");
assert_eq!(s.nodes[0].master, master);
assert_eq!(s.nodes[0].lp.as_bytes(), &lp[..]);
}
#[test]
fn a_node_redis_wrote_reads_back() {
let (master, lp, rest) = redis_node();
let lp = Listpack::from_bytes(&lp).expect("a listpack Redis wrote");
let s = Stream {
nodes: VecDeque::from(vec![Node { master, lp }]),
length: u64::from(rest[0]),
last: Id::new(u64::from(rest[1]), u64::from(rest[2])),
max_deleted: Id::new(u64::from(rest[5]), u64::from(rest[6])),
added: u64::from(rest[7]),
groups: Vec::new(),
epoch: 0,
};
assert_eq!(s.len(), 3);
assert_eq!(s.last_id(), Id::new(3, 1));
assert_eq!(s.max_deleted_id(), Id::new(1, 2));
assert_eq!(s.added(), 4);
assert_eq!(s.first_id(), Some(Id::new(1, 1)));
assert_eq!(
dump(&s),
vec![
(
Id::new(1, 1),
vec![
(b"temperature_celsius".to_vec(), b"21".to_vec()),
(b"relative_humidity".to_vec(), b"55".to_vec())
]
),
(
Id::new(2, 1),
vec![
(b"temperature_celsius".to_vec(), b"23".to_vec()),
(b"relative_humidity".to_vec(), b"57".to_vec())
]
),
(Id::new(3, 1), vec![(b"sensor".to_vec(), b"a".to_vec())]),
]
);
}
fn logged(n: u64) -> Stream {
let mut s = Stream::new();
for ms in 1..=n {
add(&mut s, ms, 0, &[("job", "x")]);
}
s
}
fn read(s: &mut Stream, group: &str, who: &str, count: Option<usize>, now: u64) -> Vec<Id> {
let mut out = Vec::new();
s.read_group(
group.as_bytes(),
who.as_bytes(),
count,
false,
now,
|id, _| {
out.push(id);
true
},
)
.expect("the group");
out
}
#[test]
fn a_group_draining_one_at_a_time_gets_every_entry_once() {
let mut s = logged(500);
s.create_group(b"workers", Id::MIN, Some(0));
let mut got = Vec::new();
for _ in 0..500 {
got.extend(read(&mut s, "workers", "alice", Some(1), 100));
}
assert_eq!(got, (1..=500).map(|ms| Id::new(ms, 0)).collect::<Vec<_>>());
assert!(read(&mut s, "workers", "alice", Some(1), 100).is_empty());
}
#[test]
fn a_group_that_has_caught_up_is_handed_the_next_append() {
let mut s = logged(3);
s.create_group(b"workers", Id::MIN, Some(0));
assert_eq!(read(&mut s, "workers", "alice", None, 100).len(), 3);
for ms in 4..=200u64 {
add(&mut s, ms, 0, &[("job", "x")]);
assert_eq!(
read(&mut s, "workers", "alice", Some(1), 100),
vec![Id::new(ms, 0)],
"the entry appended a moment ago"
);
}
}
#[test]
fn a_delete_between_two_group_reads_does_not_lose_the_rest() {
let mut s = logged(300);
s.create_group(b"workers", Id::MIN, Some(0));
let mut got = read(&mut s, "workers", "alice", Some(10), 100);
assert!(s.delete(Id::new(150, 0)));
while got.len() < 299 {
let more = read(&mut s, "workers", "alice", Some(1), 100);
assert_eq!(more.len(), 1, "at {}", got.len());
got.extend(more);
}
let want: Vec<Id> = (1..=300)
.filter(|ms| *ms != 150)
.map(|ms| Id::new(ms, 0))
.collect();
assert_eq!(got, want);
}
#[test]
fn moving_the_bookmark_back_reads_it_all_again() {
let mut s = logged(250);
s.create_group(b"workers", Id::MIN, Some(0));
for _ in 0..120 {
read(&mut s, "workers", "alice", Some(1), 100);
}
s.group_mut(b"workers")
.expect("the group")
.set_id(Id::MIN, Some(0));
let mut got = Vec::new();
for _ in 0..250 {
got.extend(read(&mut s, "workers", "alice", Some(1), 100));
}
assert_eq!(got, (1..=250).map(|ms| Id::new(ms, 0)).collect::<Vec<_>>());
}
#[test]
fn a_trim_under_a_group_does_not_hand_out_the_wrong_entries() {
let mut s = logged(500);
s.create_group(b"workers", Id::MIN, Some(0));
for _ in 0..50 {
read(&mut s, "workers", "alice", Some(1), 100);
}
assert_eq!(s.trim_maxlen(200, false, None), 300);
let mut got = Vec::new();
for _ in 0..250 {
got.extend(read(&mut s, "workers", "alice", Some(1), 100));
}
assert_eq!(
got,
(301..=500).map(|ms| Id::new(ms, 0)).collect::<Vec<_>>(),
"everything the trim left that the group had not read"
);
}
#[test]
fn a_group_is_made_once() {
let mut s = logged(3);
assert!(s.create_group(b"workers", Id::MIN, Some(0)));
assert!(!s.create_group(b"workers", Id::MIN, Some(0)));
assert!(s.group(b"workers").is_some());
assert!(s.destroy_group(b"workers"));
assert!(!s.destroy_group(b"workers"));
assert!(s.group(b"workers").is_none());
}
#[test]
fn a_group_read_hands_out_what_comes_after_the_bookmark() {
let mut s = logged(5);
s.create_group(b"workers", Id::MIN, Some(0));
assert_eq!(
read(&mut s, "workers", "alice", Some(2), 100),
vec![Id::new(1, 0), Id::new(2, 0)]
);
assert_eq!(
read(&mut s, "workers", "bob", Some(2), 100),
vec![Id::new(3, 0), Id::new(4, 0)]
);
assert_eq!(
read(&mut s, "workers", "alice", None, 100),
vec![Id::new(5, 0)]
);
assert_eq!(read(&mut s, "workers", "alice", None, 100), vec![]);
}
#[test]
fn a_group_read_fills_the_pending_list() {
let mut s = logged(3);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", None, 500);
let g = s.group(b"workers").expect("the group");
assert_eq!(g.pending_len(), 3);
assert_eq!(g.last_id(), Id::new(3, 0));
assert_eq!(g.entries_read(), Some(3));
assert_eq!(s.lag(g), Some(0));
let c = g.consumer_named(b"alice").expect("alice");
assert_eq!(c.len(), 3);
assert_eq!(c.active(), Some(500));
}
#[test]
fn a_read_that_finds_nothing_is_seen_but_not_active() {
let mut s = logged(1);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", None, 100);
read(&mut s, "workers", "alice", None, 900);
let c = s
.group(b"workers")
.expect("the group")
.consumer_named(b"alice")
.expect("alice");
assert_eq!((c.seen(), c.active()), (900, Some(100)));
}
#[test]
fn a_group_starting_at_the_end_reads_only_what_comes_next() {
let mut s = logged(3);
s.create_group(b"workers", s.last_id(), Some(s.added()));
assert_eq!(read(&mut s, "workers", "alice", None, 1), vec![]);
add(&mut s, 4, 0, &[("job", "x")]);
assert_eq!(
read(&mut s, "workers", "alice", None, 1),
vec![Id::new(4, 0)]
);
}
#[test]
fn reading_a_group_that_is_not_there_says_so() {
let mut s = logged(1);
assert!(
s.read_group(b"nope", b"alice", None, false, 1, |_, _| true)
.is_none()
);
}
#[test]
fn a_consumer_can_re_read_what_it_is_holding() {
let mut s = logged(4);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", Some(2), 1);
read(&mut s, "workers", "bob", Some(2), 1);
let mut out = Vec::new();
s.read_group_pending(b"workers", b"alice", Id::MIN, None, 2, |id, fields| {
out.push((id, fields.map(|f| f.len())));
true
})
.expect("the group");
assert_eq!(
out,
vec![(Id::new(1, 0), Some(1)), (Id::new(2, 0), Some(1))]
);
let mut after = Vec::new();
s.read_group_pending(b"workers", b"alice", Id::new(1, 0), None, 2, |id, _| {
after.push(id);
true
});
assert_eq!(after, vec![Id::new(2, 0)]);
}
#[test]
fn re_reading_counts_as_being_handed_it_again() {
let mut s = logged(2);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", None, 100);
s.read_group_pending(
b"workers",
b"alice",
Id::MIN,
None,
700,
|_: Id, _: Option<Fields<'_>>| true,
);
let g = s.group(b"workers").expect("the group");
let nack = g.nack(Id::new(1, 0)).expect("a nack");
assert_eq!((nack.count(), nack.time()), (2, 700));
assert_eq!(g.last_id(), Id::new(2, 0));
assert_eq!(g.pending_len(), 2);
}
#[test]
fn a_hole_in_front_of_a_group_takes_its_lag() {
let mut s = logged(5);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", Some(3), 1);
assert_eq!(counters(&s), (Some(3), Some(2)));
assert!(s.delete(Id::new(1, 0)));
assert_eq!(counters(&s), (Some(3), Some(2)));
assert!(s.delete(Id::new(5, 0)));
assert_eq!(counters(&s), (Some(3), None));
read(&mut s, "workers", "alice", None, 1);
assert_eq!(s.group(b"workers").expect("g").last_id(), Id::new(4, 0));
assert_eq!(counters(&s), (None, None));
}
#[test]
fn trimming_past_a_group_leaves_it_the_length() {
let mut s = logged(500);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", Some(400), 1);
assert_eq!(counters(&s), (Some(400), Some(100)));
assert_eq!(s.trim_maxlen(200, true, None), 300);
assert_eq!(counters(&s), (Some(400), Some(100)));
assert_eq!(s.trim_maxlen(10, true, None), 190);
assert_eq!(counters(&s), (Some(400), Some(10)));
}
fn counters(s: &Stream) -> (Option<u64>, Option<u64>) {
let g = s.group(b"workers").expect("the group");
(g.entries_read(), s.lag(g))
}
#[test]
fn an_entry_that_went_away_still_comes_back_as_a_hole() {
let mut s = logged(3);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", None, 1);
assert!(s.delete(Id::new(2, 0)));
let mut out = Vec::new();
s.read_group_pending(b"workers", b"alice", Id::MIN, None, 2, |id, fields| {
out.push((id, fields.is_some()));
true
});
assert_eq!(
out,
vec![
(Id::new(1, 0), true),
(Id::new(2, 0), false),
(Id::new(3, 0), true)
]
);
}
#[test]
fn acking_clears_the_pending_list() {
let mut s = logged(3);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", None, 1);
let g = s.group_mut(b"workers").expect("the group");
assert!(g.ack(Id::new(2, 0)));
assert_eq!(g.pending_len(), 2);
}
#[test]
fn a_claim_moves_work_off_a_consumer_that_stopped() {
let mut s = logged(2);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", None, 100);
let mut gone = Vec::new();
let took = s
.claim(
b"workers",
b"bob",
&[Id::new(1, 0), Id::new(2, 0)],
500,
5_000,
None,
true,
false,
5_000,
&mut gone,
)
.expect("the group");
assert_eq!(took, vec![Id::new(1, 0), Id::new(2, 0)]);
assert!(gone.is_empty());
let g = s.group(b"workers").expect("the group");
assert!(g.consumer_named(b"alice").expect("alice").is_empty());
assert_eq!(g.consumer_named(b"bob").expect("bob").len(), 2);
assert_eq!(g.nack(Id::new(1, 0)).expect("a nack").count(), 2);
}
#[test]
fn a_claim_leaves_work_that_is_not_idle_enough_alone() {
let mut s = logged(1);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", None, 100);
let mut gone = Vec::new();
let took = s
.claim(
b"workers",
b"bob",
&[Id::new(1, 0)],
5_000,
200,
None,
true,
false,
200,
&mut gone,
)
.expect("the group");
assert!(took.is_empty());
assert_eq!(
s.group(b"workers")
.expect("the group")
.consumer_named(b"alice")
.expect("alice")
.len(),
1
);
}
#[test]
fn claiming_an_entry_that_went_away_drops_it_instead() {
let mut s = logged(2);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", None, 100);
assert!(s.delete(Id::new(1, 0)));
let mut gone = Vec::new();
let took = s
.claim(
b"workers",
b"bob",
&[Id::new(1, 0), Id::new(2, 0)],
0,
5_000,
None,
true,
false,
5_000,
&mut gone,
)
.expect("the group");
assert_eq!(took, vec![Id::new(2, 0)]);
assert_eq!(gone, vec![Id::new(1, 0)]);
assert_eq!(s.group(b"workers").expect("the group").pending_len(), 1);
}
#[test]
fn force_only_works_on_an_entry_that_is_really_there() {
let mut s = logged(2);
s.create_group(b"workers", s.last_id(), Some(2));
let mut gone = Vec::new();
let took = s
.claim(
b"workers",
b"bob",
&[Id::new(1, 0), Id::new(99, 0)],
0,
100,
None,
true,
true,
100,
&mut gone,
)
.expect("the group");
assert_eq!(took, vec![Id::new(1, 0)], "99-0 is not in the stream");
assert_eq!(s.group(b"workers").expect("the group").pending_len(), 1);
}
#[test]
fn autoclaim_sweeps_the_stale_ones_and_says_where_it_stopped() {
let mut s = logged(6);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", Some(3), 100);
read(&mut s, "workers", "alice", None, 900);
let mut gone = Vec::new();
let (cursor, took) = s
.autoclaim(
b"workers",
b"bob",
Id::MIN,
500,
100,
true,
1_000,
&mut gone,
)
.expect("the group");
assert_eq!(cursor, None, "the sweep reached the end");
assert_eq!(took, vec![Id::new(1, 0), Id::new(2, 0), Id::new(3, 0)]);
assert_eq!(
s.group(b"workers")
.expect("the group")
.consumer_named(b"bob")
.expect("bob")
.len(),
3
);
}
#[test]
fn autoclaim_hands_back_a_cursor_when_it_hits_the_count() {
let mut s = logged(10);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", None, 100);
let mut gone = Vec::new();
let (cursor, took) = s
.autoclaim(b"workers", b"bob", Id::MIN, 0, 4, true, 1_000, &mut gone)
.expect("the group");
assert_eq!(took.len(), 4);
assert_eq!(cursor, Some(Id::new(5, 0)));
let (cursor, took) = s
.autoclaim(
b"workers",
b"bob",
cursor.expect("a cursor"),
0,
100,
true,
1_000,
&mut gone,
)
.expect("the group");
assert_eq!(took.len(), 6);
assert_eq!(cursor, None);
}
#[test]
fn a_group_survives_the_stream_being_trimmed_under_it() {
let mut s = logged(10);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", Some(5), 100);
assert_eq!(s.trim_maxlen(3, true, None), 7);
assert_eq!(s.group(b"workers").expect("the group").pending_len(), 5);
let mut gone = Vec::new();
let took = s
.claim(
b"workers",
b"bob",
&(1..=5).map(|ms| Id::new(ms, 0)).collect::<Vec<_>>(),
0,
1_000,
None,
true,
false,
1_000,
&mut gone,
)
.expect("the group");
assert!(took.is_empty(), "none of them are there any more");
assert_eq!(gone.len(), 5);
assert_eq!(s.group(b"workers").expect("the group").pending_len(), 0);
}
fn round_trip(s: &Stream) -> Stream {
let mut bytes = Vec::new();
s.freeze(&mut bytes);
let back = Stream::thaw(&bytes).expect("our own bytes");
assert_eq!(dump(&back), dump(s), "the entries");
assert_eq!(back.len(), s.len(), "the length");
assert_eq!(back.added(), s.added(), "the count added");
assert_eq!(back.last_id(), s.last_id(), "the last ID");
assert_eq!(back.max_deleted_id(), s.max_deleted_id(), "the max deleted");
assert_eq!(back.nodes(), s.nodes(), "the node count");
assert_eq!(back, *s, "the whole thing");
back
}
#[test]
fn a_frozen_stream_comes_back_with_every_entry_it_held() {
let mut s = Stream::new();
for ms in 1..=500u64 {
add(&mut s, ms, 0, &[("job", "x"), ("n", "1")]);
}
add(&mut s, 500, 1, &[("job", "y")]);
assert!(s.nodes() > 1, "more than one node, so the walk is tested");
round_trip(&s);
}
#[test]
fn a_frozen_stream_keeps_the_holes_and_the_counters() {
let mut s = logged(200);
for ms in [3u64, 4, 5, 100, 199] {
assert!(s.delete(Id::new(ms, 0)));
}
s.trim_minid(Id::new(20, 0), true, None);
let back = round_trip(&s);
assert_eq!(back.first_id(), Some(Id::new(20, 0)));
assert!(!back.contains(Id::new(100, 0)), "a hole is still a hole");
assert_eq!(back.max_deleted_id(), Id::new(199, 0));
}
#[test]
fn a_frozen_stream_keeps_its_groups_and_who_is_holding_what() {
let mut s = logged(20);
s.create_group(b"workers", Id::MIN, Some(0));
s.create_group(b"audit", Id::new(5, 0), None);
read(&mut s, "workers", "alice", Some(6), 1_000);
read(&mut s, "workers", "bob", Some(4), 2_000);
assert_eq!(
s.nack(b"workers", Id::new(2, 0), Retry::Keep, true),
Some(true)
);
s.group_mut(b"workers")
.expect("the group")
.create_consumer(b"carol", 3_000);
read(&mut s, "workers", "dave", Some(2), 4_000);
s.group_mut(b"workers")
.expect("the group")
.delete_consumer(b"carol");
let back = round_trip(&s);
let g = back.group(b"workers").expect("the group");
assert_eq!(g.pending_len(), 12);
assert_eq!(g.nacked_len(), 1);
assert_eq!(g.entries_read(), Some(12));
assert_eq!(g.consumer_named(b"alice").expect("alice").len(), 5);
assert_eq!(g.consumer_named(b"bob").expect("bob").len(), 4);
assert_eq!(g.consumer_named(b"dave").expect("dave").len(), 2);
assert_eq!(g.consumer_named(b"carol"), None);
assert_eq!(g.slot(b"dave"), Some(3));
assert_eq!(g.nack(Id::new(1, 0)).expect("a nack").owner(), Some(0));
assert_eq!(g.nack(Id::new(2, 0)).expect("a nack").owner(), None);
assert_eq!(
back.group(b"audit").expect("audit").last_id(),
Id::new(5, 0)
);
assert_eq!(back.group(b"audit").expect("audit").entries_read(), None);
}
#[test]
fn a_stream_that_came_back_still_takes_entries_and_reads_them() {
let mut s = logged(10);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", Some(4), 1_000);
let mut back = round_trip(&s);
add(&mut back, 11, 0, &[("job", "new")]);
assert_eq!(back.len(), 11);
assert_eq!(
read(&mut back, "workers", "alice", Some(3), 2_000),
vec![Id::new(5, 0), Id::new(6, 0), Id::new(7, 0)]
);
assert!(
back.group_mut(b"workers")
.expect("the group")
.ack(Id::new(1, 0))
);
assert_eq!(back.group(b"workers").expect("the group").pending_len(), 6);
}
#[test]
fn an_empty_stream_that_still_exists_comes_back() {
let mut s = logged(3);
for ms in 1..=3u64 {
assert!(s.delete(Id::new(ms, 0)));
}
assert_eq!(s.nodes(), 0, "the last node went with the last entry");
let back = round_trip(&s);
assert!(back.is_empty());
assert_eq!(back.last_id(), Id::new(3, 0));
round_trip(&Stream::new());
}
#[test]
fn a_frozen_stream_that_arrives_damaged_is_an_error_and_not_a_panic() {
let mut s = logged(8);
s.create_group(b"workers", Id::MIN, Some(0));
read(&mut s, "workers", "alice", Some(3), 1_000);
let mut bytes = Vec::new();
s.freeze(&mut bytes);
for cut in 0..bytes.len() {
assert!(Stream::thaw(&bytes[..cut]).is_err(), "cut at {cut}");
}
for at in 0..bytes.len().min(40) {
for bit in 0..8 {
let mut bad = bytes.clone();
bad[at] ^= 1 << bit;
let _ = Stream::thaw(&bad);
}
}
assert_eq!(Stream::thaw(&[]), Err(Broken::Short));
assert_eq!(Stream::thaw(&[9]), Err(Broken::Form));
}
#[test]
fn nothing_at_all() {
let mut s = Stream::new();
assert!(s.is_empty());
assert_eq!(s.len(), 0);
assert_eq!(s.first_id(), None);
assert_eq!(s.last_id(), Id::MIN);
assert_eq!(s.trim_maxlen(0, true, None), 0);
assert_eq!(s.trim_minid(Id::MAX, true, None), 0);
assert_eq!(dump(&s), vec![]);
}
}