reifydb_sub_flow/operator/
map.rs1use 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}