use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use arrow::array::{Array as ArrowArray, ArrayRef, Float32Array, Float64Array, UInt32Array};
use arrow::datatypes::{DataType, Field, Float64Type, Int64Type, Schema, UInt32Type};
use arrow::record_batch::RecordBatch;
use log::{debug, info};
use smartcore::linalg::basic::arrays::Array;
use smartcore::linalg::basic::matrix::DenseMatrix;
use sprs::CsMat;
use crate::graph::{GraphEdge, GraphReadOptions, GraphWriteOptions, StoredGraph};
use crate::metadata::FileInfo;
use crate::metadata::GeneMetadata;
use crate::traits::backend::StorageBackend;
use crate::traits::lance::{
LanceStorage, graph_record_batch, stored_graph_from_batch_with_options,
validate_vector_space_schema, vector_space_collection_batch,
};
use crate::{StorageError, StorageResult};
fn checked_u32_values(values: &[usize], what: &str) -> StorageResult<Vec<u32>> {
values
.iter()
.map(|&v| {
u32::try_from(v).map_err(|_| {
StorageError::Overflow(format!(
"{} value {} exceeds u32::MAX and would be silently truncated",
what, v
))
})
})
.collect()
}
async fn ensure_parent_dir(path: &Path) -> StorageResult<()> {
let parent = path
.parent()
.filter(|p| !p.as_os_str().is_empty())
.ok_or_else(|| StorageError::Invalid(format!("path has no parent directory: {path:?}")))?;
tokio::fs::create_dir_all(parent)
.await
.map_err(|e| StorageError::Io(format!("create dir {parent:?}: {e}")))
}
#[derive(Debug, Clone)]
pub struct LanceStorageGraph {
pub(crate) base: String,
pub(crate) name: String,
}
impl LanceStorageGraph {
pub fn new(base: String, name: String) -> Self {
info!("Creating LanceStorage at base={}, name={}", base, name);
Self { base, name }
}
pub async fn spawn(base_path: String) -> Result<(Self, GeneMetadata), StorageError> {
let (exists, md_path) = Self::exists(&base_path);
if !exists || md_path.is_none() {
return Err(StorageError::Invalid(format!(
"Metadata does not exist in base path: {}",
base_path
)));
}
let metadata = GeneMetadata::read(md_path.unwrap()).await?;
let storage = Self::new(base_path.clone(), metadata.name_id.clone());
Ok((storage, metadata))
}
pub fn scoped_generation(&self, generation: u64) -> Self {
let logical = crate::generations::logical_name(&self.name);
Self::new(
self.base.clone(),
crate::generations::generation_name(logical, generation),
)
}
}
impl LanceStorage for LanceStorageGraph {}
impl StorageBackend for LanceStorageGraph {
fn get_base(&self) -> String {
self.base.clone()
}
fn get_name(&self) -> String {
self.name.clone()
}
fn base_path(&self) -> PathBuf {
PathBuf::from(&self.base)
}
fn metadata_path(&self) -> PathBuf {
self.base_path()
.join(format!("{}_metadata.json", self.name))
}
fn basepath_to_uri(&self) -> StorageResult<String> {
Self::path_to_uri(PathBuf::from(self.base.clone()).as_path())
}
async fn save_dense(
&self,
key: &str,
matrix: &DenseMatrix<f64>,
md_path: &Path,
) -> StorageResult<()> {
self.validate_initialized(md_path)?;
let path = self.file_path(key);
let (n_rows, n_cols) = matrix.shape();
info!(
"Saving dense {} matrix: {} x {} at {:?}",
key, n_rows, n_cols, path
);
let batch = self.to_dense_record_batch(matrix)?;
if batch.num_rows() != n_rows {
return Err(StorageError::Invalid(format!(
"RecordBatch has {} rows but matrix has {} rows",
batch.num_rows(),
n_rows
)));
}
{
let uri = Self::path_to_uri(&path)?;
self.write_lance_batch_async(uri, batch).await?;
crate::commit::with_commit_actor(&self.metadata_path(), || async {
let mut md = self.load_metadata().await?;
md = md.add_file(
key,
FileInfo::new(
format!("{}_{}.lance", self.get_name(), key),
"dense",
matrix.shape(),
None,
None,
)?,
);
self.save_metadata(&md).await
})
.await?;
info!("Dense {} matrix saved successfully", key);
}
Ok(())
}
async fn load_dense(&self, key: &str) -> StorageResult<DenseMatrix<f64>> {
let path = self.file_path(key);
info!("Loading dense {} matrix from {:?}", key, path);
let uri = Self::path_to_uri(&path)?;
let batch = self.read_lance_all_batches_async(uri).await?;
let matrix = self.from_dense_record_batch(&batch)?;
let (n_rows, n_cols) = matrix.shape();
info!("Loaded dense {} matrix: {} x {}", key, n_rows, n_cols);
Ok(matrix)
}
async fn load_dense_from_file(&self, path: &Path) -> StorageResult<DenseMatrix<f64>> {
info!("Loading dense matrix from file (async): {:?}", path);
if !path.exists() {
return Err(StorageError::Invalid(format!(
"Dense file does not exist: {:?}",
path
)));
}
let extension = path
.extension()
.and_then(|e| e.to_str())
.ok_or_else(|| StorageError::Invalid(format!("Invalid file path: {:?}", path)))?;
match extension {
"lance" => {
let parent = path
.parent()
.ok_or_else(|| {
StorageError::Invalid(format!("Path has no parent: {:?}", path))
})?
.to_str()
.ok_or_else(|| {
StorageError::Invalid(format!("Non-UTF8 parent path for {:?}", path))
})?
.to_string();
let tmp_storage = Self::new(parent, String::from("tmp_storage"));
let uri = Self::path_to_uri(path)?;
let batch = tmp_storage.read_lance_all_batches_async(uri).await?;
let matrix = tmp_storage.from_dense_record_batch(&batch)?;
info!(
"Loaded dense matrix from Lance: {} x {}",
matrix.shape().0,
matrix.shape().1
);
Ok(matrix)
}
"parquet" => {
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
let owned_path = path.to_path_buf();
let combined =
tokio::task::spawn_blocking(move || -> StorageResult<RecordBatch> {
let file = std::fs::File::open(&owned_path).map_err(|e| {
StorageError::Io(format!("Failed to open parquet file: {}", e))
})?;
let builder =
ParquetRecordBatchReaderBuilder::try_new(file).map_err(|e| {
StorageError::Parquet(format!(
"Failed to create parquet reader: {}",
e
))
})?;
let reader = builder.build().map_err(|e| {
StorageError::Parquet(format!("Failed to build parquet reader: {}", e))
})?;
let batches: Vec<RecordBatch> =
reader.collect::<Result<Vec<_>, _>>().map_err(|e| {
StorageError::Parquet(format!(
"Failed to read parquet batch: {}",
e
))
})?;
if batches.is_empty() {
return Err(StorageError::Invalid(format!(
"Empty parquet dataset at {:?}",
owned_path
)));
}
let schema = batches[0].schema();
arrow::compute::concat_batches(&schema, &batches).map_err(|e| {
StorageError::Parquet(format!(
"Failed to concatenate parquet batches: {}",
e
))
})
})
.await
.map_err(|e| {
StorageError::Io(format!("Parquet reader task failed: {}", e))
})??;
let schema = combined.schema();
let fields = schema.fields();
let is_vector = fields.len() == 1
&& matches!(
fields[0].data_type(),
DataType::FixedSizeList(inner, _)
if matches!(inner.data_type(), DataType::Float64)
);
let is_wide_col = !is_vector
&& !fields.is_empty()
&& fields
.iter()
.all(|f| matches!(f.data_type(), DataType::Float64))
&& fields.iter().any(|f| f.name().starts_with("col_"));
let matrix = if is_vector {
let parent = path
.parent()
.ok_or_else(|| {
StorageError::Invalid(format!("Path has no parent: {:?}", path))
})?
.to_str()
.ok_or_else(|| {
StorageError::Invalid(format!("Non-UTF8 parent path for {:?}", path))
})?
.to_string();
let tmp_storage = Self::new(parent, String::from("tmp_storage"));
tmp_storage.from_dense_record_batch(&combined)?
} else if is_wide_col {
let n_rows = combined.num_rows();
let n_cols = combined.num_columns();
if n_rows == 0 || n_cols == 0 {
return Err(StorageError::Invalid(format!(
"Cannot load empty wide-column parquet at {:?}",
path
)));
}
let mut data = Vec::with_capacity(n_rows * n_cols);
for col_idx in 0..n_cols {
let col = combined.column(col_idx);
let arr = col.as_any().downcast_ref::<Float64Array>().ok_or_else(|| {
StorageError::Invalid(format!(
"Wide-column parquet expects Float64, got {:?} in column {}",
col.data_type(),
col_idx
))
})?;
for row_idx in 0..n_rows {
data.push(arr.value(row_idx));
}
}
DenseMatrix::new(n_rows, n_cols, data, true)
.map_err(|e| StorageError::Invalid(e.to_string()))?
} else {
return Err(StorageError::Invalid(format!(
"Unsupported Parquet schema at {:?}: expected FixedSizeList<Float64> \
or wide Float64 columns named col_*",
path
)));
};
info!(
"Loaded dense matrix from Parquet: {} x {}",
matrix.shape().0,
matrix.shape().1
);
Ok(matrix)
}
_ => Err(StorageError::Invalid(format!(
"Unsupported file format: {}. Only .lance and .parquet are supported",
extension
))),
}
}
fn file_path(&self, key: &str) -> PathBuf {
self.base_path()
.join(format!("{}_{}.lance", self.name, key))
}
async fn save_sparse(
&self,
key: &str,
matrix: &CsMat<f64>,
md_path: &Path,
) -> StorageResult<()> {
self.validate_initialized(md_path)?;
if matrix.nnz() == 0 {
return Err(StorageError::Invalid(format!(
"sparse '{key}' is fully disconnected: adjacency nnz=0 ({}x{}) — \
raise --eps/--k or gate the input before persisting",
matrix.rows(),
matrix.cols()
)));
}
let path = self.file_path(key);
info!(
"Saving sparse {} matrix: {} x {}, nnz={} at {:?}",
key,
matrix.rows(),
matrix.cols(),
matrix.nnz(),
path
);
let filetype = FileInfo::which_filetype(key)?;
{
crate::commit::with_commit_actor(&self.metadata_path(), || async {
let mut metadata = self.load_metadata().await?;
metadata = metadata.add_file(
key,
FileInfo::new(
format!("{}_{}.lance", self.get_name(), key),
filetype.as_str(),
(matrix.rows(), matrix.cols()),
Some(matrix.nnz()),
None,
)?,
);
self.save_metadata(&metadata).await
})
.await?;
let batch = self.to_sparse_record_batch(matrix)?;
let uri = Self::path_to_uri(&path)?;
self.write_lance_batch_async(uri, batch).await?;
}
info!("Sparse matrix {} saved successfully", filetype);
Ok(())
}
async fn load_sparse(&self, key: &str) -> StorageResult<CsMat<f64>> {
info!("Loading sparse {} matrix", key);
let metadata = self.load_metadata().await?;
let filetype = FileInfo::which_filetype(key)?;
let file_info = metadata
.files
.get(key)
.ok_or_else(|| StorageError::Invalid(format!("{key} not found in metadata")))?;
let expected_rows = file_info.rows;
let expected_cols = file_info.cols;
debug!(
"Expected dimensions from storage metadata: {} x {}",
expected_rows, expected_cols
);
let path = self.file_path(key);
let uri = Self::path_to_uri(&path)?;
let batch = self.read_lance_all_batches_async(uri).await?;
let matrix = self.from_sparse_record_batch(batch, expected_rows, expected_cols)?;
info!(
"Sparse {} matrix loaded: {} x {}, nnz={}",
filetype,
matrix.rows(),
matrix.cols(),
matrix.nnz()
);
Ok(matrix)
}
async fn save_lambdas(&self, lambdas: &[f64], md_path: &Path) -> StorageResult<()> {
info!("Saving {} lambda values", lambdas.len());
self.save_primitive_column::<Float64Type>("lambdas", "lambda", lambdas.to_vec(), md_path)
.await
}
async fn load_lambdas(&self) -> StorageResult<Vec<f64>> {
let path = self.file_path("lambdas");
info!("Loading lambda values from {:?}", path);
let uri = Self::path_to_uri(&path)?;
let batch = self.read_lance_all_batches_async(uri).await?;
let arr = batch
.column(0)
.as_any()
.downcast_ref::<Float64Array>()
.ok_or_else(|| StorageError::Invalid("lambda column type mismatch".into()))?;
let lambdas: Vec<f64> = (0..arr.len()).map(|i| arr.value(i)).collect();
info!("Loaded {} lambda values", lambdas.len());
Ok(lambdas)
}
async fn save_vector(&self, key: &str, vector: &[f64], md_path: &Path) -> StorageResult<()> {
info!("Saving {} values for vector {}", vector.len(), key);
self.save_primitive_column::<Float64Type>(key, "element", vector.to_vec(), md_path)
.await
}
async fn save_index(&self, key: &str, vector: &[usize], md_path: &Path) -> StorageResult<()> {
info!("Saving {} values for index {}", vector.len(), key);
let values = checked_u32_values(vector, "index")?;
self.save_primitive_column::<UInt32Type>(key, "id", values, md_path)
.await
}
async fn load_vector(&self, filename: &str) -> StorageResult<Vec<f64>> {
let path = self.file_path(filename);
info!("Loading vector {} from {:?}", filename, path);
let uri = Self::path_to_uri(&path)?;
let batch = self.read_lance_all_batches_async(uri).await?;
let arr = batch
.column(0)
.as_any()
.downcast_ref::<Float64Array>()
.ok_or_else(|| StorageError::Invalid("column type mismatch".into()))?;
let vector: Vec<f64> = (0..arr.len()).map(|i| arr.value(i)).collect();
info!("Loaded {} vector values for {}", vector.len(), filename);
Ok(vector)
}
async fn load_index(&self, filename: &str) -> StorageResult<Vec<usize>> {
let path = self.file_path(filename);
info!("Loading vector {} from {:?}", filename, path);
let uri = Self::path_to_uri(&path)?;
let batch = self.read_lance_all_batches_async(uri).await?;
let arr = batch
.column(0)
.as_any()
.downcast_ref::<UInt32Array>()
.ok_or_else(|| StorageError::Invalid("column type mismatch".into()))?;
let vector: Vec<usize> = (0..arr.len()).map(|i| arr.value(i) as usize).collect();
info!("Loaded {} vector values for {}", vector.len(), filename);
Ok(vector)
}
async fn save_dense_to_file(data: &DenseMatrix<f64>, path: &Path) -> StorageResult<()> {
info!("Saving dense matrix to file (async): {:?}", path);
let parent = path
.parent()
.filter(|p| !p.as_os_str().is_empty())
.ok_or_else(|| {
StorageError::Invalid(format!("path has no parent directory: {:?}", path))
})?;
tokio::fs::create_dir_all(parent)
.await
.map_err(|e| StorageError::Io(format!("create dir {:?}: {}", parent, e)))?;
let parent_str = parent
.to_str()
.ok_or_else(|| StorageError::Invalid(format!("non-UTF8 parent path for {:?}", path)))?
.to_string();
let tmp_storage = Self::new(parent_str, String::from("tmp_storage"));
let extension = path
.extension()
.and_then(|e| e.to_str())
.ok_or_else(|| StorageError::Invalid(format!("Invalid file path: {:?}", path)))?;
let (n_rows, n_cols) = data.shape();
info!("Saving matrix: {} rows x {} cols", n_rows, n_cols);
match extension {
"lance" => {
let batch = tmp_storage.to_dense_record_batch(data)?;
debug!(
"Created RecordBatch with {} rows for Lance",
batch.num_rows()
);
if batch.num_rows() != n_rows {
return Err(StorageError::Invalid(format!(
"RecordBatch has {} rows but matrix has {} rows",
batch.num_rows(),
n_rows
)));
}
let uri = Self::path_to_uri(path)?;
tmp_storage.write_lance_batch_async(uri, batch).await?;
info!("Saved dense matrix to Lance: {} x {}", n_rows, n_cols);
Ok(())
}
"parquet" => {
use parquet::arrow::ArrowWriter;
use parquet::file::properties::WriterProperties;
use std::fs::File;
let batch = tmp_storage.to_dense_record_batch(data)?;
debug!(
"Created RecordBatch with {} rows for Parquet",
batch.num_rows()
);
if batch.num_rows() != n_rows {
return Err(StorageError::Invalid(format!(
"RecordBatch has {} rows but matrix has {} rows",
batch.num_rows(),
n_rows
)));
}
let owned_path = path.to_path_buf();
tokio::task::spawn_blocking(move || -> StorageResult<()> {
let file = File::create(&owned_path).map_err(|e| {
StorageError::Io(format!("Failed to create parquet file: {}", e))
})?;
let props = WriterProperties::builder()
.set_compression(parquet::basic::Compression::SNAPPY)
.build();
let mut writer = ArrowWriter::try_new(file, batch.schema(), Some(props))
.map_err(|e| {
StorageError::Parquet(format!("Failed to create parquet writer: {}", e))
})?;
writer.write(&batch).map_err(|e| {
StorageError::Parquet(format!("Failed to write batch: {}", e))
})?;
writer.close().map_err(|e| {
StorageError::Parquet(format!("Failed to close writer: {}", e))
})?;
Ok(())
})
.await
.map_err(|e| StorageError::Io(format!("parquet writer task failed: {}", e)))??;
info!("Saved dense matrix to Parquet: {} x {}", n_rows, n_cols);
Ok(())
}
_ => Err(StorageError::Invalid(format!(
"Unsupported file format: {}. Only .lance and .parquet are supported",
extension
))),
}
}
async fn save_centroid_map(&self, map: &[usize], md_path: &Path) -> StorageResult<()> {
info!("Saving {} centroid map entries", map.len());
let values = checked_u32_values(map, "centroid map")?;
self.save_primitive_column::<UInt32Type>("centroid_map", "centroid_id", values, md_path)
.await
}
async fn load_centroid_map(&self) -> StorageResult<Vec<usize>> {
let path = self.file_path("centroid_map");
info!("Loading centroid map from {:?}", path);
let uri = Self::path_to_uri(&path)?;
let batch = self.read_lance_all_batches_async(uri).await?;
let arr = batch
.column(0)
.as_any()
.downcast_ref::<UInt32Array>()
.ok_or_else(|| StorageError::Invalid("centroid_id column type mismatch".into()))?;
let map: Vec<usize> = (0..arr.len()).map(|i| arr.value(i) as usize).collect();
info!("Loaded {} centroid map entries", map.len());
Ok(map)
}
async fn save_subcentroid_lambdas(&self, lambdas: &[f64], md_path: &Path) -> StorageResult<()> {
info!("Saving {} subcentroid lambda values", lambdas.len());
self.save_primitive_column::<Float64Type>(
"subcentroid_lambdas",
"subcentroid_lambda",
lambdas.to_vec(),
md_path,
)
.await
}
async fn load_subcentroid_lambdas(&self) -> StorageResult<Vec<f64>> {
let path = self.file_path("subcentroid_lambdas");
info!("Loading subcentroid lambda values from {:?}", path);
let uri = Self::path_to_uri(&path)?;
let batch = self.read_lance_all_batches_async(uri).await?;
let arr = batch
.column(0)
.as_any()
.downcast_ref::<Float64Array>()
.ok_or_else(|| {
StorageError::Invalid("subcentroid_lambda column type mismatch".into())
})?;
let lambdas: Vec<f64> = (0..arr.len()).map(|i| arr.value(i)).collect();
info!("Loaded {} subcentroid lambda values", lambdas.len());
Ok(lambdas)
}
async fn save_subcentroids(
&self,
subcentroids: &DenseMatrix<f64>,
md_path: &Path,
) -> StorageResult<()> {
self.validate_initialized(md_path)?;
let key = "sub_centroids";
let path = self.file_path(key);
let (n_rows, n_cols) = subcentroids.shape();
info!(
"Saving subcentroids matrix {} x {} at {:?}",
n_rows, n_cols, path
);
let batch = self.to_dense_record_batch(subcentroids)?;
{
crate::commit::with_commit_actor(&self.metadata_path(), || async {
let mut metadata = self.load_metadata().await?;
metadata = metadata.add_file(
key,
FileInfo::new(
format!("{}_{}.lance", self.get_name(), key),
"vector",
subcentroids.shape(),
None,
None,
)?,
);
self.save_metadata(&metadata).await
})
.await?;
let uri = Self::path_to_uri(&path)?;
self.write_lance_batch_async(uri, batch).await?;
}
debug!("Subcentroids matrix saved successfully");
Ok(())
}
async fn load_subcentroids(&self) -> StorageResult<Vec<Vec<f64>>> {
let path = self.file_path("sub_centroids");
info!("Loading sub_centroids from {:?}", path);
let uri = Self::path_to_uri(&path)?;
let batch = self.read_lance_all_batches_async(uri).await?;
let matrix = self.from_dense_record_batch(&batch)?;
let (n_rows, n_cols) = matrix.shape();
let mut result = Vec::with_capacity(n_rows);
for row_idx in 0..n_rows {
let row: Vec<f64> = (0..n_cols)
.map(|col_idx| *matrix.get((row_idx, col_idx)))
.collect();
result.push(row);
}
info!(
"Loaded sub_centroids: {} x {} as Vec<Vec<f64>>",
n_rows, n_cols
);
Ok(result)
}
async fn save_item_norms(&self, item_norms: &[f64], md_path: &Path) -> StorageResult<()> {
info!("Saving {} item norm values", item_norms.len());
self.save_primitive_column::<Float64Type>(
"item_norms",
"norm",
item_norms.to_vec(),
md_path,
)
.await
}
async fn load_item_norms(&self) -> StorageResult<Vec<f64>> {
let path = self.file_path("item_norms");
info!("Loading item norms from {:?}", path);
let uri = Self::path_to_uri(&path)?;
let batch = self.read_lance_all_batches_async(uri).await?;
let arr = batch
.column(0)
.as_any()
.downcast_ref::<Float64Array>()
.ok_or_else(|| StorageError::Invalid("norm column type mismatch".into()))?;
let norms: Vec<f64> = (0..arr.len()).map(|i| arr.value(i)).collect();
info!("Loaded {} item norm values", norms.len());
Ok(norms)
}
async fn save_cluster_assignments(
&self,
assignments: &[Option<usize>],
md_path: &Path,
) -> StorageResult<()> {
info!("Saving {} cluster assignments", assignments.len());
let values: Vec<i64> = assignments
.iter()
.map(|opt| opt.map(|v| v as i64).unwrap_or(-1))
.collect();
self.save_primitive_column::<Int64Type>(
"cluster_assignments",
"cluster_id",
values,
md_path,
)
.await
}
async fn load_cluster_assignments(&self) -> StorageResult<Vec<Option<usize>>> {
use arrow::array::Int64Array;
let path = self.file_path("cluster_assignments");
info!("Loading cluster assignments from {:?}", path);
let uri = Self::path_to_uri(&path)?;
let batch = self.read_lance_all_batches_async(uri).await?;
let arr = batch
.column(0)
.as_any()
.downcast_ref::<Int64Array>()
.ok_or_else(|| StorageError::Invalid("cluster_id column type mismatch".into()))?;
let assignments: Vec<Option<usize>> = (0..arr.len())
.map(|i| {
let v = arr.value(i);
if v < 0 { None } else { Some(v as usize) }
})
.collect();
info!("Loaded {} cluster assignments", assignments.len());
Ok(assignments)
}
async fn save_vectors_with(
&self,
name: &str,
batch: &RecordBatch,
properties: &BTreeMap<String, String>,
md_path: &Path,
) -> StorageResult<()> {
self.validate_initialized(md_path)?;
let (batch, dim) = vector_space_collection_batch(batch, properties)?;
let rows = batch.num_rows();
let path = self.file_path(name);
info!(
"Saving vector-space collection '{}' ({} rows, vector dim {}) at {:?}",
name, rows, dim, path
);
let uri = Self::path_to_uri(&path)?;
self.write_lance_batch_async(uri, batch).await?;
crate::commit::with_commit_actor(&self.metadata_path(), || async {
let mut md = self.load_metadata().await?;
let mut info = FileInfo::new(
format!("{}_{}.lance", self.get_name(), name),
"vectors",
(rows, dim as usize),
None,
None,
)?;
info.properties = properties.clone();
md = md.add_file(name, info);
self.save_metadata(&md).await
})
.await?;
info!("Vector-space collection '{}' saved successfully", name);
Ok(())
}
async fn load_vectors(&self, name: &str) -> StorageResult<RecordBatch> {
info!("Loading vector-space collection '{}'", name);
self.load_vectors_from_path(&self.file_path(name)).await
}
async fn save_graph_with(
&self,
name: &str,
edges: &[GraphEdge],
options: &GraphWriteOptions,
md_path: &Path,
) -> StorageResult<()> {
self.validate_initialized(md_path)?;
let (batch, num_nodes) = graph_record_batch(edges, options)?;
let num_nodes_usize = usize::try_from(num_nodes)
.map_err(|_| StorageError::Overflow(format!("node count {num_nodes} exceeds usize")))?;
let weighted = edges[0].weight.is_some();
let width = options.node_id_width;
let path = self.file_path(name);
info!(
"Saving graph collection '{}' ({} edges, {} nodes, width {:?}, weight width {:?}, weighted {}) at {:?}",
name,
edges.len(),
num_nodes,
width,
options.weight_type,
weighted,
path
);
let uri = Self::path_to_uri(&path)?;
self.write_lance_batch_async(uri, batch).await?;
crate::commit::with_commit_actor(&self.metadata_path(), || async {
let mut md = self.load_metadata().await?;
let mut info = FileInfo::new(
format!("{}_{}.lance", self.get_name(), name),
"graph",
(num_nodes_usize, num_nodes_usize),
Some(edges.len()),
None,
)?;
info.properties = options.properties.clone();
info.properties
.insert("node_id_width".to_string(), width.as_str().to_string());
info.properties
.insert("weighted".to_string(), weighted.to_string());
info.properties
.insert("num_nodes".to_string(), num_nodes.to_string());
info.properties.insert(
"weight_type".to_string(),
options.weight_type.as_str().to_string(),
);
md = md.add_file(name, info);
self.save_metadata(&md).await
})
.await?;
info!("Graph collection '{}' saved successfully", name);
Ok(())
}
async fn load_graph(&self, name: &str) -> StorageResult<StoredGraph> {
info!("Loading graph collection '{}'", name);
self.load_graph_from_path(&self.file_path(name)).await
}
async fn save_vectors_to_path(
&self,
path: &Path,
batch: &RecordBatch,
properties: &BTreeMap<String, String>,
) -> StorageResult<()> {
let (batch, dim) = vector_space_collection_batch(batch, properties)?;
ensure_parent_dir(path).await?;
info!(
"Saving vector-space collection ({} rows, vector dim {}) at {:?} (registry-free)",
batch.num_rows(),
dim,
path
);
let uri = Self::path_to_uri(path)?;
self.write_lance_batch_async(uri, batch).await?;
info!("Vector-space collection saved successfully (registry-free)");
Ok(())
}
async fn load_vectors_from_path(&self, path: &Path) -> StorageResult<RecordBatch> {
info!(
"Loading vector-space collection from {:?} (registry-free)",
path
);
let uri = Self::path_to_uri(path)?;
let batch = self.read_lance_all_batches_async(uri).await?;
validate_vector_space_schema(batch.schema().as_ref())?;
info!(
"Loaded vector-space collection: {} rows (registry-free)",
batch.num_rows()
);
Ok(batch)
}
async fn save_graph_to_path(
&self,
path: &Path,
edges: &[GraphEdge],
options: &GraphWriteOptions,
) -> StorageResult<()> {
let (batch, num_nodes) = graph_record_batch(edges, options)?;
ensure_parent_dir(path).await?;
info!(
"Saving graph collection ({} edges, {} nodes, width {:?}, weight width {:?}, weighted {}) at {:?} (registry-free)",
edges.len(),
num_nodes,
options.node_id_width,
options.weight_type,
edges.first().map(|e| e.weight.is_some()).unwrap_or(false),
path
);
let uri = Self::path_to_uri(path)?;
self.write_lance_batch_async(uri, batch).await?;
info!("Graph collection saved successfully (registry-free)");
Ok(())
}
async fn load_graph_from_path(&self, path: &Path) -> StorageResult<StoredGraph> {
self.load_graph_from_path_with_options(path, &GraphReadOptions::default())
.await
}
async fn load_graph_from_path_with_options(
&self,
path: &Path,
options: &GraphReadOptions,
) -> StorageResult<StoredGraph> {
info!(
"Loading graph collection from {:?} (registry-free, strict={})",
path, options.strict
);
let uri = Self::path_to_uri(path)?;
let batch = self.read_lance_all_batches_async(uri).await?;
let graph = stored_graph_from_batch_with_options(batch, options)?;
info!(
"Loaded graph collection: {} edges, {} nodes, width {:?}, weight width {:?} (registry-free)",
graph.edges.len(),
graph.num_nodes,
graph.node_id_width,
graph.weight_type
);
Ok(graph)
}
async fn load_graph_from_path_strict(&self, path: &Path) -> StorageResult<StoredGraph> {
self.load_graph_from_path_with_options(path, &GraphReadOptions::strict())
.await
}
async fn load_scalars(&self, name: &str) -> StorageResult<Vec<f64>> {
self.load_scalars_from_path(&self.file_path(name)).await
}
async fn save_scalars_to_path(&self, path: &Path, values: &[f64]) -> StorageResult<()> {
if values.is_empty() {
return Err(StorageError::Invalid(
"empty scalar collections are not supported".into(),
));
}
ensure_parent_dir(path).await?;
info!(
"Saving scalar collection ({} values) at {:?} (registry-free)",
values.len(),
path
);
let schema =
Schema::new(vec![Field::new("lambda", DataType::Float64, false)]).with_metadata(
std::collections::HashMap::from([("kind".to_string(), "vector-space".to_string())]),
);
let batch = RecordBatch::try_new(
Arc::new(schema),
vec![Arc::new(Float64Array::from(values.to_vec())) as ArrayRef],
)
.map_err(|e| StorageError::Lance(e.to_string()))?;
let uri = Self::path_to_uri(path)?;
self.write_lance_batch_async(uri, batch).await?;
info!("Scalar collection saved successfully (registry-free)");
Ok(())
}
async fn load_scalars_from_path(&self, path: &Path) -> StorageResult<Vec<f64>> {
info!("Loading scalar collection from {:?} (registry-free)", path);
let uri = Self::path_to_uri(path)?;
let batch = self.read_lance_all_batches_async(uri).await?;
let schema = batch.schema();
match schema.metadata().get("kind").map(String::as_str) {
Some("vector-space") => {}
other => {
return Err(StorageError::Invalid(format!(
"scalar collection at {path:?} has dataset kind {other:?}, \
expected 'vector-space'"
)));
}
}
if schema.fields().len() != 1 {
return Err(StorageError::Invalid(format!(
"scalar collection expects exactly one column, found {}",
schema.fields().len()
)));
}
let field = schema.field(0);
if field.is_nullable() {
return Err(StorageError::Invalid(format!(
"scalar collection column '{}' is nullable",
field.name()
)));
}
let column = batch.column(0);
if column.null_count() != 0 {
return Err(StorageError::Invalid(
"scalar collection column contains nulls".into(),
));
}
let values: Vec<f64> = match field.data_type() {
DataType::Float64 => {
let a = column
.as_any()
.downcast_ref::<Float64Array>()
.ok_or_else(|| StorageError::Invalid("scalar column type mismatch".into()))?;
(0..a.len()).map(|i| a.value(i)).collect()
}
DataType::Float32 => {
let a = column
.as_any()
.downcast_ref::<Float32Array>()
.ok_or_else(|| StorageError::Invalid("scalar column type mismatch".into()))?;
(0..a.len()).map(|i| f64::from(a.value(i))).collect()
}
other => {
return Err(StorageError::Invalid(format!(
"scalar collection column must be Float64|Float32, found {other:?}"
)));
}
};
info!("Loaded {} scalar values (registry-free)", values.len());
Ok(values)
}
async fn collection_schema_from_path(&self, path: &Path) -> StorageResult<Schema> {
info!("Reading collection schema from {:?} (registry-free)", path);
let path = path.to_path_buf();
tokio::task::spawn_blocking(move || crate::lancefmt::read_schema(&path))
.await
.map_err(|e| StorageError::Io(format!("lancefmt read_schema task failed: {e}")))?
}
}