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(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
109pub struct PendingSummary {
112 pub total: u64,
114 pub id_range: Option<(StreamId, StreamId)>,
116 pub by_consumer: Vec<(Vec<u8>, u64)>,
118}
119
120pub struct PendingExtended {
123 pub rows: Vec<PendingExtendedRow>,
125}
126
127pub struct PendingExtendedRow {
129 pub id: StreamId,
131 pub consumer: Vec<u8>,
133 pub idle_ms: u64,
135 pub delivery_count: u32,
137}
138
139pub struct XClaimOpts {
142 pub min_idle_ms: u64,
144 pub idle_override_ms: Option<u64>,
147 pub time_override_ms: Option<u64>,
150 pub retrycount_override: Option<u32>,
152 pub force: bool,
155 pub justid: bool,
158}
159
160impl StreamData {
161 pub fn group_create(&mut self, name: &[u8], mode: GroupCreateMode) -> Result<bool, StoreError> {
165 if self.groups.contains_key(name) {
166 return Ok(false);
167 }
168 let last_delivered_id = match mode {
169 GroupCreateMode::AtId(id) => id,
170 GroupCreateMode::AtCurrent => self.last_id,
171 };
172 self.groups.insert(
173 SmallBytes::from_slice(name),
174 Box::new(ConsumerGroup {
175 last_delivered_id,
176 pel: BTreeMap::new(),
177 consumers: KevyMap::default(),
178 }),
179 );
180 Ok(true)
181 }
182
183 pub fn group_destroy(&mut self, name: &[u8]) -> bool {
185 self.groups.remove(name).is_some()
186 }
187
188 pub fn group_setid(&mut self, name: &[u8], mode: GroupCreateMode) -> bool {
191 let Some(g) = self.groups.get_mut(name) else {
192 return false;
193 };
194 g.last_delivered_id = match mode {
195 GroupCreateMode::AtId(id) => id,
196 GroupCreateMode::AtCurrent => self.last_id,
197 };
198 true
199 }
200
201 pub fn group_create_consumer(&mut self, group: &[u8], consumer: &[u8], now_ms: u64) -> bool {
205 let Some(g) = self.groups.get_mut(group) else {
206 return false;
207 };
208 if g.consumers.contains_key(consumer) {
209 return false;
210 }
211 g.consumers.insert(
212 SmallBytes::from_slice(consumer),
213 Box::new(ConsumerState {
214 name: SmallBytes::from_slice(consumer),
215 last_seen_ms: now_ms,
216 pel_count: 0,
217 }),
218 );
219 true
220 }
221
222 pub fn group_del_consumer(&mut self, group: &[u8], consumer: &[u8]) -> u64 {
225 let Some(g) = self.groups.get_mut(group) else {
226 return 0;
227 };
228 let dropped = g.pel.len();
229 g.pel.retain(|_, p| p.consumer.as_slice() != consumer);
230 let dropped = dropped - g.pel.len();
231 g.consumers.remove(consumer);
232 dropped as u64
233 }
234
235 pub fn readgroup(
240 &mut self,
241 group: &[u8],
242 consumer: &[u8],
243 last_seen_arg: ReadGroupId,
244 count: Option<usize>,
245 noack: bool,
246 now_ms: u64,
247 ) -> Result<EntryBatch, StoreError> {
248 let Some(g) = self.groups.get_mut(group) else {
249 return Err(StoreError::NoSuchKey);
250 };
251 let consumer_smb = SmallBytes::from_slice(consumer);
252 ensure_consumer(g, &consumer_smb, now_ms);
253 if let Some(cs) = g.consumers.get_mut(consumer_smb.as_slice()) {
254 cs.last_seen_ms = now_ms;
255 }
256 match last_seen_arg {
257 ReadGroupId::New => {
258 let start = g.last_delivered_id.next();
259 let entries: Vec<(StreamId, &[(SmallBytes, SmallBytes)])> = self
260 .entries
261 .range(start..=StreamId::MAX)
262 .map(|(id, fv)| (*id, fv.as_slice()))
263 .collect();
264 let take = match count {
265 Some(n) => entries.into_iter().take(n).collect::<Vec<_>>(),
266 None => entries,
267 };
268 if take.is_empty() {
269 return Ok(Vec::new());
270 }
271 if !noack {
272 record_deliveries(g, &consumer_smb, &take, now_ms);
273 }
274 let g_mut = self.groups.get_mut(group).expect("present");
275 if let Some((last_id, _)) = take.last() {
276 g_mut.last_delivered_id = *last_id;
277 }
278 Ok(super::clone_entries(take))
279 }
280 ReadGroupId::ReplayAfter(after) => {
281 Ok(replay_pel_entries(g, &self.entries, &consumer_smb, after, count))
282 }
283 }
284 }
285
286 pub fn ack(&mut self, group: &[u8], ids: &[StreamId]) -> u64 {
288 let Some(g) = self.groups.get_mut(group) else {
289 return 0;
290 };
291 let mut n = 0u64;
292 for id in ids {
293 if let Some(p) = g.pel.remove(id) {
294 if let Some(cs) = g.consumers.get_mut(p.consumer.as_slice()) {
295 cs.pel_count = cs.pel_count.saturating_sub(1);
296 }
297 n += 1;
298 }
299 }
300 n
301 }
302
303 pub fn pending_summary(&self, group: &[u8]) -> Option<PendingSummary> {
305 let g = self.groups.get(group)?;
306 let total = g.pel.len() as u64;
307 let id_range = match (g.pel.keys().next(), g.pel.keys().next_back()) {
308 (Some(lo), Some(hi)) => Some((*lo, *hi)),
309 _ => None,
310 };
311 let mut counts: Vec<(Vec<u8>, u64)> = Vec::new();
312 for p in g.pel.values() {
313 if let Some((_, n)) = counts.iter_mut().find(|(name, _)| name == p.consumer.as_slice())
314 {
315 *n += 1;
316 } else {
317 counts.push((p.consumer.to_vec(), 1));
318 }
319 }
320 Some(PendingSummary { total, id_range, by_consumer: counts })
321 }
322
323 #[allow(clippy::too_many_arguments)]
325 pub fn pending_extended(
326 &self,
327 group: &[u8],
328 idle_min_ms: Option<u64>,
329 start: StreamId,
330 end: StreamId,
331 count: usize,
332 consumer_filter: Option<&[u8]>,
333 now_ms: u64,
334 ) -> Option<PendingExtended> {
335 let g = self.groups.get(group)?;
336 let mut rows = Vec::with_capacity(count.min(g.pel.len()));
337 for (id, p) in g.pel.range(start..=end) {
338 if rows.len() >= count {
339 break;
340 }
341 let idle = now_ms.saturating_sub(p.delivery_time_ms);
342 if let Some(min) = idle_min_ms
343 && idle < min
344 {
345 continue;
346 }
347 if let Some(c) = consumer_filter
348 && p.consumer.as_slice() != c
349 {
350 continue;
351 }
352 rows.push(PendingExtendedRow {
353 id: *id,
354 consumer: p.consumer.to_vec(),
355 idle_ms: idle,
356 delivery_count: p.delivery_count,
357 });
358 }
359 Some(PendingExtended { rows })
360 }
361}
362
363#[derive(Clone, Copy, Debug, PartialEq, Eq)]
366pub enum ReadGroupId {
367 New,
369 ReplayAfter(StreamId),
371}
372
373fn replay_pel_entries(
379 g: &ConsumerGroup,
380 entries: &alloc::collections::BTreeMap<StreamId, Vec<(SmallBytes, SmallBytes)>>,
381 consumer: &SmallBytes,
382 after: StreamId,
383 count: Option<usize>,
384) -> EntryBatch {
385 let mut hit: Vec<(StreamId, Vec<(SmallBytes, SmallBytes)>)> = Vec::new();
386 for (id, pel_entry) in g.pel.range(after.next()..=StreamId::MAX) {
387 if pel_entry.consumer != *consumer {
388 continue;
389 }
390 if let Some(fv) = entries.get(id) {
391 hit.push((*id, fv.clone()));
392 }
393 if let Some(n) = count
394 && hit.len() >= n
395 {
396 break;
397 }
398 }
399 hit.into_iter()
400 .map(|(id, fv)| (id, fv.iter().map(|(f, v)| (f.to_vec(), v.to_vec())).collect()))
401 .collect()
402}
403
404pub(super) fn ensure_consumer(g: &mut ConsumerGroup, name: &SmallBytes, now_ms: u64) {
405 if g.consumers.get(name.as_slice()).is_none() {
406 g.consumers.insert(
407 name.clone(),
408 Box::new(ConsumerState { name: name.clone(), last_seen_ms: now_ms, pel_count: 0 }),
409 );
410 }
411}
412
413fn record_deliveries(
414 g: &mut ConsumerGroup,
415 consumer: &SmallBytes,
416 entries: &[(StreamId, &[(SmallBytes, SmallBytes)])],
417 now_ms: u64,
418) {
419 let mut new_for_consumer = 0usize;
420 for (id, _) in entries {
421 let entry = g.pel.entry(*id).or_insert_with(|| {
422 new_for_consumer += 1;
423 PelEntry { consumer: consumer.clone(), delivery_time_ms: now_ms, delivery_count: 0 }
424 });
425 if entry.consumer != *consumer {
426 if let Some(prev) = g.consumers.get_mut(entry.consumer.as_slice()) {
430 prev.pel_count = prev.pel_count.saturating_sub(1);
431 }
432 entry.consumer = consumer.clone();
433 new_for_consumer += 1;
434 }
435 entry.delivery_time_ms = now_ms;
436 entry.delivery_count = entry.delivery_count.saturating_add(1);
437 }
438 if let Some(cs) = g.consumers.get_mut(consumer.as_slice()) {
439 cs.pel_count = cs.pel_count.saturating_add(new_for_consumer);
440 }
441}