1use std::collections::BTreeMap;
36use std::path::{Path, PathBuf};
37use std::sync::Arc;
38
39use arrow_array::builder::{
40 ArrayBuilder, BooleanBuilder, Float64Builder, Int64Builder, StringBuilder,
41};
42use arrow_array::{ArrayRef, RecordBatch};
43use arrow_schema::{DataType, Field, Schema};
44use cuttlefish_abi::Ty;
45
46#[derive(Debug, thiserror::Error)]
48pub enum WarehouseError {
49 #[error("creating {path}: {source}")]
51 Create {
52 path: PathBuf,
54 #[source]
56 source: std::io::Error,
57 },
58 #[error("writing {path}: {source}")]
60 Write {
61 path: PathBuf,
63 #[source]
65 source: parquet::errors::ParquetError,
66 },
67 #[error("building a record batch for {table}: {source}")]
69 Batch {
70 table: String,
72 #[source]
74 source: arrow_schema::ArrowError,
75 },
76 #[error("serializing the manifest: {0}")]
78 Manifest(#[from] serde_json::Error),
79 #[error("writing the manifest to {path}: {source}")]
81 ManifestWrite {
82 path: PathBuf,
84 #[source]
86 source: std::io::Error,
87 },
88}
89
90#[derive(Debug, Clone)]
95pub struct Lineage {
96 pub job_id: String,
98 pub spec_name: String,
100 pub spec_fingerprint: String,
105 pub model: String,
107 pub embedding_model: Option<String>,
109 pub cuttlefish_version: String,
113}
114
115fn lineage_fields() -> Vec<Field> {
120 vec![
121 Field::new("job_id", DataType::Utf8, false),
122 Field::new("node", DataType::Utf8, false),
123 Field::new("item", DataType::Int64, false),
124 Field::new("status", DataType::Utf8, false),
125 Field::new("concluded_at", DataType::Utf8, false),
126 Field::new("source_input", DataType::Utf8, true),
127 Field::new("spec_name", DataType::Utf8, false),
128 Field::new("spec_fingerprint", DataType::Utf8, false),
129 Field::new("model", DataType::Utf8, false),
130 Field::new("embedding_model", DataType::Utf8, true),
131 Field::new("cuttlefish_version", DataType::Utf8, false),
132 ]
133}
134
135#[derive(Debug, Clone)]
137pub struct Row {
138 pub node: String,
140 pub item: i64,
143 pub status: String,
147 pub concluded_at: String,
149 pub source_input: Option<String>,
152 pub output: Option<serde_json::Value>,
154 pub error: Option<String>,
156}
157
158fn is_flattenable(ty: &Ty) -> bool {
165 match ty {
166 Ty::Text | Ty::Number | Ty::Bool => true,
167 Ty::Json => true,
171 Ty::Bytes | Ty::Image | Ty::Document => false,
172 Ty::List(_) | Ty::Record(_) => false,
173 }
174}
175
176fn column_type(name: &str, ty: &Ty, rows: &[Row]) -> DataType {
184 match ty {
185 Ty::Bool => DataType::Boolean,
186 Ty::Number => {
187 let fractional = rows.iter().filter_map(|r| r.output.as_ref()).any(|out| {
188 out.get(name)
189 .and_then(|v| v.as_f64())
190 .is_some_and(|f| f.fract() != 0.0)
191 });
192 if fractional {
193 DataType::Float64
194 } else {
195 DataType::Int64
199 }
200 }
201 _ => DataType::Utf8,
202 }
203}
204
205fn silver_columns(fields: &BTreeMap<String, Ty>) -> Vec<(&String, &Ty)> {
207 fields.iter().filter(|(_, ty)| is_flattenable(ty)).collect()
208}
209
210pub fn silver_schema(item_output: &Ty, rows: &[Row]) -> Option<Schema> {
220 let Ty::Record(fields) = item_output else {
221 return None;
222 };
223 let declared = silver_columns(fields);
224 if declared.is_empty() {
225 return None;
226 }
227
228 let mut out = lineage_fields();
229 for (name, ty) in declared {
230 out.push(Field::new(
237 format!("f_{name}"),
238 column_type(name, ty, rows),
239 true,
240 ));
241 }
242 Some(Schema::new(out))
243}
244
245pub fn bronze_schema() -> Schema {
247 let mut fields = lineage_fields();
248 fields.push(Field::new("output_json", DataType::Utf8, true));
249 fields.push(Field::new("error", DataType::Utf8, true));
250 Schema::new(fields)
251}
252
253fn push_lineage(builders: &mut [Box<dyn ArrayBuilder>], row: &Row, lineage: &Lineage) {
255 macro_rules! s {
256 ($i:expr, $v:expr) => {
257 builders[$i]
258 .as_any_mut()
259 .downcast_mut::<StringBuilder>()
260 .expect("lineage column is a string column")
261 .append_option($v)
262 };
263 }
264 s!(0, Some(&lineage.job_id));
265 s!(1, Some(&row.node));
266 builders[2]
267 .as_any_mut()
268 .downcast_mut::<Int64Builder>()
269 .expect("`item` is an int column")
270 .append_value(row.item);
271 s!(3, Some(&row.status));
272 s!(4, Some(&row.concluded_at));
273 s!(5, row.source_input.as_ref());
274 s!(6, Some(&lineage.spec_name));
275 s!(7, Some(&lineage.spec_fingerprint));
276 s!(8, Some(&lineage.model));
277 s!(9, lineage.embedding_model.as_ref());
278 s!(10, Some(&lineage.cuttlefish_version));
279}
280
281fn builders_for(schema: &Schema) -> Vec<Box<dyn ArrayBuilder>> {
283 schema
284 .fields()
285 .iter()
286 .map(|f| -> Box<dyn ArrayBuilder> {
287 match f.data_type() {
288 DataType::Int64 => Box::new(Int64Builder::new()),
289 DataType::Float64 => Box::new(Float64Builder::new()),
290 DataType::Boolean => Box::new(BooleanBuilder::new()),
291 _ => Box::new(StringBuilder::new()),
292 }
293 })
294 .collect()
295}
296
297fn cell_text(value: &serde_json::Value) -> Option<String> {
303 match value {
304 serde_json::Value::Null => None,
305 serde_json::Value::String(s) => Some(s.clone()),
306 other => Some(other.to_string()),
307 }
308}
309
310fn finish(mut builders: Vec<Box<dyn ArrayBuilder>>) -> Vec<ArrayRef> {
312 builders.iter_mut().map(|b| b.finish()).collect()
313}
314
315pub fn bronze_batch(rows: &[Row], lineage: &Lineage) -> Result<RecordBatch, WarehouseError> {
317 let schema = bronze_schema();
318 let mut builders = builders_for(&schema);
319 let lineage_count = lineage_fields().len();
320
321 for row in rows {
322 push_lineage(&mut builders, row, lineage);
323 let output = row.output.as_ref().and_then(cell_text);
324 builders[lineage_count]
325 .as_any_mut()
326 .downcast_mut::<StringBuilder>()
327 .expect("`output_json` is a string column")
328 .append_option(output);
329 builders[lineage_count + 1]
330 .as_any_mut()
331 .downcast_mut::<StringBuilder>()
332 .expect("`error` is a string column")
333 .append_option(row.error.as_ref());
334 }
335
336 RecordBatch::try_new(Arc::new(schema), finish(builders)).map_err(|e| WarehouseError::Batch {
337 table: "bronze".into(),
338 source: e,
339 })
340}
341
342pub fn silver_batch(
346 rows: &[Row],
347 lineage: &Lineage,
348 item_output: &Ty,
349) -> Result<Option<RecordBatch>, WarehouseError> {
350 let Some(schema) = silver_schema(item_output, rows) else {
351 return Ok(None);
352 };
353 let Ty::Record(fields) = item_output else {
354 return Ok(None);
355 };
356 let declared = silver_columns(fields);
357
358 let mut builders = builders_for(&schema);
359 let lineage_count = lineage_fields().len();
360
361 for row in rows {
362 let Some(output) = &row.output else { continue };
366 push_lineage(&mut builders, row, lineage);
367
368 for (offset, (name, _)) in declared.iter().enumerate() {
369 let column = lineage_count + offset;
370 let value = output.get(name.as_str());
371 let builder = &mut builders[column];
372 match schema.field(column).data_type() {
373 DataType::Int64 => builder
374 .as_any_mut()
375 .downcast_mut::<Int64Builder>()
376 .expect("an Int64 column has an Int64 builder")
377 .append_option(value.and_then(|v| v.as_i64())),
378 DataType::Float64 => builder
379 .as_any_mut()
380 .downcast_mut::<Float64Builder>()
381 .expect("a Float64 column has a Float64 builder")
382 .append_option(value.and_then(|v| v.as_f64())),
383 DataType::Boolean => builder
384 .as_any_mut()
385 .downcast_mut::<BooleanBuilder>()
386 .expect("a Boolean column has a Boolean builder")
387 .append_option(value.and_then(|v| v.as_bool())),
388 _ => builder
389 .as_any_mut()
390 .downcast_mut::<StringBuilder>()
391 .expect("every other column is a string column")
392 .append_option(value.and_then(cell_text)),
393 }
394 }
395 }
396
397 let arrays = finish(builders);
398 RecordBatch::try_new(Arc::new(schema), arrays)
399 .map(Some)
400 .map_err(|e| WarehouseError::Batch {
401 table: "silver".into(),
402 source: e,
403 })
404}
405
406pub fn write_parquet(path: &Path, batch: &RecordBatch) -> Result<(), WarehouseError> {
408 if let Some(parent) = path.parent() {
409 std::fs::create_dir_all(parent).map_err(|e| WarehouseError::Create {
410 path: parent.to_path_buf(),
411 source: e,
412 })?;
413 }
414 let file = std::fs::File::create(path).map_err(|e| WarehouseError::Create {
415 path: path.to_path_buf(),
416 source: e,
417 })?;
418
419 let props = parquet::file::properties::WriterProperties::builder()
420 .set_compression(parquet::basic::Compression::SNAPPY)
424 .build();
425
426 let mut writer = parquet::arrow::ArrowWriter::try_new(file, batch.schema(), Some(props))
427 .map_err(|e| WarehouseError::Write {
428 path: path.to_path_buf(),
429 source: e,
430 })?;
431 writer.write(batch).map_err(|e| WarehouseError::Write {
432 path: path.to_path_buf(),
433 source: e,
434 })?;
435 writer.close().map_err(|e| WarehouseError::Write {
436 path: path.to_path_buf(),
437 source: e,
438 })?;
439 Ok(())
440}
441
442#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
444pub struct TableEntry {
445 pub path: String,
448 pub rows: usize,
450 pub columns: Vec<String>,
452}
453
454#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
456#[serde(untagged)]
457pub enum Layer {
458 Written(TableEntry),
460 Skipped {
463 skipped: String,
465 },
466}
467
468#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
470pub struct Manifest {
471 pub job_id: String,
473 pub spec_name: String,
475 pub spec_fingerprint: String,
477 pub model: String,
479 #[serde(skip_serializing_if = "Option::is_none")]
481 pub embedding_model: Option<String>,
482 pub cuttlefish_version: String,
484 pub written_at: String,
486 pub bronze: BTreeMap<String, Layer>,
489 pub silver: BTreeMap<String, Layer>,
491 pub gold: BTreeMap<String, Layer>,
493}
494
495pub fn now_rfc3339() -> String {
497 time::OffsetDateTime::now_utc()
498 .format(&time::format_description::well_known::Rfc3339)
499 .expect("Rfc3339 formatting cannot fail for a valid OffsetDateTime")
500}
501
502pub fn write_manifest(root: &Path, manifest: &Manifest) -> Result<PathBuf, WarehouseError> {
504 std::fs::create_dir_all(root).map_err(|e| WarehouseError::Create {
505 path: root.to_path_buf(),
506 source: e,
507 })?;
508 let path = root.join("manifest.json");
509 let body = serde_json::to_string_pretty(manifest)?;
510 std::fs::write(&path, body).map_err(|e| WarehouseError::ManifestWrite {
511 path: path.clone(),
512 source: e,
513 })?;
514 Ok(path)
515}
516
517pub fn entry_for(root: &Path, path: &Path, batch: &RecordBatch) -> TableEntry {
519 TableEntry {
520 path: path
521 .strip_prefix(root)
522 .unwrap_or(path)
523 .to_string_lossy()
524 .into_owned(),
525 rows: batch.num_rows(),
526 columns: batch
527 .schema()
528 .fields()
529 .iter()
530 .map(|f| f.name().clone())
531 .collect(),
532 }
533}
534
535#[cfg(test)]
536mod tests {
537 use super::*;
538
539 fn lineage() -> Lineage {
540 Lineage {
541 job_id: "job-1".into(),
542 spec_name: "index_corpus".into(),
543 spec_fingerprint: "abc123".into(),
544 model: "ollama:llama3.2:1b".into(),
545 embedding_model: Some("ollama:nomic-embed-text".into()),
546 cuttlefish_version: "0.8.0".into(),
547 }
548 }
549
550 fn record(fields: &[(&str, Ty)]) -> Ty {
551 Ty::Record(
552 fields
553 .iter()
554 .map(|(n, t)| (n.to_string(), t.clone()))
555 .collect(),
556 )
557 }
558
559 fn row(item: i64, output: Option<serde_json::Value>, error: Option<&str>) -> Row {
560 Row {
561 node: "extract".into(),
562 item,
563 status: if output.is_some() {
564 "completed"
565 } else {
566 "failed"
567 }
568 .into(),
569 concluded_at: "2026-08-18T00:00:00Z".into(),
570 source_input: Some(format!(r#"{{"path":"doc-{item}.pdf"}}"#)),
571 output,
572 error: error.map(str::to_string),
573 }
574 }
575
576 #[test]
577 fn a_node_declaring_json_gets_no_silver_table() {
578 assert!(silver_schema(&Ty::Json, &[]).is_none());
582 assert!(silver_schema(&Ty::Record(Default::default()), &[]).is_none());
583 assert!(silver_schema(&Ty::Text, &[]).is_none());
584 }
585
586 #[test]
587 fn silver_columns_follow_the_declared_record() {
588 let ty = record(&[("title", Ty::Text), ("body", Ty::Text)]);
589 let schema = silver_schema(&ty, &[]).expect("a declared record yields a table");
590 let names: Vec<_> = schema.fields().iter().map(|f| f.name().clone()).collect();
591
592 assert_eq!(names[0], "job_id");
595 assert!(names.contains(&"f_title".to_string()), "{names:?}");
596 assert!(names.contains(&"f_body".to_string()), "{names:?}");
597 assert!(!names.contains(&"title".to_string()), "{names:?}");
598 }
599
600 #[test]
601 fn a_record_of_only_unflattenable_fields_gets_no_table() {
602 let ty = record(&[("pages", Ty::List(Box::new(Ty::Text))), ("scan", Ty::Image)]);
606 assert!(silver_schema(&ty, &[]).is_none());
607 }
608
609 #[test]
610 fn bronze_keeps_failures_and_silver_drops_them() {
611 let rows = vec![
614 row(0, Some(serde_json::json!({"title": "A"})), None),
615 row(1, None, Some("pdf has no text layer")),
616 row(2, Some(serde_json::json!({"title": "C"})), None),
617 ];
618 let ty = record(&[("title", Ty::Text)]);
619
620 let bronze = bronze_batch(&rows, &lineage()).unwrap();
621 assert_eq!(bronze.num_rows(), 3);
622
623 let silver = silver_batch(&rows, &lineage(), &ty).unwrap().unwrap();
624 assert_eq!(silver.num_rows(), 2);
625 }
626
627 #[test]
628 fn a_string_field_is_not_re_encoded_with_its_quotes() {
629 assert_eq!(
633 cell_text(&serde_json::json!("A")),
634 Some("A".to_string()),
635 "a JSON string must become its contents"
636 );
637 assert_eq!(
638 cell_text(&serde_json::json!({"n": 1})),
639 Some(r#"{"n":1}"#.to_string()),
640 "a JSON object keeps its encoding — there is nothing else it could be"
641 );
642 assert_eq!(cell_text(&serde_json::Value::Null), None);
643 }
644
645 #[test]
646 fn a_missing_declared_field_is_null_rather_than_a_write_failure() {
647 let rows = vec![row(0, Some(serde_json::json!({"title": "A"})), None)];
650 let ty = record(&[("title", Ty::Text), ("subtitle", Ty::Text)]);
651 let silver = silver_batch(&rows, &lineage(), &ty).unwrap().unwrap();
652 assert_eq!(silver.num_rows(), 1);
653 let column = silver
654 .column_by_name("f_subtitle")
655 .expect("the declared field is a column even when unpopulated");
656 assert!(column.is_null(0), "an absent field reads as null");
657 }
658
659 #[test]
660 fn a_declared_number_becomes_a_number_column_not_a_string() {
661 let rows = vec![row(0, Some(serde_json::json!({"pages": 227})), None)];
664 let ty = record(&[("pages", Ty::Number)]);
665 let schema = silver_schema(&ty, &rows).unwrap();
666 assert_eq!(
667 schema.field_with_name("f_pages").unwrap().data_type(),
668 &DataType::Int64
669 );
670
671 let batch = silver_batch(&rows, &lineage(), &ty).unwrap().unwrap();
672 let column = batch
673 .column_by_name("f_pages")
674 .unwrap()
675 .as_any()
676 .downcast_ref::<arrow_array::Int64Array>()
677 .expect("an integral number column is Int64, not stringified");
678 assert_eq!(column.value(0), 227);
679 }
680
681 #[test]
682 fn one_fractional_value_makes_the_whole_column_a_float() {
683 let rows = vec![
687 row(0, Some(serde_json::json!({"score": 1})), None),
688 row(1, Some(serde_json::json!({"score": 0.75})), None),
689 ];
690 let ty = record(&[("score", Ty::Number)]);
691 let schema = silver_schema(&ty, &rows).unwrap();
692 assert_eq!(
693 schema.field_with_name("f_score").unwrap().data_type(),
694 &DataType::Float64
695 );
696
697 let batch = silver_batch(&rows, &lineage(), &ty).unwrap().unwrap();
698 let column = batch
699 .column_by_name("f_score")
700 .unwrap()
701 .as_any()
702 .downcast_ref::<arrow_array::Float64Array>()
703 .unwrap();
704 assert_eq!((column.value(0), column.value(1)), (1.0, 0.75));
705 }
706
707 #[test]
708 fn a_declared_bool_becomes_a_boolean_column() {
709 let rows = vec![row(0, Some(serde_json::json!({"has_text": false})), None)];
710 let ty = record(&[("has_text", Ty::Bool)]);
711 let batch = silver_batch(&rows, &lineage(), &ty).unwrap().unwrap();
712 let column = batch
713 .column_by_name("f_has_text")
714 .unwrap()
715 .as_any()
716 .downcast_ref::<arrow_array::BooleanArray>()
717 .expect("a bool column is Boolean, not the string \"false\"");
718 assert!(!column.value(0));
719 }
720
721 #[test]
722 fn a_number_field_a_block_omitted_is_null_not_zero() {
723 let rows = vec![row(0, Some(serde_json::json!({"other": 1})), None)];
726 let ty = record(&[("pages", Ty::Number), ("other", Ty::Number)]);
727 let batch = silver_batch(&rows, &lineage(), &ty).unwrap().unwrap();
728 let column = batch.column_by_name("f_pages").unwrap();
729 assert!(column.is_null(0), "an absent number is null, never 0");
730 }
731
732 #[test]
733 fn every_row_carries_its_own_lineage() {
734 let rows = vec![row(0, Some(serde_json::json!({"title": "A"})), None)];
737 let bronze = bronze_batch(&rows, &lineage()).unwrap();
738 for column in [
739 "job_id",
740 "spec_fingerprint",
741 "model",
742 "cuttlefish_version",
743 "source_input",
744 ] {
745 let c = bronze
746 .column_by_name(column)
747 .unwrap_or_else(|| panic!("bronze must carry `{column}`"));
748 assert!(!c.is_null(0), "`{column}` must be populated");
749 }
750 }
751}