Skip to main content

reifydb_engine/transaction/operation/
queue.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use reifydb_codec::row::bytes::EncodedBytes;
5use reifydb_core::{interface::catalog::queue::Queue, key::row::RowKey};
6use reifydb_transaction::{
7	change::{QueueRowInsertion, RowChange},
8	transaction::{Transaction, admin::AdminTransaction, command::CommandTransaction},
9};
10use reifydb_value::value::{datetime::DateTime, row_number::RowNumber};
11
12use crate::Result;
13
14#[derive(Debug, Clone)]
15pub struct QueueInsertRow {
16	pub row_number: RowNumber,
17	pub partition: u16,
18	pub key_hash: Option<u64>,
19	pub not_before: Option<DateTime>,
20	pub encoded: EncodedBytes,
21}
22
23pub trait QueueOperations {
24	fn insert_queue(&mut self, queue: &Queue, rows: &[QueueInsertRow]) -> Result<()>;
25}
26
27fn row_changes(queue: &Queue, rows: &[QueueInsertRow]) -> Vec<RowChange> {
28	rows.iter()
29		.map(|row| {
30			RowChange::QueueInsert(QueueRowInsertion {
31				queue_id: queue.id,
32				partition: row.partition,
33				key_hash: row.key_hash,
34				row_number: row.row_number,
35				not_before: row.not_before,
36				encoded: row.encoded.clone(),
37			})
38		})
39		.collect()
40}
41
42impl QueueOperations for CommandTransaction {
43	fn insert_queue(&mut self, queue: &Queue, rows: &[QueueInsertRow]) -> Result<()> {
44		if rows.is_empty() {
45			return Ok(());
46		}
47
48		for row in rows {
49			self.set(&RowKey::new(queue.id, row.row_number), row.encoded.clone())?;
50		}
51
52		self.track_row_change(&row_changes(queue, rows));
53
54		Ok(())
55	}
56}
57
58impl QueueOperations for AdminTransaction {
59	fn insert_queue(&mut self, queue: &Queue, rows: &[QueueInsertRow]) -> Result<()> {
60		if rows.is_empty() {
61			return Ok(());
62		}
63
64		for row in rows {
65			self.set(&RowKey::new(queue.id, row.row_number), row.encoded.clone())?;
66		}
67
68		self.track_row_change(&row_changes(queue, rows));
69
70		Ok(())
71	}
72}
73
74impl QueueOperations for Transaction<'_> {
75	fn insert_queue(&mut self, queue: &Queue, rows: &[QueueInsertRow]) -> Result<()> {
76		match self {
77			Transaction::Command(txn) => txn.insert_queue(queue, rows),
78			Transaction::Admin(txn) => txn.insert_queue(queue, rows),
79			Transaction::Test(t) => t.inner.insert_queue(queue, rows),
80			Transaction::Query(_) => panic!("Write operations not supported on Query transaction"),
81		}
82	}
83}