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 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}