reifydb_engine/queue/
hydrate.rs1use std::collections::BTreeMap;
5
6use reifydb_catalog::catalog::Catalog;
7use reifydb_codec::row::{queue::EncodedQueueRow, shape::RowShape};
8use reifydb_core::{
9 interface::{
10 catalog::{id::QueueId, queue::Queue},
11 store::SingleVersionGet,
12 },
13 internal_error,
14 key::{any::TaggedKey, queue::QueueItemStateKey, row::RowKeyRange},
15};
16use reifydb_transaction::{
17 multi::RangeScope,
18 queue::scheduling::{QueueAdmission, admit_ready_items},
19 single::SingleTransaction,
20 transaction::Transaction,
21};
22use reifydb_value::value::{identity::IdentityId, row_number::RowNumber};
23use tracing::{info, instrument};
24
25use crate::{
26 Result,
27 engine::StandardEngine,
28 queue::partition::{ordered_by_index, placement_of},
29};
30
31const HYDRATE_BATCH: usize = 1024;
32
33#[instrument(name = "queue::hydrate", level = "info", skip_all)]
34pub fn hydrate_queues(engine: &StandardEngine) -> Result<u64> {
35 let catalog = engine.catalog();
36 let single = engine.single().clone();
37
38 let mut query = engine.begin_query(IdentityId::system())?;
39 let mut txn = Transaction::Query(&mut query);
40
41 let queues = catalog.list_queues(&mut txn)?;
42
43 let mut admitted = 0u64;
44 for queue in &queues {
45 admitted += hydrate_queue(&catalog, &single, &mut txn, queue)?;
46 }
47
48 info!(queues = queues.len(), items = admitted, "queue scheduling state hydrated");
49
50 Ok(admitted)
51}
52
53fn hydrate_queue(
54 catalog: &Catalog,
55 single: &SingleTransaction,
56 txn: &mut Transaction<'_>,
57 queue: &Queue,
58) -> Result<u64> {
59 let ordered_by = ordered_by_index(queue)?;
60
61 let mut pending: BTreeMap<u16, Vec<QueueAdmission>> = BTreeMap::new();
62 let mut last_key: Option<TaggedKey> = None;
63 let mut admitted = 0u64;
64
65 loop {
66 let mut batch: Vec<(RowNumber, EncodedQueueRow)> = Vec::with_capacity(HYDRATE_BATCH);
67 let mut fetched = 0usize;
68
69 {
70 let range = RowKeyRange::scan_range_rev(queue.id.into(), last_key.as_ref());
71 let mut stream = txn.range_rev(range, RangeScope::All, HYDRATE_BATCH)?;
72
73 for _ in 0..HYDRATE_BATCH {
74 match stream.next() {
75 Some(Ok(item)) => {
76 fetched += 1;
77 if let TaggedKey::Row(key) = &item.key {
78 batch.push((key.row, EncodedQueueRow::from(item.bytes)));
79 }
80 last_key = Some(item.key.clone());
81 }
82 Some(Err(err)) => return Err(err),
83 None => break,
84 }
85 }
86 }
87
88 if !batch.is_empty() {
89 let shape = load_shape(catalog, txn, queue, &batch[0].1)?;
90 let store = single.read_store();
91
92 for (row_number, encoded) in &batch {
93 let placement = placement_of(queue, &shape, encoded, ordered_by, *row_number);
94 let state_key = QueueItemStateKey::encoded(queue.id, placement.partition, *row_number);
95 if SingleVersionGet::get(&store, &state_key)?.is_some() {
96 continue;
97 }
98
99 let items = pending.entry(placement.partition).or_default();
100 items.push(QueueAdmission {
101 row: *row_number,
102 key_hash: placement.key_hash,
103 not_before: encoded.not_before(),
104 });
105
106 if items.len() >= HYDRATE_BATCH {
107 admitted += flush(single, queue.id, placement.partition, items)?;
108 }
109 }
110 }
111
112 if fetched < HYDRATE_BATCH {
113 break;
114 }
115 }
116
117 for (partition, items) in pending.iter_mut() {
118 admitted += flush(single, queue.id, *partition, items)?;
119 }
120
121 Ok(admitted)
122}
123
124fn flush(single: &SingleTransaction, queue: QueueId, partition: u16, items: &mut Vec<QueueAdmission>) -> Result<u64> {
125 let admitted = admit_ready_items(single, queue, partition, items)?;
126 items.clear();
127 Ok(admitted)
128}
129
130fn load_shape(catalog: &Catalog, txn: &mut Transaction<'_>, queue: &Queue, row: &EncodedQueueRow) -> Result<RowShape> {
131 let fingerprint = row.fingerprint();
132 catalog.get_or_load_row_shape(fingerprint, txn)?.ok_or_else(|| {
133 internal_error!("RowShape with fingerprint {:?} not found for queue {}", fingerprint, queue.name)
134 })
135}