Skip to main content

reifydb_engine/vm/volcano/
apply_transform.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::sync::Arc;
5
6use reifydb_core::value::column::{columns::Columns, headers::ColumnHeaders};
7use reifydb_extension::transform::{Transform, context::TransformContext};
8use reifydb_transaction::transaction::Transaction;
9use reifydb_value::reifydb_assertions;
10use tracing::instrument;
11
12use crate::{
13	Result,
14	vm::volcano::query::{QueryContext, QueryNode},
15};
16
17pub(crate) struct ApplyTransformNode {
18	input: Box<dyn QueryNode>,
19	transform: Box<dyn Transform>,
20	context: Option<Arc<QueryContext>>,
21}
22
23impl ApplyTransformNode {
24	pub fn new(input: Box<dyn QueryNode>, transform: Box<dyn Transform>) -> Self {
25		Self {
26			input,
27			transform,
28			context: None,
29		}
30	}
31}
32
33impl QueryNode for ApplyTransformNode {
34	#[instrument(level = "trace", skip_all, name = "volcano::apply_transform::initialize")]
35	fn initialize<'a>(&mut self, rx: &mut Transaction<'a>, ctx: &QueryContext) -> Result<()> {
36		self.context = Some(Arc::new(ctx.clone()));
37		self.input.initialize(rx, ctx)?;
38		Ok(())
39	}
40
41	#[instrument(level = "trace", skip_all, name = "volcano::apply_transform::next")]
42	fn next<'a>(&mut self, rx: &mut Transaction<'a>, ctx: &mut QueryContext) -> Result<Option<Columns>> {
43		reifydb_assertions! {
44			assert!(self.context.is_some(), "ApplyTransformNode::next() called before initialize()");
45		}
46		let stored_ctx = self.context.as_ref().unwrap();
47
48		if let Some(columns) = self.input.next(rx, ctx)? {
49			let transform_ctx = TransformContext {
50				routines: &ctx.services.routines,
51				runtime_context: &stored_ctx.services.runtime_context,
52				params: &stored_ctx.params,
53			};
54			let result = self.transform.apply(&transform_ctx, columns)?;
55			Ok(Some(result))
56		} else {
57			Ok(None)
58		}
59	}
60
61	fn headers(&self) -> Option<ColumnHeaders> {
62		self.input.headers()
63	}
64}