akar_processor/processor/mapper/
map_update.rs1use super::ExecutionContext;
2use crate::physical_operator::*;
3use akar_common::error::ProcessorError;
4use akar_common::vector::DataChunk;
5use akar_planner::logical_operator::LogicalOperator;
6
7pub fn map_and_execute_update(
8 op: &LogicalOperator,
9 current_input: Vec<DataChunk>,
10 ctx: &mut ExecutionContext,
11) -> Result<Vec<DataChunk>, ProcessorError> {
12 match op {
13 LogicalOperator::Set(sl) => {
14 let table_catalog = ctx
15 .table_catalog
16 .clone()
17 .ok_or_else(|| "No table catalog available for SET".to_string())?;
18
19 let set_op = PhysicalSet {
20 table_name: sl.table_name.clone(),
21 table_id: sl.table_id,
22 column_name: sl.column_name.clone(),
23 column_idx: sl.column_idx,
24 value: sl.value.clone(),
25 is_node: sl.is_node,
26 table_catalog,
27 };
28 let result = set_op.execute(current_input)?;
29 record_set_writes(sl.table_id, &result, ctx);
31 Ok(result)
32 }
33 LogicalOperator::Delete(dl) => {
34 let table_catalog = ctx
35 .table_catalog
36 .clone()
37 .ok_or_else(|| "No table catalog available for DELETE".to_string())?;
38
39 let delete_op = PhysicalDelete {
40 table_name: dl.table_name.clone(),
41 table_id: dl.table_id,
42 primary_key_column: dl.primary_key_column.clone(),
43 is_node: dl.is_node,
44 detach: dl.detach,
45 row_indices: Vec::new(),
46 table_catalog,
47 };
48 let result = delete_op.execute(current_input)?;
49 record_delete_writes(dl.table_id, &result, ctx);
51 Ok(result)
52 }
53 LogicalOperator::CreateNode(cn) => {
54 let table_catalog = ctx
55 .table_catalog
56 .clone()
57 .ok_or_else(|| "No table catalog available for CREATE".to_string())?;
58
59 let create_node_op = PhysicalInsertNode {
60 table_name: cn.table_name.clone(),
61 table_id: cn.table_id,
62 out_var_name: cn.out_var_name.clone(),
63 properties: cn.properties.clone(),
64 table_catalog,
65 };
66 let result = create_node_op.execute(current_input)?;
67 record_insert_writes(cn.table_id, &result, ctx);
69 Ok(result)
70 }
71 LogicalOperator::CreateRel(cr) => {
72 let table_catalog = ctx
73 .table_catalog
74 .clone()
75 .ok_or_else(|| "No table catalog available for CREATE".to_string())?;
76
77 let create_rel_op = PhysicalInsertRel {
78 table_name: cr.table_name.clone(),
79 table_id: cr.table_id,
80 src_node_name: cr.src_node_name.clone(),
81 dst_node_name: cr.dst_node_name.clone(),
82 properties: cr.properties.clone(),
83 table_catalog,
84 };
85 let result = create_rel_op.execute(current_input)?;
86 record_insert_writes(cr.table_id, &result, ctx);
88 Ok(result)
89 }
90 LogicalOperator::Extend(ex) => {
91 let table_catalog = ctx
92 .table_catalog
93 .clone()
94 .ok_or_else(|| "No table catalog available for Extend".to_string())?;
95
96 let extend_op = PhysicalExtend {
97 rel_table_name: ex.rel_table_name.clone(),
98 rel_table_id: ex.rel_table_id,
99 rel_var: ex.rel_var.clone(),
100 bound_node_var: ex.bound_node_var.clone(),
101 direction: ex.direction.clone(),
102 dst_node_var: ex.dst_node_var.clone(),
103 dst_table_name: ex.dst_table_name.clone(),
104 dst_table_id: ex.dst_table_id,
105 table_catalog,
106 };
107 let result = extend_op.execute(current_input)?;
108 record_insert_writes(ex.rel_table_id, &result, ctx);
110 Ok(result)
111 }
112 LogicalOperator::Merge(m) => {
113 let table_catalog = ctx
114 .table_catalog
115 .clone()
116 .ok_or_else(|| "No table catalog available for MERGE".to_string())?;
117
118 let mut on_match_ops = Vec::new();
119 for set_item in &m.on_match {
120 on_match_ops.push(PhysicalSet {
121 table_name: set_item.table_name.clone(),
122 table_id: set_item.table_id,
123 column_name: set_item.column_name.clone(),
124 column_idx: set_item.column_idx,
125 value: set_item.value.clone(),
126 is_node: set_item.is_node,
127 table_catalog: table_catalog.clone(),
128 });
129 }
130
131 let mut on_create_ops = Vec::new();
132 for set_item in &m.on_create {
133 on_create_ops.push(PhysicalSet {
134 table_name: set_item.table_name.clone(),
135 table_id: set_item.table_id,
136 column_name: set_item.column_name.clone(),
137 column_idx: set_item.column_idx,
138 value: set_item.value.clone(),
139 is_node: set_item.is_node,
140 table_catalog: table_catalog.clone(),
141 });
142 }
143
144 let merge_op = PhysicalMerge {
145 table_name: m.table_name.clone(),
146 table_id: m.table_id,
147 properties: m.properties.clone(),
148 on_match: on_match_ops,
149 on_create: on_create_ops,
150 table_catalog,
151 };
152 let result = merge_op.execute(current_input)?;
153 record_insert_writes(m.table_id, &result, ctx);
155 Ok(result)
156 }
157 LogicalOperator::CopyFrom(cf) => {
158 let table_catalog = ctx
159 .table_catalog
160 .clone()
161 .ok_or_else(|| "No table catalog available for COPY FROM".to_string())?;
162
163 let columns = if let Some(node_table) = table_catalog.get_node_table_by_name(&cf.table_name) {
165 node_table.columns.clone()
166 } else if let Some(rel_table) = table_catalog.get_rel_table_by_name(&cf.table_name) {
167 rel_table.columns.clone()
168 } else {
169 return Err(format!("Table '{}' not found in storage catalog", cf.table_name).into());
170 };
171
172 let copy_op = PhysicalCopyFrom {
173 table_name: cf.table_name.clone(),
174 table_id: cf.table_id,
175 file_path: cf.file_path.clone(),
176 columns,
177 options: cf.options.clone(),
178 table_catalog,
179 vfs: ctx
180 .vfs
181 .clone()
182 .ok_or_else(|| "VFS not initialized in processor".to_string())?,
183 };
184 let result = copy_op.execute(current_input)?;
185 record_insert_writes(cf.table_id, &result, ctx);
187 Ok(result)
188 }
189 LogicalOperator::BatchInsert(bi) => {
190 let table_catalog = ctx
191 .table_catalog
192 .clone()
193 .ok_or_else(|| "No table catalog available for BATCH INSERT".to_string())?;
194
195 let batch_op = PhysicalBatchInsert {
196 table_name: bi.table_name.clone(),
197 table_id: bi.table_id,
198 rows: bi.rows.clone(),
199 table_catalog,
200 };
201 let result = batch_op.execute(current_input)?;
202 record_insert_writes(bi.table_id, &result, ctx);
204 Ok(result)
205 }
206 LogicalOperator::Insert(i) => {
207 let exec = crate::physical::misc::PhysicalInsert {
208 table_name: i.table_name.clone(),
209 table_id: i.table_id,
210 columns: i.columns.clone(),
211 values: i.values.clone(),
212 table_catalog: ctx.table_catalog.clone().unwrap(),
213 };
214 let result = exec.execute(current_input)?;
215 record_insert_writes(i.table_id, &result, ctx);
217 Ok(result)
218 }
219 _ => Err(format!("Not an update operator: {:?}", op).into()),
220 }
221}
222
223fn record_set_writes(table_id: u64, result: &[DataChunk], ctx: &mut ExecutionContext) {
226 if let Some(chunk) = result.first() {
227 for row in 0..chunk.size {
228 if !chunk.fields.is_empty() {
229 if let Some(akar_common::types::Value::Int64(row_idx)) = chunk.get_value(0, row) {
230 ctx.written_rows.push((table_id, row_idx as u64));
231 }
232 }
233 }
234 }
235}
236
237fn record_delete_writes(table_id: u64, result: &[DataChunk], ctx: &mut ExecutionContext) {
240 if let Some(chunk) = result.first() {
241 for row in 0..chunk.size {
242 if !chunk.fields.is_empty() {
243 if let Some(akar_common::types::Value::Int64(row_idx)) = chunk.get_value(0, row) {
244 ctx.written_rows.push((table_id, row_idx as u64));
245 }
246 }
247 }
248 }
249}
250
251fn record_insert_writes(table_id: u64, result: &[DataChunk], ctx: &mut ExecutionContext) {
256 if let Some(chunk) = result.first() {
257 if chunk.fields.len() > 1 {
259 for row in 0..chunk.fields[1].len() {
260 if let Some(akar_common::types::Value::Int64(row_id)) = chunk.get_value(1, row) {
261 ctx.written_rows.push((table_id, row_id as u64));
262 }
263 }
264 }
265 }
266}