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 not_before: Option<DateTime>,
19	pub encoded: EncodedBytes,
20}
21
22pub trait QueueOperations {
23	fn insert_queue(&mut self, queue: &Queue, rows: &[QueueInsertRow]) -> Result<()>;
24}
25
26fn row_changes(queue: &Queue, rows: &[QueueInsertRow]) -> Vec<RowChange> {
27	rows.iter()
28		.map(|row| {
29			RowChange::QueueInsert(QueueRowInsertion {
30				queue_id: queue.id,
31				partition: row.partition,
32				row_number: row.row_number,
33				not_before: row.not_before,
34				encoded: row.encoded.clone(),
35			})
36		})
37		.collect()
38}
39
40impl QueueOperations for CommandTransaction {
41	fn insert_queue(&mut self, queue: &Queue, rows: &[QueueInsertRow]) -> Result<()> {
42		if rows.is_empty() {
43			return Ok(());
44		}
45
46		for row in rows {
47			self.set(&RowKey::encoded(queue.id, row.row_number), row.encoded.clone())?;
48		}
49
50		self.track_row_change(&row_changes(queue, rows));
51
52		Ok(())
53	}
54}
55
56impl QueueOperations for AdminTransaction {
57	fn insert_queue(&mut self, queue: &Queue, rows: &[QueueInsertRow]) -> Result<()> {
58		if rows.is_empty() {
59			return Ok(());
60		}
61
62		for row in rows {
63			self.set(&RowKey::encoded(queue.id, row.row_number), row.encoded.clone())?;
64		}
65
66		self.track_row_change(&row_changes(queue, rows));
67
68		Ok(())
69	}
70}
71
72impl QueueOperations for Transaction<'_> {
73	fn insert_queue(&mut self, queue: &Queue, rows: &[QueueInsertRow]) -> Result<()> {
74		match self {
75			Transaction::Command(txn) => txn.insert_queue(queue, rows),
76			Transaction::Admin(txn) => txn.insert_queue(queue, rows),
77			Transaction::Test(t) => t.inner.insert_queue(queue, rows),
78			Transaction::Query(_) => panic!("Write operations not supported on Query transaction"),
79			Transaction::Replica(_) => panic!("Write operations not supported on Replica transaction"),
80		}
81	}
82}