Skip to main content

reifydb_engine/vm/
executor.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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(&params, 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, &params, 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(&params, &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}