1use std::cmp::Ordering;
108use std::collections::VecDeque;
109
110use yo_common::num::{DIGITS_MAX, i64_digits, push_u64, u64_digits};
111
112use crate::frozen::{self, Broken};
113use crate::listpack::{self, Entry, Listpack};
114
115pub mod groups;
116
117pub use groups::{Consumer, Filter, Group, Nack, Retry};
118
119pub const NODE_BYTES: usize = 4096;
126
127pub const NODE_ENTRIES: usize = 100;
131
132const LIVE: i64 = 0;
134
135const DELETED: i64 = 1;
137
138const SAME_FIELDS: i64 = 2;
140
141const MASTER_FIELDS: usize = 3;
145
146#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Default)]
152pub struct Id {
153 pub ms: u64,
155 pub seq: u64,
157}
158
159impl Id {
160 pub const MIN: Id = Id { ms: 0, seq: 0 };
162
163 pub const MAX: Id = Id {
165 ms: u64::MAX,
166 seq: u64::MAX,
167 };
168
169 #[must_use]
171 #[inline]
172 pub const fn new(ms: u64, seq: u64) -> Id {
173 Id { ms, seq }
174 }
175
176 #[must_use]
181 pub const fn next(self) -> Option<Id> {
182 if self.seq != u64::MAX {
183 Some(Id {
184 ms: self.ms,
185 seq: self.seq + 1,
186 })
187 } else if self.ms != u64::MAX {
188 Some(Id {
189 ms: self.ms + 1,
190 seq: 0,
191 })
192 } else {
193 None
194 }
195 }
196
197 #[must_use]
199 pub const fn prev(self) -> Option<Id> {
200 if self.seq != 0 {
201 Some(Id {
202 ms: self.ms,
203 seq: self.seq - 1,
204 })
205 } else if self.ms != 0 {
206 Some(Id {
207 ms: self.ms - 1,
208 seq: u64::MAX,
209 })
210 } else {
211 None
212 }
213 }
214
215 #[must_use]
220 pub fn to_bytes(self) -> [u8; 16] {
221 let mut out = [0u8; 16];
222 out[..8].copy_from_slice(&self.ms.to_be_bytes());
223 out[8..].copy_from_slice(&self.seq.to_be_bytes());
224 out
225 }
226
227 #[must_use]
229 pub fn from_bytes(bytes: [u8; 16]) -> Id {
230 let mut ms = [0u8; 8];
231 let mut seq = [0u8; 8];
232 ms.copy_from_slice(&bytes[..8]);
233 seq.copy_from_slice(&bytes[8..]);
234 Id {
235 ms: u64::from_be_bytes(ms),
236 seq: u64::from_be_bytes(seq),
237 }
238 }
239
240 pub fn write_to(self, out: &mut Vec<u8>) {
242 push_u64(out, self.ms);
243 out.push(b'-');
244 push_u64(out, self.seq);
245 }
246
247 #[must_use]
249 pub fn to_vec(self) -> Vec<u8> {
250 let mut out = Vec::with_capacity(41);
251 self.write_to(&mut out);
252 out
253 }
254
255 #[must_use]
262 pub fn parse(s: &[u8], default: u64) -> Option<Id> {
263 let (ms, seq) = match s.iter().position(|c| *c == b'-') {
264 Some(at) => (&s[..at], Some(&s[at + 1..])),
265 None => (s, None),
266 };
267 Some(Id {
268 ms: digits(ms)?,
269 seq: match seq {
270 Some(seq) => digits(seq)?,
271 None => default,
272 },
273 })
274 }
275}
276
277fn digits(s: &[u8]) -> Option<u64> {
282 if s.is_empty() || s.len() > 20 {
283 return None;
284 }
285 let mut n = 0u64;
286 for c in s {
287 let d = c.wrapping_sub(b'0');
288 if d > 9 {
289 return None;
290 }
291 n = n.checked_mul(10)?.checked_add(u64::from(d))?;
292 }
293 Some(n)
294}
295
296#[derive(Debug, Clone, Copy, PartialEq, Eq)]
298pub enum Refused {
299 NotGreater,
301 Zero,
303 Full,
305}
306
307#[derive(Debug, Clone, Copy, PartialEq, Eq)]
309pub struct Limits {
310 pub max_node_bytes: usize,
312 pub max_node_entries: usize,
314}
315
316impl Default for Limits {
317 fn default() -> Limits {
319 Limits {
320 max_node_bytes: NODE_BYTES,
321 max_node_entries: NODE_ENTRIES,
322 }
323 }
324}
325
326#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
333pub enum Refs {
334 #[default]
337 Keep,
338 Drop,
340 Acked,
342}
343
344#[derive(Debug, Clone, Copy, PartialEq, Eq)]
350pub enum Fate {
351 Missing,
353 Gone,
355 Held,
357}
358
359impl Fate {
360 #[must_use]
362 #[inline]
363 pub fn code(self) -> i64 {
364 match self {
365 Fate::Missing => -1,
366 Fate::Gone => 1,
367 Fate::Held => 2,
368 }
369 }
370}
371
372#[derive(Debug, Clone, PartialEq, Eq)]
374struct Node {
375 master: Id,
380 lp: Listpack,
381}
382
383#[derive(Debug, Clone, Copy)]
391struct Edges {
392 added: u64,
393 length: u64,
394 first: Option<Id>,
395 last: Id,
396 max_deleted: Id,
397}
398
399impl Edges {
400 fn estimate(&self, id: Id) -> Option<u64> {
412 if self.added == 0 {
413 return Some(0);
414 }
415 if id >= self.last {
416 return Some(self.added);
417 }
418 let first = self.first?;
419 if self.max_deleted != Id::MIN && self.max_deleted >= first {
423 return None;
424 }
425 let behind = self.added - self.length;
426 match id.cmp(&first) {
427 Ordering::Less => Some(behind),
428 Ordering::Equal => Some(behind + 1),
429 Ordering::Greater => None,
430 }
431 }
432
433 fn holed_from(&self, id: Id) -> bool {
439 if self.length == 0 || self.max_deleted == Id::MIN {
440 return false;
441 }
442 if self.first.is_some_and(|first| first > self.max_deleted) {
443 return false;
444 }
445 id <= self.max_deleted
446 }
447
448 fn lag(&self, group: &Group) -> Option<u64> {
450 if let Some(read) = self.estimate(group.last_id()) {
451 return Some(self.added.saturating_sub(read));
452 }
453 let read = group.entries_read()?;
454 if self.holed_from(group.last_id()) {
455 return None;
456 }
457 Some(self.added.saturating_sub(read))
458 }
459
460 fn on_deliver(&self, group: &Group, id: Id) -> Option<u64> {
467 match group.entries_read() {
468 Some(read) if !self.holed_from(id) => Some(read + 1),
469 _ => self.estimate(id),
470 }
471 }
472}
473
474const FORM_NODES: u8 = 1;
480
481#[derive(Debug, Clone, Copy)]
495pub(crate) struct Cursor {
496 epoch: u64,
498 next: Id,
502 master: Id,
504 byte: usize,
506}
507
508#[derive(Debug, Clone, Default)]
510pub struct Stream {
511 nodes: VecDeque<Node>,
512 length: u64,
514 last: Id,
516 max_deleted: Id,
518 added: u64,
523 groups: Vec<(Vec<u8>, Group)>,
529 epoch: u64,
537}
538
539impl PartialEq for Stream {
545 fn eq(&self, other: &Stream) -> bool {
546 self.nodes == other.nodes
547 && self.length == other.length
548 && self.last == other.last
549 && self.max_deleted == other.max_deleted
550 && self.added == other.added
551 && self.groups == other.groups
552 }
553}
554
555impl Eq for Stream {}
556
557impl Stream {
558 #[must_use]
560 pub fn new() -> Stream {
561 Stream::default()
562 }
563
564 pub(crate) fn raw_nodes(&self) -> impl Iterator<Item = (Id, &[u8])> + '_ {
572 self.nodes.iter().map(|n| (n.master, n.lp.as_bytes()))
573 }
574
575 pub(crate) fn push_raw_node(&mut self, master: Id, lp: Listpack) -> bool {
582 if self.nodes.back().is_some_and(|n| n.master >= master) {
583 return false;
584 }
585 self.nodes.push_back(Node { master, lp });
586 true
587 }
588
589 pub(crate) fn set_counters(
594 &mut self,
595 length: u64,
596 last: Id,
597 max_deleted: Id,
598 added: u64,
599 ) -> bool {
600 if length > added || max_deleted > last || (self.nodes.is_empty() && length != 0) {
601 return false;
602 }
603 self.length = length;
604 self.last = last;
605 self.max_deleted = max_deleted;
606 self.added = added;
607 true
608 }
609
610 pub(crate) fn push_group(&mut self, name: &[u8], group: Group) -> bool {
612 if self.groups.iter().any(|(had, _)| had == name) {
613 return false;
614 }
615 self.groups.push((name.to_vec(), group));
616 true
617 }
618
619 pub fn freeze(&self, out: &mut Vec<u8>) {
628 out.push(FORM_NODES);
629 frozen::put_uint(out, self.length);
630 frozen::put_uint(out, self.last.ms);
631 frozen::put_uint(out, self.last.seq);
632 frozen::put_uint(out, self.max_deleted.ms);
633 frozen::put_uint(out, self.max_deleted.seq);
634 frozen::put_uint(out, self.added);
635
636 frozen::put_uint(out, self.nodes.len() as u64);
637 for node in &self.nodes {
638 frozen::put_uint(out, node.master.ms);
639 frozen::put_uint(out, node.master.seq);
640 frozen::put_bytes(out, node.lp.as_bytes());
641 }
642
643 frozen::put_uint(out, self.groups.len() as u64);
644 for (name, group) in &self.groups {
645 frozen::put_bytes(out, name);
646 group.freeze(out);
647 }
648 }
649
650 pub fn thaw(bytes: &[u8]) -> Result<Stream, Broken> {
660 let mut cut = frozen::Cut::new(bytes);
661 if cut.byte()? != FORM_NODES {
662 return Err(Broken::Form);
663 }
664 let length = cut.uint()?;
665 let last = Id::new(cut.uint()?, cut.uint()?);
666 let max_deleted = Id::new(cut.uint()?, cut.uint()?);
667 let added = cut.uint()?;
668 if length > added || max_deleted > last {
672 return Err(Broken::Body);
673 }
674
675 let n = usize::try_from(cut.uint()?).map_err(|_| Broken::Short)?;
676 if n > cut.rest().len() {
679 return Err(Broken::Short);
680 }
681 let mut nodes = VecDeque::with_capacity(n);
682 let mut prev: Option<Id> = None;
683 for _ in 0..n {
684 let master = Id::new(cut.uint()?, cut.uint()?);
685 if prev.is_some_and(|p| p >= master) {
689 return Err(Broken::Body);
690 }
691 prev = Some(master);
692 let lp = Listpack::from_bytes(cut.bytes()?).map_err(|_| Broken::Body)?;
693 nodes.push_back(Node { master, lp });
694 }
695 if nodes.is_empty() && length != 0 {
696 return Err(Broken::Body);
697 }
698
699 let n = usize::try_from(cut.uint()?).map_err(|_| Broken::Short)?;
700 if n > cut.rest().len() {
701 return Err(Broken::Short);
702 }
703 let mut groups: Vec<(Vec<u8>, Group)> = Vec::with_capacity(n);
704 for _ in 0..n {
705 let name = cut.bytes()?;
706 if groups.iter().any(|(had, _)| had == name) {
709 return Err(Broken::Body);
710 }
711 groups.push((name.to_vec(), Group::thaw(&mut cut)?));
712 }
713
714 Ok(Stream {
715 nodes,
716 length,
717 last,
718 max_deleted,
719 added,
720 groups,
721 epoch: 0,
722 })
723 }
724
725 #[must_use]
727 #[inline]
728 pub fn len(&self) -> u64 {
729 self.length
730 }
731
732 #[must_use]
738 #[inline]
739 pub fn is_empty(&self) -> bool {
740 self.length == 0
741 }
742
743 #[must_use]
745 #[inline]
746 pub fn last_id(&self) -> Id {
747 self.last
748 }
749
750 #[must_use]
752 #[inline]
753 pub fn max_deleted_id(&self) -> Id {
754 self.max_deleted
755 }
756
757 #[must_use]
759 #[inline]
760 pub fn added(&self) -> u64 {
761 self.added
762 }
763
764 #[must_use]
770 pub fn first_id(&self) -> Option<Id> {
771 let mut found = None;
772 self.walk(Id::MIN, Id::MAX, Some(1), &mut |id, _| {
773 found = Some(id);
774 false
775 });
776 found
777 }
778
779 #[must_use]
787 pub fn top_id(&self) -> Option<Id> {
788 let mut found = None;
789 self.rev_range(Id::MIN, Id::MAX, Some(1), |id, _| {
790 found = Some(id);
791 false
792 });
793 found
794 }
795
796 pub fn set_id(
810 &mut self,
811 last: Id,
812 added: Option<u64>,
813 max_deleted: Option<Id>,
814 ) -> Result<(), Refused> {
815 if self.top_id().is_some_and(|top| last < top) {
816 return Err(Refused::NotGreater);
817 }
818 self.last = last;
819 if let Some(added) = added {
820 self.added = added;
821 }
822 if let Some(id) = max_deleted {
823 self.max_deleted = id;
824 }
825 Ok(())
826 }
827
828 #[must_use]
834 pub fn auto_id(&self, now: u64) -> Option<Id> {
835 if now > self.last.ms {
836 Some(Id { ms: now, seq: 0 })
837 } else {
838 self.last.next()
839 }
840 }
841
842 #[must_use]
844 pub fn auto_seq(&self, ms: u64) -> Option<Id> {
845 if ms > self.last.ms {
846 Some(Id { ms, seq: 0 })
847 } else if ms == self.last.ms {
848 self.last.next().filter(|id| id.ms == ms)
849 } else {
850 None
851 }
852 }
853
854 pub fn append(
862 &mut self,
863 id: Id,
864 fields: &[(&[u8], &[u8])],
865 limits: Limits,
866 ) -> Result<(), Refused> {
867 if id == Id::MIN {
868 return Err(Refused::Zero);
869 }
870 if id <= self.last {
871 return Err(Refused::NotGreater);
872 }
873
874 let size: usize = fields
878 .iter()
879 .map(|(f, v)| f.len() + v.len() + 11)
880 .sum::<usize>()
881 + 32;
882
883 let fits = match self.nodes.back() {
884 Some(node) => {
885 let (count, deleted) = counts(&node.lp);
886 node.lp.byte_len() + size < limits.max_node_bytes
887 && (count + deleted) < limits.max_node_entries as u64
888 }
889 None => false,
890 };
891
892 if !fits {
893 self.nodes.push_back(Node {
894 master: id,
895 lp: master_of(fields),
896 });
897 }
898 let node = self.nodes.back_mut().expect("a node was just made sure of");
899 let same = same_fields(&node.lp, fields);
900 write_entry(&mut node.lp, node.master, id, fields, same);
901 if bump(&mut node.lp, 1, 0) {
906 self.epoch = self.epoch.wrapping_add(1);
907 }
908
909 self.length += 1;
910 self.added += 1;
911 self.last = id;
912 Ok(())
913 }
914
915 pub fn delete(&mut self, id: Id) -> bool {
920 if !self.remove(id) {
921 return false;
922 }
923 self.max_deleted = self.max_deleted.max(id);
924 true
925 }
926
927 pub fn delete_ref(&mut self, id: Id, refs: Refs) -> Fate {
936 if refs == Refs::Acked && self.still_wanted(id) {
937 return Fate::Held;
938 }
939 if !self.delete(id) {
940 return Fate::Missing;
941 }
942 if refs == Refs::Drop {
943 self.drop_refs(id);
944 }
945 Fate::Gone
946 }
947
948 pub fn ack_delete(&mut self, group: &[u8], id: Id, refs: Refs) -> (Fate, bool) {
961 let Some(g) = self.group_mut(group) else {
962 return (Fate::Missing, false);
963 };
964 if !g.ack(id) {
965 return (Fate::Missing, false);
966 }
967 if refs == Refs::Acked && self.still_wanted(id) {
968 return (Fate::Held, false);
969 }
970 let gone = self.delete(id);
971 if refs == Refs::Drop {
972 self.drop_refs(id);
973 }
974 (Fate::Gone, gone)
975 }
976
977 pub fn nack(&mut self, group: &[u8], id: Id, retry: Retry, force: bool) -> Option<bool> {
984 let here = self.contains(id);
985 let g = self.group_mut(group)?;
986 if g.release(id, retry) {
987 return Some(true);
988 }
989 if force && here {
990 g.force_release(id, retry);
991 return Some(true);
992 }
993 Some(false)
994 }
995
996 fn still_wanted(&self, id: Id) -> bool {
998 self.groups
999 .iter()
1000 .any(|(_, g)| g.nack(id).is_some() || id > g.last_id())
1001 }
1002
1003 fn drop_refs(&mut self, id: Id) {
1005 for (_, g) in &mut self.groups {
1006 g.forget(id);
1007 }
1008 }
1009
1010 fn remove(&mut self, id: Id) -> bool {
1017 let Some(at) = self.node_of(id) else {
1018 return false;
1019 };
1020 let node = &self.nodes[at];
1021 let Some((offset, flags)) = find(&node.lp, node.master, id) else {
1022 return false;
1023 };
1024 if flags & DELETED != 0 {
1025 return false;
1026 }
1027
1028 let (count, _) = counts(&node.lp);
1029 if count == 1 {
1030 self.nodes.remove(at);
1031 } else {
1032 let node = &mut self.nodes[at];
1033 set_int(&mut node.lp, offset, flags | DELETED);
1034 bump(&mut node.lp, -1, 1);
1035 }
1036 self.epoch = self.epoch.wrapping_add(1);
1041 self.length -= 1;
1042 true
1043 }
1044
1045 fn edges(&self) -> Edges {
1051 Edges {
1052 added: self.added,
1053 length: self.length,
1054 first: self.first_id(),
1055 last: self.last,
1056 max_deleted: self.max_deleted,
1057 }
1058 }
1059
1060 #[must_use]
1076 pub fn lag(&self, group: &Group) -> Option<u64> {
1077 self.edges().lag(group)
1078 }
1079
1080 pub fn trim_maxlen(&mut self, len: u64, exact: bool, limit: Option<u64>) -> u64 {
1093 let mut gone = 0;
1094 while self.length > len && !limit.is_some_and(|cap| gone >= cap) {
1095 let Some(node) = self.nodes.front() else {
1096 break;
1097 };
1098 let (count, _) = counts(&node.lp);
1099 if self.length - count >= len {
1100 self.length -= count;
1101 gone += count;
1102 self.nodes.pop_front();
1103 continue;
1104 }
1105 if !exact {
1106 break;
1107 }
1108 let Some(id) = self.first_id() else { break };
1109 self.remove(id);
1110 gone += 1;
1111 }
1112 gone
1113 }
1114
1115 pub fn trim_minid(&mut self, id: Id, exact: bool, limit: Option<u64>) -> u64 {
1120 let mut gone = 0;
1121 while let Some(node) = self.nodes.front() {
1122 if limit.is_some_and(|cap| gone >= cap) {
1123 break;
1124 }
1125 let (count, _) = counts(&node.lp);
1126 if last_of(node) < id {
1127 self.length -= count;
1128 gone += count;
1129 self.nodes.pop_front();
1130 continue;
1131 }
1132 if !exact {
1133 break;
1134 }
1135 let Some(first) = self.first_id() else { break };
1136 if first >= id {
1137 break;
1138 }
1139 self.remove(first);
1140 gone += 1;
1141 }
1142 gone
1143 }
1144
1145 pub fn range<F>(&self, start: Id, end: Id, count: Option<usize>, mut f: F) -> usize
1152 where
1153 F: FnMut(Id, Fields<'_>) -> bool,
1154 {
1155 self.walk(start, end, count, &mut f)
1156 }
1157
1158 pub fn rev_range<'s, F>(&'s self, start: Id, end: Id, count: Option<usize>, mut f: F) -> usize
1164 where
1165 F: FnMut(Id, Fields<'_>) -> bool,
1166 {
1167 let mut seen = 0;
1168 let mut buf: Vec<(Id, Fields<'s>)> = Vec::new();
1173 let last = self.node_from(end);
1179 for node in self.nodes.iter().take(last + 1).rev() {
1180 if node.master > end {
1183 continue;
1184 }
1185 if last_of(node) < start {
1186 break;
1187 }
1188 buf.clear();
1189 each(&node.lp, node.master, None, &mut |id, _, fields| {
1190 if id >= start && id <= end {
1191 buf.push((id, fields));
1192 }
1193 id <= end
1194 });
1195 for (id, fields) in buf.drain(..).rev() {
1196 if count.is_some_and(|want| seen >= want) {
1197 return seen;
1198 }
1199 seen += 1;
1200 if !f(id, fields) {
1201 return seen;
1202 }
1203 }
1204 }
1205 seen
1206 }
1207
1208 #[must_use]
1213 pub fn contains(&self, id: Id) -> bool {
1214 let Some(at) = self.node_of(id) else {
1215 return false;
1216 };
1217 let node = &self.nodes[at];
1218 find(&node.lp, node.master, id).is_some_and(|(_, flags)| flags & DELETED == 0)
1219 }
1220
1221 pub fn create_group(&mut self, name: &[u8], last: Id, read: Option<u64>) -> bool {
1226 if self.group(name).is_some() {
1227 return false;
1228 }
1229 self.groups.push((name.to_vec(), Group::new(last, read)));
1230 true
1231 }
1232
1233 pub fn destroy_group(&mut self, name: &[u8]) -> bool {
1235 let Some(at) = self.groups.iter().position(|(n, _)| n == name) else {
1236 return false;
1237 };
1238 self.groups.remove(at);
1239 true
1240 }
1241
1242 #[must_use]
1244 pub fn group(&self, name: &[u8]) -> Option<&Group> {
1245 self.groups
1246 .iter()
1247 .find(|(n, _)| n.as_slice() == name)
1248 .map(|(_, g)| g)
1249 }
1250
1251 pub fn group_mut(&mut self, name: &[u8]) -> Option<&mut Group> {
1253 self.groups
1254 .iter_mut()
1255 .find(|(n, _)| n.as_slice() == name)
1256 .map(|(_, g)| g)
1257 }
1258
1259 pub fn groups(&self) -> impl Iterator<Item = (&[u8], &Group)> + '_ {
1261 self.groups.iter().map(|(n, g)| (n.as_slice(), g))
1262 }
1263
1264 pub fn read_group<F>(
1278 &mut self,
1279 group: &[u8],
1280 consumer: &[u8],
1281 count: Option<usize>,
1282 noack: bool,
1283 now: u64,
1284 mut f: F,
1285 ) -> Option<usize>
1286 where
1287 F: FnMut(Id, Fields<'_>) -> bool,
1288 {
1289 let edges = self.edges();
1292 let Stream {
1295 nodes,
1296 groups,
1297 epoch,
1298 ..
1299 } = self;
1300 let epoch = *epoch;
1301 let (_, g) = groups.iter_mut().find(|(n, _)| n.as_slice() == group)?;
1302 let slot = g.consumer_or_create(consumer, now);
1303 let Some(from) = g.last_id().next() else {
1304 g.touch(slot, now, false);
1307 return Some(0);
1308 };
1309 let resume = g.resume(epoch, from);
1314 let (seen, mark) = walk_nodes(nodes, from, Id::MAX, count, resume, &mut |id, fields| {
1315 let read = edges.on_deliver(g, id);
1318 if noack {
1319 g.skip(id);
1320 } else {
1321 g.deliver(slot, id, now);
1322 }
1323 g.set_read(read);
1324 f(id, fields)
1325 });
1326 let next = g.last_id().next();
1327 g.set_resume(match (mark, next) {
1328 (Some((master, byte)), Some(next)) => Some(Cursor {
1329 epoch,
1330 next,
1331 master,
1332 byte,
1333 }),
1334 _ => None,
1335 });
1336 g.touch(slot, now, seen > 0);
1337 Some(seen)
1338 }
1339
1340 pub fn read_group_pending<F>(
1355 &mut self,
1356 group: &[u8],
1357 consumer: &[u8],
1358 after: Id,
1359 count: Option<usize>,
1360 now: u64,
1361 mut f: F,
1362 ) -> Option<usize>
1363 where
1364 F: FnMut(Id, Option<Fields<'_>>) -> bool,
1365 {
1366 let g = self.group_mut(group)?;
1376 let slot = g.consumer_or_create(consumer, now);
1377 let ids: Vec<Id> = g
1378 .consumer(slot)
1379 .expect("the slot that was just made")
1380 .pending()
1381 .filter(|&id| id > after)
1382 .take(count.unwrap_or(usize::MAX))
1383 .collect();
1384 for &id in &ids {
1385 g.redeliver(id, now);
1386 }
1387 g.touch(slot, now, !ids.is_empty());
1388
1389 let mut seen = 0;
1390 for &id in &ids {
1391 seen += 1;
1392 let mut go = true;
1393 let mut found = false;
1394 self.walk(id, id, Some(1), &mut |got, fields| {
1395 found = true;
1396 go = f(got, Some(fields));
1397 false
1398 });
1399 if !found {
1400 go = f(id, None);
1401 }
1402 if !go {
1403 break;
1404 }
1405 }
1406 Some(seen)
1407 }
1408
1409 #[allow(clippy::too_many_arguments)]
1422 pub fn claim(
1423 &mut self,
1424 group: &[u8],
1425 consumer: &[u8],
1426 ids: &[Id],
1427 min_idle: u64,
1428 time: u64,
1429 retry: Option<u64>,
1430 bump: bool,
1431 force: bool,
1432 now: u64,
1433 gone: &mut Vec<Id>,
1434 ) -> Option<Vec<Id>> {
1435 self.group_mut(group)?.consumer_or_create(consumer, now);
1441 let mut took = Vec::new();
1442 for &id in ids {
1443 let here = self.contains(id);
1444 let g = self.group_mut(group)?;
1445 let slot = g.consumer_or_create(consumer, now);
1446 match g.nack(id) {
1447 Some(nack) => {
1448 if !here {
1449 g.forget(id);
1450 gone.push(id);
1451 continue;
1452 }
1453 if nack.idle(now) < min_idle {
1454 continue;
1455 }
1456 if g.claim(id, slot, time, retry, bump) {
1457 took.push(id);
1458 }
1459 }
1460 None => {
1461 if force && here && g.force(id, slot, time, retry.unwrap_or(1)) {
1464 took.push(id);
1465 }
1466 }
1467 }
1468 }
1469 if !took.is_empty() {
1473 let g = self.group_mut(group).expect("the group found a moment ago");
1474 let slot = g.consumer_or_create(consumer, now);
1475 g.touch(slot, now, true);
1476 }
1477 Some(took)
1478 }
1479
1480 #[allow(clippy::too_many_arguments)]
1488 pub fn autoclaim(
1489 &mut self,
1490 group: &[u8],
1491 consumer: &[u8],
1492 start: Id,
1493 min_idle: u64,
1494 count: usize,
1495 bump: bool,
1496 now: u64,
1497 gone: &mut Vec<Id>,
1498 ) -> Option<(Option<Id>, Vec<Id>)> {
1499 let mut ids = Vec::new();
1500 let cursor = self
1501 .group(group)?
1502 .claimable(start, min_idle, now, count, &mut ids);
1503 let took = self.claim(
1504 group, consumer, &ids, min_idle, now, None, bump, false, now, gone,
1505 )?;
1506 Some((cursor, took))
1507 }
1508
1509 #[must_use]
1515 pub fn memory_bytes(&self) -> usize {
1516 let nodes: usize = self
1517 .nodes
1518 .iter()
1519 .map(|node| node.lp.byte_len() + std::mem::size_of::<Node>())
1520 .sum();
1521 let groups: usize = self
1522 .groups
1523 .iter()
1524 .map(|(name, g)| {
1525 name.capacity() + std::mem::size_of::<(Vec<u8>, Group)>() + g.memory_bytes()
1526 })
1527 .sum();
1528 nodes + groups
1529 }
1530
1531 #[must_use]
1534 pub fn nodes(&self) -> usize {
1535 self.nodes.len()
1536 }
1537
1538 fn walk<F>(&self, start: Id, end: Id, count: Option<usize>, f: &mut F) -> usize
1540 where
1541 F: FnMut(Id, Fields<'_>) -> bool,
1542 {
1543 walk_nodes(&self.nodes, start, end, count, None, f).0
1544 }
1545
1546 fn node_from(&self, id: Id) -> usize {
1548 node_from(&self.nodes, id)
1549 }
1550
1551 fn node_of(&self, id: Id) -> Option<usize> {
1553 let at = self.node_from(id);
1554 let node = self.nodes.get(at)?;
1555 (node.master <= id && id <= last_of(node)).then_some(at)
1556 }
1557}
1558
1559fn node_from(nodes: &VecDeque<Node>, id: Id) -> usize {
1570 let after = nodes.partition_point(|node| node.master <= id);
1571 after.saturating_sub(1)
1572}
1573
1574fn walk_nodes<F>(
1582 nodes: &VecDeque<Node>,
1583 start: Id,
1584 end: Id,
1585 count: Option<usize>,
1586 resume: Option<(Id, usize)>,
1587 f: &mut F,
1588) -> (usize, Option<(Id, usize)>)
1589where
1590 F: FnMut(Id, Fields<'_>) -> bool,
1591{
1592 let mut seen = 0;
1593 let mut stop = false;
1594 let first = node_from(nodes, start);
1595 let mut from =
1599 resume.and_then(|(master, byte)| (nodes.get(first)?.master == master).then_some(byte));
1600 let mut mark = None;
1601 for node in nodes.iter().skip(first) {
1602 if node.master > end {
1603 break;
1604 }
1605 let at = each(&node.lp, node.master, from.take(), &mut |id, _, fields| {
1606 if id > end {
1607 stop = true;
1608 return false;
1609 }
1610 if id < start {
1611 return true;
1612 }
1613 if count.is_some_and(|want| seen >= want) {
1614 stop = true;
1615 return false;
1616 }
1617 seen += 1;
1618 if !f(id, fields) {
1619 stop = true;
1620 return false;
1621 }
1622 true
1623 });
1624 mark = Some((node.master, at));
1625 if stop {
1626 break;
1627 }
1628 }
1629 (seen, mark)
1630}
1631
1632#[derive(Debug, Clone)]
1639pub struct Fields<'a> {
1640 names: Option<listpack::Iter<'a>>,
1642 body: listpack::Iter<'a>,
1643 left: usize,
1644}
1645
1646impl<'a> Iterator for Fields<'a> {
1647 type Item = (Entry<'a>, Entry<'a>);
1648
1649 fn next(&mut self) -> Option<(Entry<'a>, Entry<'a>)> {
1650 if self.left == 0 {
1651 return None;
1652 }
1653 self.left -= 1;
1654 let name = match &mut self.names {
1655 Some(names) => names.next()?,
1656 None => self.body.next()?,
1657 };
1658 Some((name, self.body.next()?))
1659 }
1660
1661 fn size_hint(&self) -> (usize, Option<usize>) {
1662 (self.left, Some(self.left))
1663 }
1664}
1665
1666impl ExactSizeIterator for Fields<'_> {}
1667
1668impl Fields<'_> {
1669 #[must_use]
1671 #[inline]
1672 pub fn is_empty(&self) -> bool {
1673 self.left == 0
1674 }
1675}
1676
1677fn master_of(fields: &[(&[u8], &[u8])]) -> Listpack {
1679 let mut lp = Listpack::new();
1680 push_int(&mut lp, 0);
1681 push_int(&mut lp, 0);
1682 push_int(&mut lp, fields.len() as i64);
1683 for (name, _) in fields {
1684 lp.push(name);
1685 }
1686 push_int(&mut lp, 0);
1687 lp
1688}
1689
1690fn same_fields(lp: &Listpack, fields: &[(&[u8], &[u8])]) -> bool {
1692 let mut it = lp.iter();
1693 let (_, _, want) = match (it.next(), it.next(), it.next()) {
1694 (Some(_), Some(_), Some(Entry::Int(n))) => ((), (), n),
1695 _ => return false,
1696 };
1697 if want != fields.len() as i64 {
1698 return false;
1699 }
1700 fields.iter().all(|(name, _)| match it.next() {
1701 Some(Entry::Str(s)) => s == *name,
1702 Some(Entry::Int(n)) => {
1703 let mut buf = [0u8; DIGITS_MAX];
1704 i64_digits(&mut buf, n) == *name
1705 }
1706 None => false,
1707 })
1708}
1709
1710fn counts(lp: &Listpack) -> (u64, u64) {
1712 let mut it = lp.iter();
1713 let count = int_or_zero(it.next());
1714 let deleted = int_or_zero(it.next());
1715 (count.max(0) as u64, deleted.max(0) as u64)
1716}
1717
1718fn bump(lp: &mut Listpack, live: i64, dead: i64) -> bool {
1728 let was = lp.byte_len();
1729 let (count, deleted) = counts(lp);
1730 let mut buf = [0u8; DIGITS_MAX];
1731 let at = count as i64 + live;
1732 lp.replace(0, u64_digits(&mut buf, at.max(0) as u64));
1733 let at = deleted as i64 + dead;
1734 lp.replace(1, u64_digits(&mut buf, at.max(0) as u64));
1735 lp.byte_len() != was
1736}
1737
1738fn int_or_zero(entry: Option<Entry<'_>>) -> i64 {
1740 match entry {
1741 Some(Entry::Int(n)) => n,
1742 _ => 0,
1743 }
1744}
1745
1746fn push_int(lp: &mut Listpack, n: i64) {
1748 let mut buf = [0u8; DIGITS_MAX];
1749 lp.push(i64_digits(&mut buf, n));
1750}
1751
1752fn set_int(lp: &mut Listpack, index: usize, n: i64) {
1754 let mut buf = [0u8; DIGITS_MAX];
1755 lp.replace(index, i64_digits(&mut buf, n));
1756}
1757
1758fn write_entry(lp: &mut Listpack, master: Id, id: Id, fields: &[(&[u8], &[u8])], same: bool) {
1760 let flags = if same { LIVE | SAME_FIELDS } else { LIVE };
1761 push_int(lp, flags);
1762 push_int(lp, id.ms.wrapping_sub(master.ms) as i64);
1767 push_int(lp, id.seq.wrapping_sub(master.seq) as i64);
1768 if same {
1769 for (_, value) in fields {
1770 lp.push(value);
1771 }
1772 push_int(lp, fields.len() as i64 + 3);
1773 } else {
1774 push_int(lp, fields.len() as i64);
1775 for (name, value) in fields {
1776 lp.push(name);
1777 lp.push(value);
1778 }
1779 push_int(lp, fields.len() as i64 * 2 + 4);
1780 }
1781}
1782
1783fn last_of(node: &Node) -> Id {
1785 let mut last = node.master;
1786 each(&node.lp, node.master, None, &mut |id, _, _| {
1787 last = id;
1788 true
1789 });
1790 last
1791}
1792
1793fn find(lp: &Listpack, master: Id, id: Id) -> Option<(usize, i64)> {
1798 let mut at = None;
1799 walk_node(lp, master, None, &mut |mark, flags, _| {
1800 if mark.id == id {
1801 if let Some(index) = mark.index {
1805 at = Some((index, flags));
1806 }
1807 return false;
1808 }
1809 mark.id < id
1810 });
1811 at
1812}
1813
1814fn each<'a, F>(lp: &'a Listpack, master: Id, from: Option<usize>, f: &mut F) -> usize
1817where
1818 F: FnMut(Id, usize, Fields<'a>) -> bool,
1819{
1820 walk_node(lp, master, from, &mut |mark, flags, fields| {
1821 if flags & DELETED != 0 {
1822 return true;
1823 }
1824 f(mark.id, mark.byte, fields)
1825 })
1826}
1827
1828fn walk_node<'a, F>(lp: &'a Listpack, master: Id, from: Option<usize>, f: &mut F) -> usize
1834where
1835 F: FnMut(Mark, i64, Fields<'a>) -> bool,
1836{
1837 let mut it = lp.iter();
1838 let (Some(_), Some(_), Some(Entry::Int(masters))) = (it.next(), it.next(), it.next()) else {
1839 return it.offset();
1840 };
1841 let masters = masters.max(0) as usize;
1842 let names = it.clone();
1845 let mut index = MASTER_FIELDS;
1846 for _ in 0..=masters {
1847 if it.next().is_none() {
1848 return it.offset();
1849 }
1850 index += 1;
1851 }
1852 let counting = from.is_none();
1857 if let Some(byte) = from.filter(|byte| *byte > it.offset()) {
1858 it = lp.iter_at(byte);
1859 }
1860
1861 loop {
1862 let at = it.offset();
1863 let element = index;
1864 let (Some(Entry::Int(flags)), Some(Entry::Int(ms)), Some(Entry::Int(seq))) =
1865 (it.next(), it.next(), it.next())
1866 else {
1867 return at;
1868 };
1869 index += 3;
1870 let id = Id {
1871 ms: master.ms.wrapping_add(ms as u64),
1872 seq: master.seq.wrapping_add(seq as u64),
1873 };
1874
1875 let same = flags & SAME_FIELDS != 0;
1876 let (fields, skip) = if same {
1877 let fields = Fields {
1878 names: Some(names.clone()),
1879 body: it.clone(),
1880 left: masters,
1881 };
1882 (fields, masters + 1)
1883 } else {
1884 let Some(Entry::Int(n)) = it.next() else {
1885 return at;
1886 };
1887 index += 1;
1888 let fields = Fields {
1889 names: None,
1890 body: it.clone(),
1891 left: n.max(0) as usize,
1892 };
1893 (fields, n.max(0) as usize * 2 + 1)
1894 };
1895
1896 let mark = Mark {
1897 id,
1898 index: counting.then_some(element),
1899 byte: at,
1900 };
1901 if !f(mark, flags, fields) {
1902 return at;
1903 }
1904 for _ in 0..skip {
1905 if it.next().is_none() {
1906 return it.offset();
1907 }
1908 index += 1;
1909 }
1910 }
1911}
1912
1913#[derive(Debug, Clone, Copy)]
1915struct Mark {
1916 id: Id,
1919 index: Option<usize>,
1924 byte: usize,
1927}
1928
1929#[cfg(test)]
1930mod tests {
1931 use super::*;
1932 use crate::many;
1933
1934 fn pairs<'a>(of: &'a [(&'a str, &'a str)]) -> Vec<(&'a [u8], &'a [u8])> {
1935 of.iter()
1936 .map(|(f, v)| (f.as_bytes(), v.as_bytes()))
1937 .collect()
1938 }
1939
1940 type Flat = (Id, Vec<(Vec<u8>, Vec<u8>)>);
1942
1943 fn dump(s: &Stream) -> Vec<Flat> {
1945 let mut out = Vec::new();
1946 s.range(Id::MIN, Id::MAX, None, |id, fields| {
1947 out.push((id, fields.map(|(f, v)| (f.to_vec(), v.to_vec())).collect()));
1948 true
1949 });
1950 out
1951 }
1952
1953 fn add(s: &mut Stream, ms: u64, seq: u64, fields: &[(&str, &str)]) {
1954 s.append(Id::new(ms, seq), &pairs(fields), Limits::default())
1955 .expect("an append");
1956 }
1957
1958 #[test]
1959 fn an_entry_comes_back_as_it_went_in() {
1960 let mut s = Stream::new();
1961 add(&mut s, 5, 0, &[("sensor", "1"), ("reading", "23.4")]);
1962 let got = dump(&s);
1963 assert_eq!(got.len(), 1);
1964 assert_eq!(got[0].0, Id::new(5, 0));
1965 assert_eq!(
1966 got[0].1,
1967 vec![
1968 (b"sensor".to_vec(), b"1".to_vec()),
1969 (b"reading".to_vec(), b"23.4".to_vec())
1970 ]
1971 );
1972 }
1973
1974 #[test]
1975 fn entries_come_back_in_order() {
1976 let mut s = Stream::new();
1977 for ms in 1..200u64 {
1978 add(&mut s, ms, 0, &[("n", "x")]);
1979 }
1980 let got = dump(&s);
1981 assert_eq!(got.len(), 199);
1982 for (at, (id, _)) in got.iter().enumerate() {
1983 assert_eq!(*id, Id::new(at as u64 + 1, 0));
1984 }
1985 assert_eq!(s.len(), 199);
1986 assert_eq!(s.added(), 199);
1987 assert_eq!(s.last_id(), Id::new(199, 0));
1988 }
1989
1990 #[test]
1992 fn a_long_stream_is_many_nodes() {
1993 let mut s = Stream::new();
1994 for ms in 1..=1000u64 {
1995 add(&mut s, ms, 0, &[("n", "x")]);
1996 }
1997 assert_eq!(s.nodes(), 10, "a hundred entries a node");
1998 assert_eq!(dump(&s).len(), 1000);
1999 }
2000
2001 #[test]
2003 fn the_field_names_are_not_stored_twice() {
2004 let mut shared = Stream::new();
2005 let mut apart = Stream::new();
2006 for ms in 1..=100u64 {
2007 add(
2008 &mut shared,
2009 ms,
2010 0,
2011 &[("temperature_celsius", "21"), ("relative_humidity", "44")],
2012 );
2013 let a = format!("temperature_celsius{ms}");
2014 let b = format!("relative_humidity{ms}");
2015 apart
2016 .append(
2017 Id::new(ms, 0),
2018 &[(a.as_bytes(), b"21"), (b.as_bytes(), b"44")],
2019 Limits::default(),
2020 )
2021 .expect("an append");
2022 }
2023 assert!(
2024 shared.memory_bytes() * 3 < apart.memory_bytes(),
2025 "{} against {}",
2026 shared.memory_bytes(),
2027 apart.memory_bytes()
2028 );
2029 }
2030
2031 #[test]
2041 fn an_entry_costs_about_two_dozen_bytes() {
2042 let mut s = Stream::new();
2043 for ms in 1..=10_000u64 {
2044 let reading = format!("{:.3}", ms as f64 / 7.0);
2045 s.append(
2046 Id::new(ms, 0),
2047 &[(b"sensor", b"a4"), (b"reading", reading.as_bytes())],
2048 Limits::default(),
2049 )
2050 .expect("an append");
2051 }
2052 let each = s.memory_bytes() as f64 / 10_000.0;
2053 assert!(each < 32.0, "{each:.2} bytes an entry");
2054 }
2055
2056 #[test]
2057 fn an_entry_with_its_own_fields_still_reads_back() {
2058 let mut s = Stream::new();
2059 add(&mut s, 1, 0, &[("a", "1"), ("b", "2")]);
2060 add(&mut s, 2, 0, &[("c", "3")]);
2061 add(&mut s, 3, 0, &[("a", "4"), ("b", "5")]);
2062 let got = dump(&s);
2063 assert_eq!(got[1].1, vec![(b"c".to_vec(), b"3".to_vec())]);
2064 assert_eq!(
2065 got[2].1,
2066 vec![
2067 (b"a".to_vec(), b"4".to_vec()),
2068 (b"b".to_vec(), b"5".to_vec())
2069 ]
2070 );
2071 }
2072
2073 #[test]
2075 fn the_order_of_the_names_matters() {
2076 let mut s = Stream::new();
2077 add(&mut s, 1, 0, &[("a", "1"), ("b", "2")]);
2078 add(&mut s, 2, 0, &[("b", "3"), ("a", "4")]);
2079 let got = dump(&s);
2080 assert_eq!(
2081 got[1].1,
2082 vec![
2083 (b"b".to_vec(), b"3".to_vec()),
2084 (b"a".to_vec(), b"4".to_vec())
2085 ]
2086 );
2087 }
2088
2089 #[test]
2090 fn an_id_must_beat_the_last_one() {
2091 let mut s = Stream::new();
2092 add(&mut s, 5, 5, &[("n", "x")]);
2093 let f = pairs(&[("n", "x")]);
2094 for id in [Id::new(5, 5), Id::new(5, 4), Id::new(1, 0)] {
2095 assert_eq!(
2096 s.append(id, &f, Limits::default()),
2097 Err(Refused::NotGreater),
2098 "{id:?}"
2099 );
2100 }
2101 assert_eq!(s.append(Id::new(5, 6), &f, Limits::default()), Ok(()));
2102 }
2103
2104 #[test]
2105 fn nothing_can_be_added_at_zero() {
2106 let mut s = Stream::new();
2107 assert_eq!(
2108 s.append(Id::MIN, &pairs(&[("n", "x")]), Limits::default()),
2109 Err(Refused::Zero)
2110 );
2111 }
2112
2113 #[test]
2114 fn a_range_takes_both_ends() {
2115 let mut s = Stream::new();
2116 for ms in 1..=10u64 {
2117 add(&mut s, ms, 0, &[("n", "x")]);
2118 }
2119 let mut seen = Vec::new();
2120 s.range(Id::new(3, 0), Id::new(6, 0), None, |id, _| {
2121 seen.push(id.ms);
2122 true
2123 });
2124 assert_eq!(seen, vec![3, 4, 5, 6]);
2125 }
2126
2127 #[test]
2129 fn a_range_that_lands_between_entries() {
2130 let mut s = Stream::new();
2131 for ms in [10u64, 20, 30] {
2132 add(&mut s, ms, 0, &[("n", "x")]);
2133 }
2134 let mut seen = Vec::new();
2135 s.range(Id::new(11, 0), Id::new(29, 0), None, |id, _| {
2136 seen.push(id.ms);
2137 true
2138 });
2139 assert_eq!(seen, vec![20]);
2140
2141 let mut none = 0;
2142 s.range(Id::new(31, 0), Id::MAX, None, |_, _| {
2143 none += 1;
2144 true
2145 });
2146 assert_eq!(none, 0);
2147 }
2148
2149 #[test]
2150 fn a_count_stops_the_walk() {
2151 let mut s = Stream::new();
2152 for ms in 1..=500u64 {
2153 add(&mut s, ms, 0, &[("n", "x")]);
2154 }
2155 let mut seen = 0;
2156 let answered = s.range(Id::MIN, Id::MAX, Some(7), |_, _| {
2157 seen += 1;
2158 true
2159 });
2160 assert_eq!((seen, answered), (7, 7));
2161 }
2162
2163 #[test]
2164 fn the_callback_can_stop_the_walk() {
2165 let mut s = Stream::new();
2166 for ms in 1..=500u64 {
2167 add(&mut s, ms, 0, &[("n", "x")]);
2168 }
2169 let mut seen = 0;
2170 s.range(Id::MIN, Id::MAX, None, |_, _| {
2171 seen += 1;
2172 seen < 3
2173 });
2174 assert_eq!(seen, 3);
2175 }
2176
2177 #[test]
2178 fn a_reverse_range_is_the_forward_one_backwards() {
2179 let mut s = Stream::new();
2180 for ms in 1..=350u64 {
2181 add(&mut s, ms, 0, &[("n", "x")]);
2182 }
2183 let mut forward = Vec::new();
2184 s.range(Id::new(50, 0), Id::new(300, 0), None, |id, _| {
2185 forward.push(id);
2186 true
2187 });
2188 let mut back = Vec::new();
2189 s.rev_range(Id::new(50, 0), Id::new(300, 0), None, |id, _| {
2190 back.push(id);
2191 true
2192 });
2193 back.reverse();
2194 assert_eq!(forward, back);
2195 assert_eq!(forward.len(), 251);
2196 }
2197
2198 #[test]
2199 fn a_reverse_range_takes_a_count_from_the_new_end() {
2200 let mut s = Stream::new();
2201 for ms in 1..=350u64 {
2202 add(&mut s, ms, 0, &[("n", "x")]);
2203 }
2204 let mut seen = Vec::new();
2205 s.rev_range(Id::MIN, Id::MAX, Some(3), |id, _| {
2206 seen.push(id.ms);
2207 true
2208 });
2209 assert_eq!(seen, vec![350, 349, 348]);
2210 }
2211
2212 #[test]
2213 fn deleting_leaves_the_rest_readable() {
2214 let mut s = Stream::new();
2215 for ms in 1..=10u64 {
2216 add(&mut s, ms, 0, &[("n", "x")]);
2217 }
2218 assert!(s.delete(Id::new(4, 0)));
2219 assert!(!s.delete(Id::new(4, 0)), "twice is not twice");
2220 assert!(!s.delete(Id::new(99, 0)));
2221 assert_eq!(s.len(), 9);
2222 assert_eq!(s.max_deleted_id(), Id::new(4, 0));
2223 let seen: Vec<u64> = dump(&s).iter().map(|(id, _)| id.ms).collect();
2224 assert_eq!(seen, vec![1, 2, 3, 5, 6, 7, 8, 9, 10]);
2225 }
2226
2227 #[test]
2228 fn deleting_the_first_entry_moves_the_first_id() {
2229 let mut s = Stream::new();
2230 for ms in 1..=5u64 {
2231 add(&mut s, ms, 0, &[("n", "x")]);
2232 }
2233 assert_eq!(s.first_id(), Some(Id::new(1, 0)));
2234 s.delete(Id::new(1, 0));
2235 assert_eq!(s.first_id(), Some(Id::new(2, 0)));
2236 }
2237
2238 #[test]
2240 fn emptying_a_node_drops_it() {
2241 let mut s = Stream::new();
2242 for ms in 1..=250u64 {
2243 add(&mut s, ms, 0, &[("n", "x")]);
2244 }
2245 assert_eq!(s.nodes(), 3);
2246 for ms in 1..=100u64 {
2247 assert!(s.delete(Id::new(ms, 0)), "{ms}");
2248 }
2249 assert_eq!(s.nodes(), 2);
2250 assert_eq!(s.len(), 150);
2251 assert_eq!(dump(&s).len(), 150);
2252 }
2253
2254 #[test]
2256 fn an_emptied_stream_still_remembers_its_last_id() {
2257 let mut s = Stream::new();
2258 add(&mut s, 7, 0, &[("n", "x")]);
2259 s.delete(Id::new(7, 0));
2260 assert!(s.is_empty());
2261 assert_eq!(s.last_id(), Id::new(7, 0));
2262 assert_eq!(s.added(), 1);
2263 assert_eq!(s.first_id(), None);
2264 assert_eq!(
2265 s.append(Id::new(7, 0), &pairs(&[("n", "x")]), Limits::default()),
2266 Err(Refused::NotGreater),
2267 "a deleted id is still used up"
2268 );
2269 }
2270
2271 #[test]
2272 fn trimming_to_a_length_takes_the_oldest() {
2273 let mut s = Stream::new();
2274 for ms in 1..=1000u64 {
2275 add(&mut s, ms, 0, &[("n", "x")]);
2276 }
2277 assert_eq!(s.trim_maxlen(150, true, None), 850);
2278 assert_eq!(s.len(), 150);
2279 assert_eq!(s.first_id(), Some(Id::new(851, 0)));
2280 assert_eq!(s.last_id(), Id::new(1000, 0));
2281 }
2282
2283 #[test]
2285 fn an_approximate_trim_stops_at_a_node() {
2286 let mut s = Stream::new();
2287 for ms in 1..=1000u64 {
2288 add(&mut s, ms, 0, &[("n", "x")]);
2289 }
2290 assert_eq!(s.trim_maxlen(150, false, None), 800);
2291 assert_eq!(s.len(), 200, "left at the node boundary above 150");
2292 assert_eq!(s.nodes(), 2);
2293 }
2294
2295 #[test]
2296 fn trimming_to_a_length_that_is_already_met_does_nothing() {
2297 let mut s = Stream::new();
2298 for ms in 1..=10u64 {
2299 add(&mut s, ms, 0, &[("n", "x")]);
2300 }
2301 assert_eq!(s.trim_maxlen(50, true, None), 0);
2302 assert_eq!(s.len(), 10);
2303 }
2304
2305 #[test]
2306 fn trimming_to_zero_empties_it() {
2307 let mut s = Stream::new();
2308 for ms in 1..=250u64 {
2309 add(&mut s, ms, 0, &[("n", "x")]);
2310 }
2311 assert_eq!(s.trim_maxlen(0, true, None), 250);
2312 assert!(s.is_empty());
2313 assert_eq!(s.nodes(), 0);
2314 assert_eq!(s.last_id(), Id::new(250, 0));
2315 }
2316
2317 #[test]
2318 fn trimming_below_an_id_takes_everything_under_it() {
2319 let mut s = Stream::new();
2320 for ms in 1..=1000u64 {
2321 add(&mut s, ms, 0, &[("n", "x")]);
2322 }
2323 assert_eq!(s.trim_minid(Id::new(400, 0), true, None), 399);
2324 assert_eq!(s.first_id(), Some(Id::new(400, 0)));
2325 assert_eq!(s.len(), 601);
2326 }
2327
2328 #[test]
2329 fn an_approximate_minid_trim_stops_at_a_node() {
2330 let mut s = Stream::new();
2331 for ms in 1..=1000u64 {
2332 add(&mut s, ms, 0, &[("n", "x")]);
2333 }
2334 assert_eq!(s.trim_minid(Id::new(450, 0), false, None), 400);
2335 assert_eq!(s.first_id(), Some(Id::new(401, 0)));
2336 }
2337
2338 #[test]
2339 fn several_entries_share_a_millisecond() {
2340 let mut s = Stream::new();
2341 for seq in 0..250u64 {
2342 add(&mut s, 5, seq, &[("n", "x")]);
2343 }
2344 let got = dump(&s);
2345 assert_eq!(got.len(), 250);
2346 for (at, (id, _)) in got.iter().enumerate() {
2347 assert_eq!(*id, Id::new(5, at as u64), "at {at}");
2348 }
2349 let mut seen = Vec::new();
2350 s.range(Id::new(5, 100), Id::new(5, 102), None, |id, _| {
2351 seen.push(id.seq);
2352 true
2353 });
2354 assert_eq!(seen, vec![100, 101, 102]);
2355 }
2356
2357 #[test]
2360 fn a_sequence_that_crosses_a_node() {
2361 let mut s = Stream::new();
2362 for seq in 0..300u64 {
2363 add(&mut s, 1, seq, &[("n", "x")]);
2364 }
2365 assert!(s.nodes() > 1);
2366 let got = dump(&s);
2367 assert_eq!(got.len(), 300);
2368 assert_eq!(got[299].0, Id::new(1, 299));
2369 assert!(s.delete(Id::new(1, 250)));
2370 assert_eq!(dump(&s).len(), 299);
2371 }
2372
2373 #[test]
2374 fn an_entry_with_no_fields_is_still_an_entry() {
2375 let mut s = Stream::new();
2376 s.append(Id::new(1, 0), &[], Limits::default())
2377 .expect("an append");
2378 add(&mut s, 2, 0, &[("n", "x")]);
2379 let got = dump(&s);
2380 assert_eq!(got.len(), 2);
2381 assert!(got[0].1.is_empty());
2382 }
2383
2384 #[test]
2385 fn a_value_that_looks_like_a_number_comes_back_as_it_went_in() {
2386 let mut s = Stream::new();
2387 add(&mut s, 1, 0, &[("n", "007"), ("m", "7")]);
2388 let got = dump(&s);
2389 assert_eq!(got[0].1[0].1, b"007".to_vec());
2390 assert_eq!(got[0].1[1].1, b"7".to_vec());
2391 }
2392
2393 #[test]
2394 fn the_auto_id_follows_the_clock_and_never_goes_back() {
2395 let mut s = Stream::new();
2396 assert_eq!(s.auto_id(1000), Some(Id::new(1000, 0)));
2397 add(&mut s, 1000, 0, &[("n", "x")]);
2398 assert_eq!(s.auto_id(1000), Some(Id::new(1000, 1)), "same millisecond");
2399 assert_eq!(s.auto_id(900), Some(Id::new(1000, 1)), "clock went back");
2400 assert_eq!(s.auto_id(1001), Some(Id::new(1001, 0)));
2401 }
2402
2403 #[test]
2404 fn an_explicit_millisecond_takes_the_next_sequence() {
2405 let mut s = Stream::new();
2406 add(&mut s, 5, 0, &[("n", "x")]);
2407 assert_eq!(s.auto_seq(5), Some(Id::new(5, 1)));
2408 assert_eq!(s.auto_seq(6), Some(Id::new(6, 0)));
2409 assert_eq!(s.auto_seq(4), None, "below the last one");
2410 }
2411
2412 #[test]
2413 fn an_id_reads_and_writes() {
2414 for (text, default, want) in [
2415 (&b"5"[..], 0, Some(Id::new(5, 0))),
2416 (b"5", u64::MAX, Some(Id::new(5, u64::MAX))),
2417 (b"5-3", 0, Some(Id::new(5, 3))),
2418 (b"0-0", 0, Some(Id::MIN)),
2419 (b"", 0, None),
2420 (b"-1", 0, None),
2421 (b"5-", 0, None),
2422 (b"a", 0, None),
2423 (b"5-a", 0, None),
2424 (b"18446744073709551616", 0, None),
2425 ] {
2426 assert_eq!(
2427 Id::parse(text, default),
2428 want,
2429 "{:?}",
2430 String::from_utf8_lossy(text)
2431 );
2432 }
2433 assert_eq!(Id::new(5, 3).to_vec(), b"5-3".to_vec());
2434 }
2435
2436 #[test]
2437 fn an_id_round_trips_through_its_bytes() {
2438 for id in [Id::MIN, Id::MAX, Id::new(1, 2), Id::new(u64::MAX, 0)] {
2439 assert_eq!(Id::from_bytes(id.to_bytes()), id);
2440 }
2441 assert!(Id::new(1, 2).to_bytes() < Id::new(1, 3).to_bytes());
2443 assert!(Id::new(1, u64::MAX).to_bytes() < Id::new(2, 0).to_bytes());
2444 }
2445
2446 #[test]
2447 fn stepping_an_id_carries_and_stops() {
2448 assert_eq!(Id::new(1, 2).next(), Some(Id::new(1, 3)));
2449 assert_eq!(Id::new(1, u64::MAX).next(), Some(Id::new(2, 0)));
2450 assert_eq!(Id::MAX.next(), None);
2451 assert_eq!(Id::new(1, 3).prev(), Some(Id::new(1, 2)));
2452 assert_eq!(Id::new(2, 0).prev(), Some(Id::new(1, u64::MAX)));
2453 assert_eq!(Id::MIN.prev(), None);
2454 }
2455
2456 #[test]
2458 fn the_node_size_changes_nothing_but_the_node_count() {
2459 let last = many(400u64);
2462 let mut want = None;
2463 for entries in [1usize, 2, 7, 100, 4096] {
2464 let mut s = Stream::new();
2465 let limits = Limits {
2466 max_node_bytes: NODE_BYTES,
2467 max_node_entries: entries,
2468 };
2469 for ms in 1..=last {
2470 let value = format!("v{ms}");
2471 s.append(Id::new(ms, 0), &[(b"n", value.as_bytes())], limits)
2472 .expect("an append");
2473 }
2474 for ms in (1..=last).step_by(7) {
2475 s.delete(Id::new(ms, 0));
2476 }
2477 let got = dump(&s);
2478 match &want {
2479 None => want = Some(got),
2480 Some(want) => assert_eq!(&got, want, "at {entries} entries a node"),
2481 }
2482 }
2483 }
2484
2485 #[test]
2487 fn a_tiny_byte_limit_still_works() {
2488 let limits = Limits {
2489 max_node_bytes: 1,
2490 max_node_entries: NODE_ENTRIES,
2491 };
2492 let mut s = Stream::new();
2493 for ms in 1..=20u64 {
2494 s.append(Id::new(ms, 0), &[(b"n", b"x")], limits)
2495 .expect("an append");
2496 }
2497 assert_eq!(s.nodes(), 20);
2498 assert_eq!(dump(&s).len(), 20);
2499 }
2500
2501 const REDIS_DUMP: &str = "1b0110000000000000000100000000000000014070\
2518 700000001f000301010102019374656d70657261747572655f63656c736975731491\
2519 72656c61746976655f68756d69646974791200010201000100011501370105010301\
2520 0001010116013801050102010101000117013901050100010201000101018673656e\
2521 736f72078161020601ff030301010101020400406440640000000f00239a5c2c7208\
2522 ea0a";
2523
2524 fn unhex(s: &str) -> Vec<u8> {
2525 let digits: Vec<u8> = s.bytes().filter(|b| !b.is_ascii_whitespace()).collect();
2526 digits
2527 .chunks(2)
2528 .map(|pair| {
2529 let of = |b: u8| (b as char).to_digit(16).expect("a hex digit") as u8;
2530 of(pair[0]) << 4 | of(pair[1])
2531 })
2532 .collect()
2533 }
2534
2535 fn rdb_len(bytes: &[u8], at: usize) -> (usize, usize) {
2541 match bytes[at] >> 6 {
2542 0 => (usize::from(bytes[at] & 0x3F), 1),
2543 1 => (
2544 usize::from(bytes[at] & 0x3F) << 8 | usize::from(bytes[at + 1]),
2545 2,
2546 ),
2547 other => panic!("the fixture used length form {other}"),
2548 }
2549 }
2550
2551 fn redis_node() -> (Id, Vec<u8>, Vec<u8>) {
2553 let dump = unhex(REDIS_DUMP);
2554 assert_eq!(dump[0], 0x1B, "RDB_TYPE_STREAM_LISTPACKS_3");
2555 let (nodes, n) = rdb_len(&dump, 1);
2556 assert_eq!(nodes, 1, "the fixture is one node");
2557 let mut at = 1 + n;
2558 let (key, n) = rdb_len(&dump, at);
2559 assert_eq!(key, 16, "a node key is an id in sixteen bytes");
2560 at += n;
2561 let master = Id::from_bytes(dump[at..at + 16].try_into().expect("sixteen bytes"));
2562 at += 16;
2563 let (len, n) = rdb_len(&dump, at);
2564 at += n;
2565 let lp = dump[at..at + len].to_vec();
2566 let rest = dump[at + len..dump.len() - 10].to_vec();
2569 (master, lp, rest)
2570 }
2571
2572 #[test]
2573 fn a_node_is_written_the_way_redis_writes_one() {
2574 let mut s = Stream::new();
2575 add(
2576 &mut s,
2577 1,
2578 1,
2579 &[("temperature_celsius", "21"), ("relative_humidity", "55")],
2580 );
2581 add(
2582 &mut s,
2583 1,
2584 2,
2585 &[("temperature_celsius", "22"), ("relative_humidity", "56")],
2586 );
2587 add(
2588 &mut s,
2589 2,
2590 1,
2591 &[("temperature_celsius", "23"), ("relative_humidity", "57")],
2592 );
2593 add(&mut s, 3, 1, &[("sensor", "a")]);
2594 assert!(s.delete(Id::new(1, 2)));
2595
2596 let (master, lp, _) = redis_node();
2597 assert_eq!(s.nodes(), 1, "all four fit in one node");
2598 assert_eq!(s.nodes[0].master, master);
2599 assert_eq!(s.nodes[0].lp.as_bytes(), &lp[..]);
2600 }
2601
2602 #[test]
2603 fn a_node_redis_wrote_reads_back() {
2604 let (master, lp, rest) = redis_node();
2605 let lp = Listpack::from_bytes(&lp).expect("a listpack Redis wrote");
2606 let s = Stream {
2607 nodes: VecDeque::from(vec![Node { master, lp }]),
2608 length: u64::from(rest[0]),
2609 last: Id::new(u64::from(rest[1]), u64::from(rest[2])),
2610 max_deleted: Id::new(u64::from(rest[5]), u64::from(rest[6])),
2611 added: u64::from(rest[7]),
2612 groups: Vec::new(),
2613 epoch: 0,
2614 };
2615
2616 assert_eq!(s.len(), 3);
2617 assert_eq!(s.last_id(), Id::new(3, 1));
2618 assert_eq!(s.max_deleted_id(), Id::new(1, 2));
2619 assert_eq!(s.added(), 4);
2620 assert_eq!(s.first_id(), Some(Id::new(1, 1)));
2621 assert_eq!(
2622 dump(&s),
2623 vec![
2624 (
2625 Id::new(1, 1),
2626 vec![
2627 (b"temperature_celsius".to_vec(), b"21".to_vec()),
2628 (b"relative_humidity".to_vec(), b"55".to_vec())
2629 ]
2630 ),
2631 (
2632 Id::new(2, 1),
2633 vec![
2634 (b"temperature_celsius".to_vec(), b"23".to_vec()),
2635 (b"relative_humidity".to_vec(), b"57".to_vec())
2636 ]
2637 ),
2638 (Id::new(3, 1), vec![(b"sensor".to_vec(), b"a".to_vec())]),
2639 ]
2640 );
2641 }
2642
2643 fn logged(n: u64) -> Stream {
2645 let mut s = Stream::new();
2646 for ms in 1..=n {
2647 add(&mut s, ms, 0, &[("job", "x")]);
2648 }
2649 s
2650 }
2651
2652 fn read(s: &mut Stream, group: &str, who: &str, count: Option<usize>, now: u64) -> Vec<Id> {
2654 let mut out = Vec::new();
2655 s.read_group(
2656 group.as_bytes(),
2657 who.as_bytes(),
2658 count,
2659 false,
2660 now,
2661 |id, _| {
2662 out.push(id);
2663 true
2664 },
2665 )
2666 .expect("the group");
2667 out
2668 }
2669
2670 #[test]
2675 fn a_group_draining_one_at_a_time_gets_every_entry_once() {
2676 let mut s = logged(500);
2677 s.create_group(b"workers", Id::MIN, Some(0));
2678 let mut got = Vec::new();
2679 for _ in 0..500 {
2680 got.extend(read(&mut s, "workers", "alice", Some(1), 100));
2681 }
2682 assert_eq!(got, (1..=500).map(|ms| Id::new(ms, 0)).collect::<Vec<_>>());
2683 assert!(read(&mut s, "workers", "alice", Some(1), 100).is_empty());
2684 }
2685
2686 #[test]
2689 fn a_group_that_has_caught_up_is_handed_the_next_append() {
2690 let mut s = logged(3);
2691 s.create_group(b"workers", Id::MIN, Some(0));
2692 assert_eq!(read(&mut s, "workers", "alice", None, 100).len(), 3);
2693 for ms in 4..=200u64 {
2694 add(&mut s, ms, 0, &[("job", "x")]);
2695 assert_eq!(
2696 read(&mut s, "workers", "alice", Some(1), 100),
2697 vec![Id::new(ms, 0)],
2698 "the entry appended a moment ago"
2699 );
2700 }
2701 }
2702
2703 #[test]
2706 fn a_delete_between_two_group_reads_does_not_lose_the_rest() {
2707 let mut s = logged(300);
2708 s.create_group(b"workers", Id::MIN, Some(0));
2709 let mut got = read(&mut s, "workers", "alice", Some(10), 100);
2710 assert!(s.delete(Id::new(150, 0)));
2711 while got.len() < 299 {
2712 let more = read(&mut s, "workers", "alice", Some(1), 100);
2713 assert_eq!(more.len(), 1, "at {}", got.len());
2714 got.extend(more);
2715 }
2716 let want: Vec<Id> = (1..=300)
2717 .filter(|ms| *ms != 150)
2718 .map(|ms| Id::new(ms, 0))
2719 .collect();
2720 assert_eq!(got, want);
2721 }
2722
2723 #[test]
2726 fn moving_the_bookmark_back_reads_it_all_again() {
2727 let mut s = logged(250);
2728 s.create_group(b"workers", Id::MIN, Some(0));
2729 for _ in 0..120 {
2730 read(&mut s, "workers", "alice", Some(1), 100);
2731 }
2732 s.group_mut(b"workers")
2733 .expect("the group")
2734 .set_id(Id::MIN, Some(0));
2735 let mut got = Vec::new();
2736 for _ in 0..250 {
2737 got.extend(read(&mut s, "workers", "alice", Some(1), 100));
2738 }
2739 assert_eq!(got, (1..=250).map(|ms| Id::new(ms, 0)).collect::<Vec<_>>());
2740 }
2741
2742 #[test]
2746 fn a_trim_under_a_group_does_not_hand_out_the_wrong_entries() {
2747 let mut s = logged(500);
2748 s.create_group(b"workers", Id::MIN, Some(0));
2749 for _ in 0..50 {
2750 read(&mut s, "workers", "alice", Some(1), 100);
2751 }
2752 assert_eq!(s.trim_maxlen(200, false, None), 300);
2753 let mut got = Vec::new();
2754 for _ in 0..250 {
2755 got.extend(read(&mut s, "workers", "alice", Some(1), 100));
2756 }
2757 assert_eq!(
2758 got,
2759 (301..=500).map(|ms| Id::new(ms, 0)).collect::<Vec<_>>(),
2760 "everything the trim left that the group had not read"
2761 );
2762 }
2763
2764 #[test]
2765 fn a_group_is_made_once() {
2766 let mut s = logged(3);
2767 assert!(s.create_group(b"workers", Id::MIN, Some(0)));
2768 assert!(!s.create_group(b"workers", Id::MIN, Some(0)));
2769 assert!(s.group(b"workers").is_some());
2770 assert!(s.destroy_group(b"workers"));
2771 assert!(!s.destroy_group(b"workers"));
2772 assert!(s.group(b"workers").is_none());
2773 }
2774
2775 #[test]
2776 fn a_group_read_hands_out_what_comes_after_the_bookmark() {
2777 let mut s = logged(5);
2778 s.create_group(b"workers", Id::MIN, Some(0));
2779
2780 assert_eq!(
2781 read(&mut s, "workers", "alice", Some(2), 100),
2782 vec![Id::new(1, 0), Id::new(2, 0)]
2783 );
2784 assert_eq!(
2786 read(&mut s, "workers", "bob", Some(2), 100),
2787 vec![Id::new(3, 0), Id::new(4, 0)]
2788 );
2789 assert_eq!(
2790 read(&mut s, "workers", "alice", None, 100),
2791 vec![Id::new(5, 0)]
2792 );
2793 assert_eq!(read(&mut s, "workers", "alice", None, 100), vec![]);
2794 }
2795
2796 #[test]
2797 fn a_group_read_fills_the_pending_list() {
2798 let mut s = logged(3);
2799 s.create_group(b"workers", Id::MIN, Some(0));
2800 read(&mut s, "workers", "alice", None, 500);
2801
2802 let g = s.group(b"workers").expect("the group");
2803 assert_eq!(g.pending_len(), 3);
2804 assert_eq!(g.last_id(), Id::new(3, 0));
2805 assert_eq!(g.entries_read(), Some(3));
2806 assert_eq!(s.lag(g), Some(0));
2807 let c = g.consumer_named(b"alice").expect("alice");
2808 assert_eq!(c.len(), 3);
2809 assert_eq!(c.active(), Some(500));
2810 }
2811
2812 #[test]
2813 fn a_read_that_finds_nothing_is_seen_but_not_active() {
2814 let mut s = logged(1);
2815 s.create_group(b"workers", Id::MIN, Some(0));
2816 read(&mut s, "workers", "alice", None, 100);
2817 read(&mut s, "workers", "alice", None, 900);
2818
2819 let c = s
2820 .group(b"workers")
2821 .expect("the group")
2822 .consumer_named(b"alice")
2823 .expect("alice");
2824 assert_eq!((c.seen(), c.active()), (900, Some(100)));
2825 }
2826
2827 #[test]
2828 fn a_group_starting_at_the_end_reads_only_what_comes_next() {
2829 let mut s = logged(3);
2830 s.create_group(b"workers", s.last_id(), Some(s.added()));
2831 assert_eq!(read(&mut s, "workers", "alice", None, 1), vec![]);
2832 add(&mut s, 4, 0, &[("job", "x")]);
2833 assert_eq!(
2834 read(&mut s, "workers", "alice", None, 1),
2835 vec![Id::new(4, 0)]
2836 );
2837 }
2838
2839 #[test]
2840 fn reading_a_group_that_is_not_there_says_so() {
2841 let mut s = logged(1);
2842 assert!(
2843 s.read_group(b"nope", b"alice", None, false, 1, |_, _| true)
2844 .is_none()
2845 );
2846 }
2847
2848 #[test]
2849 fn a_consumer_can_re_read_what_it_is_holding() {
2850 let mut s = logged(4);
2851 s.create_group(b"workers", Id::MIN, Some(0));
2852 read(&mut s, "workers", "alice", Some(2), 1);
2853 read(&mut s, "workers", "bob", Some(2), 1);
2854
2855 let mut out = Vec::new();
2856 s.read_group_pending(b"workers", b"alice", Id::MIN, None, 2, |id, fields| {
2857 out.push((id, fields.map(|f| f.len())));
2858 true
2859 })
2860 .expect("the group");
2861 assert_eq!(
2862 out,
2863 vec![(Id::new(1, 0), Some(1)), (Id::new(2, 0), Some(1))]
2864 );
2865
2866 let mut after = Vec::new();
2868 s.read_group_pending(b"workers", b"alice", Id::new(1, 0), None, 2, |id, _| {
2869 after.push(id);
2870 true
2871 });
2872 assert_eq!(after, vec![Id::new(2, 0)]);
2873 }
2874
2875 #[test]
2884 fn re_reading_counts_as_being_handed_it_again() {
2885 let mut s = logged(2);
2886 s.create_group(b"workers", Id::MIN, Some(0));
2887 read(&mut s, "workers", "alice", None, 100);
2888 s.read_group_pending(
2889 b"workers",
2890 b"alice",
2891 Id::MIN,
2892 None,
2893 700,
2894 |_: Id, _: Option<Fields<'_>>| true,
2895 );
2896
2897 let g = s.group(b"workers").expect("the group");
2898 let nack = g.nack(Id::new(1, 0)).expect("a nack");
2899 assert_eq!((nack.count(), nack.time()), (2, 700));
2900 assert_eq!(g.last_id(), Id::new(2, 0));
2902 assert_eq!(g.pending_len(), 2);
2903 }
2904
2905 #[test]
2912 fn a_hole_in_front_of_a_group_takes_its_lag() {
2913 let mut s = logged(5);
2914 s.create_group(b"workers", Id::MIN, Some(0));
2915 read(&mut s, "workers", "alice", Some(3), 1);
2916 assert_eq!(counters(&s), (Some(3), Some(2)));
2917
2918 assert!(s.delete(Id::new(1, 0)));
2922 assert_eq!(counters(&s), (Some(3), Some(2)));
2923
2924 assert!(s.delete(Id::new(5, 0)));
2927 assert_eq!(counters(&s), (Some(3), None));
2928
2929 read(&mut s, "workers", "alice", None, 1);
2933 assert_eq!(s.group(b"workers").expect("g").last_id(), Id::new(4, 0));
2934 assert_eq!(counters(&s), (None, None));
2935 }
2936
2937 #[test]
2943 fn trimming_past_a_group_leaves_it_the_length() {
2944 let mut s = logged(500);
2945 s.create_group(b"workers", Id::MIN, Some(0));
2946 read(&mut s, "workers", "alice", Some(400), 1);
2947 assert_eq!(counters(&s), (Some(400), Some(100)));
2948
2949 assert_eq!(s.trim_maxlen(200, true, None), 300);
2951 assert_eq!(counters(&s), (Some(400), Some(100)));
2952
2953 assert_eq!(s.trim_maxlen(10, true, None), 190);
2956 assert_eq!(counters(&s), (Some(400), Some(10)));
2957 }
2958
2959 fn counters(s: &Stream) -> (Option<u64>, Option<u64>) {
2961 let g = s.group(b"workers").expect("the group");
2962 (g.entries_read(), s.lag(g))
2963 }
2964
2965 #[test]
2966 fn an_entry_that_went_away_still_comes_back_as_a_hole() {
2967 let mut s = logged(3);
2968 s.create_group(b"workers", Id::MIN, Some(0));
2969 read(&mut s, "workers", "alice", None, 1);
2970 assert!(s.delete(Id::new(2, 0)));
2971
2972 let mut out = Vec::new();
2973 s.read_group_pending(b"workers", b"alice", Id::MIN, None, 2, |id, fields| {
2974 out.push((id, fields.is_some()));
2975 true
2976 });
2977 assert_eq!(
2978 out,
2979 vec![
2980 (Id::new(1, 0), true),
2981 (Id::new(2, 0), false),
2982 (Id::new(3, 0), true)
2983 ]
2984 );
2985 }
2986
2987 #[test]
2988 fn acking_clears_the_pending_list() {
2989 let mut s = logged(3);
2990 s.create_group(b"workers", Id::MIN, Some(0));
2991 read(&mut s, "workers", "alice", None, 1);
2992 let g = s.group_mut(b"workers").expect("the group");
2993 assert!(g.ack(Id::new(2, 0)));
2994 assert_eq!(g.pending_len(), 2);
2995 }
2996
2997 #[test]
2998 fn a_claim_moves_work_off_a_consumer_that_stopped() {
2999 let mut s = logged(2);
3000 s.create_group(b"workers", Id::MIN, Some(0));
3001 read(&mut s, "workers", "alice", None, 100);
3002
3003 let mut gone = Vec::new();
3004 let took = s
3005 .claim(
3006 b"workers",
3007 b"bob",
3008 &[Id::new(1, 0), Id::new(2, 0)],
3009 500,
3010 5_000,
3011 None,
3012 true,
3013 false,
3014 5_000,
3015 &mut gone,
3016 )
3017 .expect("the group");
3018 assert_eq!(took, vec![Id::new(1, 0), Id::new(2, 0)]);
3019 assert!(gone.is_empty());
3020
3021 let g = s.group(b"workers").expect("the group");
3022 assert!(g.consumer_named(b"alice").expect("alice").is_empty());
3023 assert_eq!(g.consumer_named(b"bob").expect("bob").len(), 2);
3024 assert_eq!(g.nack(Id::new(1, 0)).expect("a nack").count(), 2);
3025 }
3026
3027 #[test]
3028 fn a_claim_leaves_work_that_is_not_idle_enough_alone() {
3029 let mut s = logged(1);
3030 s.create_group(b"workers", Id::MIN, Some(0));
3031 read(&mut s, "workers", "alice", None, 100);
3032
3033 let mut gone = Vec::new();
3034 let took = s
3035 .claim(
3036 b"workers",
3037 b"bob",
3038 &[Id::new(1, 0)],
3039 5_000,
3040 200,
3041 None,
3042 true,
3043 false,
3044 200,
3045 &mut gone,
3046 )
3047 .expect("the group");
3048 assert!(took.is_empty());
3049 assert_eq!(
3050 s.group(b"workers")
3051 .expect("the group")
3052 .consumer_named(b"alice")
3053 .expect("alice")
3054 .len(),
3055 1
3056 );
3057 }
3058
3059 #[test]
3060 fn claiming_an_entry_that_went_away_drops_it_instead() {
3061 let mut s = logged(2);
3062 s.create_group(b"workers", Id::MIN, Some(0));
3063 read(&mut s, "workers", "alice", None, 100);
3064 assert!(s.delete(Id::new(1, 0)));
3065
3066 let mut gone = Vec::new();
3067 let took = s
3068 .claim(
3069 b"workers",
3070 b"bob",
3071 &[Id::new(1, 0), Id::new(2, 0)],
3072 0,
3073 5_000,
3074 None,
3075 true,
3076 false,
3077 5_000,
3078 &mut gone,
3079 )
3080 .expect("the group");
3081 assert_eq!(took, vec![Id::new(2, 0)]);
3082 assert_eq!(gone, vec![Id::new(1, 0)]);
3083 assert_eq!(s.group(b"workers").expect("the group").pending_len(), 1);
3084 }
3085
3086 #[test]
3087 fn force_only_works_on_an_entry_that_is_really_there() {
3088 let mut s = logged(2);
3089 s.create_group(b"workers", s.last_id(), Some(2));
3090
3091 let mut gone = Vec::new();
3092 let took = s
3093 .claim(
3094 b"workers",
3095 b"bob",
3096 &[Id::new(1, 0), Id::new(99, 0)],
3097 0,
3098 100,
3099 None,
3100 true,
3101 true,
3102 100,
3103 &mut gone,
3104 )
3105 .expect("the group");
3106 assert_eq!(took, vec![Id::new(1, 0)], "99-0 is not in the stream");
3107 assert_eq!(s.group(b"workers").expect("the group").pending_len(), 1);
3108 }
3109
3110 #[test]
3111 fn autoclaim_sweeps_the_stale_ones_and_says_where_it_stopped() {
3112 let mut s = logged(6);
3113 s.create_group(b"workers", Id::MIN, Some(0));
3114 read(&mut s, "workers", "alice", Some(3), 100);
3115 read(&mut s, "workers", "alice", None, 900);
3116
3117 let mut gone = Vec::new();
3118 let (cursor, took) = s
3119 .autoclaim(
3120 b"workers",
3121 b"bob",
3122 Id::MIN,
3123 500,
3124 100,
3125 true,
3126 1_000,
3127 &mut gone,
3128 )
3129 .expect("the group");
3130 assert_eq!(cursor, None, "the sweep reached the end");
3131 assert_eq!(took, vec![Id::new(1, 0), Id::new(2, 0), Id::new(3, 0)]);
3132 assert_eq!(
3133 s.group(b"workers")
3134 .expect("the group")
3135 .consumer_named(b"bob")
3136 .expect("bob")
3137 .len(),
3138 3
3139 );
3140 }
3141
3142 #[test]
3143 fn autoclaim_hands_back_a_cursor_when_it_hits_the_count() {
3144 let mut s = logged(10);
3145 s.create_group(b"workers", Id::MIN, Some(0));
3146 read(&mut s, "workers", "alice", None, 100);
3147
3148 let mut gone = Vec::new();
3149 let (cursor, took) = s
3150 .autoclaim(b"workers", b"bob", Id::MIN, 0, 4, true, 1_000, &mut gone)
3151 .expect("the group");
3152 assert_eq!(took.len(), 4);
3153 assert_eq!(cursor, Some(Id::new(5, 0)));
3154
3155 let (cursor, took) = s
3157 .autoclaim(
3158 b"workers",
3159 b"bob",
3160 cursor.expect("a cursor"),
3161 0,
3162 100,
3163 true,
3164 1_000,
3165 &mut gone,
3166 )
3167 .expect("the group");
3168 assert_eq!(took.len(), 6);
3169 assert_eq!(cursor, None);
3170 }
3171
3172 #[test]
3173 fn a_group_survives_the_stream_being_trimmed_under_it() {
3174 let mut s = logged(10);
3175 s.create_group(b"workers", Id::MIN, Some(0));
3176 read(&mut s, "workers", "alice", Some(5), 100);
3177 assert_eq!(s.trim_maxlen(3, true, None), 7);
3179
3180 assert_eq!(s.group(b"workers").expect("the group").pending_len(), 5);
3181 let mut gone = Vec::new();
3182 let took = s
3183 .claim(
3184 b"workers",
3185 b"bob",
3186 &(1..=5).map(|ms| Id::new(ms, 0)).collect::<Vec<_>>(),
3187 0,
3188 1_000,
3189 None,
3190 true,
3191 false,
3192 1_000,
3193 &mut gone,
3194 )
3195 .expect("the group");
3196 assert!(took.is_empty(), "none of them are there any more");
3197 assert_eq!(gone.len(), 5);
3198 assert_eq!(s.group(b"workers").expect("the group").pending_len(), 0);
3199 }
3200
3201 fn round_trip(s: &Stream) -> Stream {
3203 let mut bytes = Vec::new();
3204 s.freeze(&mut bytes);
3205 let back = Stream::thaw(&bytes).expect("our own bytes");
3206 assert_eq!(dump(&back), dump(s), "the entries");
3207 assert_eq!(back.len(), s.len(), "the length");
3208 assert_eq!(back.added(), s.added(), "the count added");
3209 assert_eq!(back.last_id(), s.last_id(), "the last ID");
3210 assert_eq!(back.max_deleted_id(), s.max_deleted_id(), "the max deleted");
3211 assert_eq!(back.nodes(), s.nodes(), "the node count");
3212 assert_eq!(back, *s, "the whole thing");
3213 back
3214 }
3215
3216 #[test]
3217 fn a_frozen_stream_comes_back_with_every_entry_it_held() {
3218 let mut s = Stream::new();
3219 for ms in 1..=500u64 {
3220 add(&mut s, ms, 0, &[("job", "x"), ("n", "1")]);
3221 }
3222 add(&mut s, 500, 1, &[("job", "y")]);
3223 assert!(s.nodes() > 1, "more than one node, so the walk is tested");
3224 round_trip(&s);
3225 }
3226
3227 #[test]
3228 fn a_frozen_stream_keeps_the_holes_and_the_counters() {
3229 let mut s = logged(200);
3230 for ms in [3u64, 4, 5, 100, 199] {
3231 assert!(s.delete(Id::new(ms, 0)));
3232 }
3233 s.trim_minid(Id::new(20, 0), true, None);
3234 let back = round_trip(&s);
3235 assert_eq!(back.first_id(), Some(Id::new(20, 0)));
3236 assert!(!back.contains(Id::new(100, 0)), "a hole is still a hole");
3237 assert_eq!(back.max_deleted_id(), Id::new(199, 0));
3238 }
3239
3240 #[test]
3241 fn a_frozen_stream_keeps_its_groups_and_who_is_holding_what() {
3242 let mut s = logged(20);
3243 s.create_group(b"workers", Id::MIN, Some(0));
3244 s.create_group(b"audit", Id::new(5, 0), None);
3245 read(&mut s, "workers", "alice", Some(6), 1_000);
3246 read(&mut s, "workers", "bob", Some(4), 2_000);
3247 assert_eq!(
3250 s.nack(b"workers", Id::new(2, 0), Retry::Keep, true),
3251 Some(true)
3252 );
3253 s.group_mut(b"workers")
3256 .expect("the group")
3257 .create_consumer(b"carol", 3_000);
3258 read(&mut s, "workers", "dave", Some(2), 4_000);
3259 s.group_mut(b"workers")
3260 .expect("the group")
3261 .delete_consumer(b"carol");
3262
3263 let back = round_trip(&s);
3264 let g = back.group(b"workers").expect("the group");
3265 assert_eq!(g.pending_len(), 12);
3266 assert_eq!(g.nacked_len(), 1);
3267 assert_eq!(g.entries_read(), Some(12));
3268 assert_eq!(g.consumer_named(b"alice").expect("alice").len(), 5);
3270 assert_eq!(g.consumer_named(b"bob").expect("bob").len(), 4);
3271 assert_eq!(g.consumer_named(b"dave").expect("dave").len(), 2);
3272 assert_eq!(g.consumer_named(b"carol"), None);
3273 assert_eq!(g.slot(b"dave"), Some(3));
3276 assert_eq!(g.nack(Id::new(1, 0)).expect("a nack").owner(), Some(0));
3277 assert_eq!(g.nack(Id::new(2, 0)).expect("a nack").owner(), None);
3278 assert_eq!(
3279 back.group(b"audit").expect("audit").last_id(),
3280 Id::new(5, 0)
3281 );
3282 assert_eq!(back.group(b"audit").expect("audit").entries_read(), None);
3283 }
3284
3285 #[test]
3286 fn a_stream_that_came_back_still_takes_entries_and_reads_them() {
3287 let mut s = logged(10);
3288 s.create_group(b"workers", Id::MIN, Some(0));
3289 read(&mut s, "workers", "alice", Some(4), 1_000);
3290
3291 let mut back = round_trip(&s);
3292 add(&mut back, 11, 0, &[("job", "new")]);
3293 assert_eq!(back.len(), 11);
3294 assert_eq!(
3295 read(&mut back, "workers", "alice", Some(3), 2_000),
3296 vec![Id::new(5, 0), Id::new(6, 0), Id::new(7, 0)]
3297 );
3298 assert!(
3299 back.group_mut(b"workers")
3300 .expect("the group")
3301 .ack(Id::new(1, 0))
3302 );
3303 assert_eq!(back.group(b"workers").expect("the group").pending_len(), 6);
3304 }
3305
3306 #[test]
3307 fn an_empty_stream_that_still_exists_comes_back() {
3308 let mut s = logged(3);
3309 for ms in 1..=3u64 {
3310 assert!(s.delete(Id::new(ms, 0)));
3311 }
3312 assert_eq!(s.nodes(), 0, "the last node went with the last entry");
3313 let back = round_trip(&s);
3314 assert!(back.is_empty());
3315 assert_eq!(back.last_id(), Id::new(3, 0));
3318 round_trip(&Stream::new());
3319 }
3320
3321 #[test]
3322 fn a_frozen_stream_that_arrives_damaged_is_an_error_and_not_a_panic() {
3323 let mut s = logged(8);
3324 s.create_group(b"workers", Id::MIN, Some(0));
3325 read(&mut s, "workers", "alice", Some(3), 1_000);
3326 let mut bytes = Vec::new();
3327 s.freeze(&mut bytes);
3328
3329 for cut in 0..bytes.len() {
3330 assert!(Stream::thaw(&bytes[..cut]).is_err(), "cut at {cut}");
3331 }
3332 for at in 0..bytes.len().min(40) {
3335 for bit in 0..8 {
3336 let mut bad = bytes.clone();
3337 bad[at] ^= 1 << bit;
3338 let _ = Stream::thaw(&bad);
3342 }
3343 }
3344 assert_eq!(Stream::thaw(&[]), Err(Broken::Short));
3345 assert_eq!(Stream::thaw(&[9]), Err(Broken::Form));
3346 }
3347
3348 #[test]
3349 fn nothing_at_all() {
3350 let mut s = Stream::new();
3351 assert!(s.is_empty());
3352 assert_eq!(s.len(), 0);
3353 assert_eq!(s.first_id(), None);
3354 assert_eq!(s.last_id(), Id::MIN);
3355 assert_eq!(s.trim_maxlen(0, true, None), 0);
3356 assert_eq!(s.trim_minid(Id::MAX, true, None), 0);
3357 assert_eq!(dump(&s), vec![]);
3358 }
3359}