use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use arrow_array::builder::{
ArrayBuilder, BooleanBuilder, Float64Builder, Int64Builder, StringBuilder,
};
use arrow_array::{ArrayRef, RecordBatch};
use arrow_schema::{DataType, Field, Schema};
use cuttlefish_abi::Ty;
#[derive(Debug, thiserror::Error)]
pub enum WarehouseError {
#[error("creating {path}: {source}")]
Create {
path: PathBuf,
#[source]
source: std::io::Error,
},
#[error("writing {path}: {source}")]
Write {
path: PathBuf,
#[source]
source: parquet::errors::ParquetError,
},
#[error("building a record batch for {table}: {source}")]
Batch {
table: String,
#[source]
source: arrow_schema::ArrowError,
},
#[error("serializing the manifest: {0}")]
Manifest(#[from] serde_json::Error),
#[error("writing the manifest to {path}: {source}")]
ManifestWrite {
path: PathBuf,
#[source]
source: std::io::Error,
},
}
#[derive(Debug, Clone)]
pub struct Lineage {
pub job_id: String,
pub spec_name: String,
pub spec_fingerprint: String,
pub model: String,
pub embedding_model: Option<String>,
pub cuttlefish_version: String,
}
fn lineage_fields() -> Vec<Field> {
vec![
Field::new("job_id", DataType::Utf8, false),
Field::new("node", DataType::Utf8, false),
Field::new("item", DataType::Int64, false),
Field::new("status", DataType::Utf8, false),
Field::new("concluded_at", DataType::Utf8, false),
Field::new("source_input", DataType::Utf8, true),
Field::new("spec_name", DataType::Utf8, false),
Field::new("spec_fingerprint", DataType::Utf8, false),
Field::new("model", DataType::Utf8, false),
Field::new("embedding_model", DataType::Utf8, true),
Field::new("cuttlefish_version", DataType::Utf8, false),
]
}
#[derive(Debug, Clone)]
pub struct Row {
pub node: String,
pub item: i64,
pub status: String,
pub concluded_at: String,
pub source_input: Option<String>,
pub output: Option<serde_json::Value>,
pub error: Option<String>,
}
fn is_flattenable(ty: &Ty) -> bool {
match ty {
Ty::Text | Ty::Number | Ty::Bool => true,
Ty::Json => true,
Ty::Bytes | Ty::Image | Ty::Document => false,
Ty::List(_) | Ty::Record(_) => false,
}
}
fn column_type(name: &str, ty: &Ty, rows: &[Row]) -> DataType {
match ty {
Ty::Bool => DataType::Boolean,
Ty::Number => {
let fractional = rows.iter().filter_map(|r| r.output.as_ref()).any(|out| {
out.get(name)
.and_then(|v| v.as_f64())
.is_some_and(|f| f.fract() != 0.0)
});
if fractional {
DataType::Float64
} else {
DataType::Int64
}
}
_ => DataType::Utf8,
}
}
fn silver_columns(fields: &BTreeMap<String, Ty>) -> Vec<(&String, &Ty)> {
fields.iter().filter(|(_, ty)| is_flattenable(ty)).collect()
}
pub fn silver_schema(item_output: &Ty, rows: &[Row]) -> Option<Schema> {
let Ty::Record(fields) = item_output else {
return None;
};
let declared = silver_columns(fields);
if declared.is_empty() {
return None;
}
let mut out = lineage_fields();
for (name, ty) in declared {
out.push(Field::new(
format!("f_{name}"),
column_type(name, ty, rows),
true,
));
}
Some(Schema::new(out))
}
pub fn bronze_schema() -> Schema {
let mut fields = lineage_fields();
fields.push(Field::new("output_json", DataType::Utf8, true));
fields.push(Field::new("error", DataType::Utf8, true));
Schema::new(fields)
}
fn push_lineage(builders: &mut [Box<dyn ArrayBuilder>], row: &Row, lineage: &Lineage) {
macro_rules! s {
($i:expr, $v:expr) => {
builders[$i]
.as_any_mut()
.downcast_mut::<StringBuilder>()
.expect("lineage column is a string column")
.append_option($v)
};
}
s!(0, Some(&lineage.job_id));
s!(1, Some(&row.node));
builders[2]
.as_any_mut()
.downcast_mut::<Int64Builder>()
.expect("`item` is an int column")
.append_value(row.item);
s!(3, Some(&row.status));
s!(4, Some(&row.concluded_at));
s!(5, row.source_input.as_ref());
s!(6, Some(&lineage.spec_name));
s!(7, Some(&lineage.spec_fingerprint));
s!(8, Some(&lineage.model));
s!(9, lineage.embedding_model.as_ref());
s!(10, Some(&lineage.cuttlefish_version));
}
fn builders_for(schema: &Schema) -> Vec<Box<dyn ArrayBuilder>> {
schema
.fields()
.iter()
.map(|f| -> Box<dyn ArrayBuilder> {
match f.data_type() {
DataType::Int64 => Box::new(Int64Builder::new()),
DataType::Float64 => Box::new(Float64Builder::new()),
DataType::Boolean => Box::new(BooleanBuilder::new()),
_ => Box::new(StringBuilder::new()),
}
})
.collect()
}
fn cell_text(value: &serde_json::Value) -> Option<String> {
match value {
serde_json::Value::Null => None,
serde_json::Value::String(s) => Some(s.clone()),
other => Some(other.to_string()),
}
}
fn finish(mut builders: Vec<Box<dyn ArrayBuilder>>) -> Vec<ArrayRef> {
builders.iter_mut().map(|b| b.finish()).collect()
}
pub fn bronze_batch(rows: &[Row], lineage: &Lineage) -> Result<RecordBatch, WarehouseError> {
let schema = bronze_schema();
let mut builders = builders_for(&schema);
let lineage_count = lineage_fields().len();
for row in rows {
push_lineage(&mut builders, row, lineage);
let output = row.output.as_ref().and_then(cell_text);
builders[lineage_count]
.as_any_mut()
.downcast_mut::<StringBuilder>()
.expect("`output_json` is a string column")
.append_option(output);
builders[lineage_count + 1]
.as_any_mut()
.downcast_mut::<StringBuilder>()
.expect("`error` is a string column")
.append_option(row.error.as_ref());
}
RecordBatch::try_new(Arc::new(schema), finish(builders)).map_err(|e| WarehouseError::Batch {
table: "bronze".into(),
source: e,
})
}
pub fn silver_batch(
rows: &[Row],
lineage: &Lineage,
item_output: &Ty,
) -> Result<Option<RecordBatch>, WarehouseError> {
let Some(schema) = silver_schema(item_output, rows) else {
return Ok(None);
};
let Ty::Record(fields) = item_output else {
return Ok(None);
};
let declared = silver_columns(fields);
let mut builders = builders_for(&schema);
let lineage_count = lineage_fields().len();
for row in rows {
let Some(output) = &row.output else { continue };
push_lineage(&mut builders, row, lineage);
for (offset, (name, _)) in declared.iter().enumerate() {
let column = lineage_count + offset;
let value = output.get(name.as_str());
let builder = &mut builders[column];
match schema.field(column).data_type() {
DataType::Int64 => builder
.as_any_mut()
.downcast_mut::<Int64Builder>()
.expect("an Int64 column has an Int64 builder")
.append_option(value.and_then(|v| v.as_i64())),
DataType::Float64 => builder
.as_any_mut()
.downcast_mut::<Float64Builder>()
.expect("a Float64 column has a Float64 builder")
.append_option(value.and_then(|v| v.as_f64())),
DataType::Boolean => builder
.as_any_mut()
.downcast_mut::<BooleanBuilder>()
.expect("a Boolean column has a Boolean builder")
.append_option(value.and_then(|v| v.as_bool())),
_ => builder
.as_any_mut()
.downcast_mut::<StringBuilder>()
.expect("every other column is a string column")
.append_option(value.and_then(cell_text)),
}
}
}
let arrays = finish(builders);
RecordBatch::try_new(Arc::new(schema), arrays)
.map(Some)
.map_err(|e| WarehouseError::Batch {
table: "silver".into(),
source: e,
})
}
pub fn write_parquet(path: &Path, batch: &RecordBatch) -> Result<(), WarehouseError> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).map_err(|e| WarehouseError::Create {
path: parent.to_path_buf(),
source: e,
})?;
}
let file = std::fs::File::create(path).map_err(|e| WarehouseError::Create {
path: path.to_path_buf(),
source: e,
})?;
let props = parquet::file::properties::WriterProperties::builder()
.set_compression(parquet::basic::Compression::SNAPPY)
.build();
let mut writer = parquet::arrow::ArrowWriter::try_new(file, batch.schema(), Some(props))
.map_err(|e| WarehouseError::Write {
path: path.to_path_buf(),
source: e,
})?;
writer.write(batch).map_err(|e| WarehouseError::Write {
path: path.to_path_buf(),
source: e,
})?;
writer.close().map_err(|e| WarehouseError::Write {
path: path.to_path_buf(),
source: e,
})?;
Ok(())
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct TableEntry {
pub path: String,
pub rows: usize,
pub columns: Vec<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(untagged)]
pub enum Layer {
Written(TableEntry),
Skipped {
skipped: String,
},
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct Manifest {
pub job_id: String,
pub spec_name: String,
pub spec_fingerprint: String,
pub model: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub embedding_model: Option<String>,
pub cuttlefish_version: String,
pub written_at: String,
pub bronze: BTreeMap<String, Layer>,
pub silver: BTreeMap<String, Layer>,
pub gold: BTreeMap<String, Layer>,
}
pub fn now_rfc3339() -> String {
time::OffsetDateTime::now_utc()
.format(&time::format_description::well_known::Rfc3339)
.expect("Rfc3339 formatting cannot fail for a valid OffsetDateTime")
}
pub fn write_manifest(root: &Path, manifest: &Manifest) -> Result<PathBuf, WarehouseError> {
std::fs::create_dir_all(root).map_err(|e| WarehouseError::Create {
path: root.to_path_buf(),
source: e,
})?;
let path = root.join("manifest.json");
let body = serde_json::to_string_pretty(manifest)?;
std::fs::write(&path, body).map_err(|e| WarehouseError::ManifestWrite {
path: path.clone(),
source: e,
})?;
Ok(path)
}
pub fn entry_for(root: &Path, path: &Path, batch: &RecordBatch) -> TableEntry {
TableEntry {
path: path
.strip_prefix(root)
.unwrap_or(path)
.to_string_lossy()
.into_owned(),
rows: batch.num_rows(),
columns: batch
.schema()
.fields()
.iter()
.map(|f| f.name().clone())
.collect(),
}
}
#[cfg(test)]
mod tests {
use super::*;
fn lineage() -> Lineage {
Lineage {
job_id: "job-1".into(),
spec_name: "index_corpus".into(),
spec_fingerprint: "abc123".into(),
model: "ollama:llama3.2:1b".into(),
embedding_model: Some("ollama:nomic-embed-text".into()),
cuttlefish_version: "0.8.0".into(),
}
}
fn record(fields: &[(&str, Ty)]) -> Ty {
Ty::Record(
fields
.iter()
.map(|(n, t)| (n.to_string(), t.clone()))
.collect(),
)
}
fn row(item: i64, output: Option<serde_json::Value>, error: Option<&str>) -> Row {
Row {
node: "extract".into(),
item,
status: if output.is_some() {
"completed"
} else {
"failed"
}
.into(),
concluded_at: "2026-08-18T00:00:00Z".into(),
source_input: Some(format!(r#"{{"path":"doc-{item}.pdf"}}"#)),
output,
error: error.map(str::to_string),
}
}
#[test]
fn a_node_declaring_json_gets_no_silver_table() {
assert!(silver_schema(&Ty::Json, &[]).is_none());
assert!(silver_schema(&Ty::Record(Default::default()), &[]).is_none());
assert!(silver_schema(&Ty::Text, &[]).is_none());
}
#[test]
fn silver_columns_follow_the_declared_record() {
let ty = record(&[("title", Ty::Text), ("body", Ty::Text)]);
let schema = silver_schema(&ty, &[]).expect("a declared record yields a table");
let names: Vec<_> = schema.fields().iter().map(|f| f.name().clone()).collect();
assert_eq!(names[0], "job_id");
assert!(names.contains(&"f_title".to_string()), "{names:?}");
assert!(names.contains(&"f_body".to_string()), "{names:?}");
assert!(!names.contains(&"title".to_string()), "{names:?}");
}
#[test]
fn a_record_of_only_unflattenable_fields_gets_no_table() {
let ty = record(&[("pages", Ty::List(Box::new(Ty::Text))), ("scan", Ty::Image)]);
assert!(silver_schema(&ty, &[]).is_none());
}
#[test]
fn bronze_keeps_failures_and_silver_drops_them() {
let rows = vec![
row(0, Some(serde_json::json!({"title": "A"})), None),
row(1, None, Some("pdf has no text layer")),
row(2, Some(serde_json::json!({"title": "C"})), None),
];
let ty = record(&[("title", Ty::Text)]);
let bronze = bronze_batch(&rows, &lineage()).unwrap();
assert_eq!(bronze.num_rows(), 3);
let silver = silver_batch(&rows, &lineage(), &ty).unwrap().unwrap();
assert_eq!(silver.num_rows(), 2);
}
#[test]
fn a_string_field_is_not_re_encoded_with_its_quotes() {
assert_eq!(
cell_text(&serde_json::json!("A")),
Some("A".to_string()),
"a JSON string must become its contents"
);
assert_eq!(
cell_text(&serde_json::json!({"n": 1})),
Some(r#"{"n":1}"#.to_string()),
"a JSON object keeps its encoding — there is nothing else it could be"
);
assert_eq!(cell_text(&serde_json::Value::Null), None);
}
#[test]
fn a_missing_declared_field_is_null_rather_than_a_write_failure() {
let rows = vec![row(0, Some(serde_json::json!({"title": "A"})), None)];
let ty = record(&[("title", Ty::Text), ("subtitle", Ty::Text)]);
let silver = silver_batch(&rows, &lineage(), &ty).unwrap().unwrap();
assert_eq!(silver.num_rows(), 1);
let column = silver
.column_by_name("f_subtitle")
.expect("the declared field is a column even when unpopulated");
assert!(column.is_null(0), "an absent field reads as null");
}
#[test]
fn a_declared_number_becomes_a_number_column_not_a_string() {
let rows = vec![row(0, Some(serde_json::json!({"pages": 227})), None)];
let ty = record(&[("pages", Ty::Number)]);
let schema = silver_schema(&ty, &rows).unwrap();
assert_eq!(
schema.field_with_name("f_pages").unwrap().data_type(),
&DataType::Int64
);
let batch = silver_batch(&rows, &lineage(), &ty).unwrap().unwrap();
let column = batch
.column_by_name("f_pages")
.unwrap()
.as_any()
.downcast_ref::<arrow_array::Int64Array>()
.expect("an integral number column is Int64, not stringified");
assert_eq!(column.value(0), 227);
}
#[test]
fn one_fractional_value_makes_the_whole_column_a_float() {
let rows = vec![
row(0, Some(serde_json::json!({"score": 1})), None),
row(1, Some(serde_json::json!({"score": 0.75})), None),
];
let ty = record(&[("score", Ty::Number)]);
let schema = silver_schema(&ty, &rows).unwrap();
assert_eq!(
schema.field_with_name("f_score").unwrap().data_type(),
&DataType::Float64
);
let batch = silver_batch(&rows, &lineage(), &ty).unwrap().unwrap();
let column = batch
.column_by_name("f_score")
.unwrap()
.as_any()
.downcast_ref::<arrow_array::Float64Array>()
.unwrap();
assert_eq!((column.value(0), column.value(1)), (1.0, 0.75));
}
#[test]
fn a_declared_bool_becomes_a_boolean_column() {
let rows = vec![row(0, Some(serde_json::json!({"has_text": false})), None)];
let ty = record(&[("has_text", Ty::Bool)]);
let batch = silver_batch(&rows, &lineage(), &ty).unwrap().unwrap();
let column = batch
.column_by_name("f_has_text")
.unwrap()
.as_any()
.downcast_ref::<arrow_array::BooleanArray>()
.expect("a bool column is Boolean, not the string \"false\"");
assert!(!column.value(0));
}
#[test]
fn a_number_field_a_block_omitted_is_null_not_zero() {
let rows = vec![row(0, Some(serde_json::json!({"other": 1})), None)];
let ty = record(&[("pages", Ty::Number), ("other", Ty::Number)]);
let batch = silver_batch(&rows, &lineage(), &ty).unwrap().unwrap();
let column = batch.column_by_name("f_pages").unwrap();
assert!(column.is_null(0), "an absent number is null, never 0");
}
#[test]
fn every_row_carries_its_own_lineage() {
let rows = vec![row(0, Some(serde_json::json!({"title": "A"})), None)];
let bronze = bronze_batch(&rows, &lineage()).unwrap();
for column in [
"job_id",
"spec_fingerprint",
"model",
"cuttlefish_version",
"source_input",
] {
let c = bronze
.column_by_name(column)
.unwrap_or_else(|| panic!("bronze must carry `{column}`"));
assert!(!c.is_null(0), "`{column}` must be populated");
}
}
}