reifydb_engine/vm/volcano/
query.rs1use std::sync::Arc;
5
6use reifydb_core::{
7 error::diagnostic::{operation, query},
8 interface::{
9 catalog::config::{ConfigKey, GetConfig},
10 resolved::ResolvedObject,
11 },
12 sort::SortKey,
13 util::budget::MemoryBudget,
14 value::column::{buffer::ColumnBuffer, columns::Columns, headers::ColumnHeaders},
15};
16use reifydb_evaluate::{expression::context::EvalContext, stack::SymbolTable};
17use reifydb_extension::transform::context::TransformContext;
18use reifydb_transaction::transaction::Transaction;
19use reifydb_value::{
20 byte_size::ByteSize,
21 error,
22 params::Params,
23 value::{identity::IdentityId, value_type::ValueType},
24};
25
26use crate::{Result, vm::services::Services};
27
28pub fn query_budget(services: &Services) -> Arc<MemoryBudget> {
29 let limit = services.catalog.get_config_uint8(ConfigKey::QueryMemoryLimit);
30 Arc::new(MemoryBudget::new(ByteSize::from_bytes(limit)))
31}
32
33pub(crate) fn is_scalar_type(ty: &ValueType) -> bool {
34 !matches!(ty.inner_type(), ValueType::Any | ValueType::List(_) | ValueType::Record(_) | ValueType::Tuple(_))
35}
36
37pub(crate) fn ensure_sort_key_orderable(key: &SortKey, data: &ColumnBuffer) -> Result<()> {
38 let ty = data.get_type();
39 if is_scalar_type(&ty) {
40 Ok(())
41 } else {
42 Err(error!(operation::sort_key_not_orderable(key.column.clone(), ty)))
43 }
44}
45
46pub(crate) fn charge_query_memory(budget: &MemoryBudget, charged: &mut usize, buffer: &Columns) -> Result<()> {
47 charge_query_memory_bytes(budget, charged, buffer.heap_size())
48}
49
50pub(crate) fn charge_query_memory_bytes(budget: &MemoryBudget, charged: &mut usize, total: usize) -> Result<()> {
51 if total > *charged {
52 let delta = (total - *charged) as u64;
53 if !budget.try_charge(ByteSize::from_bytes(delta)) {
54 return Err(error!(query::memory_limit_exceeded(budget.used(), budget.limit())));
55 }
56 *charged = total;
57 }
58 Ok(())
59}
60
61pub trait QueryNode: Send + Sync {
62 fn initialize<'a>(&mut self, rx: &mut Transaction<'a>, ctx: &QueryContext) -> Result<()>;
63
64 fn next<'a>(&mut self, rx: &mut Transaction<'a>, ctx: &mut QueryContext) -> Result<Option<Columns>>;
65
66 fn headers(&self) -> Option<ColumnHeaders>;
67}
68
69#[derive(Clone)]
70pub struct QueryContext {
71 pub services: Arc<Services>,
72 pub source: Option<ResolvedObject>,
73 pub batch_size: u64,
74 pub params: Params,
75 pub symbols: SymbolTable,
76 pub identity: IdentityId,
77 pub memory: Arc<MemoryBudget>,
78}
79
80impl QueryNode for Box<dyn QueryNode> {
81 fn initialize<'a>(&mut self, rx: &mut Transaction<'a>, ctx: &QueryContext) -> Result<()> {
82 (**self).initialize(rx, ctx)
83 }
84
85 fn next<'a>(&mut self, rx: &mut Transaction<'a>, ctx: &mut QueryContext) -> Result<Option<Columns>> {
86 let result = (**self).next(rx, ctx)?;
87 if let Some(ref columns) = result {
88 columns.assert_invariants("QueryNode::next output");
89 }
90 Ok(result)
91 }
92
93 fn headers(&self) -> Option<ColumnHeaders> {
94 (**self).headers()
95 }
96}
97
98#[cfg(test)]
99mod tests {
100 use reifydb_core::{
101 util::budget::MemoryBudget,
102 value::column::{ColumnWithName, columns::Columns},
103 };
104 use reifydb_value::byte_size::ByteSize;
105
106 use super::charge_query_memory;
107
108 #[test]
109 fn charge_query_memory_delta_charges_and_rejects_over_budget() {
110 let budget = MemoryBudget::new(ByteSize::from_kib(1));
111 let mut charged = 0usize;
112
113 let small = Columns::new(vec![ColumnWithName::int4("c", [1i32, 2, 3, 4])]);
114 charge_query_memory(&budget, &mut charged, &small).expect("small buffer fits under 1 KiB");
115 let after_first = budget.used().as_bytes();
116 assert!(after_first > 0, "charging a non-empty buffer must consume budget");
117 assert_eq!(charged as u64, after_first, "charged must track exactly what the budget recorded");
118
119 charge_query_memory(&budget, &mut charged, &small).expect("re-charge of the same buffer is free");
120 assert_eq!(
121 budget.used().as_bytes(),
122 after_first,
123 "delta charging must not double count an unchanged buffer"
124 );
125
126 let big = Columns::new(vec![ColumnWithName::int4("c", 0..4000i32)]);
127 let mut big_charged = 0usize;
128 let err = charge_query_memory(&budget, &mut big_charged, &big).unwrap_err();
129 assert_eq!(err.0.code, "QUERY_006", "an over-budget charge must raise the memory-limit diagnostic");
130 }
131}
132
133pub fn eval_context_from_query<'a>(ctx: &'a QueryContext) -> EvalContext<'a> {
134 EvalContext {
135 target: None,
136 columns: Columns::empty(),
137 row_count: 1,
138 take: None,
139 params: &ctx.params,
140 symbols: &ctx.symbols,
141 is_aggregate_context: false,
142 routines: &ctx.services.routines,
143 runtime_context: &ctx.services.runtime_context,
144 identity: ctx.identity,
145 }
146}
147
148pub fn eval_context_from_transform<'a>(ctx: &'a TransformContext<'a>, stored: &'a QueryContext) -> EvalContext<'a> {
149 EvalContext {
150 target: None,
151 columns: Columns::empty(),
152 row_count: 1,
153 take: None,
154 params: ctx.params,
155 symbols: &stored.symbols,
156 is_aggregate_context: false,
157 routines: &stored.services.routines,
158 runtime_context: ctx.runtime_context,
159 identity: stored.identity,
160 }
161}