1use std::{ops::Deref, result::Result as StdResult, sync::Arc};
5
6use bumpalo::Bump;
7use reifydb_catalog::{catalog::Catalog, vtable::system::flow_operator_store::SystemFlowOperatorStore};
8use reifydb_core::{
9 error::diagnostic::subscription,
10 execution::ExecutionResult,
11 interface::catalog::policy::SessionOp,
12 metric::{ExecutionMetrics, StatementMetric},
13 value::column::columns::Columns,
14};
15use reifydb_metric::storage::metric::MetricReader;
16use reifydb_policy::inject_from_policies;
17use reifydb_rql::{
18 ast::parse_str,
19 compiler::{CompilationResult, Compiled, IncrementalCompilation, constrain_policy},
20 fingerprint::request::fingerprint_request,
21};
22use reifydb_runtime::context::clock::Instant;
23use reifydb_store_single::SingleStore;
24use reifydb_transaction::transaction::{
25 RqlExecutor, TestTransaction, Transaction, admin::AdminTransaction, command::CommandTransaction,
26 query::QueryTransaction,
27};
28#[cfg(not(reifydb_single_threaded))]
29use reifydb_value::error::Diagnostic;
30use reifydb_value::{
31 error::Error,
32 params::Params,
33 value::{Value, duration::Duration, frame::frame::Frame, value_type::ValueType},
34};
35use tracing::instrument;
36
37#[cfg(not(reifydb_single_threaded))]
38use crate::remote;
39use crate::{
40 Result,
41 policy::PolicyEvaluator,
42 vm::{
43 Admin, Command, Query, Subscription, Test,
44 services::{EngineConfig, Services},
45 stack::{SymbolTable, Variable},
46 vm::Vm,
47 },
48};
49
50pub struct Executor(Arc<Services>);
51
52impl Clone for Executor {
53 fn clone(&self) -> Self {
54 Self(self.0.clone())
55 }
56}
57
58impl Deref for Executor {
59 type Target = Services;
60
61 fn deref(&self) -> &Self::Target {
62 &self.0
63 }
64}
65
66impl Executor {
67 pub fn new(
68 catalog: Catalog,
69 config: EngineConfig,
70 flow_operator_store: SystemFlowOperatorStore,
71 stats_reader: MetricReader<SingleStore>,
72 ) -> Self {
73 Self(Arc::new(Services::new(catalog, config, flow_operator_store, stats_reader)))
74 }
75
76 pub fn services(&self) -> &Arc<Services> {
77 &self.0
78 }
79
80 pub fn from_services(services: Arc<Services>) -> Self {
81 Self(services)
82 }
83
84 #[allow(dead_code)]
85 pub fn testing() -> Self {
86 Self(Services::testing())
87 }
88
89 #[cfg(not(reifydb_single_threaded))]
90 fn try_forward_remote_query(&self, err: &Error, rql: &str, params: Params) -> Result<Option<Vec<Frame>>> {
91 if let Some(ref registry) = self.0.remote_registry
92 && remote::is_remote_query(err)
93 && let Some(address) = remote::extract_remote_address(err)
94 {
95 let token = remote::extract_remote_token(err);
96 return registry.forward_query(&address, rql, params, token.as_deref()).map(Some);
97 }
98 Ok(None)
99 }
100}
101
102impl RqlExecutor for Executor {
103 fn rql(&self, tx: &mut Transaction<'_>, rql: &str, params: Params) -> ExecutionResult {
104 Executor::rql(self, tx, rql, params)
105 }
106}
107
108fn populate_symbols(symbols: &mut SymbolTable, params: &Params) -> Result<()> {
109 match params {
110 Params::Positional(values) => {
111 for (index, value) in values.iter().enumerate() {
112 let param_name = (index + 1).to_string();
113 symbols.set(param_name, Variable::scalar(value.clone()), false)?;
114 }
115 }
116 Params::Named(map) => {
117 for (name, value) in map.iter() {
118 symbols.set(name.clone(), Variable::scalar(value.clone()), false)?;
119 }
120 }
121 Params::None => {}
122 }
123 Ok(())
124}
125
126fn populate_identity(symbols: &mut SymbolTable, catalog: &Catalog, tx: &mut Transaction<'_>) -> Result<()> {
127 let identity = tx.identity();
128 if identity.is_privileged() {
129 return Ok(());
130 }
131 if identity.is_anonymous() {
132 let columns = Columns::single_row([
133 ("id", Value::IdentityId(identity)),
134 ("name", Value::none_of(ValueType::Utf8)),
135 ("roles", Value::List(vec![])),
136 ]);
137 symbols.set("identity".to_string(), Variable::columns(columns), false)?;
138 return Ok(());
139 }
140 if let Some(user) = catalog.find_identity(tx, identity)? {
141 let roles = catalog.find_role_names_for_identity(tx, identity)?;
142 let role_values: Vec<Value> = roles.into_iter().map(Value::Utf8).collect();
143 let columns = Columns::single_row([
144 ("id", Value::IdentityId(identity)),
145 ("name", Value::Utf8(user.name)),
146 ("roles", Value::List(role_values)),
147 ]);
148 symbols.set("identity".to_string(), Variable::columns(columns), false)?;
149 }
150 Ok(())
151}
152
153type CompiledUnitsResult = (Vec<Frame>, Vec<Frame>, SymbolTable, Vec<StatementMetric>);
154
155struct ExecutionFailure {
156 error: Error,
157 partial_metrics: Vec<StatementMetric>,
158}
159
160fn build_metrics(statements: Vec<StatementMetric>) -> ExecutionMetrics {
161 let fps: Vec<_> = statements.iter().map(|m| m.fingerprint).collect();
162 ExecutionMetrics {
163 fingerprint: fingerprint_request(&fps),
164 statements,
165 ..Default::default()
166 }
167}
168
169struct RunUnitOutcome {
170 symbols: SymbolTable,
171 run_result: Result<()>,
172 execute_duration_us: u64,
173}
174
175#[instrument(
176 name = "vm::run",
177 level = "debug",
178 skip_all,
179 fields(fingerprint = ?compiled.fingerprint, instr_count = compiled.instructions.len()),
180)]
181fn run_compiled_unit(
182 services: &Arc<Services>,
183 tx: &mut Transaction<'_>,
184 compiled: &Compiled,
185 params: &Params,
186 symbols: SymbolTable,
187 result: &mut Vec<Frame>,
188) -> RunUnitOutcome {
189 let mut vm = Vm::from_services(symbols, services, params, tx.identity());
190 let start = services.runtime_context.clock.instant();
191 let run_result = vm.run(services, tx, &compiled.instructions, result);
192 let execute_duration = start.elapsed();
193 RunUnitOutcome {
194 symbols: vm.symbols,
195 run_result,
196 execute_duration_us: execute_duration.as_micros() as u64,
197 }
198}
199
200#[instrument(
201 name = "executor::execute_units",
202 level = "debug",
203 skip_all,
204 fields(unit_count = compiled_list.len()),
205)]
206fn execute_compiled_units(
207 services: &Arc<Services>,
208 tx: &mut Transaction<'_>,
209 compiled_list: &[Compiled],
210 params: &Params,
211 mut symbols: SymbolTable,
212 compile_duration: Duration,
213) -> StdResult<CompiledUnitsResult, ExecutionFailure> {
214 let compile_duration_us = compile_duration.to_std().as_micros() as u64 / compiled_list.len().max(1) as u64;
215 let mut result = vec![];
216 let mut output_results: Vec<Frame> = Vec::new();
217 let mut metrics = Vec::new();
218
219 for compiled in compiled_list.iter() {
220 result.clear();
221 let outcome = run_compiled_unit(services, tx, compiled, params, symbols, &mut result);
222 symbols = outcome.symbols;
223
224 metrics.push(StatementMetric {
225 fingerprint: compiled.fingerprint,
226 normalized_rql: compiled.normalized_rql.clone(),
227 compile_duration_us,
228 execute_duration_us: outcome.execute_duration_us,
229 rows_affected: if outcome.run_result.is_ok() {
230 extract_rows_affected(&result)
231 } else {
232 0
233 },
234 });
235
236 if let Err(error) = outcome.run_result {
237 return Err(ExecutionFailure {
238 error,
239 partial_metrics: metrics,
240 });
241 }
242
243 if compiled.is_output {
244 output_results.append(&mut result);
245 }
246 }
247
248 Ok((output_results, result, symbols, metrics))
249}
250
251fn merge_results(mut output_results: Vec<Frame>, mut remaining: Vec<Frame>) -> Vec<Frame> {
252 output_results.append(&mut remaining);
253 output_results
254}
255
256#[inline]
257fn error_result(error: Error, metrics: ExecutionMetrics) -> ExecutionResult {
258 ExecutionResult {
259 frames: vec![],
260 error: Some(error),
261 metrics,
262 }
263}
264
265fn extract_rows_affected(result: &[Frame]) -> u64 {
266 if result.len() == 1 {
267 let frame = &result[0];
268 for col in &frame.columns {
269 match col.name.as_str() {
270 "inserted" | "updated" | "deleted" => {
271 if col.data.len() == 1
272 && let Value::Uint8(n) = col.data.get_value(0)
273 {
274 return n;
275 }
276 }
277 _ => {}
278 }
279 }
280 }
281 result.len() as u64
282}
283
284impl Executor {
285 #[instrument(name = "executor::setup_symbols", level = "debug", skip_all)]
286 fn setup_symbols(&self, params: &Params, tx: &mut Transaction<'_>) -> Result<SymbolTable> {
287 let mut symbols = SymbolTable::new();
288 populate_symbols(&mut symbols, params)?;
289 populate_identity(&mut symbols, &self.catalog, tx)?;
290 Ok(symbols)
291 }
292
293 #[instrument(name = "executor::compile", level = "debug", skip(self, tx), fields(rql = %rql))]
294 fn compile_query(&self, tx: &mut Transaction<'_>, rql: &str) -> Result<CompilationResult> {
295 self.compiler.compile_with_policy(tx, rql, inject_from_policies)
296 }
297
298 #[instrument(name = "executor::rql", level = "debug", skip(self, tx, params), fields(rql = %rql))]
299 pub fn rql(&self, tx: &mut Transaction<'_>, rql: &str, params: Params) -> ExecutionResult {
300 let symbols = match self.setup_symbols(¶ms, tx) {
301 Ok(s) => s,
302 Err(e) => return error_result(e, ExecutionMetrics::default()),
303 };
304
305 let start_compile = self.0.runtime_context.clock.instant();
306 let compiled_list = match self.compile_query(tx, rql) {
307 Ok(CompilationResult::Ready(compiled)) => compiled,
308 Ok(CompilationResult::Incremental(_)) => {
309 unreachable!("incremental compilation not supported in rql()")
310 }
311 Err(err) => return self.handle_rql_compile_error(err, rql, params),
312 };
313 let compile_duration = Duration::from_std(start_compile.elapsed());
314
315 match self.run_units_collecting_last(tx, &compiled_list, ¶ms, symbols, compile_duration) {
316 Ok((frames, metrics)) => ExecutionResult {
317 frames,
318 error: None,
319 metrics: build_metrics(metrics),
320 },
321 Err(f) => error_result(f.error, build_metrics(f.partial_metrics)),
322 }
323 }
324
325 #[inline]
326 #[cfg_attr(reifydb_single_threaded, allow(unused_variables))]
327 fn handle_rql_compile_error(&self, err: Error, rql: &str, params: Params) -> ExecutionResult {
328 #[cfg(not(reifydb_single_threaded))]
329 if let Ok(Some(frames)) = self.try_forward_remote_query(&err, rql, params) {
330 return ExecutionResult {
331 frames,
332 error: None,
333 metrics: ExecutionMetrics::default(),
334 };
335 }
336 error_result(err, ExecutionMetrics::default())
337 }
338
339 #[inline]
340 fn run_units_collecting_last(
341 &self,
342 tx: &mut Transaction<'_>,
343 compiled_list: &[Compiled],
344 params: &Params,
345 mut symbols: SymbolTable,
346 compile_duration: Duration,
347 ) -> StdResult<(Vec<Frame>, Vec<StatementMetric>), ExecutionFailure> {
348 let compile_duration_us =
349 compile_duration.to_std().as_micros() as u64 / compiled_list.len().max(1) as u64;
350 let mut result = vec![];
351 let mut metrics = Vec::new();
352 for compiled in compiled_list.iter() {
353 result.clear();
354 let outcome = run_compiled_unit(&self.0, tx, compiled, params, symbols, &mut result);
355 symbols = outcome.symbols;
356
357 metrics.push(StatementMetric {
358 fingerprint: compiled.fingerprint,
359 normalized_rql: compiled.normalized_rql.clone(),
360 compile_duration_us,
361 execute_duration_us: outcome.execute_duration_us,
362 rows_affected: if outcome.run_result.is_ok() {
363 extract_rows_affected(&result)
364 } else {
365 0
366 },
367 });
368
369 if let Err(error) = outcome.run_result {
370 return Err(ExecutionFailure {
371 error,
372 partial_metrics: metrics,
373 });
374 }
375 }
376
377 Ok((result, metrics))
378 }
379
380 #[instrument(name = "executor::admin", level = "debug", skip(self, txn, cmd), fields(rql = %cmd.rql))]
381 pub fn admin(&self, txn: &mut AdminTransaction, cmd: Admin<'_>) -> ExecutionResult {
382 let symbols = match self.setup_symbols(&cmd.params, &mut Transaction::Admin(&mut *txn)) {
383 Ok(s) => s,
384 Err(e) => return error_result(e, ExecutionMetrics::default()),
385 };
386 if let Err(e) = self.enforce_admin_policy(&symbols, txn) {
387 return error_result(e, ExecutionMetrics::default());
388 }
389 let start_compile = self.0.runtime_context.clock.instant();
390 match self.compile_query(&mut Transaction::Admin(txn), cmd.rql) {
391 Err(err) => self.handle_admin_compile_error(err, cmd.rql, cmd.params),
392 Ok(CompilationResult::Ready(compiled)) => {
393 self.execute_admin_ready(txn, compiled, &cmd.params, symbols, start_compile)
394 }
395 Ok(CompilationResult::Incremental(state)) => {
396 self.execute_admin_incremental(txn, state, &cmd.params, symbols)
397 }
398 }
399 }
400
401 #[inline]
402 fn enforce_admin_policy(&self, symbols: &SymbolTable, txn: &mut AdminTransaction) -> Result<()> {
403 PolicyEvaluator::new(&self.0, symbols).enforce_session_policy(
404 &mut Transaction::Admin(txn),
405 SessionOp::Admin,
406 true,
407 )
408 }
409
410 #[inline]
411 #[cfg_attr(reifydb_single_threaded, allow(unused_variables))]
412 fn handle_admin_compile_error(&self, err: Error, rql: &str, params: Params) -> ExecutionResult {
413 #[cfg(not(reifydb_single_threaded))]
414 if let Ok(Some(frames)) = self.try_forward_remote_query(&err, rql, params) {
415 return ExecutionResult {
416 frames,
417 error: None,
418 metrics: ExecutionMetrics::default(),
419 };
420 }
421 error_result(err, ExecutionMetrics::default())
422 }
423
424 #[inline]
425 fn execute_admin_ready(
426 &self,
427 txn: &mut AdminTransaction,
428 compiled: Arc<Vec<Compiled>>,
429 params: &Params,
430 symbols: SymbolTable,
431 start_compile: Instant,
432 ) -> ExecutionResult {
433 let compile_duration = Duration::from_std(start_compile.elapsed());
434 match execute_compiled_units(
435 &self.0,
436 &mut Transaction::Admin(txn),
437 &compiled,
438 params,
439 symbols,
440 compile_duration,
441 ) {
442 Ok((output, remaining, _, metrics)) => ExecutionResult {
443 frames: merge_results(output, remaining),
444 error: None,
445 metrics: build_metrics(metrics),
446 },
447 Err(f) => ExecutionResult {
448 frames: vec![],
449 error: Some(f.error),
450 metrics: build_metrics(f.partial_metrics),
451 },
452 }
453 }
454
455 fn execute_admin_incremental(
456 &self,
457 txn: &mut AdminTransaction,
458 mut state: IncrementalCompilation,
459 params: &Params,
460 symbols: SymbolTable,
461 ) -> ExecutionResult {
462 let policy = constrain_policy(inject_from_policies);
463 let mut result = vec![];
464 let mut output_results: Vec<Frame> = Vec::new();
465 let mut symbols = symbols;
466 let mut metrics = Vec::new();
467 loop {
468 let start_incr = self.0.runtime_context.clock.instant();
469 let next = match self.compiler.compile_next_with_policy(
470 &mut Transaction::Admin(txn),
471 &mut state,
472 &policy,
473 ) {
474 Ok(n) => n,
475 Err(e) => return error_result(e, build_metrics(metrics)),
476 };
477 let compile_duration = start_incr.elapsed();
478
479 let Some(compiled) = next else {
480 break;
481 };
482
483 result.clear();
484 let mut tx = Transaction::Admin(txn);
485 let mut vm = Vm::from_services(symbols, &self.0, params, tx.identity());
486 let start_execute = self.0.runtime_context.clock.instant();
487 let run_result = vm.run(&self.0, &mut tx, &compiled.instructions, &mut result);
488 let execute_duration = start_execute.elapsed();
489 symbols = vm.symbols;
490
491 metrics.push(StatementMetric {
492 fingerprint: compiled.fingerprint,
493 normalized_rql: compiled.normalized_rql,
494 compile_duration_us: compile_duration.as_micros() as u64,
495 execute_duration_us: execute_duration.as_micros() as u64,
496 rows_affected: if run_result.is_ok() {
497 extract_rows_affected(&result)
498 } else {
499 0
500 },
501 });
502
503 if let Err(e) = run_result {
504 return error_result(e, build_metrics(metrics));
505 }
506
507 if compiled.is_output {
508 output_results.append(&mut result);
509 }
510 }
511 ExecutionResult {
512 frames: merge_results(output_results, result),
513 error: None,
514 metrics: build_metrics(metrics),
515 }
516 }
517
518 #[instrument(name = "executor::test", level = "debug", skip(self, txn, cmd), fields(rql = %cmd.rql))]
519 pub fn test(&self, txn: &mut TestTransaction<'_>, cmd: Test<'_>) -> ExecutionResult {
520 let symbols = match self.setup_symbols(&cmd.params, &mut Transaction::Test(Box::new(txn.reborrow()))) {
521 Ok(s) => s,
522 Err(e) => return error_result(e, ExecutionMetrics::default()),
523 };
524 if let Err(e) = self.enforce_test_policy(&symbols, txn) {
525 return error_result(e, ExecutionMetrics::default());
526 }
527 let start_compile = self.0.runtime_context.clock.instant();
528 match self.compiler.compile_with_policy(
529 &mut Transaction::Test(Box::new(txn.reborrow())),
530 cmd.rql,
531 inject_from_policies,
532 ) {
533 Err(err) => self.handle_test_compile_error(err, cmd.rql, cmd.params),
534 Ok(CompilationResult::Ready(compiled)) => {
535 self.execute_test_ready(txn, compiled, &cmd.params, symbols, start_compile)
536 }
537 Ok(CompilationResult::Incremental(state)) => {
538 self.execute_test_incremental(txn, state, &cmd.params, symbols)
539 }
540 }
541 }
542
543 #[inline]
544 fn enforce_test_policy(&self, symbols: &SymbolTable, txn: &mut TestTransaction<'_>) -> Result<()> {
545 let session_type = txn.session_type;
546 let session_default_deny = txn.session_default_deny;
547 PolicyEvaluator::new(&self.0, symbols).enforce_session_policy(
548 &mut Transaction::Test(Box::new(txn.reborrow())),
549 session_type,
550 session_default_deny,
551 )
552 }
553
554 #[inline]
555 #[cfg_attr(reifydb_single_threaded, allow(unused_variables))]
556 fn handle_test_compile_error(&self, err: Error, rql: &str, params: Params) -> ExecutionResult {
557 #[cfg(not(reifydb_single_threaded))]
558 if let Ok(Some(frames)) = self.try_forward_remote_query(&err, rql, params) {
559 return ExecutionResult {
560 frames,
561 error: None,
562 metrics: ExecutionMetrics::default(),
563 };
564 }
565 error_result(err, ExecutionMetrics::default())
566 }
567
568 #[inline]
569 fn execute_test_ready(
570 &self,
571 txn: &mut TestTransaction<'_>,
572 compiled: Arc<Vec<Compiled>>,
573 params: &Params,
574 symbols: SymbolTable,
575 start_compile: Instant,
576 ) -> ExecutionResult {
577 let compile_duration = Duration::from_std(start_compile.elapsed());
578 match execute_compiled_units(
579 &self.0,
580 &mut Transaction::Test(Box::new(txn.reborrow())),
581 &compiled,
582 params,
583 symbols,
584 compile_duration,
585 ) {
586 Ok((output, remaining, _, metrics)) => ExecutionResult {
587 frames: merge_results(output, remaining),
588 error: None,
589 metrics: build_metrics(metrics),
590 },
591 Err(f) => ExecutionResult {
592 frames: vec![],
593 error: Some(f.error),
594 metrics: build_metrics(f.partial_metrics),
595 },
596 }
597 }
598
599 fn execute_test_incremental(
600 &self,
601 txn: &mut TestTransaction<'_>,
602 mut state: IncrementalCompilation,
603 params: &Params,
604 symbols: SymbolTable,
605 ) -> ExecutionResult {
606 let policy = constrain_policy(inject_from_policies);
607 let mut result = vec![];
608 let mut output_results: Vec<Frame> = Vec::new();
609 let mut symbols = symbols;
610 let mut metrics = Vec::new();
611 loop {
612 let start_incr = self.0.runtime_context.clock.instant();
613 let next = match self.compiler.compile_next_with_policy(
614 &mut Transaction::Test(Box::new(txn.reborrow())),
615 &mut state,
616 &policy,
617 ) {
618 Ok(n) => n,
619 Err(e) => return error_result(e, build_metrics(metrics)),
620 };
621 let compile_duration = start_incr.elapsed();
622
623 let Some(compiled) = next else {
624 break;
625 };
626
627 result.clear();
628 let mut tx = Transaction::Test(Box::new(txn.reborrow()));
629 let mut vm = Vm::from_services(symbols, &self.0, params, tx.identity());
630 let start_execute = self.0.runtime_context.clock.instant();
631 let run_result = vm.run(&self.0, &mut tx, &compiled.instructions, &mut result);
632 let execute_duration = start_execute.elapsed();
633 symbols = vm.symbols;
634
635 metrics.push(StatementMetric {
636 fingerprint: compiled.fingerprint,
637 normalized_rql: compiled.normalized_rql,
638 compile_duration_us: compile_duration.as_micros() as u64,
639 execute_duration_us: execute_duration.as_micros() as u64,
640 rows_affected: if run_result.is_ok() {
641 extract_rows_affected(&result)
642 } else {
643 0
644 },
645 });
646
647 if let Err(e) = run_result {
648 return error_result(e, build_metrics(metrics));
649 }
650
651 if compiled.is_output {
652 output_results.append(&mut result);
653 }
654 }
655 ExecutionResult {
656 frames: merge_results(output_results, result),
657 error: None,
658 metrics: build_metrics(metrics),
659 }
660 }
661
662 #[instrument(name = "executor::subscription", level = "debug", skip(self, txn, cmd), fields(rql = %cmd.rql))]
663 pub fn subscription(&self, txn: &mut QueryTransaction, cmd: Subscription<'_>) -> ExecutionResult {
664 let bump = Bump::new();
665 let statements = match parse_str(&bump, cmd.rql) {
666 Ok(s) => s,
667 Err(e) => {
668 return ExecutionResult {
669 frames: vec![],
670 error: Some(e),
671 metrics: ExecutionMetrics::default(),
672 };
673 }
674 };
675
676 if statements.len() != 1 {
677 return ExecutionResult {
678 frames: vec![],
679 error: Some(Error(Box::new(subscription::single_statement_required(
680 "Subscription endpoint requires exactly one statement",
681 )))),
682 metrics: ExecutionMetrics::default(),
683 };
684 }
685
686 let statement = &statements[0];
687 if statement.nodes.len() != 1 || !statement.nodes[0].is_subscription_ddl() {
688 return ExecutionResult {
689 frames: vec![],
690 error: Some(Error(Box::new(subscription::invalid_statement(
691 "Subscription endpoint only supports CREATE SUBSCRIPTION or DROP SUBSCRIPTION",
692 )))),
693 metrics: ExecutionMetrics::default(),
694 };
695 }
696
697 let symbols = match self.setup_symbols(&cmd.params, &mut Transaction::Query(&mut *txn)) {
698 Ok(s) => s,
699 Err(e) => {
700 return ExecutionResult {
701 frames: vec![],
702 error: Some(e),
703 metrics: ExecutionMetrics::default(),
704 };
705 }
706 };
707
708 if let Err(e) = PolicyEvaluator::new(&self.0, &symbols).enforce_session_policy(
709 &mut Transaction::Query(&mut *txn),
710 SessionOp::Subscription,
711 true,
712 ) {
713 return ExecutionResult {
714 frames: vec![],
715 error: Some(e),
716 metrics: ExecutionMetrics::default(),
717 };
718 }
719
720 let start_compile = self.0.runtime_context.clock.instant();
721 let compiled = match self.compiler.compile_with_policy(
722 &mut Transaction::Query(txn),
723 cmd.rql,
724 inject_from_policies,
725 ) {
726 Ok(CompilationResult::Ready(compiled)) => compiled,
727 Ok(CompilationResult::Incremental(_)) => {
728 unreachable!("Single subscription statement should not require incremental compilation")
729 }
730 Err(err) => {
731 return ExecutionResult {
732 frames: vec![],
733 error: Some(err),
734 metrics: ExecutionMetrics::default(),
735 };
736 }
737 };
738 let compile_duration = Duration::from_std(start_compile.elapsed());
739
740 match execute_compiled_units(
741 &self.0,
742 &mut Transaction::Query(txn),
743 &compiled,
744 &cmd.params,
745 symbols,
746 compile_duration,
747 ) {
748 Ok((output, remaining, _, metrics)) => ExecutionResult {
749 frames: merge_results(output, remaining),
750 error: None,
751 metrics: build_metrics(metrics),
752 },
753 Err(f) => ExecutionResult {
754 frames: vec![],
755 error: Some(f.error),
756 metrics: build_metrics(f.partial_metrics),
757 },
758 }
759 }
760
761 #[instrument(name = "executor::command", level = "debug", skip(self, txn, cmd), fields(rql = %cmd.rql))]
762 pub fn command(&self, txn: &mut CommandTransaction, cmd: Command<'_>) -> ExecutionResult {
763 let symbols = match self.setup_symbols(&cmd.params, &mut Transaction::Command(&mut *txn)) {
764 Ok(s) => s,
765 Err(e) => {
766 return ExecutionResult {
767 frames: vec![],
768 error: Some(e),
769 metrics: ExecutionMetrics::default(),
770 };
771 }
772 };
773
774 if let Err(e) = PolicyEvaluator::new(&self.0, &symbols).enforce_session_policy(
775 &mut Transaction::Command(&mut *txn),
776 SessionOp::Command,
777 false,
778 ) {
779 return ExecutionResult {
780 frames: vec![],
781 error: Some(e),
782 metrics: ExecutionMetrics::default(),
783 };
784 }
785
786 let start_compile = self.0.runtime_context.clock.instant();
787 let compiled = match self.compile_query(&mut Transaction::Command(txn), cmd.rql) {
788 Ok(CompilationResult::Ready(compiled)) => compiled,
789 Ok(CompilationResult::Incremental(_)) => {
790 unreachable!("DDL statements require admin transactions, not command transactions")
791 }
792 Err(err) => {
793 #[cfg(not(reifydb_single_threaded))]
794 if self.0.remote_registry.is_some() && remote::is_remote_query(&err) {
795 return ExecutionResult {
796 frames: vec![],
797 error: Some(Error(Box::new(Diagnostic {
798 code: "REMOTE_002".to_string(),
799 message: "Write operations on remote namespaces are not supported"
800 .to_string(),
801 help: Some("Use the remote instance directly for write operations"
802 .to_string()),
803 ..Default::default()
804 }))),
805 metrics: ExecutionMetrics::default(),
806 };
807 }
808 return ExecutionResult {
809 frames: vec![],
810 error: Some(err),
811 metrics: ExecutionMetrics::default(),
812 };
813 }
814 };
815 let compile_duration = Duration::from_std(start_compile.elapsed());
816
817 match execute_compiled_units(
818 &self.0,
819 &mut Transaction::Command(txn),
820 &compiled,
821 &cmd.params,
822 symbols,
823 compile_duration,
824 ) {
825 Ok((output, remaining, _, metrics)) => ExecutionResult {
826 frames: merge_results(output, remaining),
827 error: None,
828 metrics: build_metrics(metrics),
829 },
830 Err(f) => ExecutionResult {
831 frames: vec![],
832 error: Some(f.error),
833 metrics: build_metrics(f.partial_metrics),
834 },
835 }
836 }
837
838 #[instrument(name = "executor::call_procedure", level = "debug", skip(self, txn, params), fields(name = %name))]
839 pub fn call_procedure(&self, txn: &mut CommandTransaction, name: &str, params: &Params) -> ExecutionResult {
840 let rql = format!("CALL {}()", name);
841 let symbols = match self.setup_symbols(params, &mut Transaction::Command(&mut *txn)) {
842 Ok(s) => s,
843 Err(e) => return error_result(e, ExecutionMetrics::default()),
844 };
845
846 let start_compile = self.0.runtime_context.clock.instant();
847 let compiled = match self.compiler.compile(&mut Transaction::Command(txn), &rql) {
848 Ok(CompilationResult::Ready(compiled)) => compiled,
849 Ok(CompilationResult::Incremental(_)) => {
850 unreachable!("CALL statements should not require incremental compilation")
851 }
852 Err(e) => return error_result(e, ExecutionMetrics::default()),
853 };
854 let compile_duration = Duration::from_std(start_compile.elapsed());
855
856 match self.run_command_units_collecting_last(txn, &compiled, params, symbols, compile_duration) {
857 Ok((frames, metrics)) => ExecutionResult {
858 frames,
859 error: None,
860 metrics: build_metrics(metrics),
861 },
862 Err(f) => error_result(f.error, build_metrics(f.partial_metrics)),
863 }
864 }
865
866 #[inline]
867 fn run_command_units_collecting_last(
868 &self,
869 txn: &mut CommandTransaction,
870 compiled: &[Compiled],
871 params: &Params,
872 mut symbols: SymbolTable,
873 compile_duration: Duration,
874 ) -> StdResult<(Vec<Frame>, Vec<StatementMetric>), ExecutionFailure> {
875 let compile_duration_us = compile_duration.to_std().as_micros() as u64 / compiled.len().max(1) as u64;
876 let mut result = vec![];
877 let mut metrics = Vec::new();
878 for compiled in compiled.iter() {
879 result.clear();
880 let mut tx = Transaction::Command(txn);
881 let mut vm = Vm::from_services(symbols, &self.0, params, tx.identity());
882 let start_execute = self.0.runtime_context.clock.instant();
883 let run_result = vm.run(&self.0, &mut tx, &compiled.instructions, &mut result);
884 let execute_duration = start_execute.elapsed();
885 symbols = vm.symbols;
886
887 metrics.push(StatementMetric {
888 fingerprint: compiled.fingerprint,
889 normalized_rql: compiled.normalized_rql.clone(),
890 compile_duration_us,
891 execute_duration_us: execute_duration.as_micros() as u64,
892 rows_affected: if run_result.is_ok() {
893 extract_rows_affected(&result)
894 } else {
895 0
896 },
897 });
898
899 if let Err(error) = run_result {
900 return Err(ExecutionFailure {
901 error,
902 partial_metrics: metrics,
903 });
904 }
905 }
906
907 Ok((result, metrics))
908 }
909
910 #[instrument(name = "executor::query", level = "debug", skip(self, txn, qry), fields(rql = %qry.rql))]
911 pub fn query(&self, txn: &mut QueryTransaction, qry: Query<'_>) -> ExecutionResult {
912 let symbols = match self.setup_symbols(&qry.params, &mut Transaction::Query(&mut *txn)) {
913 Ok(s) => s,
914 Err(e) => {
915 return ExecutionResult {
916 frames: vec![],
917 error: Some(e),
918 metrics: ExecutionMetrics::default(),
919 };
920 }
921 };
922
923 if let Err(e) = PolicyEvaluator::new(&self.0, &symbols).enforce_session_policy(
924 &mut Transaction::Query(&mut *txn),
925 SessionOp::Query,
926 false,
927 ) {
928 return ExecutionResult {
929 frames: vec![],
930 error: Some(e),
931 metrics: ExecutionMetrics::default(),
932 };
933 }
934
935 let start_compile = self.0.runtime_context.clock.instant();
936 let compiled = match self.compile_query(&mut Transaction::Query(txn), qry.rql) {
937 Ok(CompilationResult::Ready(compiled)) => compiled,
938 Ok(CompilationResult::Incremental(_)) => {
939 unreachable!("DDL statements require admin transactions, not query transactions")
940 }
941 Err(err) => {
942 #[cfg(not(reifydb_single_threaded))]
943 if let Ok(Some(frames)) = self.try_forward_remote_query(&err, qry.rql, qry.params) {
944 return ExecutionResult {
945 frames,
946 error: None,
947 metrics: ExecutionMetrics::default(),
948 };
949 }
950 return ExecutionResult {
951 frames: vec![],
952 error: Some(err),
953 metrics: ExecutionMetrics::default(),
954 };
955 }
956 };
957 let compile_duration = Duration::from_std(start_compile.elapsed());
958
959 let exec_result = execute_compiled_units(
960 &self.0,
961 &mut Transaction::Query(txn),
962 &compiled,
963 &qry.params,
964 symbols,
965 compile_duration,
966 );
967
968 match exec_result {
969 Ok((output, remaining, _, metrics)) => ExecutionResult {
970 frames: merge_results(output, remaining),
971 error: None,
972 metrics: build_metrics(metrics),
973 },
974 Err(f) => ExecutionResult {
975 frames: vec![],
976 error: Some(f.error),
977 metrics: build_metrics(f.partial_metrics),
978 },
979 }
980 }
981}