reifydb_transaction/queue/
chain.rs1use std::collections::BTreeSet;
5
6use reifydb_codec::row::bytes::EncodedBytes;
7use reifydb_core::{
8 interface::{catalog::id::QueueId, store::SingleVersionRangeRev},
9 key::queue::QueueKeyActiveKey,
10};
11use reifydb_value::{Result, util::cowvec::CowVec, value::row_number::RowNumber};
12
13use crate::single::{SingleTransaction, write::SingleWriteTransaction};
14
15#[derive(Debug, Default)]
16pub struct ChainOverlay {
17 added: BTreeSet<(u64, RowNumber)>,
18 removed: BTreeSet<(u64, RowNumber)>,
19}
20
21impl ChainOverlay {
22 fn of_key(set: &BTreeSet<(u64, RowNumber)>, key_hash: u64) -> impl Iterator<Item = RowNumber> + '_ {
23 set.range((key_hash, RowNumber(0))..=(key_hash, RowNumber(u64::MAX))).map(|(_, row)| *row)
24 }
25}
26
27#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub enum ChainHead {
29 Empty,
30 Single(RowNumber),
31 Multiple(RowNumber),
32}
33
34pub fn chain_peek(
35 single: &SingleTransaction,
36 overlay: &ChainOverlay,
37 queue: QueueId,
38 partition: u16,
39 key_hash: u64,
40) -> Result<ChainHead> {
41 let store = single.read_store();
42 let budget = 2 + ChainOverlay::of_key(&overlay.removed, key_hash).count() as u64;
43 let batch = SingleVersionRangeRev::range_rev_batch(
44 &store,
45 QueueKeyActiveKey::key_scan(queue, partition, key_hash).encode(),
46 budget,
47 )?;
48
49 let mut rows: BTreeSet<RowNumber> = batch
50 .items
51 .iter()
52 .filter_map(|item| QueueKeyActiveKey::decode(&item.key))
53 .map(|key| key.row)
54 .filter(|row| !overlay.removed.contains(&(key_hash, *row)))
55 .collect();
56 rows.extend(ChainOverlay::of_key(&overlay.added, key_hash));
57
58 let mut ascending = rows.into_iter();
59
60 Ok(match (ascending.next(), ascending.next()) {
61 (None, _) => ChainHead::Empty,
62 (Some(row), None) => ChainHead::Single(row),
63 (Some(row), Some(_)) => ChainHead::Multiple(row),
64 })
65}
66
67pub fn chain_add(
68 tx: &mut SingleWriteTransaction<'_>,
69 overlay: &mut ChainOverlay,
70 queue: QueueId,
71 partition: u16,
72 key_hash: u64,
73 row: RowNumber,
74) -> Result<()> {
75 tx.set(&QueueKeyActiveKey::new(queue, partition, key_hash, row), EncodedBytes(CowVec::new(vec![])))?;
76 overlay.removed.remove(&(key_hash, row));
77 overlay.added.insert((key_hash, row));
78
79 Ok(())
80}
81
82pub fn chain_remove(
83 tx: &mut SingleWriteTransaction<'_>,
84 overlay: &mut ChainOverlay,
85 queue: QueueId,
86 partition: u16,
87 key_hash: u64,
88 row: RowNumber,
89) -> Result<()> {
90 tx.remove(&QueueKeyActiveKey::new(queue, partition, key_hash, row))?;
91 overlay.added.remove(&(key_hash, row));
92 overlay.removed.insert((key_hash, row));
93
94 Ok(())
95}