1#[cfg(not(feature = "std"))]
8use crate::nostd_prelude::*;
9use alloc::collections::BTreeMap;
10
11use kevy_map::KevyMap;
12
13pub(super) use super::claim::AutoclaimResult;
14use super::{EntryBatch, StreamData, StreamId};
15use crate::StoreError;
16use crate::value::SmallBytes;
17
18#[derive(Debug, Clone)]
21pub struct ConsumerGroup {
22 pub last_delivered_id: StreamId,
25 pub pel: BTreeMap<StreamId, PelEntry>,
28 pub consumers: KevyMap<SmallBytes, Box<ConsumerState>>,
30}
31
32impl ConsumerGroup {
33 pub fn last_delivered_id(&self) -> StreamId {
35 self.last_delivered_id
36 }
37 pub fn pending_count(&self) -> usize {
39 self.pel.len()
40 }
41 pub fn consumer_count(&self) -> usize {
43 self.consumers.len()
44 }
45 pub fn consumers_iter(&self) -> impl Iterator<Item = (&[u8], &ConsumerState)> {
47 self.consumers.iter().map(|(k, v)| (k.as_slice(), v.as_ref()))
48 }
49}
50
51impl ConsumerState {
52 pub fn pending_count(&self) -> usize {
54 self.pel_count
55 }
56 pub fn last_seen_ms(&self) -> u64 {
58 self.last_seen_ms
59 }
60}
61
62impl Default for ConsumerGroup {
63 fn default() -> Self {
64 Self {
65 last_delivered_id: StreamId::MIN,
66 pel: BTreeMap::new(),
67 consumers: KevyMap::default(),
68 }
69 }
70}
71
72#[derive(Clone, Debug)]
74pub struct PelEntry {
75 pub consumer: SmallBytes,
78 pub delivery_time_ms: u64,
81 pub delivery_count: u32,
84}
85
86#[derive(Clone, Debug)]
88#[allow(dead_code)]
89pub struct ConsumerState {
90 pub name: SmallBytes,
92 pub last_seen_ms: u64,
95 pub pel_count: usize,
97}
98
99#[derive(Clone, Copy, Debug, PartialEq, Eq)]
102pub enum GroupCreateMode {
103 AtId(StreamId),
105 AtCurrent,
107}
108
109#[derive(Debug)]
112pub struct PendingSummary {
113 pub total: u64,
115 pub id_range: Option<(StreamId, StreamId)>,
117 pub by_consumer: Vec<(Vec<u8>, u64)>,
119}
120
121#[derive(Debug)]
124pub struct PendingExtended {
125 pub rows: Vec<PendingExtendedRow>,
127}
128
129#[derive(Debug)]
131pub struct PendingExtendedRow {
132 pub id: StreamId,
134 pub consumer: Vec<u8>,
136 pub idle_ms: u64,
138 pub delivery_count: u32,
140}
141
142#[derive(Debug)]
145pub struct XClaimOpts {
146 pub min_idle_ms: u64,
148 pub idle_override_ms: Option<u64>,
151 pub time_override_ms: Option<u64>,
154 pub retrycount_override: Option<u32>,
156 pub force: bool,
159 pub justid: bool,
162}
163
164impl StreamData {
165 pub fn group_create(&mut self, name: &[u8], mode: GroupCreateMode) -> Result<bool, StoreError> {
169 if self.groups.contains_key(name) {
170 return Ok(false);
171 }
172 let last_delivered_id = match mode {
173 GroupCreateMode::AtId(id) => id,
174 GroupCreateMode::AtCurrent => self.last_id,
175 };
176 self.groups.insert(
177 SmallBytes::from_slice(name),
178 Box::new(ConsumerGroup {
179 last_delivered_id,
180 pel: BTreeMap::new(),
181 consumers: KevyMap::default(),
182 }),
183 );
184 Ok(true)
185 }
186
187 pub fn group_destroy(&mut self, name: &[u8]) -> bool {
189 self.groups.remove(name).is_some()
190 }
191
192 pub fn group_setid(&mut self, name: &[u8], mode: GroupCreateMode) -> bool {
195 let Some(g) = self.groups.get_mut(name) else {
196 return false;
197 };
198 g.last_delivered_id = match mode {
199 GroupCreateMode::AtId(id) => id,
200 GroupCreateMode::AtCurrent => self.last_id,
201 };
202 true
203 }
204
205 pub fn group_create_consumer(&mut self, group: &[u8], consumer: &[u8], now_ms: u64) -> bool {
209 let Some(g) = self.groups.get_mut(group) else {
210 return false;
211 };
212 if g.consumers.contains_key(consumer) {
213 return false;
214 }
215 g.consumers.insert(
216 SmallBytes::from_slice(consumer),
217 Box::new(ConsumerState {
218 name: SmallBytes::from_slice(consumer),
219 last_seen_ms: now_ms,
220 pel_count: 0,
221 }),
222 );
223 true
224 }
225
226 pub fn group_del_consumer(&mut self, group: &[u8], consumer: &[u8]) -> u64 {
229 let Some(g) = self.groups.get_mut(group) else {
230 return 0;
231 };
232 let dropped = g.pel.len();
233 g.pel.retain(|_, p| p.consumer.as_slice() != consumer);
234 let dropped = dropped - g.pel.len();
235 g.consumers.remove(consumer);
236 dropped as u64
237 }
238
239 pub fn readgroup(
244 &mut self,
245 group: &[u8],
246 consumer: &[u8],
247 last_seen_arg: ReadGroupId,
248 count: Option<usize>,
249 noack: bool,
250 now_ms: u64,
251 ) -> Result<EntryBatch, StoreError> {
252 let Some(g) = self.groups.get_mut(group) else {
253 return Err(StoreError::NoSuchKey);
254 };
255 let consumer_smb = SmallBytes::from_slice(consumer);
256 ensure_consumer(g, &consumer_smb, now_ms);
257 if let Some(cs) = g.consumers.get_mut(consumer_smb.as_slice()) {
258 cs.last_seen_ms = now_ms;
259 }
260 match last_seen_arg {
261 ReadGroupId::New => {
262 let start = g.last_delivered_id.next();
263 let entries: Vec<(StreamId, &[(SmallBytes, SmallBytes)])> = self
264 .entries
265 .range(start..=StreamId::MAX)
266 .map(|(id, fv)| (*id, fv.as_slice()))
267 .collect();
268 let take = match count {
269 Some(n) => entries.into_iter().take(n).collect::<Vec<_>>(),
270 None => entries,
271 };
272 if take.is_empty() {
273 return Ok(Vec::new());
274 }
275 if !noack {
276 record_deliveries(g, &consumer_smb, &take, now_ms);
277 }
278 let g_mut = self.groups.get_mut(group).expect("present");
279 if let Some((last_id, _)) = take.last() {
280 g_mut.last_delivered_id = *last_id;
281 }
282 Ok(super::clone_entries(take))
283 }
284 ReadGroupId::ReplayAfter(after) => {
285 Ok(replay_pel_entries(g, &self.entries, &consumer_smb, after, count))
286 }
287 }
288 }
289
290 pub fn ack(&mut self, group: &[u8], ids: &[StreamId]) -> u64 {
292 let Some(g) = self.groups.get_mut(group) else {
293 return 0;
294 };
295 let mut n = 0u64;
296 for id in ids {
297 if let Some(p) = g.pel.remove(id) {
298 if let Some(cs) = g.consumers.get_mut(p.consumer.as_slice()) {
299 cs.pel_count = cs.pel_count.saturating_sub(1);
300 }
301 n += 1;
302 }
303 }
304 n
305 }
306
307 pub fn pending_summary(&self, group: &[u8]) -> Option<PendingSummary> {
309 let g = self.groups.get(group)?;
310 let total = g.pel.len() as u64;
311 let id_range = match (g.pel.keys().next(), g.pel.keys().next_back()) {
312 (Some(lo), Some(hi)) => Some((*lo, *hi)),
313 _ => None,
314 };
315 let mut counts: Vec<(Vec<u8>, u64)> = Vec::new();
316 for p in g.pel.values() {
317 if let Some((_, n)) = counts.iter_mut().find(|(name, _)| name == p.consumer.as_slice())
318 {
319 *n += 1;
320 } else {
321 counts.push((p.consumer.to_vec(), 1));
322 }
323 }
324 Some(PendingSummary { total, id_range, by_consumer: counts })
325 }
326
327 #[allow(clippy::too_many_arguments)]
329 pub fn pending_extended(
330 &self,
331 group: &[u8],
332 idle_min_ms: Option<u64>,
333 start: StreamId,
334 end: StreamId,
335 count: usize,
336 consumer_filter: Option<&[u8]>,
337 now_ms: u64,
338 ) -> Option<PendingExtended> {
339 let g = self.groups.get(group)?;
340 let mut rows = Vec::with_capacity(count.min(g.pel.len()));
341 for (id, p) in g.pel.range(start..=end) {
342 if rows.len() >= count {
343 break;
344 }
345 let idle = now_ms.saturating_sub(p.delivery_time_ms);
346 if let Some(min) = idle_min_ms
347 && idle < min
348 {
349 continue;
350 }
351 if let Some(c) = consumer_filter
352 && p.consumer.as_slice() != c
353 {
354 continue;
355 }
356 rows.push(PendingExtendedRow {
357 id: *id,
358 consumer: p.consumer.to_vec(),
359 idle_ms: idle,
360 delivery_count: p.delivery_count,
361 });
362 }
363 Some(PendingExtended { rows })
364 }
365}
366
367#[derive(Clone, Copy, Debug, PartialEq, Eq)]
370pub enum ReadGroupId {
371 New,
373 ReplayAfter(StreamId),
375}
376
377fn replay_pel_entries(
383 g: &ConsumerGroup,
384 entries: &alloc::collections::BTreeMap<StreamId, Vec<(SmallBytes, SmallBytes)>>,
385 consumer: &SmallBytes,
386 after: StreamId,
387 count: Option<usize>,
388) -> EntryBatch {
389 let mut hit: Vec<(StreamId, Vec<(SmallBytes, SmallBytes)>)> = Vec::new();
390 for (id, pel_entry) in g.pel.range(after.next()..=StreamId::MAX) {
391 if pel_entry.consumer != *consumer {
392 continue;
393 }
394 if let Some(fv) = entries.get(id) {
395 hit.push((*id, fv.clone()));
396 }
397 if let Some(n) = count
398 && hit.len() >= n
399 {
400 break;
401 }
402 }
403 hit.into_iter()
404 .map(|(id, fv)| (id, fv.iter().map(|(f, v)| (f.to_vec(), v.to_vec())).collect()))
405 .collect()
406}
407
408pub(super) fn ensure_consumer(g: &mut ConsumerGroup, name: &SmallBytes, now_ms: u64) {
409 if g.consumers.get(name.as_slice()).is_none() {
410 g.consumers.insert(
411 name.clone(),
412 Box::new(ConsumerState { name: name.clone(), last_seen_ms: now_ms, pel_count: 0 }),
413 );
414 }
415}
416
417fn record_deliveries(
418 g: &mut ConsumerGroup,
419 consumer: &SmallBytes,
420 entries: &[(StreamId, &[(SmallBytes, SmallBytes)])],
421 now_ms: u64,
422) {
423 let mut new_for_consumer = 0usize;
424 for (id, _) in entries {
425 let entry = g.pel.entry(*id).or_insert_with(|| {
426 new_for_consumer += 1;
427 PelEntry { consumer: consumer.clone(), delivery_time_ms: now_ms, delivery_count: 0 }
428 });
429 if entry.consumer != *consumer {
430 if let Some(prev) = g.consumers.get_mut(entry.consumer.as_slice()) {
434 prev.pel_count = prev.pel_count.saturating_sub(1);
435 }
436 entry.consumer = consumer.clone();
437 new_for_consumer += 1;
438 }
439 entry.delivery_time_ms = now_ms;
440 entry.delivery_count = entry.delivery_count.saturating_add(1);
441 }
442 if let Some(cs) = g.consumers.get_mut(consumer.as_slice()) {
443 cs.pel_count = cs.pel_count.saturating_add(new_for_consumer);
444 }
445}