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