reifydb_engine/vm/volcano/
apply_transform.rs1use 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}