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