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