use crate::datatypes::values::Value;
use crate::graph::constraints::{NamedConstraint, UniqueConstraintKey};
use crate::graph::features::timeseries::{NodeTimeseries, TimeseriesConfig};
use crate::graph::property_types::DeclaredType;
use crate::graph::schema::{
CompositeIndexKey, ConnectionTypeInfo, ConnectivityTriple, DirGraph, EmbeddingStore, IndexKey,
PropertyStorage, SaveMetadata, SchemaDefinition, SerdeDeserializeGuard, SerdeSerializeGuard,
SpatialConfig, StringInterner, StripPropertiesGuard, TemporalConfig,
};
use crate::graph::storage::column_store::ColumnStore;
use crate::graph::storage::property_storage::ColumnarRow;
use crate::graph::storage::{GraphRead, GraphWrite};
use memmap2::Mmap;
use rustc_hash::FxHashMap;
use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::fs::File;
use std::io::{self, BufWriter, Read, Write};
use std::path::Path;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::{Duration, SystemTime};
use crate::graph::io::magic::{
newer_portable_format_error, unrecognized_magic_error, V3_HARD_BREAK_MSG, V3_MAGIC, V4_MAGIC,
V5_MAGIC, V6_MAGIC,
};
use crate::serde_codec;
const MAX_CODEC_BYTES: u64 = 2 * 1024 * 1024 * 1024;
const DISK_SERDE_MAGIC: &[u8; 8] = b"KGLDSC1\0";
const CURRENT_CORE_DATA_VERSION: u32 = 3;
const EMBED_PROVENANCE_MIN_VERSION: u32 = 3;
const TOPOLOGY_SECTION: &str = "topology";
const EMBEDDINGS_SECTION: &str = "embeddings";
const TIMESERIES_SECTION: &str = "timeseries";
const SECONDARY_LABELS_SECTION: &str = "secondary_labels";
const VECTOR_INDEX_SECTION: &str = "vector_index";
fn column_section_key(type_name: &str) -> String {
format!("columns:{type_name}")
}
fn section_digest(compressed: &[u8]) -> u32 {
crate::graph::wal::crc32(compressed)
}
#[derive(Debug, Clone, Serialize, Deserialize)]
struct PortableColumnSection {
type_name: String,
compressed_size: u64,
row_count: u32,
columns: HashMap<String, String>, }
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub(crate) struct FileMetadata {
#[serde(default)]
core_data_version: u32,
#[serde(default)]
library_version: String,
#[serde(default)]
schema_definition: Option<SchemaDefinition>,
#[serde(default)]
property_index_keys: Vec<IndexKey>,
#[serde(default)]
composite_index_keys: Vec<CompositeIndexKey>,
#[serde(default)]
range_index_keys: Vec<IndexKey>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
unique_constraint_keys: Vec<UniqueConstraintKey>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
constraint_names: HashMap<String, NamedConstraint>,
#[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
ddl_not_null_constraints: BTreeSet<(String, String)>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
ddl_property_type_constraints: BTreeMap<String, BTreeMap<String, DeclaredType>>,
#[serde(default, skip_serializing_if = "BTreeSet::is_empty")]
rel_ddl_not_null_constraints: BTreeSet<(String, String)>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
rel_ddl_property_type_constraints: BTreeMap<String, BTreeMap<String, DeclaredType>>,
#[serde(default)]
node_type_metadata: HashMap<String, HashMap<String, String>>,
#[serde(default)]
connection_type_metadata: HashMap<String, ConnectionTypeInfo>,
#[serde(default)]
id_field_aliases: FxHashMap<String, String>,
#[serde(default)]
title_field_aliases: FxHashMap<String, String>,
#[serde(default = "crate::graph::dir_graph::default_auto_vacuum_threshold")]
auto_vacuum_threshold: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
storage_mode: Option<String>,
#[serde(default)]
parent_types: HashMap<String, String>,
#[serde(default, skip_serializing_if = "HashMap::is_empty")]
graph_instructions: HashMap<String, String>,
#[serde(default, skip_serializing_if = "is_zero")]
user_schema_version: u32,
#[serde(default, skip_serializing_if = "is_zero")]
checkpoint_lsn: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
cdc_handoff: Option<crate::graph::cdc::CdcHandoff>,
#[serde(default)]
spatial_configs: HashMap<String, SpatialConfig>,
#[serde(default)]
timeseries_configs: HashMap<String, TimeseriesConfig>,
#[serde(default)]
temporal_node_configs: HashMap<String, TemporalConfig>,
#[serde(default)]
temporal_edge_configs: HashMap<String, Vec<TemporalConfig>>,
#[serde(default = "default_ts_data_version")]
timeseries_data_version: u32,
#[serde(default)]
topology_compressed_size: u64,
#[serde(default)]
column_sections: Vec<PortableColumnSection>,
#[serde(default)]
embeddings_compressed_size: u64,
#[serde(default)]
timeseries_compressed_size: u64,
#[serde(default)]
secondary_labels_compressed_size: u64,
#[serde(default)]
vector_index_compressed_size: u64,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
section_digests: BTreeMap<String, u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
edge_type_counts: Option<HashMap<String, usize>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
type_connectivity: Option<Vec<ConnectivityTriple>>,
}
fn default_ts_data_version() -> u32 {
2
}
fn is_zero<T: Default + PartialEq>(value: &T) -> bool {
*value == T::default()
}
impl FileMetadata {
pub(crate) fn from_graph(graph: &DirGraph) -> Self {
FileMetadata {
core_data_version: CURRENT_CORE_DATA_VERSION,
library_version: env!("CARGO_PKG_VERSION").to_string(),
schema_definition: graph.schema_definition.clone(),
property_index_keys: graph.property_index_keys.clone(),
composite_index_keys: graph.composite_index_keys.clone(),
range_index_keys: graph.range_index_keys.clone(),
unique_constraint_keys: graph.unique_constraint_keys.clone(),
constraint_names: graph.constraint_names.clone(),
ddl_not_null_constraints: graph.ddl_not_null_constraints.clone(),
ddl_property_type_constraints: graph.ddl_property_type_constraints.clone(),
rel_ddl_not_null_constraints: graph.rel_ddl_not_null_constraints.clone(),
rel_ddl_property_type_constraints: graph.rel_ddl_property_type_constraints.clone(),
node_type_metadata: (*graph.node_type_metadata).clone(),
connection_type_metadata: (*graph.connection_type_metadata).clone(),
id_field_aliases: (*graph.id_field_aliases).clone(),
title_field_aliases: (*graph.title_field_aliases).clone(),
auto_vacuum_threshold: graph.auto_vacuum_threshold,
storage_mode: recorded_storage_mode_tag(graph),
parent_types: (*graph.parent_types).clone(),
graph_instructions: graph.graph_instructions.clone(),
user_schema_version: graph.user_schema_version,
checkpoint_lsn: graph.checkpoint_lsn,
cdc_handoff: crate::graph::cdc::status(graph)
.map(|status| crate::graph::cdc::CdcHandoff {
epoch: status.epoch,
last_seq: status.current,
})
.or(graph.cdc_handoff),
spatial_configs: graph.spatial_configs.clone(),
timeseries_configs: graph.timeseries_configs.clone(),
temporal_node_configs: graph.temporal_node_configs.clone(),
temporal_edge_configs: graph.temporal_edge_configs.clone(),
timeseries_data_version: 2,
topology_compressed_size: 0,
column_sections: Vec::new(),
embeddings_compressed_size: 0,
timeseries_compressed_size: 0,
secondary_labels_compressed_size: 0,
vector_index_compressed_size: 0,
section_digests: BTreeMap::new(),
edge_type_counts: if graph.has_edge_type_counts_cache() {
Some((*graph.get_edge_type_counts()).clone())
} else {
None
},
type_connectivity: graph.get_type_connectivity(),
}
}
pub(crate) fn apply_to(self, graph: &mut DirGraph) {
self.apply_to_with(graph, true)
}
pub(crate) fn apply_to_with(self, graph: &mut DirGraph, derive_type_connectivity: bool) {
graph.schema_definition = self.schema_definition;
graph.property_index_keys = self.property_index_keys;
graph.composite_index_keys = self.composite_index_keys;
graph.range_index_keys = self.range_index_keys;
graph.unique_constraint_keys = self.unique_constraint_keys;
graph.constraint_names = self.constraint_names;
graph.ddl_not_null_constraints = self.ddl_not_null_constraints;
graph.ddl_property_type_constraints = self.ddl_property_type_constraints;
graph.rel_ddl_not_null_constraints = self.rel_ddl_not_null_constraints;
graph.rel_ddl_property_type_constraints = self.rel_ddl_property_type_constraints;
graph.node_type_metadata = Arc::new(self.node_type_metadata);
graph.connection_type_metadata = Arc::new(self.connection_type_metadata);
graph.id_field_aliases = Arc::new(self.id_field_aliases);
graph.title_field_aliases = Arc::new(self.title_field_aliases);
graph.auto_vacuum_threshold = self.auto_vacuum_threshold;
graph.parent_types = Arc::new(self.parent_types);
graph.graph_instructions = self.graph_instructions;
graph.user_schema_version = self.user_schema_version;
graph.checkpoint_lsn = self.checkpoint_lsn;
graph.cdc_handoff = self.cdc_handoff;
graph.spatial_configs = self.spatial_configs;
graph.timeseries_configs = self.timeseries_configs;
graph.temporal_node_configs = self.temporal_node_configs;
graph.temporal_edge_configs = self.temporal_edge_configs;
graph.save_metadata = SaveMetadata {
format_version: crate::graph::schema::KGL_FORMAT_VERSION,
library_version: self.library_version,
};
if let Some(counts) = self.edge_type_counts {
*graph.edge_type_counts_cache.write().unwrap() = Some(std::sync::Arc::new(counts));
}
if let Some(triples) = self.type_connectivity {
*graph.type_connectivity_cache.write().unwrap() = Some(triples);
} else if derive_type_connectivity && !graph.connection_type_metadata.is_empty() {
let edge_counts = graph.edge_type_counts_cache.read().unwrap();
let mut triples = Vec::new();
for (conn_type, info) in graph.connection_type_metadata.iter() {
let count = edge_counts
.as_ref()
.and_then(|c| c.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,
});
}
}
}
if !triples.is_empty() {
*graph.type_connectivity_cache.write().unwrap() = Some(triples);
}
}
}
}
pub(crate) fn build_disk_metadata(graph: &DirGraph) -> FileMetadata {
FileMetadata::from_graph(graph)
}
pub(crate) fn strip_type_connectivity(meta: &mut FileMetadata) {
meta.type_connectivity = None;
}
pub(crate) fn strip_heavy_metadata(meta: &mut FileMetadata) {
meta.node_type_metadata.clear();
meta.connection_type_metadata.clear();
}
mod metadata_sidecars;
mod storage_mode;
pub(crate) use metadata_sidecars::{
read_connection_type_metadata_bin, read_node_type_metadata_bin,
write_connection_type_metadata_bin, write_node_type_metadata_bin,
};
use storage_mode::recorded_storage_mode_tag;
mod fast_load_sidecars;
use fast_load_sidecars::{decode_secondary_label_index, encode_secondary_label_index};
pub(crate) use fast_load_sidecars::{
read_id_indices_bin, read_interner_bin, read_secondary_labels_bin, read_type_connectivity_bin,
read_type_indices_bin, write_interner_bin, write_secondary_labels_bin,
write_type_connectivity_bin,
};
#[cfg(test)]
pub(crate) use fast_load_sidecars::{
ID_INDICES_MAGIC, ID_INDICES_VERSION, TYPE_INDICES_MAGIC, TYPE_INDICES_VERSION,
};
pub fn prepare_save(graph: &mut Arc<DirGraph>) {
let g = crate::graph::handle::make_dir_graph_mut_preserving_lineage(graph);
g.save_metadata = SaveMetadata::current();
g.populate_index_keys();
}
fn zstd_compress(data: &[u8]) -> io::Result<Vec<u8>> {
let mut encoder = zstd::Encoder::new(Vec::new(), 1)?;
encoder.include_checksum(true)?;
encoder.write_all(data)?;
encoder.finish()
}
fn zstd_decompress(data: &[u8]) -> io::Result<Vec<u8>> {
zstd_decompress_limited(data, MAX_DECOMPRESSED_SECTION_BYTES)
}
pub(crate) fn encode_disk_serde<T: Serialize + ?Sized>(value: &T) -> io::Result<Vec<u8>> {
let payload = serde_codec::encode_versioned(serde_codec::CURRENT_CODEC, value, MAX_CODEC_BYTES)
.map_err(io::Error::other)?;
let mut framed = Vec::with_capacity(DISK_SERDE_MAGIC.len() + 1 + payload.len());
framed.extend_from_slice(DISK_SERDE_MAGIC);
framed.push(serde_codec::CURRENT_CODEC.tag());
framed.extend_from_slice(&payload);
Ok(framed)
}
pub(crate) fn decode_disk_serde<'de, T: Deserialize<'de>>(
bytes: &'de [u8],
allocated_bytes: u64,
) -> io::Result<T> {
if bytes.starts_with(DISK_SERDE_MAGIC) {
let codec_tag = *bytes
.get(DISK_SERDE_MAGIC.len())
.ok_or_else(|| invalid_data("disk codec frame is truncated"))?;
let payload = &bytes[DISK_SERDE_MAGIC.len() + 1..];
return serde_codec::decode_exact_with(
serde_codec::CodecVersion::from_tag(codec_tag).map_err(io::Error::other)?,
payload,
allocated_bytes,
serde_codec::DecodeLimits::new(MAX_CODEC_BYTES, MAX_CODEC_BYTES),
)
.map_err(io::Error::other);
}
Err(pre_014_bincode_error("unframed disk sidecar"))
}
fn corrupt_sidecar_error(file_name: &str, cause: &io::Error) -> io::Error {
io::Error::new(
io::ErrorKind::InvalidData,
format!(
"disk graph sidecar '{file_name}' exists but is corrupt ({cause}); refusing to \
load the graph with this data silently missing. Restore '{file_name}' from a \
backup, rebuild the graph, or delete the file to load without it."
),
)
}
fn zstd_decompress_limited(data: &[u8], limit: u64) -> io::Result<Vec<u8>> {
let decoder = zstd::Decoder::new(std::io::Cursor::new(data))
.map_err(|e| invalid_data(format!("invalid zstd section: {e}")))?;
let mut bounded = decoder.take(limit.saturating_add(1));
let mut decoded = Vec::new();
bounded
.read_to_end(&mut decoded)
.map_err(|e| invalid_data(format!("invalid zstd section: {e}")))?;
if decoded.len() as u64 > limit {
return Err(invalid_data(format!(
"decompressed section exceeds the {} byte load limit",
limit
)));
}
Ok(decoded)
}
fn codec_ser<T: Serialize>(codec: serde_codec::CodecVersion, val: &T) -> io::Result<Vec<u8>> {
serde_codec::encode_versioned(codec, val, MAX_CODEC_BYTES).map_err(io::Error::other)
}
fn codec_deser<'a, T: Deserialize<'a>>(
codec: serde_codec::CodecVersion,
buf: &'a [u8],
allocated_bytes: u64,
) -> io::Result<T> {
let envelope = serde_codec::PayloadEnvelope::from_tag(
codec.tag(),
buf,
allocated_bytes,
serde_codec::DecodeLimits::new(MAX_CODEC_BYTES, MAX_CODEC_BYTES),
)
.map_err(|e| invalid_data(format!("binary payload envelope is invalid: {e}")))?;
let decoded = serde_codec::decode_versioned_exact(envelope);
decoded.map_err(|e| invalid_data(format!("binary deserialization failed: {e}")))
}
fn validate_column_keys_registered(graph: &DirGraph) -> io::Result<()> {
for (type_name, store) in graph.column_stores_by_name() {
let schema = store.schema();
for (_slot, key) in schema.iter() {
if graph.interner.try_resolve(key).is_none() {
return Err(invalid_data(format!(
"ColumnStore for type '{type_name}' contains unregistered InternedKey {}; \
refusing to serialize an unknown property name",
key.as_u64()
)));
}
}
}
Ok(())
}
fn save_temp_prefix(dest: &Path) -> String {
format!(
"{}.tmp.",
dest.file_name()
.map(|n| n.to_string_lossy().into_owned())
.unwrap_or_else(|| "graph.kgl".to_string()),
)
}
const UNIDENTIFIED_TEMP_MAX_AGE: Duration = Duration::from_secs(24 * 60 * 60);
pub fn reap_stale_save_temps(path: &Path) -> usize {
let prefix = save_temp_prefix(path);
let dir = match path.parent().filter(|p| !p.as_os_str().is_empty()) {
Some(d) => d.to_path_buf(),
None => Path::new(".").to_path_buf(),
};
let entries = match std::fs::read_dir(&dir) {
Ok(entries) => entries,
Err(_) => return 0,
};
let mut reaped = 0;
for entry in entries.flatten() {
let name = entry.file_name();
let Some(name) = name.to_str() else { continue };
let Some(pid) = temp_owner_pid(name, &prefix) else {
continue;
};
if pid == std::process::id() {
continue;
}
let abandoned = match process_is_alive(pid) {
Some(alive) => !alive,
None => temp_is_older_than(&entry, UNIDENTIFIED_TEMP_MAX_AGE),
};
if abandoned && std::fs::remove_file(entry.path()).is_ok() {
reaped += 1;
}
}
reaped
}
fn temp_owner_pid(name: &str, prefix: &str) -> Option<u32> {
let rest = name.strip_prefix(prefix)?;
let (pid, nonce) = rest.split_once('.')?;
nonce.parse::<u64>().ok()?;
let pid: u32 = pid.parse().ok()?;
(pid > 0).then_some(pid)
}
fn temp_is_older_than(entry: &std::fs::DirEntry, max_age: Duration) -> bool {
entry
.metadata()
.and_then(|m| m.modified())
.ok()
.and_then(|modified| SystemTime::now().duration_since(modified).ok())
.is_some_and(|age| age > max_age)
}
#[cfg(unix)]
fn process_is_alive(pid: u32) -> Option<bool> {
if unsafe { libc::kill(pid as libc::pid_t, 0) } == 0 {
return Some(true);
}
match std::io::Error::last_os_error().raw_os_error() {
Some(libc::ESRCH) => Some(false),
Some(libc::EPERM) => Some(true),
_ => None,
}
}
#[cfg(not(unix))]
fn process_is_alive(_pid: u32) -> Option<bool> {
None
}
pub fn write_kgl_with(graph: &DirGraph, path: &str, fsync: bool) -> io::Result<()> {
let dest = Path::new(path);
let dir = dest.parent().filter(|p| !p.as_os_str().is_empty());
static SAVE_COUNTER: AtomicU64 = AtomicU64::new(0);
let nonce = SAVE_COUNTER.fetch_add(1, Ordering::Relaxed);
let tmp_name = format!("{}{}.{}", save_temp_prefix(dest), std::process::id(), nonce);
let tmp = match dir {
Some(d) => d.join(&tmp_name),
None => Path::new(&tmp_name).to_path_buf(),
};
let write_result = (|| -> io::Result<()> {
let file = File::create(&tmp)?;
let mut writer = BufWriter::new(file);
write_kgl_to(graph, &mut writer)?;
writer.flush()?;
let file = writer
.into_inner()
.map_err(|e| io::Error::other(e.to_string()))?;
if fsync {
file.sync_all()?;
}
Ok(())
})();
if let Err(e) = write_result {
let _ = std::fs::remove_file(&tmp);
return Err(e);
}
if let Err(e) = std::fs::rename(&tmp, dest) {
let _ = std::fs::remove_file(&tmp);
return Err(e);
}
if fsync {
if let Some(d) = dir {
if let Ok(dirfile) = File::open(d) {
let _ = dirfile.sync_all();
}
}
}
Ok(())
}
pub fn write_kgl(graph: &DirGraph, path: &str) -> io::Result<()> {
write_kgl_with(graph, path, true)
}
fn build_section_digests(
topology: &[u8],
column_meta: &[PortableColumnSection],
column_data: &[Vec<u8>],
optional: [(&str, Option<&[u8]>); 4],
) -> BTreeMap<String, u32> {
let mut digests = BTreeMap::new();
digests.insert(TOPOLOGY_SECTION.to_string(), section_digest(topology));
for (meta, data) in column_meta.iter().zip(column_data.iter()) {
digests.insert(column_section_key(&meta.type_name), section_digest(data));
}
for (key, data) in optional {
if let Some(bytes) = data {
digests.insert(key.to_string(), section_digest(bytes));
}
}
digests
}
pub fn write_kgl_to<W: Write>(graph: &DirGraph, writer: &mut W) -> io::Result<()> {
validate_column_keys_registered(graph)?;
let codec = serde_codec::CodecVersion::PostcardV1;
let topology_raw = {
let _strip = StripPropertiesGuard::new();
let _guard = SerdeSerializeGuard::new(&graph.interner);
codec_ser(codec, &graph.graph)?
};
let topology_compressed = zstd_compress(&topology_raw)?;
drop(topology_raw);
let mut column_sections_meta: Vec<PortableColumnSection> = Vec::new();
let mut column_sections_data: Vec<Vec<u8>> = Vec::new();
let mut column_stores_sorted: Vec<(&str, &Arc<ColumnStore>)> = graph.column_stores_by_name();
column_stores_sorted.sort_by(|a, b| a.0.cmp(b.0));
for (type_name, store) in column_stores_sorted {
let packed = store.write_packed_with_codec(
&graph.interner,
codec,
crate::graph::storage::packed_codec::IntColumnEncoding::Auto,
)?;
let compressed = zstd_compress(&packed)?;
drop(packed);
let mut cols = HashMap::new();
for (slot, ik) in store.schema().iter() {
let prop_name = graph.interner.resolve(ik);
if let Some(col) = store.column(slot as usize) {
cols.insert(prop_name.to_string(), col.type_tag().to_string());
}
}
column_sections_meta.push(PortableColumnSection {
type_name: type_name.to_string(),
compressed_size: compressed.len() as u64,
row_count: store.row_count(),
columns: cols,
});
column_sections_data.push(compressed);
}
let embedding_compressed = if !graph.embeddings.is_empty() {
let ordered: std::collections::BTreeMap<_, _> = graph.embeddings.iter().collect();
let raw = codec_ser(codec, &ordered)?;
Some(zstd_compress(&raw)?)
} else {
None
};
let timeseries_compressed = if !graph.timeseries_store.is_empty() {
let ordered: std::collections::BTreeMap<_, _> = graph.timeseries_store.iter().collect();
let raw = codec_ser(codec, &ordered)?;
Some(zstd_compress(&raw)?)
} else {
None
};
let secondary_labels_compressed = match encode_secondary_label_index(graph) {
Some(payload) => Some(zstd_compress(&payload)?),
None => None,
};
let vector_index_compressed = match encode_vector_indexes(graph)? {
Some(payload) => Some(zstd_compress(&payload)?),
None => None,
};
let section_digests = build_section_digests(
&topology_compressed,
&column_sections_meta,
&column_sections_data,
[
(EMBEDDINGS_SECTION, embedding_compressed.as_deref()),
(TIMESERIES_SECTION, timeseries_compressed.as_deref()),
(
SECONDARY_LABELS_SECTION,
secondary_labels_compressed.as_deref(),
),
(VECTOR_INDEX_SECTION, vector_index_compressed.as_deref()),
],
);
let mut metadata = FileMetadata::from_graph(graph);
metadata.section_digests = section_digests;
metadata.topology_compressed_size = topology_compressed.len() as u64;
metadata.column_sections = column_sections_meta;
metadata.embeddings_compressed_size = embedding_compressed
.as_ref()
.map(|b| b.len() as u64)
.unwrap_or(0);
metadata.timeseries_compressed_size = timeseries_compressed
.as_ref()
.map(|b| b.len() as u64)
.unwrap_or(0);
metadata.secondary_labels_compressed_size = secondary_labels_compressed
.as_ref()
.map(|b| b.len() as u64)
.unwrap_or(0);
metadata.vector_index_compressed_size = vector_index_compressed
.as_ref()
.map(|b| b.len() as u64)
.unwrap_or(0);
let metadata_value = serde_json::to_value(&metadata).map_err(io::Error::other)?;
let metadata_json = serde_json::to_vec(&metadata_value).map_err(io::Error::other)?;
writer.write_all(&V6_MAGIC)?;
writer.write_all(&[codec.tag()])?;
writer.write_all(&CURRENT_CORE_DATA_VERSION.to_le_bytes())?;
writer.write_all(&(metadata_json.len() as u32).to_le_bytes())?;
writer.write_all(&metadata_json)?;
writer.write_all(&topology_compressed)?;
for section_data in &column_sections_data {
writer.write_all(section_data)?;
}
if let Some(emb_data) = &embedding_compressed {
writer.write_all(emb_data)?;
}
if let Some(ts_data) = ×eries_compressed {
writer.write_all(ts_data)?;
}
if let Some(sl_data) = &secondary_labels_compressed {
writer.write_all(sl_data)?;
}
if let Some(vi_data) = &vector_index_compressed {
writer.write_all(vi_data)?;
}
writer.flush()?;
Ok(())
}
pub fn prepare_kgl_write(graph: &mut Arc<DirGraph>) {
prepare_save(graph);
let dir = crate::graph::handle::make_dir_graph_mut_preserving_lineage(graph);
dir.enable_columnar();
}
pub fn save_inmemory_with(graph: &mut Arc<DirGraph>, path: &str, fsync: bool) -> io::Result<()> {
prepare_kgl_write(graph);
write_kgl_with(graph, path, fsync)
}
pub fn save_graph(graph: &mut Arc<DirGraph>, path: &str) -> Result<(), SaveError> {
save_graph_with(graph, path, true)
}
pub fn save_graph_with(
graph: &mut Arc<DirGraph>,
path: &str,
fsync: bool,
) -> Result<(), SaveError> {
save_guard::ensure_target_recovered(graph, path)?;
if graph.graph.is_disk() {
let dir = crate::graph::handle::make_dir_graph_mut_preserving_lineage(graph);
return dir.save_disk(path).map_err(SaveError::Io);
}
save_inmemory_with(graph, path, fsync).map_err(|e| SaveError::Io(e.to_string()))
}
pub fn materialize_disk_graph(graph: &mut Arc<DirGraph>, path: &str) -> Result<(), SaveError> {
save_guard::ensure_target_recovered(graph, path)?;
let dir = crate::graph::handle::make_dir_graph_mut(graph);
dir.enable_disk_mode_at(path).map_err(SaveError::Io)
}
const FILE_MMAP_THRESHOLD: u64 = 65_536;
const MAX_METADATA_BYTES: usize = 64 * 1024 * 1024;
const MAX_DECOMPRESSED_SECTION_BYTES: u64 = 16 * 1024 * 1024 * 1024;
fn invalid_data(message: impl Into<String>) -> io::Error {
io::Error::new(io::ErrorKind::InvalidData, message.into())
}
fn validate_and_rebuild_embedding_norms(
embeddings: &mut HashMap<(String, String), EmbeddingStore>,
) -> io::Result<()> {
for store in embeddings.values_mut() {
store.validate_shape().map_err(invalid_data)?;
store.rebuild_norms();
}
Ok(())
}
pub(crate) fn pre_014_bincode_error(artifact: &str) -> io::Error {
invalid_data(format!(
"Unsupported pre-0.14 bincode persistence: {artifact}. This build reads Postcard \
persistence only. Open the artifact with kglite 0.13.4 and re-save or re-export it, \
then retry; alternatively rebuild it from the original source."
))
}
struct SectionCursor<'a> {
bytes: &'a [u8],
offset: usize,
digests: BTreeMap<String, u32>,
}
impl<'a> SectionCursor<'a> {
fn new(bytes: &'a [u8], offset: usize, digests: BTreeMap<String, u32>) -> io::Result<Self> {
if offset > bytes.len() {
return Err(invalid_data("section cursor starts past end of file"));
}
Ok(Self {
bytes,
offset,
digests,
})
}
fn take(&mut self, encoded_len: u64, section: &str) -> io::Result<&'a [u8]> {
let len = usize::try_from(encoded_len)
.map_err(|_| invalid_data(format!("{section} section size does not fit usize")))?;
let end = self
.offset
.checked_add(len)
.ok_or_else(|| invalid_data(format!("{section} section offset overflow")))?;
let bytes = self.bytes.get(self.offset..end).ok_or_else(|| {
invalid_data(format!(
"file is truncated — {section} section needs {len} bytes at offset {}",
self.offset
))
})?;
self.verify(section, bytes)?;
self.offset = end;
Ok(bytes)
}
fn verify(&self, section: &str, bytes: &[u8]) -> io::Result<()> {
let Some(&expected) = self.digests.get(section) else {
return Ok(());
};
let actual = section_digest(bytes);
if actual != expected {
return Err(invalid_data(format!(
"the '{section}' section of this .kgl file is corrupt — it does not match the \
CRC32 digest recorded when the file was written (recorded {expected:#010x}, \
computed {actual:#010x}). Restore the file from a backup or rebuild the graph \
from its source."
)));
}
Ok(())
}
}
pub fn load_file(path: &str) -> io::Result<Arc<DirGraph>> {
let p = std::path::Path::new(path);
if p.is_dir() {
return load_disk_dir(p);
}
let file = File::open(path)?;
let file_len = file.metadata()?.len();
if file_len >= FILE_MMAP_THRESHOLD {
let mmap = unsafe { Mmap::map(&file)? };
if mmap.len() < 4 {
return Err(io::Error::other(
"File is too small to be a valid kglite file.",
));
}
if mmap[..4] == V6_MAGIC {
return load_portable_container(&mmap, "v6");
}
if mmap[..4] == V5_MAGIC {
return load_portable_container(&mmap, "v5");
}
if mmap[..4] == V4_MAGIC {
return Err(pre_014_bincode_error(".kgl container v4"));
}
if mmap[..4] == V3_MAGIC {
return Err(io::Error::other(V3_HARD_BREAK_MSG));
}
if mmap[..3] == V6_MAGIC[..3] && mmap[3] > V6_MAGIC[3] {
return Err(newer_portable_format_error(mmap[3]));
}
return Err(unrecognized_magic_error(&mmap[..4], &format!("'{path}'")));
}
let buf = std::fs::read(path)?;
if buf.len() < 4 {
return Err(io::Error::other(
"File is too small to be a valid kglite file.",
));
}
if buf[..4] == V6_MAGIC {
load_portable_container(&buf, "v6")
} else if buf[..4] == V5_MAGIC {
load_portable_container(&buf, "v5")
} else if buf[..4] == V4_MAGIC {
Err(pre_014_bincode_error(".kgl container v4"))
} else if buf[..4] == V3_MAGIC {
Err(io::Error::other(V3_HARD_BREAK_MSG))
} else if buf[..3] == V6_MAGIC[..3] && buf[3] > V6_MAGIC[3] {
Err(newer_portable_format_error(buf[3]))
} else {
Err(unrecognized_magic_error(&buf[..4], &format!("'{path}'")))
}
}
pub fn load_kgl_bytes(data: &[u8]) -> io::Result<Arc<DirGraph>> {
if data.len() < 4 {
return Err(io::Error::other(
"Byte buffer is too small to be a valid kglite graph.",
));
}
if data[..4] == V6_MAGIC {
load_portable_container(data, "v6")
} else if data[..4] == V5_MAGIC {
load_portable_container(data, "v5")
} else if data[..4] == V4_MAGIC {
Err(pre_014_bincode_error(".kgl container v4"))
} else if data[..4] == V3_MAGIC {
Err(io::Error::other(V3_HARD_BREAK_MSG))
} else if data[..3] == V6_MAGIC[..3] && data[3] > V6_MAGIC[3] {
Err(newer_portable_format_error(data[3]))
} else {
Err(unrecognized_magic_error(&data[..4], "the byte buffer"))
}
}
const EMBED_FORMAT_BREAK_MSG: &str =
"This .kgl was saved with an older embedding format (before per-vector model \
id + text-hash provenance, kglite 0.10.29). Its embeddings can't be loaded by \
this binary. The graph's nodes/edges are fine — reload, re-run \
embed_texts()/add_embeddings() to rebuild the vectors, and save again. \
(Embeddings are a rebuildable cache; only the vector section broke.)";
fn rebuild_disk_type_schemas(graph: &mut DirGraph) -> io::Result<()> {
let metadata = std::sync::Arc::clone(&graph.node_type_metadata);
for (node_type, props) in metadata.iter() {
let mut schema = crate::graph::schema::TypeSchema::new();
let mut prop_names: Vec<&String> = props.keys().collect();
prop_names.sort();
for prop_name in prop_names {
let key = graph
.interner
.try_get_or_intern(prop_name)
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
schema.add_key(key);
}
graph
.type_schemas_mut()
.insert(node_type.clone(), std::sync::Arc::new(schema));
}
Ok(())
}
fn load_disk_dir(dir: &std::path::Path) -> io::Result<Arc<DirGraph>> {
use crate::graph::io::load_timing::{log_stage, stage_timer};
use crate::graph::schema::GraphBackend;
let _load_t = stage_timer();
let resolved = crate::graph::storage::disk::generation::resolve_snapshot(dir)?;
let logical_root = resolved.logical_root;
let snapshot_dir = resolved.snapshot_dir;
let dir = snapshot_dir.as_path();
if !dir.join("disk_graph_meta.json").exists() {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"Directory does not contain a valid disk graph (missing disk_graph_meta.json)",
));
}
let mut graph = DirGraph::new();
let t = stage_timer();
apply_disk_metadata(dir, &mut graph)?;
log_stage("metadata_json", t);
let t = stage_timer();
let loaded_from_bin = read_interner_bin(dir, &mut graph)?;
if !loaded_from_bin && dir.join("interner.json").exists() {
let interner_str = std::fs::read_to_string(dir.join("interner.json"))?;
let interner_map: std::collections::HashMap<String, String> =
serde_json::from_str(&interner_str)
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
for original in interner_map.values() {
graph
.interner
.try_get_or_intern(original)
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
}
}
log_stage("interner_load", t);
let t = stage_timer();
let (mut disk_graph, temp_dir) =
crate::graph::storage::disk::graph::DiskGraph::load_from_dir(dir, &mut graph.interner)?;
disk_graph.set_logical_root(logical_root);
log_stage("disk_graph_load", t);
if std::env::var_os("KGLITE_PREFETCH").is_some() {
let t = stage_timer();
disk_graph.prefetch_hot_regions();
log_stage("prefetch_hot_regions", t);
}
graph.graph = GraphBackend::Disk(Box::new(disk_graph));
if let Ok(mut dirs) = graph.temp_dirs.lock() {
dirs.push(temp_dir);
}
let t = stage_timer();
if let GraphBackend::Disk(ref dg) = graph.graph {
let mut loaded = false;
if let Some(base) =
crate::graph::storage::disk::type_index::TypeIndexBase::load_from(dir, &graph.interner)?
{
graph.type_indices =
crate::graph::storage::disk::type_index::TypeIndexStore::from_base(base);
loaded = true;
}
if !loaded {
let ti_path = dir.join("type_indices.bin.zst");
if ti_path.exists() {
if let Ok(compressed) = std::fs::read(&ti_path) {
if let Ok(bytes) = zstd_decompress(&compressed) {
if let Ok(Some(indices)) = read_type_indices_bin(&bytes, &graph.interner) {
graph.type_indices.replace_with(indices);
loaded = true;
}
}
}
}
}
if !loaded {
let mut new_type_indices: std::collections::HashMap<
String,
Vec<petgraph::graph::NodeIndex>,
> = std::collections::HashMap::new();
for i in 0..dg.node_slot_len() {
let slot = dg.node_slot(i);
if slot.is_alive() {
let key = crate::graph::schema::InternedKey::from_u64(slot.node_type);
if let Some(type_name) = graph.interner.try_resolve(key) {
new_type_indices
.entry(type_name.to_string())
.or_default()
.push(petgraph::graph::NodeIndex::new(i));
}
}
}
graph.type_indices.replace_with(new_type_indices);
}
}
log_stage("type_indices_load", t);
rebuild_disk_type_schemas(&mut graph)?;
load_disk_column_stores(dir, &mut graph)?;
let t = stage_timer();
if crate::graph::storage::GraphRead::is_disk(&graph.graph) {
if let Some(base) =
crate::graph::storage::disk::id_index::IdIndexBase::load_from(dir, &graph.interner)?
{
graph.id_indices = crate::graph::storage::disk::id_index::IdIndexStore::from_base(base);
} else {
let id_indices_path = dir.join("id_indices.bin.zst");
if id_indices_path.exists() {
if let Ok(compressed) = std::fs::read(&id_indices_path) {
if let Ok(bytes) = zstd_decompress(&compressed) {
if let Ok(Some(indices)) = read_id_indices_bin(&bytes, &graph.interner) {
graph.id_indices.replace_with(indices);
}
}
}
}
}
}
log_stage("id_indices_load", t);
let t = stage_timer();
if std::env::var_os("KGLITE_EAGER_TYPE_CONNECTIVITY").is_some()
&& !graph.has_type_connectivity_cache()
{
if let Ok(Some(triples)) = read_type_connectivity_bin(dir, &graph) {
if !triples.is_empty() {
*graph.type_connectivity_cache.write().unwrap() = Some(triples);
}
}
}
log_stage("type_connectivity_load", t);
load_disk_sidecars(dir, &mut graph)?;
graph.build_connection_types_cache();
log_stage("load_disk_dir_total", _load_t);
Ok(Arc::new(graph))
}
fn load_disk_column_stores(dir: &std::path::Path, graph: &mut DirGraph) -> io::Result<()> {
use crate::graph::io::load_timing::{log_stage, stage_timer};
let mmap_path = {
let seg0 = dir.join("seg_000/columns.bin");
if seg0.exists() {
seg0
} else {
dir.join("columns.bin")
}
};
let meta_bin_path = {
let seg0 = dir.join("seg_000/columns_meta.bin.zst");
if seg0.exists() {
seg0
} else {
dir.join("columns_meta.bin.zst")
}
};
let meta_json_path = {
let seg0 = dir.join("seg_000/columns_meta.json");
if seg0.exists() {
seg0
} else {
dir.join("columns_meta.json")
}
};
let has_mmap = mmap_path.exists() && (meta_bin_path.exists() || meta_json_path.exists());
let t = stage_timer();
if has_mmap {
use crate::graph::io::ntriples::ColumnTypeMeta;
use memmap2::MmapMut;
let file = std::fs::OpenOptions::new()
.read(true)
.write(true)
.open(&mmap_path)?;
let mmap = unsafe { MmapMut::map_mut(&file)? };
let mmap_arc = std::sync::Arc::new(mmap);
let type_metas: Vec<ColumnTypeMeta> = if meta_bin_path.exists() {
let compressed = std::fs::read(&meta_bin_path)?;
let bytes = zstd_decompress(&compressed)?;
decode_disk_serde(&bytes, bytes.capacity() as u64)?
} else {
let meta_json = std::fs::read_to_string(&meta_json_path)?;
serde_json::from_str(&meta_json).map_err(io::Error::other)?
};
let skip_utf8 = std::env::var_os("KGLITE_SKIP_UTF8_VALIDATION").is_some();
for tm in type_metas {
let store = tm.to_mmap_store(std::sync::Arc::clone(&mmap_arc));
if !skip_utf8 {
store.validate_utf8(&tm.type_name)?;
}
let cs = crate::graph::storage::column_store::ColumnStore::from_mmap_store(
std::sync::Arc::new(store),
);
graph.install_column_store(&tm.type_name, Arc::new(cs));
}
load_column_sidecars(dir, graph)?;
} else {
load_column_sidecars(dir, graph)?;
}
log_stage("column_stores_load", t);
Ok(())
}
fn load_disk_sidecars(dir: &std::path::Path, graph: &mut DirGraph) -> io::Result<()> {
let emb_path = dir.join("embeddings.bin.zst");
if emb_path.exists() {
let mut embeddings = (|| -> io::Result<HashMap<(String, String), EmbeddingStore>> {
let compressed = std::fs::read(&emb_path)?;
let bytes = zstd_decompress(&compressed)?;
decode_disk_serde(&bytes, bytes.capacity() as u64)
.map_err(|e| invalid_data(e.to_string()))
})()
.map_err(|e| corrupt_sidecar_error("embeddings.bin.zst", &e))?;
validate_and_rebuild_embedding_norms(&mut embeddings)
.map_err(|e| corrupt_sidecar_error("embeddings.bin.zst", &e))?;
graph.embeddings = embeddings;
}
let ts_path = dir.join("timeseries.bin.zst");
if ts_path.exists() {
graph.timeseries_store = (|| -> io::Result<HashMap<usize, NodeTimeseries>> {
let compressed = std::fs::read(&ts_path)?;
let bytes = zstd_decompress(&compressed)?;
decode_disk_serde(&bytes, bytes.capacity() as u64)
.map_err(|e| invalid_data(e.to_string()))
})()
.map_err(|e| corrupt_sidecar_error("timeseries.bin.zst", &e))?;
}
read_secondary_labels_bin(dir, graph)
.map_err(|e| corrupt_sidecar_error("secondary_labels.bin.zst", &e))?;
Ok(())
}
fn apply_disk_metadata(dir: &std::path::Path, graph: &mut DirGraph) -> io::Result<()> {
if !dir.join("metadata.json").exists() {
return Ok(());
}
let meta_bytes = std::fs::read(dir.join("metadata.json"))?;
let mut meta: FileMetadata = serde_json::from_slice(&meta_bytes)
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
meta.validate_disk_storage_mode()?;
if let Some(ntm) = read_node_type_metadata_bin(dir)? {
meta.node_type_metadata = ntm;
}
if let Some(ctm) = read_connection_type_metadata_bin(dir)? {
meta.connection_type_metadata = ctm;
}
meta.apply_to_with(graph, false);
Ok(())
}
fn load_portable_container(buf: &[u8], format_name: &str) -> io::Result<Arc<DirGraph>> {
if buf.len() < 13 {
return Err(invalid_data(format!(
"{format_name} file is truncated — header incomplete"
)));
}
let codec = serde_codec::CodecVersion::from_tag(buf[4]).map_err(|e| {
invalid_data(format!(
"{format_name} header has an invalid codec tag: {e}"
))
})?;
if codec != serde_codec::CodecVersion::PostcardV1 {
return Err(invalid_data(format!(
"{format_name} header selects codec {}, but {format_name} requires Postcard codec {}",
codec.tag(),
serde_codec::CodecVersion::PostcardV1.tag()
)));
}
let core_version = u32::from_le_bytes([buf[5], buf[6], buf[7], buf[8]]);
let metadata_len = u32::from_le_bytes([buf[9], buf[10], buf[11], buf[12]]) as usize;
load_portable_columnar(buf, format_name, codec, core_version, metadata_len, 13)
}
struct PortableSectionPlan {
columns: Vec<PortableColumnSection>,
embeddings: u64,
timeseries: u64,
secondary_labels: u64,
vector_index: u64,
}
fn parse_portable_metadata<'a>(
buf: &'a [u8],
format_name: &str,
metadata_len: usize,
metadata_start: usize,
) -> io::Result<(FileMetadata, SectionCursor<'a>)> {
let metadata_end = metadata_start
.checked_add(metadata_len)
.ok_or_else(|| invalid_data(format!("{format_name} metadata offset overflow")))?;
let metadata_bytes = buf.get(metadata_start..metadata_end).ok_or_else(|| {
invalid_data(format!(
"{format_name} file is truncated — metadata incomplete"
))
})?;
let metadata: FileMetadata = serde_json::from_slice(metadata_bytes)
.map_err(|e| invalid_data(format!("failed to parse {format_name} metadata: {e}")))?;
if metadata.column_sections.len() > 1_000_000 {
return Err(invalid_data(format!(
"{format_name} metadata declares too many column sections"
)));
}
let digests = metadata.section_digests.clone();
Ok((metadata, SectionCursor::new(buf, metadata_end, digests)?))
}
fn decode_portable_topology(
codec: serde_codec::CodecVersion,
sections: &mut SectionCursor<'_>,
metadata: FileMetadata,
) -> io::Result<(DirGraph, PortableSectionPlan)> {
let topology_compressed = sections.take(metadata.topology_compressed_size, TOPOLOGY_SECTION)?;
let topology_raw = zstd_decompress(topology_compressed)?;
let mut interner = StringInterner::new();
let graph: crate::graph::schema::GraphBackend = {
let _guard = SerdeDeserializeGuard::new(&mut interner);
codec_deser(codec, &topology_raw, topology_raw.capacity() as u64)?
};
let plan = PortableSectionPlan {
columns: metadata.column_sections.clone(),
embeddings: metadata.embeddings_compressed_size,
timeseries: metadata.timeseries_compressed_size,
secondary_labels: metadata.secondary_labels_compressed_size,
vector_index: metadata.vector_index_compressed_size,
};
let mut dir_graph = DirGraph::from_graph(graph);
dir_graph.interner = interner;
metadata.apply_to(&mut dir_graph);
dir_graph.rebuild_type_indices_and_schemas();
dir_graph.build_connection_types_cache();
Ok((dir_graph, plan))
}
fn portable_temp_dir() -> std::path::PathBuf {
std::env::temp_dir().join(format!(
"kglite_portable_{}_{:x}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos()
))
}
fn load_portable_column_section(
codec: serde_codec::CodecVersion,
dir_graph: &mut DirGraph,
sections: &mut SectionCursor<'_>,
section_meta: &PortableColumnSection,
section_index: usize,
temp_dir: &Path,
) -> io::Result<()> {
let compressed = sections.take(
section_meta.compressed_size,
&column_section_key(§ion_meta.type_name),
)?;
let packed = zstd_decompress(compressed)?;
let expected_rows = dir_graph
.type_indices
.get(§ion_meta.type_name)
.map_or(0, |nodes| nodes.len());
if section_meta.row_count as usize != expected_rows {
return Err(invalid_data(format!(
"column section {section_index} for '{}' declares {} rows; topology has {expected_rows}",
section_meta.type_name, section_meta.row_count
)));
}
let mut ordered_names = ColumnStore::packed_column_names(&packed)?;
let in_payload: std::collections::HashSet<&str> =
ordered_names.iter().map(String::as_str).collect();
let mut orphans: Vec<&String> = section_meta
.columns
.keys()
.filter(|name| !in_payload.contains(name.as_str()))
.collect();
orphans.sort();
let orphans: Vec<String> = orphans.into_iter().cloned().collect();
ordered_names.extend(orphans);
let col_keys = ordered_names
.iter()
.map(|name| {
dir_graph
.interner
.try_get_or_intern(name)
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))
})
.collect::<io::Result<Vec<_>>>()?;
let column_schema = Arc::new(crate::graph::schema::TypeSchema::from_keys(col_keys));
let type_meta = dir_graph
.node_type_metadata
.get(§ion_meta.type_name)
.cloned()
.unwrap_or_default();
let type_temp_dir = temp_dir.join(format!("type_{section_index}"));
std::fs::create_dir_all(&type_temp_dir)?;
let store = ColumnStore::load_packed_with_codec(
column_schema,
&type_meta,
&dir_graph.interner,
&packed,
section_meta.row_count,
Some(&type_temp_dir),
codec,
)?;
dir_graph.install_column_store(§ion_meta.type_name, Arc::new(store));
Ok(())
}
fn load_portable_columns(
codec: serde_codec::CodecVersion,
dir_graph: &mut DirGraph,
sections: &mut SectionCursor<'_>,
columns: &[PortableColumnSection],
) -> io::Result<()> {
let temp_dir = portable_temp_dir();
if let Ok(mut dirs) = dir_graph.temp_dirs.lock() {
dirs.push(temp_dir.clone());
}
for (index, metadata) in columns.iter().enumerate() {
load_portable_column_section(codec, dir_graph, sections, metadata, index, &temp_dir)?;
}
attach_portable_column_stores(dir_graph);
Ok(())
}
fn load_portable_optional_sections(
codec: serde_codec::CodecVersion,
core_version: u32,
dir_graph: &mut DirGraph,
sections: &mut SectionCursor<'_>,
plan: &PortableSectionPlan,
) -> io::Result<()> {
if plan.embeddings > 0 {
if core_version < EMBED_PROVENANCE_MIN_VERSION {
return Err(io::Error::other(EMBED_FORMAT_BREAK_MSG));
}
let compressed = sections.take(plan.embeddings, EMBEDDINGS_SECTION)?;
let raw = zstd_decompress(compressed)?;
let mut embeddings: HashMap<(String, String), EmbeddingStore> =
codec_deser(codec, &raw, raw.capacity() as u64)?;
validate_and_rebuild_embedding_norms(&mut embeddings)?;
dir_graph.embeddings = embeddings;
}
if plan.timeseries > 0 {
let compressed = sections.take(plan.timeseries, TIMESERIES_SECTION)?;
let raw = zstd_decompress(compressed)?;
dir_graph.timeseries_store = codec_deser(codec, &raw, raw.capacity() as u64)?;
}
if plan.secondary_labels > 0 {
let compressed = sections.take(plan.secondary_labels, SECONDARY_LABELS_SECTION)?;
let raw = zstd_decompress(compressed)?;
decode_secondary_label_index(&raw, dir_graph)?;
}
if plan.vector_index > 0 {
let compressed = sections.take(plan.vector_index, VECTOR_INDEX_SECTION)?;
if let Ok(raw) = zstd_decompress(compressed) {
decode_vector_indexes(&raw, dir_graph);
}
}
Ok(())
}
fn load_portable_columnar(
buf: &[u8],
format_name: &str,
codec: serde_codec::CodecVersion,
core_version: u32,
metadata_len: usize,
metadata_start: usize,
) -> io::Result<Arc<DirGraph>> {
if metadata_len > MAX_METADATA_BYTES {
return Err(invalid_data(format!(
"{format_name} metadata is {metadata_len} bytes; limit is {MAX_METADATA_BYTES}"
)));
}
if core_version > CURRENT_CORE_DATA_VERSION {
return Err(io::Error::other(format!(
"File uses core data version {} but this library only supports up to version {}. \
Please upgrade kglite.",
core_version, CURRENT_CORE_DATA_VERSION,
)));
}
let (metadata, mut sections) =
parse_portable_metadata(buf, format_name, metadata_len, metadata_start)?;
let recorded_mode = metadata.portable_storage_mode()?;
let (mut dir_graph, plan) = decode_portable_topology(codec, &mut sections, metadata)?;
load_portable_columns(codec, &mut dir_graph, &mut sections, &plan.columns)?;
load_portable_optional_sections(codec, core_version, &mut dir_graph, &mut sections, &plan)?;
crate::graph::storage::mode::convert_dir_graph_to_mode(&mut dir_graph, recorded_mode)
.map_err(io::Error::other)?;
Ok(Arc::new(dir_graph))
}
mod columns;
use columns::{attach_portable_column_stores, load_column_sidecars};
mod save_guard;
pub use save_guard::SaveError;
mod vector_persistence;
#[allow(unused_imports)]
pub use vector_persistence::ExportStats;
use vector_persistence::{decode_vector_indexes, encode_vector_indexes};
pub use vector_persistence::{
export_embeddings_to_file, import_embeddings_from_file, EmbeddingExportFilter, ImportStats,
};
#[cfg(test)]
#[path = "file_tests.rs"]
mod file_tests;