use crate::datatypes::values::Value;
use crate::graph::schema::{DirGraph, InternedKey, NodeData, PropertyStorage};
use crate::graph::storage::mapped::mmap_vec::MmapOrVec;
use crate::graph::storage::type_build_meta::TypeBuildMeta;
use crate::graph::storage::{GraphRead, GraphWrite};
use flate2::read::GzDecoder;
use std::collections::HashMap;
#[cfg(test)]
use std::collections::HashSet;
use std::fs::File;
use std::io::{BufRead, BufReader, Read};
use std::path::Path;
use std::sync::Arc;
use std::time::Instant;
use super::column_builder::ColumnTypeMeta;
use super::parser::{
extract_lang_text, language_matches, parse_line, parse_qcode_number, typed_literal_to_value,
CompactNTripleEdge, EdgeBuffer, EntityAccumulator, Object, Predicate, Subject,
};
use super::writer::{create_edges_from_buffer, create_edges_with_qnum_map};
use super::{
Cancelled, NTriplesConfig, NTriplesStats, ProgressEvent, ProgressSink, ProgressValue as PV,
};
macro_rules! eplog {
($($arg:tt)*) => {
eprintln!("[{}] {}", chrono::Local::now().format("%H:%M:%S"), format_args!($($arg)*))
};
}
const CANCELLED_TOKEN: &str = "<cancelled>";
const READER_BATCH_SIZE: usize = 200_000;
const READER_TARGET_BATCH_BYTES: usize = 16 * 1024 * 1024;
type ReaderBatch = Result<super::parser::LineBuffer, String>;
fn spawn_reader(
reader: Box<dyn Read + Send>,
) -> (
std::sync::mpsc::Receiver<ReaderBatch>,
std::thread::JoinHandle<Result<(), String>>,
) {
let (tx, rx) = std::sync::mpsc::sync_channel::<ReaderBatch>(32);
let handle = std::thread::spawn(move || {
let mut reader = BufReader::with_capacity(8 * 1024 * 1024, reader);
let mut raw = Vec::with_capacity(512);
let mut batch =
super::parser::LineBuffer::with_capacity(READER_BATCH_SIZE, READER_TARGET_BATCH_BYTES);
let prefix: &[u8] = b"<http://www.wikidata.org/entity/Q";
loop {
raw.clear();
let bytes_read = match reader.read_until(b'\n', &mut raw) {
Ok(bytes_read) => bytes_read,
Err(error) => {
let message = format!("N-Triples reader error: {error}");
if !batch.is_empty() {
let next = super::parser::LineBuffer::with_capacity(
READER_BATCH_SIZE,
READER_TARGET_BATCH_BYTES,
);
let full = std::mem::replace(&mut batch, next);
if tx.send(Ok(full)).is_err() {
return Ok(());
}
}
if tx.send(Err(message.clone())).is_err() {
return Ok(());
}
return Err(message);
}
};
if bytes_read == 0 {
if !batch.is_empty() && tx.send(Ok(batch)).is_err() {
return Ok(());
}
return Ok(());
}
if !raw.starts_with(prefix) {
continue;
}
batch.push_line(&raw);
if batch.offsets.len() >= READER_BATCH_SIZE
|| batch.data.len() >= READER_TARGET_BATCH_BYTES
{
let next = super::parser::LineBuffer::with_capacity(
READER_BATCH_SIZE,
READER_TARGET_BATCH_BYTES,
);
let full = std::mem::replace(&mut batch, next);
if tx.send(Ok(full)).is_err() {
return Ok(());
}
}
}
});
(rx, handle)
}
fn join_reader(handle: std::thread::JoinHandle<Result<(), String>>) -> Result<(), String> {
match handle.join() {
Ok(result) => result,
Err(payload) => {
let detail = payload
.downcast_ref::<&str>()
.copied()
.or_else(|| payload.downcast_ref::<String>().map(String::as_str))
.unwrap_or("unknown panic");
Err(format!("N-Triples reader thread panicked: {detail}"))
}
}
}
#[inline]
fn validated_line(bytes: &[u8]) -> Result<&str, std::str::Utf8Error> {
if bytes.is_ascii() {
Ok(unsafe { std::str::from_utf8_unchecked(bytes) })
} else {
std::str::from_utf8(bytes)
}
}
#[inline]
fn emit(sink: Option<&dyn ProgressSink>, event: ProgressEvent<'_>) -> Result<(), String> {
if let Some(s) = sink {
s.emit(event)
.map_err(|Cancelled| CANCELLED_TOKEN.to_string())?;
}
Ok(())
}
fn apply_type_renames(
backend: &mut crate::graph::schema::GraphBackend,
rename_map: &HashMap<u64, u64>,
) {
use crate::graph::schema::GraphBackend;
match backend {
GraphBackend::Disk(dg) => {
for i in 0..dg.node_slot_len() {
let slot = dg.node_slot(i);
if slot.is_alive() {
if let Some(&new_type) = rename_map.get(&slot.node_type) {
let mut new_slot = slot;
new_slot.node_type = new_type;
dg.set_node_slot(i, new_slot);
}
}
}
}
GraphBackend::Memory(g) => retype_heap_nodes(
crate::graph::storage::backend::unique_heap_backend(g),
rename_map,
),
GraphBackend::Mapped(g) => retype_heap_nodes(
crate::graph::storage::backend::unique_heap_backend(g),
rename_map,
),
GraphBackend::Recording(rg) => apply_type_renames(rg.inner_mut(), rename_map),
GraphBackend::Forked(_) => {
backend.flatten_fork();
apply_type_renames(backend, rename_map);
}
}
}
fn retype_heap_nodes(g: &mut impl GraphWrite, rename_map: &HashMap<u64, u64>) {
for i in 0..g.node_bound() {
let idx = petgraph::graph::NodeIndex::new(i);
if let Some(node) = g.node_weight_mut(idx) {
if let Some(&new_key) = rename_map.get(&node.node_type.as_u64()) {
node.node_type = InternedKey::from_u64(new_key);
}
}
}
}
fn open_ntriples_reader(path: &Path, display_path: &str) -> Result<Box<dyn Read + Send>, String> {
if display_path.ends_with(".bz2") {
return super::parallel_bz2::open(path)
.map_err(|error| format!("Cannot open {display_path}: {error}"));
}
let file = File::open(path).map_err(|error| format!("Cannot open {display_path}: {error}"))?;
if display_path.ends_with(".gz") {
Ok(Box::new(GzDecoder::new(BufReader::new(file))))
} else if display_path.ends_with(".zst") || display_path.ends_with(".zstd") {
Ok(Box::new(
zstd::Decoder::new(BufReader::new(file))
.map_err(|error| format!("zstd decoder error: {error}"))?,
))
} else {
Ok(Box::new(file))
}
}
fn finalize_disk_graph(
graph: &mut DirGraph,
config: &NTriplesConfig,
sink: Option<&dyn ProgressSink>,
) -> Result<(), String> {
let finalising_start = Instant::now();
if config.verbose && graph.graph.is_disk() {
eplog!("[Finalising] Building auxiliary indexes + saving metadata");
}
if graph.graph.is_disk() {
emit(
sink,
ProgressEvent::Start {
phase: "finalising",
label: "Finalising: Auxiliary indexes + metadata",
total: None,
unit: "step",
},
)?;
}
if graph.graph.is_disk() {
if let crate::graph::schema::GraphBackend::Disk(ref mut dg) = graph.graph {
if let Some(raw_counts) = dg.edge_type_counts_raw.take() {
let string_counts: HashMap<String, usize> = raw_counts
.into_iter()
.map(|(key_u64, count)| {
let key = InternedKey::from_u64(key_u64);
let name = graph.interner.resolve(key).to_string();
(name, count)
})
.collect();
if build_debug() {
eplog!(
" Cached {} edge type counts from CSR build",
string_counts.len()
);
}
*graph.edge_type_counts_cache.write().unwrap() = Some(Arc::new(string_counts));
}
}
}
if graph.graph.is_disk() {
let rebuild_start = Instant::now();
if let crate::graph::schema::GraphBackend::Disk(ref dg) = graph.graph {
for i in 0..dg.node_slot_len() {
let slot = dg.node_slot(i);
if slot.is_alive() {
let type_key = InternedKey::from_u64(slot.node_type);
let type_name = graph.interner.resolve(type_key).to_string();
graph
.type_indices
.entry_or_default(type_name)
.push(petgraph::graph::NodeIndex::new(i));
}
}
}
if build_debug() {
eplog!(
" Rebuilt {} type indices ({})",
graph.type_indices.len(),
fmt_dur(rebuild_start.elapsed().as_secs_f64()),
);
}
}
if graph.graph.is_disk() {
let data_dir = if let crate::graph::schema::GraphBackend::Disk(ref dg) = graph.graph {
dg.active_write_dir().to_path_buf()
} else {
std::path::PathBuf::new()
};
let mmap_path = data_dir.join("columns.bin");
let meta_path = data_dir.join("columns_meta.json");
if mmap_path.exists() && meta_path.exists() {
let reload_start = Instant::now();
let reload_result: Result<(), String> = (|| {
let meta_json = std::fs::read_to_string(&meta_path)
.map_err(|e| format!("read columns_meta.json: {}", e))?;
let columns_meta: Vec<ColumnTypeMeta> = serde_json::from_str(&meta_json)
.map_err(|e| format!("parse columns_meta.json: {}", e))?;
let file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.open(&mmap_path)
.map_err(|e| format!("open columns.bin: {}", e))?;
let mmap = unsafe {
memmap2::MmapMut::map_mut(&file)
.map_err(|e| format!("mmap columns.bin: {}", e))?
};
let mmap_arc = Arc::new(mmap);
for type_meta in &columns_meta {
let mmap_store = type_meta.to_mmap_store(Arc::clone(&mmap_arc));
let store = crate::graph::storage::column_store::ColumnStore::from_mmap_store(
Arc::new(mmap_store),
);
graph.install_column_store(&type_meta.type_name, Arc::new(store));
}
Ok(())
})();
if let Err(e) = reload_result {
eplog!(" Warning: failed to reload column stores: {}", e);
}
if build_debug() {
eplog!(
" Reloaded {} column stores from mmap ({})",
graph.column_store_count(),
fmt_dur(reload_start.elapsed().as_secs_f64()),
);
}
}
}
if graph.graph.is_disk() {
let id_start = Instant::now();
let type_names: Vec<String> = graph.type_indices.keys().map(|s| s.to_string()).collect();
for type_name in &type_names {
graph.build_id_index(type_name);
}
if build_debug() {
eplog!(
" Built {} id indices ({})",
type_names.len(),
fmt_dur(id_start.elapsed().as_secs_f64()),
);
}
}
if graph.graph.is_disk() {
if let crate::graph::schema::GraphBackend::Disk(ref mut dg) = graph.graph {
let flush_step = Instant::now();
dg.flush_node_slots()
.map_err(|e| format!("Failed to flush node slots: {e}"))?;
if build_debug() {
eplog!(
" node_slots.bin flush: {}",
fmt_dur(flush_step.elapsed().as_secs_f64())
);
}
}
build_type_connectivity_cache(graph);
let generation_root = match &graph.graph {
crate::graph::schema::GraphBackend::Disk(dg) => dg
.bulk_build_generation_root()
.map(std::path::Path::to_path_buf),
_ => None,
};
match generation_root {
Some(root) => publish_build_as_generation(graph, &root)?,
None => publish_build_in_place(graph)?,
}
let finalising_elapsed = finalising_start.elapsed().as_secs_f64();
if config.verbose {
eplog!("[Finalising] Complete in {}", fmt_dur(finalising_elapsed));
}
emit(
sink,
ProgressEvent::Complete {
phase: "finalising",
elapsed_s: finalising_elapsed,
fields: &[],
},
)?;
}
Ok(())
}
fn publish_flat_root(root_dir: &Path, metadata_json: &str) -> Result<(), String> {
std::fs::write(root_dir.join("metadata.json"), metadata_json)
.map_err(|e| format!("Failed to publish graph metadata: {e}"))?;
match std::fs::remove_file(root_dir.join("CURRENT")) {
Ok(()) => {
if let Ok(handle) = std::fs::File::open(root_dir) {
let _ = handle.sync_all();
}
Ok(())
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(e) => Err(format!(
"Failed to retire the previous generation pointer: {e}"
)),
}
}
fn build_type_connectivity_cache(graph: &mut DirGraph) {
let mut triples = Vec::new();
for (conn_type, info) in graph.connection_type_metadata.iter() {
let edge_count = graph
.edge_type_counts_cache
.read()
.unwrap()
.as_ref()
.and_then(|counts| counts.get(conn_type).copied())
.unwrap_or(0);
for src in &info.source_types {
for tgt in &info.target_types {
triples.push(crate::graph::schema::ConnectivityTriple {
src: src.clone(),
conn: conn_type.clone(),
tgt: tgt.clone(),
count: edge_count,
});
}
}
}
if !triples.is_empty() {
if build_debug() {
eplog!(
" Built type connectivity cache ({} triples)",
triples.len()
);
}
*graph.type_connectivity_cache.write().unwrap() = Some(triples);
}
}
fn publish_build_in_place(graph: &DirGraph) -> Result<(), String> {
let crate::graph::schema::GraphBackend::Disk(ref dg) = graph.graph else {
return Ok(());
};
let data_dir = dg.active_write_dir().to_path_buf();
let root_dir = data_dir
.parent()
.map(std::path::Path::to_path_buf)
.unwrap_or_else(|| data_dir.clone());
let save_step = Instant::now();
let interner_map: HashMap<String, String> = graph
.interner
.iter()
.map(|(k, v)| (k.as_u64().to_string(), v.to_string()))
.collect();
let json = serde_json::to_string(&interner_map)
.map_err(|e| format!("Failed to serialize interner: {e}"))?;
std::fs::write(root_dir.join("interner.json"), json)
.map_err(|e| format!("Failed to write interner: {e}"))?;
if build_debug() {
eplog!(
" interner.json ({} entries): {}",
interner_map.len(),
fmt_dur(save_step.elapsed().as_secs_f64())
);
}
let save_step = Instant::now();
crate::graph::io::file::write_node_type_metadata_bin(&root_dir, graph)
.map_err(|e| format!("Failed to write node-type metadata: {e}"))?;
crate::graph::io::file::write_connection_type_metadata_bin(&root_dir, graph)
.map_err(|e| format!("Failed to write connection-type metadata: {e}"))?;
let mut meta = crate::graph::io::file::build_disk_metadata(graph);
crate::graph::io::file::strip_heavy_metadata(&mut meta);
let json = serde_json::to_string_pretty(&meta)
.map_err(|e| format!("Failed to serialize graph metadata: {e}"))?;
if build_debug() {
eplog!(" metadata: {}", fmt_dur(save_step.elapsed().as_secs_f64()));
}
let save_step = Instant::now();
if !graph.id_indices.is_empty() {
crate::graph::storage::disk::id_index::write_id_indices_bin(
&root_dir,
&graph.id_indices,
&graph.interner,
)
.map_err(|e| format!("Failed to write id indexes: {e}"))?;
}
if build_debug() {
eplog!(
" id_indices.bin ({} types): {}",
graph.id_indices.len(),
fmt_dur(save_step.elapsed().as_secs_f64())
);
}
let save_step = Instant::now();
if !graph.type_indices.is_empty() {
crate::graph::storage::disk::type_index::write_type_indices_bin(
&root_dir,
&graph.type_indices,
&graph.interner,
)
.map_err(|e| format!("Failed to write type indexes: {e}"))?;
}
if build_debug() {
eplog!(
" type_indices.bin ({} types): {}",
graph.type_indices.len(),
fmt_dur(save_step.elapsed().as_secs_f64())
);
}
publish_flat_root(&root_dir, &json)?;
Ok(())
}
fn publish_build_as_generation(graph: &mut DirGraph, root: &Path) -> Result<(), String> {
let publish_step = Instant::now();
let root_str = root.to_str().ok_or_else(|| {
format!(
"Cannot publish the disk build: graph root is not valid UTF-8: {}",
root.display()
)
})?;
graph.save_disk(root_str)?;
if build_debug() {
eplog!(
" published build as a new generation: {}",
fmt_dur(publish_step.elapsed().as_secs_f64())
);
}
Ok(())
}
fn build_columns(
graph: &mut DirGraph,
config: &NTriplesConfig,
sink: Option<&dyn ProgressSink>,
mut prop_log: Option<crate::graph::storage::memory::property_log::PropertyLogWriter>,
type_meta: HashMap<String, TypeBuildMeta>,
type_rename_map: &HashMap<String, String>,
) -> Result<(), String> {
if let Some(log_writer) = prop_log.take() {
let phase1b_total = log_writer.count();
let phase1b_label = format!(
"Phase 1b: Building columnar storage ({} entities, {} types)",
format_count(phase1b_total),
type_meta.len(),
);
if config.verbose {
eplog!("[Phase 1b] {}", phase1b_label);
}
emit(
sink,
ProgressEvent::Start {
phase: "phase1b",
label: &phase1b_label,
total: Some(phase1b_total),
unit: "ent",
},
)?;
let conv_start = Instant::now();
let log_path = log_writer
.finish()
.map_err(|e| format!("Failed to finish property log: {}", e))?;
let build_result = super::column_builder::build_columns_direct(
graph,
&log_path,
&type_meta,
type_rename_map,
build_debug(),
sink,
);
let _ = std::fs::remove_file(&log_path);
if let Some(parent) = log_path.parent() {
let _ = std::fs::remove_dir(parent);
}
match build_result {
Ok(()) => {}
Err(super::column_builder::BuildColumnsError::Cancelled) => {
return Err(CANCELLED_TOKEN.to_string());
}
Err(super::column_builder::BuildColumnsError::Io(error)) => {
return Err(format!("Failed to build columns: {error}"));
}
}
let phase1b_elapsed = conv_start.elapsed().as_secs_f64();
if config.verbose {
eplog!("[Phase 1b] Complete in {}", fmt_dur(phase1b_elapsed));
}
emit(
sink,
ProgressEvent::Complete {
phase: "phase1b",
elapsed_s: phase1b_elapsed,
fields: &[("entities", PV::U64(phase1b_total))],
},
)?;
} else {
debug_assert!(
!graph.graph.is_mapped(),
"mapped load_ntriples must populate prop_log and use build_columns_direct"
);
graph.enable_columnar();
}
if graph.graph.is_disk() {
let dropped_stores = graph.column_store_count();
graph.clear_column_stores();
drop(type_meta);
let type_indices_count = graph.type_indices.len();
graph.type_indices.clear();
if build_debug() {
eplog!(
" Freed {} column stores + {} type indices before Phase 2",
dropped_stores,
type_indices_count,
);
}
}
Ok(())
}
fn build_edges(
graph: &mut DirGraph,
config: &NTriplesConfig,
sink: Option<&dyn ProgressSink>,
stats: &mut NTriplesStats,
edge_buffer: EdgeBuffer,
mut qnum_to_idx: Option<MmapOrVec<u32>>,
) -> Result<(), String> {
let phase2_total = edge_buffer.len() as u64;
if config.verbose {
eplog!("[Phase 2] Creating edges");
let _ = std::io::Write::flush(&mut std::io::stderr());
}
emit(
sink,
ProgressEvent::Start {
phase: "phase2",
label: "Phase 2: Creating edges",
total: Some(phase2_total),
unit: "edge",
},
)?;
let edge_start = Instant::now();
if let Some(ref qt) = qnum_to_idx {
create_edges_with_qnum_map(graph, &edge_buffer, stats, qt, sink)?;
} else {
create_edges_from_buffer(graph, &edge_buffer, stats, sink)?;
}
let phase2_elapsed = edge_start.elapsed().as_secs_f64();
if config.verbose {
eplog!(
"[Phase 2] Complete: {} edges created ({} skipped) in {}",
format_count(stats.edges_created),
format_count(stats.edges_skipped),
fmt_dur(phase2_elapsed),
);
}
emit(
sink,
ProgressEvent::Complete {
phase: "phase2",
elapsed_s: phase2_elapsed,
fields: &[
("edges_created", PV::U64(stats.edges_created)),
("edges_skipped", PV::U64(stats.edges_skipped)),
],
},
)?;
if let Some(qt) = qnum_to_idx.take() {
let qt_path = qt.file_path().map(|p| p.to_path_buf());
drop(qt);
if let Some(path) = qt_path {
let _ = std::fs::remove_file(path);
}
}
let edge_file_path = match &edge_buffer {
EdgeBuffer::Compact(buf) => buf.file_path().map(|p| p.to_path_buf()),
_ => None,
};
drop(edge_buffer);
if let Some(path) = edge_file_path {
let _ = std::fs::remove_file(&path);
}
Ok(())
}
fn build_disk_csr(
graph: &mut DirGraph,
config: &NTriplesConfig,
sink: Option<&dyn ProgressSink>,
) -> Result<(), String> {
if let crate::graph::schema::GraphBackend::Disk(ref mut dg) = graph.graph {
if config.verbose {
eplog!("[Phase 3] Building CSR edge index");
}
emit(
sink,
ProgressEvent::Start {
phase: "phase3",
label: "Phase 3: Building CSR edge index",
total: None,
unit: "step",
},
)?;
let csr_start = Instant::now();
dg.build_csr_from_pending()
.map_err(|e| format!("Failed to build disk CSR: {e}"))?;
let phase3_elapsed = csr_start.elapsed().as_secs_f64();
if config.verbose {
eplog!("[Phase 3] Complete in {}", fmt_dur(phase3_elapsed));
}
emit(
sink,
ProgressEvent::Complete {
phase: "phase3",
elapsed_s: phase3_elapsed,
fields: &[],
},
)?;
}
Ok(())
}
fn resolve_type_labels(
graph: &mut DirGraph,
config: &NTriplesConfig,
mut label_writer: Option<super::label_spill::LabelSpillWriter>,
type_meta: &mut HashMap<String, TypeBuildMeta>,
) -> Result<HashMap<String, String>, String> {
let mut type_rename_map: HashMap<String, String> = HashMap::new();
let label_journal_path = if let Some(writer) = label_writer.take() {
let spill_dir = graph.spill_dir.clone().unwrap_or_else(|| {
std::env::temp_dir().join(format!("kglite_build_{}", std::process::id()))
});
let path = spill_dir.join("labels.bin");
let journal_size = writer
.finish()
.map_err(|e| format!("Failed to finish label journal: {e}"))?;
if build_debug() {
eplog!(
" Label journal: {} bytes on disk",
format_count(journal_size)
);
}
Some(path)
} else {
None
};
if config.auto_type {
let wanted: std::collections::HashSet<u32> = graph
.type_indices
.keys()
.filter_map(parse_qcode_number)
.collect();
let label_lookup: HashMap<u32, String> = if let Some(ref path) = label_journal_path {
super::label_spill::read_labels_for(path, &wanted).unwrap_or_else(|e| {
eplog!(" WARN: failed to read label journal: {}", e);
HashMap::new()
})
} else {
HashMap::new()
};
let mut renames: Vec<(String, String)> = Vec::new();
for type_name in graph.type_indices.keys() {
if let Some(qnum) = parse_qcode_number(type_name) {
if let Some(label) = label_lookup.get(&qnum) {
if label != type_name {
renames.push((type_name.to_string(), label.clone()));
}
}
}
}
if !renames.is_empty() {
let rename_start = std::time::Instant::now();
if build_debug() {
eplog!(
" Resolving {} Q-code type names to labels...",
renames.len()
);
}
graph
.interner
.validate_names(
renames
.iter()
.flat_map(|(old, new)| [old.as_str(), new.as_str()]),
)
.map_err(|e| e.to_string())?;
let old_key_to_new_key: Vec<(InternedKey, InternedKey)> = renames
.iter()
.map(|(old, new)| {
let old_key = graph.interner.get_or_intern(old);
let new_key = graph.interner.get_or_intern(new);
(old_key, new_key)
})
.collect();
for (old_name, new_name) in &renames {
if let Some(indices) = graph.type_indices.remove(old_name) {
graph
.type_indices
.entry_or_default(new_name.clone())
.extend(indices);
}
if let Some(old_meta) = graph.node_type_metadata_mut().remove(old_name) {
let entry = graph
.node_type_metadata_mut()
.entry(new_name.clone())
.or_default();
for (k, v) in old_meta {
entry.entry(k).or_insert(v);
}
}
if let Some(old_schema) = graph.type_schemas_mut().remove(old_name) {
if let Some(existing) = graph.type_schemas.get(new_name) {
let merged = existing.merge(&old_schema);
graph
.type_schemas_mut()
.insert(new_name.clone(), Arc::new(merged));
} else {
graph
.type_schemas_mut()
.insert(new_name.clone(), old_schema);
}
}
if let Some(old_build) = type_meta.remove(old_name) {
let entry = type_meta
.entry(new_name.clone())
.or_insert_with(TypeBuildMeta::new);
entry.merge_from(&old_build);
}
}
for (old_name, new_name) in &renames {
type_rename_map.insert(old_name.clone(), new_name.clone());
}
let rename_map: HashMap<u64, u64> = old_key_to_new_key
.iter()
.map(|(old, new)| (old.as_u64(), new.as_u64()))
.collect();
apply_type_renames(&mut graph.graph, &rename_map);
if build_debug() {
eplog!(
" Resolved {} Q-code types ({})",
renames.len(),
fmt_dur(rename_start.elapsed().as_secs_f64())
);
}
}
}
Ok(type_rename_map)
}
struct LoadSpills {
prop_log: Option<crate::graph::storage::memory::property_log::PropertyLogWriter>,
edge_buffer: EdgeBuffer,
label_writer: Option<super::label_spill::LabelSpillWriter>,
type_meta: HashMap<String, TypeBuildMeta>,
qnum_to_idx: Option<MmapOrVec<u32>>,
}
fn initialize_spills(graph: &mut DirGraph, config: &NTriplesConfig) -> Result<LoadSpills, String> {
let use_streaming_build = graph.graph.is_disk() || graph.graph.is_mapped();
let use_compact = use_streaming_build;
let prop_log: Option<crate::graph::storage::memory::property_log::PropertyLogWriter> =
if use_streaming_build {
let spill_dir = graph.spill_dir.clone().unwrap_or_else(|| {
std::env::temp_dir().join(format!("kglite_build_{}", std::process::id()))
});
const STALE_AFTER_SECS: u64 = 3600;
if let Some(parent) = spill_dir.parent() {
if let Ok(entries) = std::fs::read_dir(parent) {
let now = std::time::SystemTime::now();
for entry in entries.flatten() {
let name = entry.file_name();
let name = name.to_string_lossy();
if !name.starts_with("kglite_build_") {
continue;
}
if entry.path() == spill_dir {
continue;
}
let is_stale = entry
.metadata()
.and_then(|m| m.modified())
.ok()
.and_then(|m| now.duration_since(m).ok())
.map(|d| d.as_secs() > STALE_AFTER_SECS)
.unwrap_or(false);
if is_stale {
let _ = std::fs::remove_dir_all(entry.path());
}
}
}
}
if let crate::graph::schema::GraphBackend::Disk(ref mut dg) = graph.graph {
let stale = dg.active_write_dir().join("_pending_edges.bin");
let mapped_live = dg.pending_edges.get_mut().file_path() == Some(stale.as_path());
if stale.exists() && !mapped_live {
let _ = std::fs::remove_file(&stale);
}
}
let log_path = spill_dir.join("properties.log.zst");
if build_debug() {
eplog!(" Property log: {}", log_path.display());
}
Some(
crate::graph::storage::memory::property_log::PropertyLogWriter::new(&log_path, 1)
.map_err(|e| format!("Failed to create property log: {}", e))?,
)
} else {
None
};
let edge_buffer = if use_compact {
if use_streaming_build {
let spill_dir = graph.spill_dir.clone().unwrap_or_else(|| {
std::env::temp_dir().join(format!("kglite_build_{}", std::process::id()))
});
std::fs::create_dir_all(&spill_dir)
.map_err(|e| format!("Failed to create spill dir: {}", e))?;
let edge_path = spill_dir.join("edges.bin");
EdgeBuffer::Compact(
MmapOrVec::mapped(&edge_path, 1 << 20)
.map_err(|e| format!("Failed to create edge buffer file: {}", e))?,
)
} else {
EdgeBuffer::Compact(MmapOrVec::new())
}
} else {
EdgeBuffer::Strings(Vec::new())
};
let label_writer: Option<super::label_spill::LabelSpillWriter> = if config.auto_type {
let spill_dir = graph.spill_dir.clone().unwrap_or_else(|| {
std::env::temp_dir().join(format!("kglite_build_{}", std::process::id()))
});
std::fs::create_dir_all(&spill_dir)
.map_err(|e| format!("Failed to create label spill directory: {e}"))?;
Some(
super::label_spill::LabelSpillWriter::new(&spill_dir.join("labels.bin"))
.map_err(|e| format!("Failed to create label journal: {}", e))?,
)
} else {
None
};
let type_meta: HashMap<String, TypeBuildMeta> = HashMap::new();
let qnum_to_idx: Option<MmapOrVec<u32>> = if graph.graph.is_disk() {
let spill_dir = graph.spill_dir.clone().unwrap_or_else(|| {
std::env::temp_dir().join(format!("kglite_build_{}", std::process::id()))
});
std::fs::create_dir_all(&spill_dir)
.map_err(|e| format!("Failed to create id spill directory: {e}"))?;
Some(
MmapOrVec::mapped_prefilled(&spill_dir.join("qnum_to_idx.bin"), 150_000_000)
.unwrap_or_else(|_| MmapOrVec::from_vec(vec![0u32; 150_000_000])),
)
} else {
None
};
Ok(LoadSpills {
prop_log,
edge_buffer,
label_writer,
type_meta,
qnum_to_idx,
})
}
struct Phase1Ingest<'a> {
graph: &'a mut DirGraph,
path: &'a str,
path_obj: &'a Path,
config: &'a NTriplesConfig,
stats: &'a mut NTriplesStats,
edge_buffer: &'a mut EdgeBuffer,
prop_log: &'a mut Option<crate::graph::storage::memory::property_log::PropertyLogWriter>,
label_writer: &'a mut Option<super::label_spill::LabelSpillWriter>,
type_meta: &'a mut HashMap<String, TypeBuildMeta>,
qnum_to_idx: &'a mut Option<MmapOrVec<u32>>,
rx: std::sync::mpsc::Receiver<ReaderBatch>,
reader_handle: std::thread::JoinHandle<Result<(), String>>,
started: Instant,
}
fn ingest_phase1(context: Phase1Ingest<'_>) -> Result<(), String> {
let Phase1Ingest {
graph,
path,
path_obj,
config,
stats,
edge_buffer,
prop_log,
label_writer,
type_meta,
qnum_to_idx,
rx,
reader_handle,
started: start,
} = context;
let mut current: Option<EntityAccumulator> = None;
let mut entity_limit_reached = false;
let mut scratch_props: Vec<(InternedKey, Value)> = Vec::with_capacity(64);
let mut progress_countdown: u64 = 5_000_000;
let mut last_progress_log = Instant::now();
const PROGRESS_BUCKET: u64 = 5_000_000;
const PROGRESS_INTERVAL_SECS: f64 = 60.0;
if config.verbose {
eplog!("[Phase 1] Streaming and parsing N-triples ({})", path);
}
let sink = config.progress.as_deref();
let phase1_label = format!(
"Phase 1: Streaming N-triples ({})",
path_obj
.file_name()
.and_then(|s| s.to_str())
.unwrap_or(path)
);
emit(
sink,
ProgressEvent::Start {
phase: "phase1",
label: &phase1_label,
total: config.max_triples,
unit: "tri",
},
)?;
let mut reader_error = None;
'outer: while let Ok(message) = rx.recv() {
let batch = match message {
Ok(batch) => batch,
Err(error) => {
reader_error = Some(error);
break;
}
};
let n_lines = batch.offsets.len();
for i in 0..n_lines {
let line = match validated_line(batch.line(i)) {
Ok(line) => line,
Err(error) => {
reader_error = Some(format!(
"N-Triples input contains invalid UTF-8 in an accepted entity line: {error}"
));
break 'outer;
}
};
if entity_limit_reached && current.is_none() {
break 'outer;
}
if let Some(cap) = config.max_triples {
if stats.triples_scanned >= cap {
break 'outer;
}
}
stats.triples_scanned += 1;
progress_countdown -= 1;
if progress_countdown == 0 {
progress_countdown = PROGRESS_BUCKET;
let buf_len = edge_buffer.len() as u64;
emit(
sink,
ProgressEvent::Update {
phase: "phase1",
current: stats.triples_scanned,
fields: &[
("entities", PV::U64(stats.entities_created)),
("edges_buffered", PV::U64(buf_len)),
],
},
)?;
if config.verbose
&& last_progress_log.elapsed().as_secs_f64() >= PROGRESS_INTERVAL_SECS
{
let elapsed = start.elapsed().as_secs_f64();
let rate = stats.triples_scanned as f64 / elapsed;
eplog!(
"[Phase 1] {} triples, {} entities, {} edges buffered — {:.0}k triples/s",
format_count(stats.triples_scanned),
format_count(stats.entities_created),
format_count(buf_len),
rate / 1000.0,
);
last_progress_log = Instant::now();
}
}
let (subject, predicate, object) = match parse_line(line) {
Some(parsed) => parsed,
None => continue,
};
let subj_id = match subject {
Subject::Entity(id) => id,
Subject::Other => continue,
};
if current.as_ref().is_some_and(|c| c.id != subj_id) {
if let Some(acc) = current.take() {
flush_entity(
graph,
acc,
config,
edge_buffer,
stats,
prop_log,
label_writer,
type_meta,
qnum_to_idx,
&mut scratch_props,
)?;
}
if entity_limit_reached {
break;
}
}
if current.is_none() {
if let Some(max) = config.max_entities {
if stats.entities_created >= max as u64 {
entity_limit_reached = true;
continue;
}
}
current = Some(EntityAccumulator::new(subj_id.to_string()));
}
let acc = current.as_mut().unwrap();
match predicate {
Predicate::Label => {
if let Some(text) = extract_lang_text(&object, &config.languages) {
acc.label = Some(text);
}
}
Predicate::Description => {
if let Some(text) = extract_lang_text(&object, &config.languages) {
acc.description = Some(text);
}
}
Predicate::AltLabel => {}
Predicate::Type => {}
Predicate::WikidataDirect(pcode) => {
if let Some(ref allowed) = config.predicates {
if !allowed.contains(pcode) {
continue;
}
}
let pred_label = config
.predicate_labels
.get(pcode)
.cloned()
.unwrap_or_else(|| pcode.to_string());
match &object {
Object::Entity(target_qcode) => {
if pcode == "P31" && acc.type_qcode.is_none() {
acc.type_qcode = Some(target_qcode.to_string());
}
acc.outgoing_edges
.push((pred_label, target_qcode.to_string()));
}
Object::Literal(text) => {
acc.properties
.insert(pred_label, Value::String(text.clone()));
}
Object::LangLiteral(text, lang) => {
if language_matches(lang, &config.languages) {
acc.properties
.insert(pred_label, Value::String(text.clone()));
}
}
Object::TypedLiteral(text, type_uri) => {
acc.properties
.insert(pred_label, typed_literal_to_value(text, type_uri));
}
Object::Other => {}
}
}
Predicate::Other => {}
}
} } drop(rx);
let join_result = join_reader(reader_handle);
if let Some(error) = reader_error {
return Err(error);
}
join_result?;
if let Some(acc) = current.take() {
flush_entity(
graph,
acc,
config,
edge_buffer,
stats,
prop_log,
label_writer,
type_meta,
qnum_to_idx,
&mut scratch_props,
)?;
}
Ok(())
}
fn open_reader_for_load(
graph: &mut DirGraph,
path: &Path,
display_path: &str,
) -> Result<Box<dyn Read + Send>, String> {
let reader = open_ntriples_reader(path, display_path)?;
if let crate::graph::schema::GraphBackend::Disk(ref mut disk) = graph.graph {
disk.prepare_bulk_load_workspace()
.map_err(|error| format!("Failed to prepare the disk build workspace: {error}"))?;
}
Ok(reader)
}
pub fn load_ntriples(
graph: &mut DirGraph,
path: &str,
config: &NTriplesConfig,
) -> Result<NTriplesStats, String> {
let start = Instant::now();
let path_obj = Path::new(path);
let reader = open_reader_for_load(graph, path_obj, path)?;
let (rx, reader_handle) = spawn_reader(reader);
let mut stats = NTriplesStats {
triples_scanned: 0,
entities_created: 0,
edges_created: 0,
edges_skipped: 0,
seconds: 0.0,
};
let LoadSpills {
mut prop_log,
mut edge_buffer,
mut label_writer,
mut type_meta,
mut qnum_to_idx,
} = initialize_spills(graph, config)?;
ingest_phase1(Phase1Ingest {
graph,
path,
path_obj,
config,
stats: &mut stats,
edge_buffer: &mut edge_buffer,
prop_log: &mut prop_log,
label_writer: &mut label_writer,
type_meta: &mut type_meta,
qnum_to_idx: &mut qnum_to_idx,
rx,
reader_handle,
started: start,
})?;
let sink = config.progress.as_deref();
let type_rename_map = resolve_type_labels(graph, config, label_writer, &mut type_meta)?;
let phase1_elapsed = start.elapsed().as_secs_f64();
let phase1_buf_len = edge_buffer.len() as u64;
let phase1_num_types = type_meta.len() as u64;
let phase1_total_cols: u64 = type_meta.values().map(|m| m.columns.len() as u64).sum();
if config.verbose {
eplog!(
"[Phase 1] Complete: {} entities, {} types ({} columns), {} edges buffered in {}",
format_count(stats.entities_created),
format_count(phase1_num_types),
format_count(phase1_total_cols),
format_count(phase1_buf_len),
fmt_dur(phase1_elapsed),
);
}
emit(
sink,
ProgressEvent::Complete {
phase: "phase1",
elapsed_s: phase1_elapsed,
fields: &[
("entities", PV::U64(stats.entities_created)),
("edges_buffered", PV::U64(phase1_buf_len)),
("types", PV::U64(phase1_num_types)),
("columns", PV::U64(phase1_total_cols)),
("triples_scanned", PV::U64(stats.triples_scanned)),
],
},
)?;
build_columns(graph, config, sink, prop_log, type_meta, &type_rename_map)?;
build_edges(graph, config, sink, &mut stats, edge_buffer, qnum_to_idx)?;
build_disk_csr(graph, config, sink)?;
finalize_disk_graph(graph, config, sink)?;
stats.seconds = start.elapsed().as_secs_f64();
if config.verbose {
eplog!("[Build] Total elapsed: {}", fmt_dur(stats.seconds));
}
Ok(stats)
}
#[allow(clippy::too_many_arguments)]
fn flush_entity(
graph: &mut DirGraph,
acc: EntityAccumulator,
config: &NTriplesConfig,
edge_buffer: &mut EdgeBuffer,
stats: &mut NTriplesStats,
prop_log: &mut Option<crate::graph::storage::memory::property_log::PropertyLogWriter>,
label_writer: &mut Option<super::label_spill::LabelSpillWriter>,
type_meta: &mut HashMap<String, TypeBuildMeta>,
qnum_to_idx: &mut Option<MmapOrVec<u32>>,
scratch_props: &mut Vec<(InternedKey, Value)>,
) -> Result<(), String> {
let use_compact_ids = graph.graph.is_disk() || graph.graph.is_mapped();
if use_compact_ids && parse_qcode_number(&acc.id).is_none() {
return Ok(());
}
let title = acc.label.unwrap_or_else(|| acc.id.clone());
if let Some(ref mut w) = label_writer {
if let Some(qnum) = parse_qcode_number(&acc.id) {
w.append(qnum, &title)
.map_err(|e| format!("Failed to append label journal: {e}"))?;
}
}
let node_type = if let Some(ref tq) = acc.type_qcode {
if let Some(mapped_name) = config.node_types.get(tq) {
mapped_name.clone()
} else if config.auto_type {
tq.clone()
} else {
"Entity".to_string()
}
} else {
"Entity".to_string()
};
let mut names = vec![node_type.as_str(), "nid", "description", "P31"];
names.extend(acc.properties.keys().map(String::as_str));
names.extend(
acc.outgoing_edges
.iter()
.map(|(predicate, _)| predicate.as_str()),
);
graph
.interner
.validate_names(names)
.map_err(|e| e.to_string())?;
let mut properties = acc.properties;
properties.insert("nid".to_string(), Value::String(acc.id.clone()));
if let Some(desc) = acc.description {
properties.insert("description".to_string(), Value::String(desc));
}
if let Some(ref tq) = acc.type_qcode {
properties.insert("P31".to_string(), Value::String(tq.clone()));
}
let id_value = parse_qcode_number(&acc.id)
.map(Value::UniqueId)
.unwrap_or_else(|| Value::String(acc.id.clone()));
let title_value = Value::String(title);
let mut node_data = NodeData::new(
id_value.clone(),
title_value,
node_type.clone(),
properties,
&mut graph.interner,
);
let saved_id = node_data.id.clone();
let saved_title = node_data.title.clone();
if prop_log.is_some() {
scratch_props.clear();
scratch_props.extend(
node_data
.properties
.drain_to_interned_pairs(&graph.interner),
);
node_data.properties = PropertyStorage::Map(HashMap::new());
node_data.id = Value::Null;
node_data.title = Value::Null;
}
let node_idx = GraphWrite::add_node(&mut graph.graph, node_data);
if let Some(ref mut log) = prop_log {
let node_type_key = graph.interner.get_or_intern(&node_type);
log.write_entity(
node_type_key,
node_idx,
&saved_id,
&saved_title,
scratch_props,
)
.map_err(|e| format!("Property log write failed: {e}"))?;
type_meta
.entry(node_type.clone())
.or_insert_with(TypeBuildMeta::new)
.record_entity(&saved_id, &saved_title, scratch_props);
}
graph
.type_indices
.entry_or_default(node_type.clone())
.push(node_idx);
if let Some(ref mut qt) = qnum_to_idx {
if let Some(qnum) = parse_qcode_number(&acc.id) {
if (qnum as usize) < qt.len() {
qt.set(qnum as usize, node_idx.index() as u32 + 1); }
}
} else {
graph
.id_indices
.entry_or_default(node_type)
.insert(id_value, node_idx);
}
stats.entities_created += 1;
match edge_buffer {
EdgeBuffer::Compact(buf) => {
if let Some(src_num) = parse_qcode_number(&acc.id) {
for (pred_label, target_qcode) in acc.outgoing_edges {
if let Some(tgt_num) = parse_qcode_number(&target_qcode) {
let pred_key = graph.interner.get_or_intern(&pred_label);
buf.try_push(CompactNTripleEdge {
source_qnum: src_num,
target_qnum: tgt_num,
predicate: pred_key,
})
.map_err(|error| format!("append compact N-Triples edge: {error}"))?;
}
}
}
}
EdgeBuffer::Strings(buf) => {
for (pred_label, target_qcode) in acc.outgoing_edges {
buf.push((acc.id.clone(), target_qcode, pred_label));
}
}
}
Ok(())
}
pub(super) fn format_count(n: u64) -> String {
if n >= 1_000_000 {
format!("{:.1}M", n as f64 / 1_000_000.0)
} else if n >= 1_000 {
format!("{:.1}K", n as f64 / 1_000.0)
} else {
format!("{}", n)
}
}
pub(super) fn fmt_dur(secs: f64) -> String {
if secs < 60.0 {
return format!("{:.1}s", secs);
}
let total = secs as u64;
let h = total / 3600;
let m = (total % 3600) / 60;
let s = total % 60;
if h > 0 {
format!("{}h{:02}m{:02}s", h, m, s)
} else {
format!("{}m{:02}s", m, s)
}
}
fn build_debug() -> bool {
std::env::var("KGLITE_BUILD_DEBUG").is_ok()
}
#[cfg(test)]
#[path = "loader_tests.rs"]
mod tests;