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