Skip to main content

reifydb_sub_flow/operator/
map.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::sync::LazyLock;
5
6use reifydb_abi::operator::capabilities::OperatorCapability;
7use reifydb_core::{
8	interface::{
9		catalog::flow::FlowNodeId,
10		change::{Change, Diff},
11	},
12	value::column::{ColumnWithName, columns::Columns},
13};
14use reifydb_engine::{
15	expression::{
16		compile::{CompiledExpr, compile_expression},
17		context::{CompileContext, EvalContext},
18	},
19	vm::stack::SymbolTable,
20};
21use reifydb_routine::routine::registry::Routines;
22use reifydb_rql::expression::{Expression, name::display_label};
23use reifydb_runtime::context::RuntimeContext;
24use reifydb_value::{Result, fragment::Fragment, params::Params, value::identity::IdentityId};
25
26use crate::{Operator, operator::OperatorCell, transaction::FlowTransaction};
27
28static EMPTY_PARAMS: Params = Params::None;
29static EMPTY_SYMBOL_TABLE: LazyLock<SymbolTable> = LazyLock::new(SymbolTable::new);
30
31pub struct MapOperator {
32	parent: OperatorCell,
33	node: FlowNodeId,
34	expressions: Vec<Expression>,
35	compiled_expressions: Vec<CompiledExpr>,
36	routines: Routines,
37	runtime_context: RuntimeContext,
38}
39
40impl MapOperator {
41	pub fn new(
42		parent: OperatorCell,
43		node: FlowNodeId,
44		expressions: Vec<Expression>,
45		routines: Routines,
46		runtime_context: RuntimeContext,
47	) -> Self {
48		let compile_ctx = CompileContext {
49			symbols: &EMPTY_SYMBOL_TABLE,
50		};
51		let compiled_expressions: Vec<CompiledExpr> = expressions
52			.iter()
53			.map(|e| compile_expression(&compile_ctx, e))
54			.collect::<Result<Vec<_>>>()
55			.expect("Failed to compile expressions");
56
57		Self {
58			parent,
59			node,
60			expressions,
61			compiled_expressions,
62			routines,
63			runtime_context,
64		}
65	}
66
67	pub(crate) fn output_schema(&self) -> Option<Columns> {
68		self.parent.output_schema()
69	}
70
71	fn project(&self, columns: &Columns) -> Result<Columns> {
72		let row_count = columns.row_count();
73		if row_count == 0 {
74			return Ok(Columns::empty());
75		}
76
77		let session = EvalContext {
78			params: &EMPTY_PARAMS,
79			symbols: &EMPTY_SYMBOL_TABLE,
80			routines: &self.routines,
81			runtime_context: &self.runtime_context,
82			arena: None,
83			identity: IdentityId::root(),
84			is_aggregate_context: false,
85			columns: Columns::empty(),
86			row_count: 1,
87			target: None,
88			take: None,
89		};
90		let exec_ctx = session.with_eval(columns.clone(), row_count);
91
92		let mut result_columns = Vec::with_capacity(self.expressions.len());
93
94		for (i, compiled_expr) in self.compiled_expressions.iter().enumerate() {
95			let evaluated_col = compiled_expr.execute(&exec_ctx)?;
96
97			let expr = &self.expressions[i];
98			let field_name = display_label(expr).text().to_string();
99
100			let named_column =
101				ColumnWithName::new(Fragment::internal(field_name), evaluated_col.data().clone());
102
103			result_columns.push(named_column);
104		}
105
106		let row_numbers = if columns.row_numbers.is_empty() {
107			Vec::new()
108		} else {
109			columns.row_numbers.iter().cloned().collect()
110		};
111
112		Ok(Columns::with_system_columns(
113			result_columns,
114			row_numbers,
115			columns.created_at.to_vec(),
116			columns.updated_at.to_vec(),
117		))
118	}
119}
120
121impl Operator for MapOperator {
122	fn id(&self) -> FlowNodeId {
123		self.node
124	}
125
126	fn capabilities(&self) -> &[OperatorCapability] {
127		OperatorCapability::STANDARD
128	}
129
130	fn apply(&self, _txn: &mut FlowTransaction, change: Change) -> Result<Change> {
131		let mut result = Vec::new();
132
133		for diff in change.diffs.into_iter() {
134			match diff {
135				Diff::Insert {
136					post,
137					..
138				} => {
139					let projected = match self.project(&post) {
140						Ok(projected) => projected,
141						Err(err) => {
142							panic!("{:#?}", err)
143						}
144					};
145
146					if !projected.is_empty() {
147						result.push(Diff::insert(projected));
148					}
149				}
150				Diff::Update {
151					pre,
152					post,
153					..
154				} => {
155					let projected_post = self.project(&post)?;
156					let projected_pre = self.project(&pre)?;
157
158					if !projected_post.is_empty() {
159						result.push(Diff::update(projected_pre, projected_post));
160					}
161				}
162				Diff::Remove {
163					pre,
164					..
165				} => {
166					let projected_pre = self.project(&pre)?;
167					if !projected_pre.is_empty() {
168						result.push(Diff::remove(projected_pre));
169					}
170				}
171			}
172		}
173
174		Ok(Change::from_flow(self.node, change.version, result, change.changed_at))
175	}
176}