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