reifydb_engine/transaction/operation/
queue.rs1use 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}