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::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(&params, 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, &params, 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}