Skip to main content

reifydb_transaction/queue/
chain.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}