Skip to main content

reifydb_engine/queue/
hydrate.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}