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