formualizer_eval/engine/arrow_ingest.rs
1use crate::arrow_store::{ArrowSheet, IngestBuilder};
2use crate::engine::Engine;
3use crate::traits::EvaluationContext;
4use formualizer_common::{ExcelError, LiteralValue};
5use rustc_hash::FxHashMap;
6
7#[derive(Debug, Clone, Default)]
8pub struct ArrowBulkIngestSummary {
9 pub sheets: usize,
10 pub total_rows: usize,
11}
12
13/// Bulk Arrow ingest builder for Phase A base values.
14pub struct ArrowBulkIngestBuilder<'e, R: EvaluationContext> {
15 engine: &'e mut Engine<R>,
16 builders: FxHashMap<String, IngestBuilder>,
17 rows: FxHashMap<String, usize>,
18}
19
20impl<'e, R: EvaluationContext> ArrowBulkIngestBuilder<'e, R> {
21 pub fn new(engine: &'e mut Engine<R>) -> Self {
22 Self {
23 engine,
24 builders: FxHashMap::default(),
25 rows: FxHashMap::default(),
26 }
27 }
28
29 /// Add a sheet ingest target. Creates or replaces any existing Arrow sheet on finish.
30 pub fn add_sheet(&mut self, name: &str, ncols: usize, chunk_rows: usize) {
31 let ib = IngestBuilder::new(name, ncols, chunk_rows, self.engine.config.date_system);
32 self.builders.insert(name.to_string(), ib);
33 self.rows.insert(name.to_string(), 0);
34 self.engine.sheet_id_mut(name);
35 }
36
37 /// Append a single row of values for the given sheet (0-based col order, length == ncols).
38 pub fn append_row(&mut self, name: &str, row: &[LiteralValue]) -> Result<(), ExcelError> {
39 let ib = self
40 .builders
41 .get_mut(name)
42 .expect("sheet must be added before append_row");
43 ib.append_row(row)?;
44 *self.rows.get_mut(name).unwrap() += 1;
45 Ok(())
46 }
47
48 /// Finish all sheet builders, installing ArrowSheets into the engine store.
49 pub fn finish(mut self) -> Result<ArrowBulkIngestSummary, ExcelError> {
50 let mut sheets: Vec<(String, ArrowSheet)> = Vec::with_capacity(self.builders.len());
51 for (name, builder) in self.builders.drain() {
52 let sheet = builder.finish();
53 sheets.push((name, sheet));
54 }
55 // Insert or replace by name
56 for (name, sheet) in sheets {
57 let store = self.engine.sheet_store_mut();
58 if let Some(pos) = store.sheets.iter().position(|s| s.name.as_ref() == name) {
59 store.sheets[pos] = sheet;
60 } else {
61 store.sheets.push(sheet);
62 }
63 }
64 let total_rows = self.rows.values().copied().sum();
65 Ok(ArrowBulkIngestSummary {
66 sheets: self.rows.len(),
67 total_rows,
68 })
69 }
70}
71
72/// Bulk Arrow update builder for Phase C. Chooses overlay vs rebuild per chunk.
73pub struct ArrowBulkUpdateBuilder<'e, R: EvaluationContext> {
74 engine: &'e mut Engine<R>,
75 // sheet -> col0 -> row0 -> value
76 updates: FxHashMap<String, FxHashMap<usize, FxHashMap<usize, LiteralValue>>>,
77}
78
79impl<'e, R: EvaluationContext> ArrowBulkUpdateBuilder<'e, R> {
80 pub fn new(engine: &'e mut Engine<R>) -> Self {
81 Self {
82 engine,
83 updates: FxHashMap::default(),
84 }
85 }
86
87 pub fn update_cell(&mut self, sheet: &str, row: u32, col: u32, value: LiteralValue) {
88 let s = self.updates.entry(sheet.to_string()).or_default();
89 let c = s.entry(col.saturating_sub(1) as usize).or_default();
90 c.insert(row.saturating_sub(1) as usize, value);
91 }
92
93 pub fn finish(mut self) -> Result<usize, ExcelError> {
94 use std::sync::Arc;
95 let date_system = self.engine.config.date_system;
96 let mut total = 0usize;
97 for (sheet_name, by_col) in self.updates.drain() {
98 let maybe_sheet = self.engine.sheet_store_mut().sheet_mut(&sheet_name);
99 if maybe_sheet.is_none() {
100 continue;
101 }
102 let sheet = maybe_sheet.unwrap();
103 for (col0, rows_map) in by_col {
104 total += rows_map.len();
105 if col0 >= sheet.columns.len() {
106 continue;
107 }
108 // Partition by chunk
109 let mut by_chunk: FxHashMap<usize, Vec<(usize, LiteralValue)>> =
110 FxHashMap::default();
111 for (row0, v) in rows_map {
112 if row0 >= sheet.nrows as usize {
113 sheet.ensure_row_capacity(row0 + 1);
114 }
115 if let Some((ch_idx, in_off)) = sheet.chunk_of_row(row0) {
116 by_chunk.entry(ch_idx).or_default().push((in_off, v));
117 }
118 }
119 for (ch_idx, mut items) in by_chunk {
120 let Some(ch) = sheet.ensure_column_chunk_mut(col0, ch_idx) else {
121 continue;
122 };
123 let len = ch.type_tag.len();
124 // heuristic: rebuild if > 2% or > 1024 updates in this chunk
125 let rebuild = items.len() > len / 50 || items.len() > 1024;
126 if !rebuild {
127 // overlay
128 for (off, v) in items {
129 let ov = match v {
130 LiteralValue::Empty => crate::arrow_store::OverlayValue::Empty,
131 LiteralValue::Int(i) => {
132 crate::arrow_store::OverlayValue::Number(i as f64)
133 }
134 LiteralValue::Number(n) => {
135 crate::arrow_store::OverlayValue::Number(n)
136 }
137 LiteralValue::Boolean(b) => {
138 crate::arrow_store::OverlayValue::Boolean(b)
139 }
140 LiteralValue::Text(s) => {
141 crate::arrow_store::OverlayValue::Text(Arc::from(s))
142 }
143 LiteralValue::Error(e) => crate::arrow_store::OverlayValue::Error(
144 crate::arrow_store::map_error_code(e.kind),
145 ),
146 LiteralValue::Date(d) => {
147 let dt = d.and_hms_opt(0, 0, 0).unwrap();
148 let serial = formualizer_common::datetime_to_serial_for(
149 date_system,
150 &dt,
151 );
152 crate::arrow_store::OverlayValue::DateTime(serial)
153 }
154 LiteralValue::DateTime(dt) => {
155 let serial = formualizer_common::datetime_to_serial_for(
156 date_system,
157 &dt,
158 );
159 crate::arrow_store::OverlayValue::DateTime(serial)
160 }
161 LiteralValue::Time(t) => {
162 let serial = formualizer_common::time_to_fraction(&t);
163 crate::arrow_store::OverlayValue::DateTime(serial)
164 }
165 LiteralValue::Duration(d) => {
166 let serial = d.num_seconds() as f64 / 86_400.0;
167 crate::arrow_store::OverlayValue::Duration(serial)
168 }
169 LiteralValue::Pending => crate::arrow_store::OverlayValue::Pending,
170 LiteralValue::Array(_) => crate::arrow_store::OverlayValue::Error(
171 crate::arrow_store::map_error_code(
172 formualizer_common::ExcelErrorKind::Value,
173 ),
174 ),
175 };
176 let _ = ch.overlay.set(off, ov);
177 }
178 } else {
179 // rebuild chunk with updates applied
180 use arrow_array::Array as _;
181 use arrow_array::builder::{
182 BooleanBuilder, Float64Builder, StringBuilder, UInt8Builder,
183 };
184 items.sort_by_key(|(o, _)| *o);
185 let mut tag_b = UInt8Builder::with_capacity(len);
186 let mut nb = Float64Builder::with_capacity(len);
187 let mut bb = BooleanBuilder::with_capacity(len);
188 let mut sb = StringBuilder::with_capacity(len, len * 8);
189 let mut eb = UInt8Builder::with_capacity(len);
190 let mut non_num = 0usize;
191 let mut non_bool = 0usize;
192 let mut non_text = 0usize;
193 let mut non_err = 0usize;
194 let mut it = items.into_iter().peekable();
195 for i in 0..len {
196 let upd = if it.peek().map(|(o, _)| *o == i).unwrap_or(false) {
197 Some(it.next().unwrap().1)
198 } else {
199 None
200 };
201 let val = if let Some(v) = upd {
202 v
203 } else {
204 // read from base tag/lane
205 let t = crate::arrow_store::TypeTag::from_u8(ch.type_tag.value(i));
206 match t {
207 crate::arrow_store::TypeTag::Empty => LiteralValue::Empty,
208 crate::arrow_store::TypeTag::Number => {
209 if let Some(a) = &ch.numbers {
210 let fa = a
211 .as_any()
212 .downcast_ref::<arrow_array::Float64Array>()
213 .unwrap();
214 if fa.is_null(i) {
215 LiteralValue::Empty
216 } else {
217 LiteralValue::Number(fa.value(i))
218 }
219 } else {
220 LiteralValue::Empty
221 }
222 }
223 crate::arrow_store::TypeTag::DateTime => {
224 if let Some(a) = &ch.numbers {
225 let fa = a
226 .as_any()
227 .downcast_ref::<arrow_array::Float64Array>()
228 .unwrap();
229 if fa.is_null(i) {
230 LiteralValue::Empty
231 } else {
232 LiteralValue::try_from_serial_number_for(
233 date_system,
234 fa.value(i),
235 )
236 .unwrap_or_else(LiteralValue::Error)
237 }
238 } else {
239 LiteralValue::Empty
240 }
241 }
242 crate::arrow_store::TypeTag::Duration => {
243 if let Some(a) = &ch.numbers {
244 let fa = a
245 .as_any()
246 .downcast_ref::<arrow_array::Float64Array>()
247 .unwrap();
248 if fa.is_null(i) {
249 LiteralValue::Empty
250 } else {
251 let serial = fa.value(i);
252 let nanos_f = serial * 86_400.0 * 1_000_000_000.0;
253 let nanos = nanos_f
254 .round()
255 .clamp(i64::MIN as f64, i64::MAX as f64)
256 as i64;
257 LiteralValue::Duration(
258 chrono::Duration::nanoseconds(nanos),
259 )
260 }
261 } else {
262 LiteralValue::Empty
263 }
264 }
265 crate::arrow_store::TypeTag::Boolean => {
266 if let Some(a) = &ch.booleans {
267 let ba = a
268 .as_any()
269 .downcast_ref::<arrow_array::BooleanArray>()
270 .unwrap();
271 if ba.is_null(i) {
272 LiteralValue::Empty
273 } else {
274 LiteralValue::Boolean(ba.value(i))
275 }
276 } else {
277 LiteralValue::Empty
278 }
279 }
280 crate::arrow_store::TypeTag::Text => {
281 if let Some(a) = &ch.text {
282 let sa = a
283 .as_any()
284 .downcast_ref::<arrow_array::StringArray>()
285 .unwrap();
286 if sa.is_null(i) {
287 LiteralValue::Empty
288 } else {
289 LiteralValue::Text(sa.value(i).to_string())
290 }
291 } else {
292 LiteralValue::Empty
293 }
294 }
295 crate::arrow_store::TypeTag::Error => {
296 if let Some(a) = &ch.errors {
297 let ea = a
298 .as_any()
299 .downcast_ref::<arrow_array::UInt8Array>()
300 .unwrap();
301 if ea.is_null(i) {
302 LiteralValue::Empty
303 } else {
304 LiteralValue::Error(ExcelError::new(
305 crate::arrow_store::unmap_error_code(
306 ea.value(i),
307 ),
308 ))
309 }
310 } else {
311 LiteralValue::Empty
312 }
313 }
314 crate::arrow_store::TypeTag::Pending => LiteralValue::Pending,
315 }
316 };
317 match val {
318 LiteralValue::Empty => {
319 tag_b.append_value(crate::arrow_store::TypeTag::Empty as u8);
320 nb.append_null();
321 bb.append_null();
322 sb.append_null();
323 eb.append_null();
324 }
325 LiteralValue::Int(i) => {
326 tag_b.append_value(crate::arrow_store::TypeTag::Number as u8);
327 nb.append_value(i as f64);
328 non_num += 1;
329 bb.append_null();
330 sb.append_null();
331 eb.append_null();
332 }
333 LiteralValue::Number(n) => {
334 tag_b.append_value(crate::arrow_store::TypeTag::Number as u8);
335 nb.append_value(n);
336 non_num += 1;
337 bb.append_null();
338 sb.append_null();
339 eb.append_null();
340 }
341 LiteralValue::Boolean(b) => {
342 tag_b.append_value(crate::arrow_store::TypeTag::Boolean as u8);
343 nb.append_null();
344 bb.append_value(b);
345 non_bool += 1;
346 sb.append_null();
347 eb.append_null();
348 }
349 LiteralValue::Text(s) => {
350 tag_b.append_value(crate::arrow_store::TypeTag::Text as u8);
351 nb.append_null();
352 bb.append_null();
353 sb.append_value(&s);
354 non_text += 1;
355 eb.append_null();
356 }
357 LiteralValue::Error(e) => {
358 tag_b.append_value(crate::arrow_store::TypeTag::Error as u8);
359 nb.append_null();
360 bb.append_null();
361 sb.append_null();
362 eb.append_value(crate::arrow_store::map_error_code(e.kind));
363 non_err += 1;
364 }
365 LiteralValue::Date(d) => {
366 tag_b.append_value(crate::arrow_store::TypeTag::Number as u8);
367 let dt = d.and_hms_opt(0, 0, 0).unwrap();
368 let serial = formualizer_common::datetime_to_serial_for(
369 date_system,
370 &dt,
371 );
372 nb.append_value(serial);
373 non_num += 1;
374 bb.append_null();
375 sb.append_null();
376 eb.append_null();
377 }
378 LiteralValue::DateTime(dt) => {
379 tag_b.append_value(crate::arrow_store::TypeTag::Number as u8);
380 let serial = formualizer_common::datetime_to_serial_for(
381 date_system,
382 &dt,
383 );
384 nb.append_value(serial);
385 non_num += 1;
386 bb.append_null();
387 sb.append_null();
388 eb.append_null();
389 }
390 LiteralValue::Time(t) => {
391 tag_b.append_value(crate::arrow_store::TypeTag::Number as u8);
392 let serial = formualizer_common::time_to_fraction(&t);
393 nb.append_value(serial);
394 non_num += 1;
395 bb.append_null();
396 sb.append_null();
397 eb.append_null();
398 }
399 LiteralValue::Duration(d) => {
400 tag_b.append_value(crate::arrow_store::TypeTag::Number as u8);
401 let serial = d.num_seconds() as f64 / 86_400.0;
402 nb.append_value(serial);
403 non_num += 1;
404 bb.append_null();
405 sb.append_null();
406 eb.append_null();
407 }
408 LiteralValue::Pending => {
409 tag_b.append_value(crate::arrow_store::TypeTag::Pending as u8);
410 nb.append_null();
411 bb.append_null();
412 sb.append_null();
413 eb.append_null();
414 }
415 LiteralValue::Array(_) => {
416 tag_b.append_value(crate::arrow_store::TypeTag::Error as u8);
417 nb.append_null();
418 bb.append_null();
419 sb.append_null();
420 eb.append_value(crate::arrow_store::map_error_code(
421 formualizer_common::ExcelErrorKind::Value,
422 ));
423 non_err += 1;
424 }
425 }
426 }
427 ch.type_tag = Arc::new(tag_b.finish());
428 ch.numbers = if non_num == 0 {
429 None
430 } else {
431 Some(Arc::new(nb.finish()))
432 };
433 ch.booleans = if non_bool == 0 {
434 None
435 } else {
436 Some(Arc::new(bb.finish()))
437 };
438 ch.text = if non_text == 0 {
439 None
440 } else {
441 Some(Arc::new(sb.finish()))
442 };
443 ch.errors = if non_err == 0 {
444 None
445 } else {
446 Some(Arc::new(eb.finish()))
447 };
448 ch.meta.len = len;
449 ch.meta.non_null_num = non_num;
450 ch.meta.non_null_bool = non_bool;
451 ch.meta.non_null_text = non_text;
452 ch.meta.non_null_err = non_err;
453 let _ = ch.overlay.clear();
454 }
455 }
456 }
457 }
458 // Advance snapshot and mark edited
459 self.engine.mark_data_edited();
460 Ok(total)
461 }
462}
463
464#[cfg(test)]
465mod tests {
466 use super::*;
467 use crate::engine::EvalConfig;
468 use crate::test_workbook::TestWorkbook;
469
470 #[test]
471 fn arrow_bulk_ingest_basic() {
472 let mut engine = Engine::new(TestWorkbook::default(), EvalConfig::default());
473 let mut ab = engine.begin_bulk_ingest_arrow();
474 ab.add_sheet("S", 3, 2);
475 ab.append_row(
476 "S",
477 &[
478 LiteralValue::Number(1.0),
479 LiteralValue::Text("a".into()),
480 LiteralValue::Empty,
481 ],
482 )
483 .unwrap();
484 ab.append_row(
485 "S",
486 &[
487 LiteralValue::Boolean(true),
488 LiteralValue::Text("".into()),
489 LiteralValue::Error(formualizer_common::ExcelError::new_value()),
490 ],
491 )
492 .unwrap();
493 let summary = ab.finish().unwrap();
494 assert_eq!(summary.sheets, 1);
495 assert_eq!(summary.total_rows, 2);
496
497 let sheet = engine
498 .sheet_store()
499 .sheet("S")
500 .expect("arrow sheet present");
501 assert_eq!(sheet.columns.len(), 3);
502 assert_eq!(sheet.nrows, 2);
503 // Validate chunking (chunk_rows=2 => single chunk)
504 for col in &sheet.columns {
505 assert_eq!(col.chunks.len(), 1);
506 assert_eq!(col.chunks[0].len(), 2);
507 }
508 }
509}