use super::super::filter::apply_filter;
use super::super::input::InputRegistry;
use super::super::table::{ListMisparseTally, RawCsv};
use super::super::timeseries as ts;
use super::super::typing::map_blueprint_type;
use super::cache::{CsvCache, IdTypeCache};
use super::nodes::{fk_id_columns, node_chunk_size, should_stream_spec};
use super::prepass;
use super::specs::FlatSpec;
use super::table_ops::subset_rows;
use super::BuildReport;
use crate::datatypes::values::DataFrame;
use crate::graph::mutation::maintain;
use crate::graph::schema::DirGraph;
use indexmap::IndexMap;
use std::collections::{HashMap, HashSet};
struct PreppedFkEdges {
source_type: String,
pk: String,
edges: Vec<PreppedFkEdge>,
errors: Vec<String>,
warnings: Vec<String>,
}
struct PreppedFkEdge {
edge_type: String,
target_type: String,
target_col: String,
df: DataFrame,
}
fn prep_fk_edges(
spec: &FlatSpec,
registry: &InputRegistry,
cache: &CsvCache,
) -> Option<PreppedFkEdges> {
let input = spec.input.as_deref()?;
let mut fk_edges: IndexMap<String, super::super::schema::FkEdge> = spec
.spec
.connections
.fk_edges
.iter()
.map(|(k, v)| (k.clone(), v.clone()))
.collect();
if let (Some(parent_type), Some(parent_fk)) = (&spec.spec.parent, &spec.spec.parent_fk) {
let edge_type = format!("OF_{}", parent_type.to_uppercase());
fk_edges.entry(edge_type).or_insert_with(|| {
super::super::schema::FkEdge::plain(parent_type.clone(), parent_fk.clone())
});
}
if fk_edges.is_empty() {
return None;
}
let known_types = registry
.get(input)
.map(|s| s.known_column_types())
.unwrap_or_default();
let raw_rc = cache.get(registry, input).ok()?;
let mut raw: RawCsv = (*raw_rc).clone_raw();
if !spec.spec.filter.is_empty() {
apply_filter(&mut raw, &spec.spec.filter);
}
if let Some(tspec) = &spec.spec.timeseries {
ts::drop_zero_time_components(&mut raw, tspec);
}
let raw_pk = spec.spec.pk.clone().unwrap_or_else(|| "id".to_string());
let pk = if raw_pk == "auto" {
let synth = format!("_{}_id", spec.node_type);
let n = raw.row_count();
let values: Vec<String> = (1..=n).map(|i| i.to_string()).collect();
raw.headers.push(synth.clone());
for (r, row) in raw.rows.iter_mut().enumerate() {
row.push(values[r].clone());
raw.nulls[r].push(false);
}
synth
} else {
raw_pk
};
let mut built = Vec::new();
let mut errors = Vec::new();
let mut warnings = Vec::new();
for (edge_type, edge) in &fk_edges {
let Some(fk_idx) = raw.col_index(&edge.fk) else {
errors.push(format!(
"[{}] FK column '{}' not found for edge {}",
spec.node_type, edge.fk, edge_type
));
continue;
};
let Some(pk_idx) = raw.col_index(&pk) else {
errors.push(format!(
"[{}] pk column '{}' not found for edge {}",
spec.node_type, pk, edge_type
));
continue;
};
let props = match fk_edge_properties(edge_type, &spec.node_type, edge, &pk) {
Ok(mut p) => {
super::super::typing::overlay_known_types(&mut p.declared, &known_types);
p
}
Err(e) => {
errors.push(e);
continue;
}
};
let mut misparses = ListMisparseTally::default();
let frame = match fk_edge_frame(
&raw,
&pk,
edge,
IdColumnIdx {
pk: pk_idx,
fk: fk_idx,
},
&props,
&IdTypes(None),
&mut misparses,
) {
Ok(Some(frame)) => frame,
Ok(None) => continue,
Err(e) => {
errors.push(format!(
"[{}] failed to build edge DataFrame for {}: {}",
spec.node_type, edge_type, e
));
continue;
}
};
warnings.extend(misparses.into_warnings(&format!(
"fk_edge '{edge_type}' (node '{}')",
spec.node_type
)));
for col in &frame.missing_properties {
errors.push(missing_fk_property_error(&spec.node_type, edge_type, col));
}
built.push(PreppedFkEdge {
edge_type: edge_type.clone(),
target_type: edge.target.clone(),
target_col: frame.target_col,
df: frame.df,
});
}
Some(PreppedFkEdges {
source_type: spec.node_type.clone(),
pk,
edges: built,
errors,
warnings,
})
}
fn missing_fk_property_error(node_type: &str, edge_type: &str, column: &str) -> String {
format!(
"[{node_type}] fk_edge {edge_type}: property column '{column}' not found in the \
source CSV — the edge is built without it"
)
}
struct FkEdgeProperties {
columns: Vec<String>,
declared: HashMap<String, String>,
rename: HashMap<String, String>,
}
fn fk_edge_properties(
edge_type: &str,
node_type: &str,
edge: &super::super::schema::FkEdge,
pk: &str,
) -> Result<FkEdgeProperties, String> {
let target_col = fk_target_col(pk, &edge.fk);
let mut columns: Vec<String> = Vec::new();
for col in &edge.properties {
if col == pk || col == &edge.fk {
return Err(format!(
"[{node_type}] fk_edge {edge_type}: property '{col}' is an id column \
(pk '{pk}', fk '{}'); the edge already carries it",
edge.fk
));
}
if !columns.contains(col) {
columns.push(col.clone());
}
}
let mut rename: HashMap<String, String> = HashMap::new();
for (col, new_name) in &edge.rename {
if col == pk || col == &edge.fk {
return Err(format!(
"[{node_type}] fk_edge {edge_type}: rename of fk column '{col}' is not \
supported — 'fk' and 'pk' name the CSV columns"
));
}
if !columns.contains(col) {
return Err(format!(
"[{node_type}] fk_edge {edge_type}: rename key '{col}' is not in 'properties'"
));
}
let collides = new_name == pk
|| new_name == &target_col
|| columns.iter().any(|c| c == new_name && c != col)
|| rename.values().any(|v| v == new_name);
if collides {
return Err(format!(
"[{node_type}] fk_edge {edge_type}: rename target '{new_name}' collides with \
another column"
));
}
rename.insert(col.clone(), new_name.clone());
}
let declared = edge
.property_types
.iter()
.filter(|(_, ty)| map_blueprint_type(ty).is_some())
.map(|(col, ty)| (col.clone(), ty.clone()))
.collect();
Ok(FkEdgeProperties {
columns,
declared,
rename,
})
}
#[derive(Clone, Copy)]
struct IdColumnIdx {
pk: usize,
fk: usize,
}
struct FkEdgeFrame {
target_col: String,
df: DataFrame,
missing_properties: Vec<String>,
}
fn fk_edge_frame(
raw: &RawCsv,
pk: &str,
edge: &super::super::schema::FkEdge,
idx: IdColumnIdx,
props: &FkEdgeProperties,
id_types: &IdTypes<'_>,
misparses: &mut ListMisparseTally,
) -> Result<Option<FkEdgeFrame>, String> {
let cols = build_fk_columns(raw, pk, &edge.fk, idx.pk, idx.fk);
if cols.src.is_empty() {
return Ok(None);
}
let mut df = build_edge_df(
pk,
&cols.target_col,
cols.src,
cols.tgt,
id_types.for_columns(pk, &edge.fk),
)?;
let mut missing_properties = Vec::new();
let mut present = Vec::new();
for col in &props.columns {
if raw.col_index(col).is_some() {
present.push(col.clone());
} else {
missing_properties.push(col.clone());
}
}
if !present.is_empty() {
let subset;
let source: &RawCsv = if cols.rows.len() == raw.row_count() {
raw
} else {
subset = subset_rows(raw, &cols.rows);
&subset
};
super::super::typing::append_typed_columns(
&mut df,
source,
&present,
&props.declared,
&props.rename,
misparses,
)?;
}
Ok(Some(FkEdgeFrame {
target_col: cols.target_col,
df,
missing_properties,
}))
}
fn fk_target_col(pk: &str, fk: &str) -> String {
if pk == fk {
format!("_target_{}", fk)
} else {
fk.to_string()
}
}
struct FkColumns {
target_col: String,
src: Vec<Option<String>>,
tgt: Vec<Option<String>>,
rows: Vec<usize>,
}
fn build_fk_columns(raw: &RawCsv, pk: &str, fk: &str, pk_idx: usize, fk_idx: usize) -> FkColumns {
let target_col = fk_target_col(pk, fk);
let mut src = Vec::new();
let mut tgt = Vec::new();
let mut rows = Vec::new();
if pk == fk {
for (r, row) in raw.rows.iter().enumerate() {
if raw.nulls[r][pk_idx] {
continue;
}
src.push(Some(row[pk_idx].clone()));
tgt.push(Some(row[pk_idx].clone()));
rows.push(r);
}
} else {
for (r, row) in raw.rows.iter().enumerate() {
if raw.nulls[r][fk_idx] {
continue;
}
let src_val = if raw.nulls[r][pk_idx] {
None
} else {
Some(row[pk_idx].clone())
};
src.push(src_val);
tgt.push(Some(row[fk_idx].clone()));
rows.push(r);
}
}
FkColumns {
target_col,
src,
tgt,
rows,
}
}
pub(super) fn load_fk_edges(
graph: &mut DirGraph,
specs: &[&FlatSpec],
registry: &InputRegistry,
cache: &CsvCache,
id_types: &IdTypeCache,
report: &mut BuildReport,
) -> Result<(), String> {
use rayon::prelude::*;
let profile = std::env::var("KGLITE_BLUEPRINT_PROFILE").is_ok();
let (streamable, buffered): (Vec<&FlatSpec>, Vec<&FlatSpec>) = specs
.iter()
.copied()
.partition(|s| should_stream_spec(s, registry));
let t_par = std::time::Instant::now();
let prepped: Vec<Option<PreppedFkEdges>> = buffered
.par_iter()
.map(|spec| prep_fk_edges(spec, registry, cache))
.collect();
let t_par_ms = t_par.elapsed().as_millis();
let t_serial = std::time::Instant::now();
let mut t_connect = std::time::Duration::ZERO;
for result in prepped {
let Some(pfx) = result else { continue };
for err in pfx.errors {
report.errors.push(err);
}
report.warnings.extend(pfx.warnings);
for edge in pfx.edges {
let t_c = std::time::Instant::now();
let count = connect(
graph,
edge.df,
&edge.edge_type,
&pfx.source_type,
&pfx.pk,
&edge.target_type,
&edge.target_col,
report,
maintain::InitialLoad::Detect,
)?;
t_connect += t_c.elapsed();
*report
.edges_by_type
.entry(edge.edge_type.clone())
.or_insert(0) += count;
}
}
let t_stream = std::time::Instant::now();
for spec in &streamable {
if let Err(e) = load_streamed_fk_edges(graph, spec, registry, id_types, report) {
report.errors.push(e);
}
}
let t_stream_ms = t_stream.elapsed().as_millis();
if profile {
eprintln!(
" fk parallel prep: {} ms | serial connect: {} ms | streaming ({} specs): {} ms | serial total: {} ms",
t_par_ms,
t_connect.as_millis(),
streamable.len(),
t_stream_ms,
t_serial.elapsed().as_millis(),
);
}
Ok(())
}
type PreparedFkChunks<'a> = (
Box<dyn Iterator<Item = Result<RawCsv, String>> + 'a>,
IndexMap<String, crate::datatypes::values::ColumnType>,
);
fn resolve_fk_property_types<'a>(
source: &'a dyn super::super::input::Source,
chunk_size: usize,
spec: &FlatSpec,
id_columns: &[String],
edge_props: &mut IndexMap<String, FkEdgeProperties>,
report: &mut BuildReport,
) -> Result<PreparedFkChunks<'a>, String> {
let mut seen: HashSet<&str> = HashSet::new();
let mut wanted: Vec<String> = Vec::new();
for props in edge_props.values() {
for col in &props.columns {
if !props.declared.contains_key(col) && seen.insert(col.as_str()) {
wanted.push(col.clone());
}
}
}
let filtered = !spec.spec.filter.is_empty();
let prepared = prepass::prepare_chunks(
source,
chunk_size,
&HashMap::new(),
id_columns,
!filtered,
|raw| {
if filtered {
apply_filter(raw, &spec.spec.filter);
}
wanted.clone()
},
)
.map_err(|e| format!("[{}] {}", spec.node_type, e))?;
if let Some(w) = prepass::prepass_warning(
&format!("fk_edge properties (node '{}')", spec.node_type),
&prepared,
) {
report.warnings.push(w);
}
for props in edge_props.values_mut() {
for (col, keyword) in &prepared.resolved {
if props.columns.contains(col) && !props.declared.contains_key(col) {
props.declared.insert(col.clone(), keyword.clone());
}
}
}
Ok((prepared.chunks, prepared.resolved_ids))
}
fn load_streamed_fk_edges(
graph: &mut DirGraph,
spec: &FlatSpec,
registry: &InputRegistry,
id_types: &IdTypeCache,
report: &mut BuildReport,
) -> Result<(), String> {
let Some(input) = spec.input.as_deref() else {
return Ok(());
};
let mut fk_edges: IndexMap<String, super::super::schema::FkEdge> = spec
.spec
.connections
.fk_edges
.iter()
.map(|(k, v)| (k.clone(), v.clone()))
.collect();
if let (Some(parent_type), Some(parent_fk)) = (&spec.spec.parent, &spec.spec.parent_fk) {
let edge_type = format!("OF_{}", parent_type.to_uppercase());
fk_edges.entry(edge_type).or_insert_with(|| {
super::super::schema::FkEdge::plain(parent_type.clone(), parent_fk.clone())
});
}
if fk_edges.is_empty() {
return Ok(());
}
let chunk_size = node_chunk_size();
let initial_load: HashMap<String, maintain::InitialLoad> = fk_edges
.keys()
.map(|edge_type| {
let unseen = !graph.connection_type_metadata.contains_key(edge_type);
(edge_type.clone(), maintain::InitialLoad::Preset(unseen))
})
.collect();
let source = registry
.get(input)
.map_err(|e| format!("[{}] {}", spec.node_type, e))?;
let raw_pk = spec.spec.pk.clone().unwrap_or_else(|| "id".to_string());
let (pk, is_auto_pk) = if raw_pk == "auto" {
(format!("_{}_id", spec.node_type), true)
} else {
(raw_pk, false)
};
let mut auto_pk_counter: u64 = 1;
let known_types = source.known_column_types();
let mut edge_props: IndexMap<String, FkEdgeProperties> = IndexMap::new();
for (edge_type, edge) in &fk_edges {
match fk_edge_properties(edge_type, &spec.node_type, edge, &pk) {
Ok(mut props) => {
super::super::typing::overlay_known_types(&mut props.declared, &known_types);
edge_props.insert(edge_type.clone(), props);
}
Err(e) => report.errors.push(e),
}
}
let id_columns = fk_id_columns(spec, &pk);
let known = id_types.get(input, &id_columns);
let (chunks, resolved_ids) = resolve_fk_property_types(
source,
chunk_size,
spec,
if known.is_some() { &[] } else { &id_columns },
&mut edge_props,
report,
)?;
let resolved_ids = known.unwrap_or(resolved_ids);
let mut reported_missing_fk: HashSet<String> = HashSet::new();
let mut reported_missing_pk: HashSet<String> = HashSet::new();
let mut reported_missing_prop: HashSet<(String, String)> = HashSet::new();
let mut misparses: IndexMap<String, ListMisparseTally> = IndexMap::new();
for chunk_result in chunks {
let mut raw = chunk_result.map_err(|e| format!("[{}] {}", spec.node_type, e))?;
if !spec.spec.filter.is_empty() {
apply_filter(&mut raw, &spec.spec.filter);
}
if raw.row_count() == 0 {
continue;
}
if is_auto_pk {
raw.headers.push(pk.clone());
for r in 0..raw.row_count() {
raw.rows[r].push(auto_pk_counter.to_string());
raw.nulls[r].push(false);
auto_pk_counter += 1;
}
}
let Some(pk_idx) = raw.col_index(&pk) else {
for edge_type in fk_edges.keys() {
if reported_missing_pk.insert(edge_type.clone()) {
report.errors.push(format!(
"[{}] pk column '{}' not found for edge {}",
spec.node_type, pk, edge_type
));
}
}
continue;
};
for (edge_type, edge) in &fk_edges {
let Some(props) = edge_props.get(edge_type) else {
continue;
};
let Some(fk_idx) = raw.col_index(&edge.fk) else {
if reported_missing_fk.insert(edge_type.clone()) {
report.errors.push(format!(
"[{}] FK column '{}' not found for edge {}",
spec.node_type, edge.fk, edge_type
));
}
continue;
};
let tally = misparses.entry(edge_type.clone()).or_default();
let frame = match fk_edge_frame(
&raw,
&pk,
edge,
IdColumnIdx {
pk: pk_idx,
fk: fk_idx,
},
props,
&IdTypes(Some(&resolved_ids)),
tally,
) {
Ok(Some(frame)) => frame,
Ok(None) => continue,
Err(e) => {
report.errors.push(format!(
"[{}] failed to build edge DataFrame for {}: {}",
spec.node_type, edge_type, e
));
continue;
}
};
for col in &frame.missing_properties {
if reported_missing_prop.insert((edge_type.clone(), col.clone())) {
report
.errors
.push(missing_fk_property_error(&spec.node_type, edge_type, col));
}
}
let (target_col, df) = (frame.target_col, frame.df);
let count = connect(
graph,
df,
edge_type,
&spec.node_type,
&pk,
&edge.target,
&target_col,
report,
initial_load[edge_type],
)?;
*report.edges_by_type.entry(edge_type.clone()).or_insert(0) += count;
}
}
for (edge_type, tally) in misparses {
report.warnings.extend(tally.into_warnings(&format!(
"fk_edge '{edge_type}' (node '{}')",
spec.node_type
)));
}
Ok(())
}
struct IdTypes<'a>(Option<&'a IndexMap<String, crate::datatypes::values::ColumnType>>);
impl IdTypes<'_> {
fn for_columns(
&self,
src: &str,
tgt: &str,
) -> (
Option<crate::datatypes::values::ColumnType>,
Option<crate::datatypes::values::ColumnType>,
) {
match self.0 {
Some(map) => (map.get(src).cloned(), map.get(tgt).cloned()),
None => (None, None),
}
}
}
fn build_edge_df(
src_name: &str,
tgt_name: &str,
src: Vec<Option<String>>,
tgt: Vec<Option<String>>,
resolved: (
Option<crate::datatypes::values::ColumnType>,
Option<crate::datatypes::values::ColumnType>,
),
) -> Result<DataFrame, String> {
let src_type = resolved.0.unwrap_or_else(|| infer_id_type(&src));
let tgt_type = resolved.1.unwrap_or_else(|| infer_id_type(&tgt));
let mut df = DataFrame::new(Vec::new());
add_id_column(&mut df, src_name, src, src_type)?;
add_id_column(&mut df, tgt_name, tgt, tgt_type)?;
Ok(df)
}
pub(super) fn infer_id_type(vals: &[Option<String>]) -> crate::datatypes::values::ColumnType {
let mut inference = super::super::typing::IdInference::default();
for v in vals {
if inference.is_settled() {
break;
}
if let Some(s) = v {
inference.observe(s);
}
}
inference.resolve()
}
pub(super) fn add_id_column(
df: &mut DataFrame,
name: &str,
vals: Vec<Option<String>>,
col_type: crate::datatypes::values::ColumnType,
) -> Result<(), String> {
use crate::datatypes::values::{ColumnData, ColumnType};
let data = match col_type {
ColumnType::Int64 => {
let ints: Vec<Option<i64>> = vals
.iter()
.map(|v| {
v.as_ref().and_then(|s| {
let t = s.trim();
if t.is_empty() {
None
} else if let Ok(i) = t.parse::<i64>() {
Some(i)
} else if let Ok(f) = t.parse::<f64>() {
if f.is_finite()
&& f.fract() == 0.0
&& f >= i64::MIN as f64
&& f <= i64::MAX as f64
{
Some(f as i64)
} else {
None
}
} else {
None
}
})
})
.collect();
ColumnData::Int64(ints)
}
_ => ColumnData::String(
vals.into_iter()
.map(|v| v.filter(|s| !s.is_empty()))
.collect(),
),
};
df.add_column(name.to_string(), col_type, data)
}
#[allow(clippy::too_many_arguments)]
pub(super) fn connect(
graph: &mut DirGraph,
df: DataFrame,
connection_type: &str,
source_type: &str,
source_id_field: &str,
target_type: &str,
target_id_field: &str,
report: &mut BuildReport,
initial_load: maintain::InitialLoad,
) -> Result<usize, String> {
match maintain::add_connections_with_initial_load(
graph,
df,
connection_type.to_string(),
source_type.to_string(),
source_id_field.to_string(),
target_type.to_string(),
target_id_field.to_string(),
None,
None,
None,
initial_load,
) {
Ok(r) => {
if r.connections_skipped > 0 {
let detail = r.errors.join("; ");
report.warnings.push(format!(
"[{}] -[{}]-> {}: {} skipped ({})",
source_type, connection_type, target_type, r.connections_skipped, detail
));
}
if r.stubs_vivified > 0 {
report.warnings.push(format!(
"[{}] -[{}]-> {}: {} stub node(s) vivified for missing endpoints",
source_type, connection_type, target_type, r.stubs_vivified
));
}
Ok(r.connections_created)
}
Err(e) => {
report
.errors
.push(format!("[{}] edge {}: {}", source_type, connection_type, e));
Ok(0)
}
}
}