From d4dc9badd449053255c7824e9f7e1f81dcb1bf7c Mon Sep 17 00:00:00 2001 From: Weiyao Luo <9347182+SeliMeli@users.noreply.github.com> Date: Wed, 29 Jul 2026 02:33:42 +0000 Subject: [PATCH 1/6] build: add PiPNN algorithm and memory selection --- Cargo.lock | 3 + diskann-disk/Cargo.toml | 4 + .../build/configuration/build_algorithm.rs | 123 ++++++++++ .../disk_index_build_parameter.rs | 211 +++++++++++++++++- diskann-disk/src/build/configuration/mod.rs | 5 + diskann-disk/src/build/mod.rs | 5 +- diskann-disk/src/lib.rs | 5 +- 7 files changed, 351 insertions(+), 5 deletions(-) create mode 100644 diskann-disk/src/build/configuration/build_algorithm.rs diff --git a/Cargo.lock b/Cargo.lock index 6a88b5550..ec189c7ca 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -475,6 +475,7 @@ dependencies = [ "diskann-disk", "diskann-inmem", "diskann-label-filter", + "diskann-pipnn", "diskann-providers", "diskann-quantization", "diskann-tools", @@ -580,6 +581,7 @@ dependencies = [ "criterion", "diskann", "diskann-linalg", + "diskann-pipnn", "diskann-providers", "diskann-quantization", "diskann-utils", @@ -596,6 +598,7 @@ dependencies = [ "rayon", "rstest", "serde", + "serde_json", "tempfile", "thiserror 2.0.17", "tokio", diff --git a/diskann-disk/Cargo.toml b/diskann-disk/Cargo.toml index 693614381..0c53172e9 100644 --- a/diskann-disk/Cargo.toml +++ b/diskann-disk/Cargo.toml @@ -43,6 +43,7 @@ vfs = { workspace = true } # Optional dependencies opentelemetry = { workspace = true, optional = true } +diskann-pipnn = { workspace = true, optional = true } [target.'cfg(target_os = "linux")'.dependencies] io-uring = "0.6.4" @@ -68,6 +69,7 @@ features = [ rstest.workspace = true tempfile.workspace = true vfs.workspace = true +serde_json.workspace = true diskann-providers = { workspace = true, default-features = false, features = [ "testing", "virtual_storage", @@ -82,6 +84,8 @@ proptest.workspace = true [features] default = [] perf_test = ["dep:opentelemetry"] +pipnn = ["dep:diskann-pipnn"] +virtual_storage = ["diskann-providers/virtual_storage"] experimental_diversity_search = [ "diskann/experimental_diversity_search", "diskann-providers/experimental_diversity_search", diff --git a/diskann-disk/src/build/configuration/build_algorithm.rs b/diskann-disk/src/build/configuration/build_algorithm.rs new file mode 100644 index 000000000..f798138a1 --- /dev/null +++ b/diskann-disk/src/build/configuration/build_algorithm.rs @@ -0,0 +1,123 @@ +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT license. + */ + +//! Graph-build algorithm selection and its JSON-facing configuration. + +use std::fmt; + +use serde::{Deserialize, Serialize}; + +/// JSON-facing PiPNN parameters. +/// +/// Graph degree, build-L, alpha, metric, threads, and memory limits remain in +/// the outer index configuration shared with Vamana. +#[cfg(feature = "pipnn")] +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[serde(default, deny_unknown_fields)] +pub struct PiPNNParameters { + /// Maximum number of points in a leaf. + pub c_max: usize, + /// Minimum leaf size used by global small-leaf merging. + pub c_min: usize, + /// Fraction of a cluster sampled as leaders. + pub p_samp: f64, + /// Number of nearest leaders retained at each partition level. + pub fanout: Vec, + /// Number of nearest neighbors selected within each leaf. + pub k: usize, + /// Number of independent partition passes. + pub replicas: usize, +} + +#[cfg(feature = "pipnn")] +impl Default for PiPNNParameters { + fn default() -> Self { + Self { + c_max: 256, + c_min: 16, + p_samp: 0.005, + fanout: vec![8, 3], + k: 2, + replicas: 1, + } + } +} + +#[cfg(feature = "pipnn")] +impl From<&PiPNNParameters> for diskann_pipnn::PiPNNConfig { + fn from(config: &PiPNNParameters) -> Self { + Self { + c_max: config.c_max, + c_min: config.c_min, + p_samp: config.p_samp, + fanout: config.fanout.clone(), + k: config.k, + replicas: config.replicas, + } + } +} + +/// Selects the graph construction algorithm for index building. +#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)] +#[serde(tag = "algorithm")] +#[non_exhaustive] +pub enum BuildAlgorithm { + /// Default Vamana graph construction. + #[default] + Vamana, + + /// PiPNN one-shot partition-based graph construction. + #[cfg(feature = "pipnn")] + PiPNN(PiPNNParameters), +} + +impl fmt::Display for BuildAlgorithm { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Vamana => write!(f, "Vamana"), + #[cfg(feature = "pipnn")] + Self::PiPNN(config) => write!(f, "PiPNN({config:?})"), + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn default_is_vamana() { + assert_eq!(BuildAlgorithm::default(), BuildAlgorithm::Vamana); + } + + #[test] + fn vamana_serde_roundtrip() { + let json = serde_json::to_string(&BuildAlgorithm::Vamana).unwrap(); + assert_eq!( + serde_json::from_str::(&json).unwrap(), + BuildAlgorithm::Vamana + ); + } + + #[cfg(feature = "pipnn")] + #[test] + fn pipnn_serde_uses_inline_defaults_and_rejects_unknown_fields() { + let algorithm: BuildAlgorithm = serde_json::from_str( + r#"{"algorithm":"PiPNN","c_max":512,"c_min":64,"fanout":[10,3],"k":3}"#, + ) + .unwrap(); + let BuildAlgorithm::PiPNN(config) = algorithm else { + panic!("expected PiPNN"); + }; + assert_eq!(config.c_max, 512); + assert_eq!(config.c_min, 64); + assert_eq!(config.fanout, [10, 3]); + assert_eq!(config.k, 3); + assert_eq!(config.replicas, 1); + assert!( + serde_json::from_str::(r#"{"algorithm":"PiPNN","l_max":72}"#).is_err() + ); + } +} diff --git a/diskann-disk/src/build/configuration/disk_index_build_parameter.rs b/diskann-disk/src/build/configuration/disk_index_build_parameter.rs index 077c5090d..939c44742 100644 --- a/diskann-disk/src/build/configuration/disk_index_build_parameter.rs +++ b/diskann-disk/src/build/configuration/disk_index_build_parameter.rs @@ -8,9 +8,13 @@ use std::num::NonZeroUsize; use diskann::ANNError; +#[cfg(feature = "pipnn")] +use diskann::ANNResult; use thiserror::Error; -use super::QuantizationType; +#[cfg(feature = "pipnn")] +use super::PiPNNParameters; +use super::{BuildAlgorithm, QuantizationType}; /// GB to bytes ratio. pub const BYTES_IN_GB: f64 = 1024_f64 * 1024_f64 * 1024_f64; @@ -105,9 +109,10 @@ impl NumPQChunks { } /// Parameters specific for disk index construction. -#[derive(Clone, Copy, PartialEq, Debug)] +#[derive(Clone, PartialEq, Debug)] pub struct DiskIndexBuildParameters { - /// Limit on the memory allowed for building the index. + /// Limit on graph-construction memory. PiPNN falls back to Vamana when its + /// estimated one-shot peak exceeds this value. build_memory_limit: MemoryBudget, /// Number of PQ chunks stored in-memory for search and to be generated during build. @@ -118,6 +123,9 @@ pub struct DiskIndexBuildParameters { /// Number of vectors processed per data-compression chunk. data_compression_chunk_vector_count: usize, + + /// Which graph construction algorithm to use. + build_algorithm: BuildAlgorithm, } impl DiskIndexBuildParameters { @@ -132,6 +140,26 @@ impl DiskIndexBuildParameters { search_pq_chunks, build_quantization, data_compression_chunk_vector_count: DEFAULT_DATA_COMPRESSION_CHUNK_VECTOR_COUNT, + build_algorithm: BuildAlgorithm::default(), + } + } + + /// Create parameters for one-shot PiPNN graph construction. + /// + /// PiPNN uses the common search-PQ and disk-layout pipeline. The memory + /// budget selects Vamana when the estimated PiPNN peak does not fit. + #[cfg(feature = "pipnn")] + pub fn new_pipnn( + build_memory_limit: MemoryBudget, + search_pq_chunks: NumPQChunks, + config: PiPNNParameters, + ) -> Self { + Self { + build_memory_limit, + search_pq_chunks, + build_quantization: QuantizationType::FP, + data_compression_chunk_vector_count: DEFAULT_DATA_COMPRESSION_CHUNK_VECTOR_COUNT, + build_algorithm: BuildAlgorithm::PiPNN(config), } } @@ -163,6 +191,129 @@ impl DiskIndexBuildParameters { pub fn data_compression_chunk_vector_count(&self) -> usize { self.data_compression_chunk_vector_count } + + /// Get the graph-construction algorithm. + pub fn build_algorithm(&self) -> &BuildAlgorithm { + &self.build_algorithm + } + + #[cfg(feature = "pipnn")] + pub(crate) fn pipnn_config(&self) -> Option { + match &self.build_algorithm { + BuildAlgorithm::PiPNN(config) => Some(config.into()), + BuildAlgorithm::Vamana => None, + } + } + + #[cfg(feature = "pipnn")] + pub(crate) fn use_vamana_if_pipnn_exceeds( + &mut self, + npoints: usize, + dimensions: usize, + element_size: usize, + num_threads: usize, + ) -> ANNResult { + let parameters = match &self.build_algorithm { + BuildAlgorithm::PiPNN(parameters) => parameters, + BuildAlgorithm::Vamana => { + return Err(ANNError::log_index_error( + "memory selection requires PiPNN parameters", + )); + } + }; + let estimate = + estimate_pipnn_peak_memory(parameters, npoints, dimensions, element_size, num_threads); + if estimate.is_none_or(|bytes| bytes > self.build_memory_limit.in_bytes()) { + self.build_algorithm = BuildAlgorithm::Vamana; + } + Ok(estimate.unwrap_or(usize::MAX)) + } +} + +#[cfg(feature = "pipnn")] +fn estimate_pipnn_peak_memory( + config: &PiPNNParameters, + npoints: usize, + dimensions: usize, + element_size: usize, + num_threads: usize, +) -> Option { + const PROPORTIONAL_HEADROOM_PERCENT: u128 = 108; + const PROCESS_HEADROOM: u128 = 16 * 1024 * 1024; + const PER_WORKER_HEADROOM: u128 = 28 * 1024 * 1024; + + let workers = if num_threads == 0 { + std::thread::available_parallelism().map_or(1, NonZeroUsize::get) + } else { + num_threads + } as u128; + let copies = config + .fanout + .iter() + .try_fold(config.replicas as u128, |copies, &fanout| { + copies.checked_mul(fanout as u128) + })?; + let dataset = (dimensions as u128).checked_mul(element_size as u128)?; + let leaf_ids = copies.checked_mul(size_of::() as u128)?; + let offered = copies.checked_mul(config.k as u128)?.checked_mul(2)?; + let candidate_capacity = offered.checked_next_power_of_two()?; + let candidate_storage = candidate_capacity.checked_mul(size_of::() as u128)?; + let candidate_metadata = + size_of::>>() as u128; + + let partition_per_point = dataset + .checked_add(16)? + .checked_add(leaf_ids.checked_mul(2)?)?; + let leaf_per_point = dataset + .checked_add(candidate_metadata)? + .checked_add(candidate_storage)? + .checked_add(leaf_ids)?; + let finalization_per_point = dataset + .checked_add(candidate_metadata)? + .checked_add(size_of::>() as u128)? + .checked_add(candidate_storage)?; + + let leaf_size = config.c_max.min(npoints).max(1) as u128; + let leaf_scratch = leaf_size + .checked_mul(dimensions as u128)? + .checked_mul(size_of::() as u128)? + .checked_add(leaf_size.checked_mul(leaf_size)?.checked_mul(5)?)? + .checked_add(leaf_size.checked_mul(config.k as u128)?.checked_mul(24)?)? + .checked_add(leaf_size.checked_mul(20)?)? + .checked_add(68)?; + let partition_rows = (npoints as u128).min(1_024); + let partition_leaders = (npoints as u128).min(1_000); + let partition_scratch = partition_leaders + .checked_mul(dimensions as u128)? + .checked_mul(size_of::() as u128)? + .checked_add( + partition_rows + .checked_mul(dimensions as u128)? + .checked_mul(size_of::() as u128)?, + )? + .checked_add( + partition_rows + .checked_mul(partition_leaders)? + .checked_mul(size_of::() as u128)?, + )?; + + let points = npoints as u128; + let structural = points + .checked_mul(partition_per_point)? + .checked_add(workers.checked_mul(partition_scratch)?)? + .max( + points + .checked_mul(leaf_per_point)? + .checked_add(workers.checked_mul(leaf_scratch)?)?, + ) + .max(points.checked_mul(finalization_per_point)?); + let proportional = structural + .checked_mul(PROPORTIONAL_HEADROOM_PERCENT)? + .div_ceil(100); + let allocator_floor = structural + .checked_add(PROCESS_HEADROOM)? + .checked_add(workers.checked_mul(PER_WORKER_HEADROOM)?)?; + usize::try_from(proportional.max(allocator_floor)).ok() } #[cfg(test)] @@ -245,4 +396,58 @@ mod dataset_test { let chunks = NumPQChunks::new_with(64, 128).unwrap(); assert_eq!(chunks.get(), 64); } + + #[cfg(feature = "pipnn")] + #[test] + fn new_pipnn_uses_the_common_disk_pipeline_parameters() { + let pq = NumPQChunks::new_with(1, 128).unwrap(); + let parameters = PiPNNParameters::default(); + let config = diskann_pipnn::PiPNNConfig::from(¶meters); + let budget = MemoryBudget::try_from_gb(2.0).unwrap(); + let params = DiskIndexBuildParameters::new_pipnn(budget, pq, parameters); + + assert_eq!(params.pipnn_config(), Some(config)); + assert_eq!(params.search_pq_chunks(), pq); + assert_eq!( + params.data_compression_chunk_vector_count(), + DEFAULT_DATA_COMPRESSION_CHUNK_VECTOR_COUNT + ); + assert!(matches!(params.build_algorithm(), BuildAlgorithm::PiPNN(_))); + } + + #[cfg(feature = "pipnn")] + #[test] + fn pipnn_memory_estimate_covers_measured_bigann10m_peak() { + const MEASURED_PEAK: usize = 8_656_004 * 1024; + let parameters = PiPNNParameters { + c_max: 512, + c_min: 64, + p_samp: 0.01, + fanout: vec![10, 3], + k: 2, + replicas: 1, + }; + + let estimate = estimate_pipnn_peak_memory(¶meters, 10_000_000, 128, 2, 16).unwrap(); + + assert!(estimate >= MEASURED_PEAK); + assert!(estimate <= MEASURED_PEAK * 120 / 100); + } + + #[cfg(feature = "pipnn")] + #[test] + fn pipnn_falls_back_to_vamana_above_budget() { + let budget = MemoryBudget::try_from_gb(1.0).unwrap(); + let pq = NumPQChunks::new_with(1, 128).unwrap(); + let mut params = + DiskIndexBuildParameters::new_pipnn(budget, pq, PiPNNParameters::default()); + + let estimate = params + .use_vamana_if_pipnn_exceeds(10_000_000, 128, 2, 16) + .unwrap(); + + assert!(estimate > budget.in_bytes()); + assert!(matches!(params.build_algorithm(), BuildAlgorithm::Vamana)); + assert_eq!(params.build_quantization(), &QuantizationType::FP); + } } diff --git a/diskann-disk/src/build/configuration/mod.rs b/diskann-disk/src/build/configuration/mod.rs index 25453abd0..a7e343fb5 100644 --- a/diskann-disk/src/build/configuration/mod.rs +++ b/diskann-disk/src/build/configuration/mod.rs @@ -2,6 +2,11 @@ * Copyright (c) Microsoft Corporation. * Licensed under the MIT license. */ +pub mod build_algorithm; +pub use build_algorithm::BuildAlgorithm; +#[cfg(feature = "pipnn")] +pub use build_algorithm::PiPNNParameters; + pub mod disk_index_build_parameter; pub use disk_index_build_parameter::{DiskIndexBuildParameters, MemoryBudget, NumPQChunks}; diff --git a/diskann-disk/src/build/mod.rs b/diskann-disk/src/build/mod.rs index 20e6e4b38..27f4c124a 100644 --- a/diskann-disk/src/build/mod.rs +++ b/diskann-disk/src/build/mod.rs @@ -12,6 +12,9 @@ pub mod builder; pub mod configuration; // Re-export key types for convenience +#[cfg(feature = "pipnn")] +pub use configuration::PiPNNParameters; pub use configuration::{ - disk_index_build_parameter, filter_parameter, DiskIndexBuildParameters, QuantizationType, + disk_index_build_parameter, filter_parameter, BuildAlgorithm, DiskIndexBuildParameters, + QuantizationType, }; diff --git a/diskann-disk/src/lib.rs b/diskann-disk/src/lib.rs index 845b774df..a7d3d29a5 100644 --- a/diskann-disk/src/lib.rs +++ b/diskann-disk/src/lib.rs @@ -12,8 +12,11 @@ pub(crate) mod test_utils; pub mod build; +#[cfg(feature = "pipnn")] +pub use build::PiPNNParameters; pub use build::{ - disk_index_build_parameter, filter_parameter, DiskIndexBuildParameters, QuantizationType, + disk_index_build_parameter, filter_parameter, BuildAlgorithm, DiskIndexBuildParameters, + QuantizationType, }; pub mod data_model; From cc31fa8c1e4db7074a7f922017e9eda0da4aac80 Mon Sep 17 00:00:00 2001 From: Weiyao Luo <9347182+SeliMeli@users.noreply.github.com> Date: Wed, 29 Jul 2026 02:34:10 +0000 Subject: [PATCH 2/6] disk: adapt PiPNN graph into common build pipeline --- diskann-disk/src/build/builder/build.rs | 34 +++- diskann-disk/src/build/builder/build/pipnn.rs | 75 +++++++ .../src/build/builder/build/pipnn/tests.rs | 190 ++++++++++++++++++ diskann-providers/src/storage/bin.rs | 55 +++++ diskann-providers/src/storage/mod.rs | 1 + diskann-providers/src/utils/rayon_util.rs | 5 + 6 files changed, 358 insertions(+), 2 deletions(-) create mode 100644 diskann-disk/src/build/builder/build/pipnn.rs create mode 100644 diskann-disk/src/build/builder/build/pipnn/tests.rs diff --git a/diskann-disk/src/build/builder/build.rs b/diskann-disk/src/build/builder/build.rs index 432c577fc..054244c12 100644 --- a/diskann-disk/src/build/builder/build.rs +++ b/diskann-disk/src/build/builder/build.rs @@ -30,6 +30,9 @@ use diskann_providers::{ use tokio::task::JoinSet; use tracing::{debug, info}; +#[cfg(feature = "pipnn")] +mod pipnn; + use crate::{ build::builder::{ core::{determine_build_strategy, IndexBuildStrategy, MergedVamanaIndexBuilder}, @@ -72,6 +75,28 @@ where index_configuration: IndexConfiguration, index_writer: DiskIndexWriter, ) -> ANNResult { + #[cfg(feature = "pipnn")] + let disk_build_param = { + let mut disk_build_param = disk_build_param; + if let Some(config) = disk_build_param.pipnn_config() { + config.validate()?; + let estimate = disk_build_param.use_vamana_if_pipnn_exceeds( + index_configuration.max_points, + index_configuration.dim, + std::mem::size_of::(), + index_configuration.num_threads, + )?; + let selected = disk_build_param.build_algorithm(); + info!( + estimated_peak_bytes = estimate, + memory_limit_bytes = disk_build_param.build_memory_limit().in_bytes(), + algorithm = %selected, + "Selected graph build algorithm" + ); + } + disk_build_param + }; + let pq_storage = PQStorage::new( &(index_writer.get_index_path_prefix() + "_pq_pivots.bin"), &(index_writer.get_index_path_prefix() + "_pq_compressed.bin"), @@ -122,7 +147,7 @@ where self.generate_compressed_data(pool.as_ref())?; logger.log_checkpoint(DiskIndexBuildCheckpoint::PqConstruction); - self.build_inmem_index(pool.as_ref()).await?; + self.build_graph(pool.as_ref()).await?; logger.log_checkpoint(DiskIndexBuildCheckpoint::InmemIndexBuild); // Use physical file to pass the memory index to the disk writer @@ -171,7 +196,12 @@ where ) } - async fn build_inmem_index(&mut self, pool: RayonThreadPoolRef<'_>) -> ANNResult<()> { + async fn build_graph(&mut self, pool: RayonThreadPoolRef<'_>) -> ANNResult<()> { + #[cfg(feature = "pipnn")] + if let Some(config) = self.disk_build_param.pipnn_config() { + return pipnn::build_graph(self, pool, config); + } + match determine_build_strategy::( &self.index_configuration, self.disk_build_param.build_memory_limit().in_bytes() as f64, diff --git a/diskann-disk/src/build/builder/build/pipnn.rs b/diskann-disk/src/build/builder/build/pipnn.rs new file mode 100644 index 000000000..a224adc02 --- /dev/null +++ b/diskann-disk/src/build/builder/build/pipnn.rs @@ -0,0 +1,75 @@ +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT license. + */ + +//! PiPNN graph adapter for the common disk-build pipeline. + +use diskann::{utils::VectorRepr, ANNError, ANNResult}; +use diskann_pipnn::{PiPNNBuildContext, PiPNNConfig}; +use diskann_providers::{ + storage::{save_adjacency_graph, StorageReadProvider, StorageWriteProvider}, + utils::{find_medoid_with_sampling, RayonThreadPoolRef, MAX_MEDOID_SAMPLE_SIZE}, +}; +use diskann_utils::io::{read_bin, Metadata}; + +use super::{u32_try_from, DiskIndexBuilder}; +use crate::data_model::GraphDataType; + +pub(super) fn build_graph( + builder: &DiskIndexBuilder<'_, Data, StorageProvider>, + pool: RayonThreadPoolRef<'_>, + config: PiPNNConfig, +) -> ANNResult<()> +where + Data: GraphDataType, + Data::VectorDataType: VectorRepr, + StorageProvider: StorageReadProvider + StorageWriteProvider, +{ + let data_path = builder.index_writer.get_dataset_file(); + let (points, dimensions) = + Metadata::read(&mut builder.storage_provider.open_reader(&data_path)?)?.into_dims(); + if dimensions != builder.index_configuration.dim { + return Err(ANNError::log_dimension_mismatch_error(format!( + "configured dimension {} does not match dataset dimension {dimensions}", + builder.index_configuration.dim + ))); + } + if points != builder.index_configuration.max_points { + return Err(ANNError::log_index_error(format!( + "configured point count {} does not match dataset point count {points}", + builder.index_configuration.max_points + ))); + } + + let data = + read_bin::(&mut builder.storage_provider.open_reader(&data_path)?)?; + let context = PiPNNBuildContext::new( + config, + &builder.index_configuration.config, + builder.index_configuration.dist_metric, + pool.as_rayon(), + )?; + let adjacency = diskann_pipnn::build_graph(data.as_view(), &context)?; + + let mut rng = diskann_providers::utils::create_rnd_from_optional_seed( + builder.index_configuration.random_seed, + ); + let (_, start_id) = find_medoid_with_sampling::( + &data_path, + builder.storage_provider, + MAX_MEDOID_SAMPLE_SIZE, + &mut rng, + )?; + save_adjacency_graph( + &adjacency, + u32_try_from(builder.index_configuration.config.pruned_degree().get())?, + builder.storage_provider, + u32_try_from(start_id)?, + &builder.index_writer.get_mem_index_file(), + )?; + Ok(()) +} + +#[cfg(test)] +mod tests; diff --git a/diskann-disk/src/build/builder/build/pipnn/tests.rs b/diskann-disk/src/build/builder/build/pipnn/tests.rs new file mode 100644 index 000000000..db2662171 --- /dev/null +++ b/diskann-disk/src/build/builder/build/pipnn/tests.rs @@ -0,0 +1,190 @@ +/* + * Copyright (c) Microsoft Corporation. + * Licensed under the MIT license. + */ + +use diskann::{graph::config, utils::ONE}; +use diskann_providers::utils::create_thread_pool; +use diskann_providers::{ + model::IndexConfiguration, + storage::{ + get_disk_index_file, StorageReadProvider, StorageWriteProvider, VirtualStorageProvider, + }, +}; +use diskann_utils::{io::write_bin, views::MatrixView}; +use diskann_vector::distance::Metric; +use vfs::MemoryFS; + +use crate::{ + build::{ + builder::build::DiskIndexBuilder, + configuration::{MemoryBudget, NumPQChunks, PiPNNParameters}, + }, + data_model::AdHoc, + storage::DiskIndexWriter, + DiskIndexBuildParameters, +}; + +fn pipnn() -> PiPNNParameters { + PiPNNParameters { + c_max: 512, + c_min: 64, + p_samp: 0.01, + fanout: vec![10, 3], + k: 2, + replicas: 1, + } +} + +fn write_data(storage: &VirtualStorageProvider, points: usize, dimensions: usize) { + let data: Vec = (0..points * dimensions) + .map(|index| ((index * 17) % 251) as f32) + .collect(); + write_bin( + MatrixView::try_from(data.as_slice(), points, dimensions).unwrap(), + &mut storage.create_for_write("/data.fbin").unwrap(), + ) + .unwrap(); +} + +fn graph_config(degree: usize, alpha: f32) -> diskann::graph::Config { + config::Builder::new_with( + degree, + config::MaxDegree::default_slack(), + 50, + Metric::L2.into(), + |builder| { + builder.alpha(alpha); + }, + ) + .build() + .unwrap() +} + +fn builder<'a>( + storage: &'a VirtualStorageProvider, + points: usize, + dimensions: usize, + budget_gib: f64, + alpha: f32, + parameters: PiPNNParameters, +) -> DiskIndexBuilder<'a, AdHoc, VirtualStorageProvider> { + let params = DiskIndexBuildParameters::new_pipnn( + MemoryBudget::try_from_gb(budget_gib).unwrap(), + NumPQChunks::new_with(dimensions, dimensions).unwrap(), + parameters, + ); + let config = IndexConfiguration::new( + Metric::L2, + dimensions, + points, + ONE, + 1, + graph_config(32, alpha), + ) + .with_pseudo_rng_from_seed(42); + let writer = DiskIndexWriter::new("/data.fbin".into(), "/index".into(), None, 4096).unwrap(); + DiskIndexBuilder::new(storage, params, config, writer).unwrap() +} + +#[test] +fn pipnn_disk_build_rejects_configuration_dataset_mismatch() { + let storage = VirtualStorageProvider::new_memory(); + write_data(&storage, 2, 8); + let params = DiskIndexBuildParameters::new_pipnn( + MemoryBudget::try_from_gb(10_000.0).unwrap(), + NumPQChunks::new_with(4, 4).unwrap(), + PiPNNParameters::default(), + ); + let config = IndexConfiguration::new(Metric::L2, 4, 3, ONE, 1, graph_config(4, 1.2)); + let writer = DiskIndexWriter::new("/data.fbin".into(), "/index".into(), None, 4096).unwrap(); + let mut builder = + DiskIndexBuilder::, _>::new(&storage, params, config, writer).unwrap(); + + let error = builder.build().unwrap_err(); + assert!(format!("{error:?}").contains("configured dimension 4")); + assert!(storage.exists("/index_pq_compressed.bin")); +} + +#[test] +fn pipnn_disk_build_uses_common_pipeline() { + let storage = VirtualStorageProvider::new_memory(); + let (points, dimensions) = (256, 8); + write_data(&storage, points, dimensions); + let mut builder = builder(&storage, points, dimensions, 1.0, 1.2, pipnn()); + + builder.build().unwrap(); + + assert!(storage.exists(&get_disk_index_file("/index"))); + assert!(storage.exists("/index_pq_compressed.bin")); +} + +#[test] +fn pipnn_graph_adapter_writes_real_point_header() { + let storage = VirtualStorageProvider::new_memory(); + let (points, dimensions) = (256, 8); + write_data(&storage, points, dimensions); + let parameters = pipnn(); + let builder = builder(&storage, points, dimensions, 1.0, 1.2, parameters.clone()); + let pool = create_thread_pool(1).unwrap(); + + super::build_graph(&builder, pool.as_ref(), (¶meters).into()).unwrap(); + + let mut header = [0_u8; 24]; + std::io::Read::read_exact( + &mut storage + .open_reader(&builder.index_writer.get_mem_index_file()) + .unwrap(), + &mut header, + ) + .unwrap(); + assert_eq!(u32::from_le_bytes(header[8..12].try_into().unwrap()), 32); + assert!(u32::from_le_bytes(header[12..16].try_into().unwrap()) < points as u32); + assert_eq!(u64::from_le_bytes(header[16..24].try_into().unwrap()), 0); +} + +#[test] +fn pipnn_disk_build_falls_back_to_complete_vamana_pipeline() { + let storage = VirtualStorageProvider::new_memory(); + let (points, dimensions) = (256, 8); + write_data(&storage, points, dimensions); + let mut builder = builder(&storage, points, dimensions, 0.0001, 1.3, pipnn()); + + assert!(matches!( + builder.disk_build_param.build_algorithm(), + crate::BuildAlgorithm::Vamana + )); + assert_eq!( + builder.disk_build_param.build_quantization(), + &crate::QuantizationType::FP + ); + assert_eq!(builder.index_configuration.config.pruned_degree().get(), 32); + assert_eq!(builder.index_configuration.config.l_build().get(), 50); + assert_eq!(builder.index_configuration.config.alpha(), 1.3); + + builder.build().unwrap(); + assert!(storage.exists(&get_disk_index_file("/index"))); +} + +#[test] +fn pipnn_disk_build_rejects_invalid_config_before_fallback() { + let storage = VirtualStorageProvider::new_memory(); + let invalid = PiPNNParameters { + c_max: 0, + ..PiPNNParameters::default() + }; + let params = DiskIndexBuildParameters::new_pipnn( + MemoryBudget::try_from_gb(0.0001).unwrap(), + NumPQChunks::new_with(1, 1).unwrap(), + invalid, + ); + let config = IndexConfiguration::new(Metric::L2, 1, 1, ONE, 1, graph_config(4, 1.2)); + let writer = DiskIndexWriter::new("/data.fbin".into(), "/index".into(), None, 4096).unwrap(); + + let error = match DiskIndexBuilder::, _>::new(&storage, params, config, writer) { + Ok(_) => panic!("invalid PiPNN config must not silently fall back to Vamana"), + Err(error) => error, + }; + + assert!(format!("{error:?}").contains("c_max must be greater than zero")); +} diff --git a/diskann-providers/src/storage/bin.rs b/diskann-providers/src/storage/bin.rs index 9b607e715..62992d82e 100644 --- a/diskann-providers/src/storage/bin.rs +++ b/diskann-providers/src/storage/bin.rs @@ -9,6 +9,7 @@ use super::{StorageReadProvider, StorageWriteProvider}; use byteorder::{LittleEndian, ReadBytesExt}; use diskann::{ ANNError, ANNResult, + graph::AdjacencyList, utils::{IntoUsize, VectorRepr}, }; use diskann_utils::io::Metadata; @@ -378,3 +379,57 @@ where out.flush()?; Ok(index_size.into_usize()) } + +/// Save real-point adjacency lists in the canonical graph layout. +pub fn save_adjacency_graph

( + adjacency: &[AdjacencyList], + max_degree: u32, + provider: &P, + start_point: u32, + path: &str, +) -> ANNResult +where + P: StorageWriteProvider, +{ + save_graph( + &AdjacencyGraph { + adjacency, + max_degree, + }, + provider, + start_point, + path, + ) +} + +struct AdjacencyGraph<'a> { + adjacency: &'a [AdjacencyList], + max_degree: u32, +} + +impl GetAdjacencyList for AdjacencyGraph<'_> { + type Element = u32; + type Item<'a> + = &'a [u32] + where + Self: 'a; + + fn get_adjacency_list(&self, index: usize) -> ANNResult> { + self.adjacency + .get(index) + .map(|row| &**row) + .ok_or_else(|| ANNError::log_index_error(format_args!("missing graph row {index}"))) + } + + fn total(&self) -> usize { + self.adjacency.len() + } + + fn additional_points(&self) -> u64 { + 0 + } + + fn max_degree(&self) -> Option { + Some(self.max_degree) + } +} diff --git a/diskann-providers/src/storage/mod.rs b/diskann-providers/src/storage/mod.rs index 1233b11f6..f2add5ff2 100644 --- a/diskann-providers/src/storage/mod.rs +++ b/diskann-providers/src/storage/mod.rs @@ -17,6 +17,7 @@ mod api; pub use api::{AsyncIndexMetadata, AsyncQuantLoadContext, DiskGraphOnly, LoadWith, SaveWith}; pub(crate) mod bin; +pub use bin::save_adjacency_graph; pub(crate) mod file_storage_provider; // Use VirtualStorageProvider in tests to avoid filesystem side-effects diff --git a/diskann-providers/src/utils/rayon_util.rs b/diskann-providers/src/utils/rayon_util.rs index ac5c2101d..cad31bba2 100644 --- a/diskann-providers/src/utils/rayon_util.rs +++ b/diskann-providers/src/utils/rayon_util.rs @@ -74,6 +74,11 @@ impl<'a> RayonThreadPoolRef<'a> { { self.0.install(op) } + + /// Borrow the underlying pool for APIs that retain a caller-owned pool. + pub fn as_rayon(self) -> &'a rayon::ThreadPool { + self.0 + } } // Allow use of disallowed methods within this trait to provide custom From 8805b8f26aacd40fdcf6159dec68474478ba3ebc Mon Sep 17 00:00:00 2001 From: Weiyao Luo <9347182+SeliMeli@users.noreply.github.com> Date: Wed, 29 Jul 2026 02:39:17 +0000 Subject: [PATCH 3/6] test(disk): assert sharded Vamana fallback --- .../src/build/builder/build/pipnn/tests.rs | 15 +++++++++++++-- 1 file changed, 13 insertions(+), 2 deletions(-) diff --git a/diskann-disk/src/build/builder/build/pipnn/tests.rs b/diskann-disk/src/build/builder/build/pipnn/tests.rs index db2662171..a39c1f2da 100644 --- a/diskann-disk/src/build/builder/build/pipnn/tests.rs +++ b/diskann-disk/src/build/builder/build/pipnn/tests.rs @@ -17,7 +17,10 @@ use vfs::MemoryFS; use crate::{ build::{ - builder::build::DiskIndexBuilder, + builder::{ + build::DiskIndexBuilder, + core::{determine_build_strategy, IndexBuildStrategy}, + }, configuration::{MemoryBudget, NumPQChunks, PiPNNParameters}, }, data_model::AdHoc, @@ -148,7 +151,7 @@ fn pipnn_disk_build_falls_back_to_complete_vamana_pipeline() { let storage = VirtualStorageProvider::new_memory(); let (points, dimensions) = (256, 8); write_data(&storage, points, dimensions); - let mut builder = builder(&storage, points, dimensions, 0.0001, 1.3, pipnn()); + let mut builder = builder(&storage, points, dimensions, 0.000001, 1.3, pipnn()); assert!(matches!( builder.disk_build_param.build_algorithm(), @@ -161,6 +164,14 @@ fn pipnn_disk_build_falls_back_to_complete_vamana_pipeline() { assert_eq!(builder.index_configuration.config.pruned_degree().get(), 32); assert_eq!(builder.index_configuration.config.l_build().get(), 50); assert_eq!(builder.index_configuration.config.alpha(), 1.3); + assert!(matches!( + determine_build_strategy::>( + &builder.index_configuration, + builder.disk_build_param.build_memory_limit().in_bytes() as f64, + builder.disk_build_param.build_quantization(), + ), + IndexBuildStrategy::Merged + )); builder.build().unwrap(); assert!(storage.exists(&get_disk_index_file("/index"))); From 6960ffd36bf6bc065663355f3f826418b397d63b Mon Sep 17 00:00:00 2001 From: Weiyao Luo <9347182+SeliMeli@users.noreply.github.com> Date: Wed, 29 Jul 2026 13:42:09 +0000 Subject: [PATCH 4/6] fix(disk): honor explicit PiPNN build selection --- Cargo.lock | 1 - diskann-disk/src/build/builder/build.rs | 23 +-- .../src/build/builder/build/pipnn/tests.rs | 22 +-- .../disk_index_build_parameter.rs | 157 +----------------- 4 files changed, 13 insertions(+), 190 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index ec189c7ca..c065d7e7f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -475,7 +475,6 @@ dependencies = [ "diskann-disk", "diskann-inmem", "diskann-label-filter", - "diskann-pipnn", "diskann-providers", "diskann-quantization", "diskann-tools", diff --git a/diskann-disk/src/build/builder/build.rs b/diskann-disk/src/build/builder/build.rs index 054244c12..20a01949d 100644 --- a/diskann-disk/src/build/builder/build.rs +++ b/diskann-disk/src/build/builder/build.rs @@ -76,26 +76,9 @@ where index_writer: DiskIndexWriter, ) -> ANNResult { #[cfg(feature = "pipnn")] - let disk_build_param = { - let mut disk_build_param = disk_build_param; - if let Some(config) = disk_build_param.pipnn_config() { - config.validate()?; - let estimate = disk_build_param.use_vamana_if_pipnn_exceeds( - index_configuration.max_points, - index_configuration.dim, - std::mem::size_of::(), - index_configuration.num_threads, - )?; - let selected = disk_build_param.build_algorithm(); - info!( - estimated_peak_bytes = estimate, - memory_limit_bytes = disk_build_param.build_memory_limit().in_bytes(), - algorithm = %selected, - "Selected graph build algorithm" - ); - } - disk_build_param - }; + if let Some(config) = disk_build_param.pipnn_config() { + config.validate()?; + } let pq_storage = PQStorage::new( &(index_writer.get_index_path_prefix() + "_pq_pivots.bin"), diff --git a/diskann-disk/src/build/builder/build/pipnn/tests.rs b/diskann-disk/src/build/builder/build/pipnn/tests.rs index a39c1f2da..a53a1a6e0 100644 --- a/diskann-disk/src/build/builder/build/pipnn/tests.rs +++ b/diskann-disk/src/build/builder/build/pipnn/tests.rs @@ -17,10 +17,7 @@ use vfs::MemoryFS; use crate::{ build::{ - builder::{ - build::DiskIndexBuilder, - core::{determine_build_strategy, IndexBuildStrategy}, - }, + builder::build::DiskIndexBuilder, configuration::{MemoryBudget, NumPQChunks, PiPNNParameters}, }, data_model::AdHoc, @@ -147,7 +144,7 @@ fn pipnn_graph_adapter_writes_real_point_header() { } #[test] -fn pipnn_disk_build_falls_back_to_complete_vamana_pipeline() { +fn explicit_pipnn_selection_is_not_replaced_by_memory_budget() { let storage = VirtualStorageProvider::new_memory(); let (points, dimensions) = (256, 8); write_data(&storage, points, dimensions); @@ -155,7 +152,7 @@ fn pipnn_disk_build_falls_back_to_complete_vamana_pipeline() { assert!(matches!( builder.disk_build_param.build_algorithm(), - crate::BuildAlgorithm::Vamana + crate::BuildAlgorithm::PiPNN(_) )); assert_eq!( builder.disk_build_param.build_quantization(), @@ -164,21 +161,12 @@ fn pipnn_disk_build_falls_back_to_complete_vamana_pipeline() { assert_eq!(builder.index_configuration.config.pruned_degree().get(), 32); assert_eq!(builder.index_configuration.config.l_build().get(), 50); assert_eq!(builder.index_configuration.config.alpha(), 1.3); - assert!(matches!( - determine_build_strategy::>( - &builder.index_configuration, - builder.disk_build_param.build_memory_limit().in_bytes() as f64, - builder.disk_build_param.build_quantization(), - ), - IndexBuildStrategy::Merged - )); - builder.build().unwrap(); assert!(storage.exists(&get_disk_index_file("/index"))); } #[test] -fn pipnn_disk_build_rejects_invalid_config_before_fallback() { +fn pipnn_disk_build_rejects_invalid_config() { let storage = VirtualStorageProvider::new_memory(); let invalid = PiPNNParameters { c_max: 0, @@ -193,7 +181,7 @@ fn pipnn_disk_build_rejects_invalid_config_before_fallback() { let writer = DiskIndexWriter::new("/data.fbin".into(), "/index".into(), None, 4096).unwrap(); let error = match DiskIndexBuilder::, _>::new(&storage, params, config, writer) { - Ok(_) => panic!("invalid PiPNN config must not silently fall back to Vamana"), + Ok(_) => panic!("invalid PiPNN config must be rejected"), Err(error) => error, }; diff --git a/diskann-disk/src/build/configuration/disk_index_build_parameter.rs b/diskann-disk/src/build/configuration/disk_index_build_parameter.rs index 939c44742..f09af881d 100644 --- a/diskann-disk/src/build/configuration/disk_index_build_parameter.rs +++ b/diskann-disk/src/build/configuration/disk_index_build_parameter.rs @@ -8,8 +8,6 @@ use std::num::NonZeroUsize; use diskann::ANNError; -#[cfg(feature = "pipnn")] -use diskann::ANNResult; use thiserror::Error; #[cfg(feature = "pipnn")] @@ -111,8 +109,8 @@ impl NumPQChunks { /// Parameters specific for disk index construction. #[derive(Clone, PartialEq, Debug)] pub struct DiskIndexBuildParameters { - /// Limit on graph-construction memory. PiPNN falls back to Vamana when its - /// estimated one-shot peak exceeds this value. + /// Memory budget for disk-index pipeline stages that support bounded work. + /// Explicit one-shot PiPNN selection is never silently replaced. build_memory_limit: MemoryBudget, /// Number of PQ chunks stored in-memory for search and to be generated during build. @@ -146,8 +144,8 @@ impl DiskIndexBuildParameters { /// Create parameters for one-shot PiPNN graph construction. /// - /// PiPNN uses the common search-PQ and disk-layout pipeline. The memory - /// budget selects Vamana when the estimated PiPNN peak does not fit. + /// PiPNN uses the common search-PQ and disk-layout pipeline. Its one-shot + /// graph build is not governed by the pipeline memory budget. #[cfg(feature = "pipnn")] pub fn new_pipnn( build_memory_limit: MemoryBudget, @@ -204,116 +202,6 @@ impl DiskIndexBuildParameters { BuildAlgorithm::Vamana => None, } } - - #[cfg(feature = "pipnn")] - pub(crate) fn use_vamana_if_pipnn_exceeds( - &mut self, - npoints: usize, - dimensions: usize, - element_size: usize, - num_threads: usize, - ) -> ANNResult { - let parameters = match &self.build_algorithm { - BuildAlgorithm::PiPNN(parameters) => parameters, - BuildAlgorithm::Vamana => { - return Err(ANNError::log_index_error( - "memory selection requires PiPNN parameters", - )); - } - }; - let estimate = - estimate_pipnn_peak_memory(parameters, npoints, dimensions, element_size, num_threads); - if estimate.is_none_or(|bytes| bytes > self.build_memory_limit.in_bytes()) { - self.build_algorithm = BuildAlgorithm::Vamana; - } - Ok(estimate.unwrap_or(usize::MAX)) - } -} - -#[cfg(feature = "pipnn")] -fn estimate_pipnn_peak_memory( - config: &PiPNNParameters, - npoints: usize, - dimensions: usize, - element_size: usize, - num_threads: usize, -) -> Option { - const PROPORTIONAL_HEADROOM_PERCENT: u128 = 108; - const PROCESS_HEADROOM: u128 = 16 * 1024 * 1024; - const PER_WORKER_HEADROOM: u128 = 28 * 1024 * 1024; - - let workers = if num_threads == 0 { - std::thread::available_parallelism().map_or(1, NonZeroUsize::get) - } else { - num_threads - } as u128; - let copies = config - .fanout - .iter() - .try_fold(config.replicas as u128, |copies, &fanout| { - copies.checked_mul(fanout as u128) - })?; - let dataset = (dimensions as u128).checked_mul(element_size as u128)?; - let leaf_ids = copies.checked_mul(size_of::() as u128)?; - let offered = copies.checked_mul(config.k as u128)?.checked_mul(2)?; - let candidate_capacity = offered.checked_next_power_of_two()?; - let candidate_storage = candidate_capacity.checked_mul(size_of::() as u128)?; - let candidate_metadata = - size_of::>>() as u128; - - let partition_per_point = dataset - .checked_add(16)? - .checked_add(leaf_ids.checked_mul(2)?)?; - let leaf_per_point = dataset - .checked_add(candidate_metadata)? - .checked_add(candidate_storage)? - .checked_add(leaf_ids)?; - let finalization_per_point = dataset - .checked_add(candidate_metadata)? - .checked_add(size_of::>() as u128)? - .checked_add(candidate_storage)?; - - let leaf_size = config.c_max.min(npoints).max(1) as u128; - let leaf_scratch = leaf_size - .checked_mul(dimensions as u128)? - .checked_mul(size_of::() as u128)? - .checked_add(leaf_size.checked_mul(leaf_size)?.checked_mul(5)?)? - .checked_add(leaf_size.checked_mul(config.k as u128)?.checked_mul(24)?)? - .checked_add(leaf_size.checked_mul(20)?)? - .checked_add(68)?; - let partition_rows = (npoints as u128).min(1_024); - let partition_leaders = (npoints as u128).min(1_000); - let partition_scratch = partition_leaders - .checked_mul(dimensions as u128)? - .checked_mul(size_of::() as u128)? - .checked_add( - partition_rows - .checked_mul(dimensions as u128)? - .checked_mul(size_of::() as u128)?, - )? - .checked_add( - partition_rows - .checked_mul(partition_leaders)? - .checked_mul(size_of::() as u128)?, - )?; - - let points = npoints as u128; - let structural = points - .checked_mul(partition_per_point)? - .checked_add(workers.checked_mul(partition_scratch)?)? - .max( - points - .checked_mul(leaf_per_point)? - .checked_add(workers.checked_mul(leaf_scratch)?)?, - ) - .max(points.checked_mul(finalization_per_point)?); - let proportional = structural - .checked_mul(PROPORTIONAL_HEADROOM_PERCENT)? - .div_ceil(100); - let allocator_floor = structural - .checked_add(PROCESS_HEADROOM)? - .checked_add(workers.checked_mul(PER_WORKER_HEADROOM)?)?; - usize::try_from(proportional.max(allocator_floor)).ok() } #[cfg(test)] @@ -407,6 +295,7 @@ mod dataset_test { let params = DiskIndexBuildParameters::new_pipnn(budget, pq, parameters); assert_eq!(params.pipnn_config(), Some(config)); + assert_eq!(params.build_memory_limit(), budget); assert_eq!(params.search_pq_chunks(), pq); assert_eq!( params.data_compression_chunk_vector_count(), @@ -414,40 +303,4 @@ mod dataset_test { ); assert!(matches!(params.build_algorithm(), BuildAlgorithm::PiPNN(_))); } - - #[cfg(feature = "pipnn")] - #[test] - fn pipnn_memory_estimate_covers_measured_bigann10m_peak() { - const MEASURED_PEAK: usize = 8_656_004 * 1024; - let parameters = PiPNNParameters { - c_max: 512, - c_min: 64, - p_samp: 0.01, - fanout: vec![10, 3], - k: 2, - replicas: 1, - }; - - let estimate = estimate_pipnn_peak_memory(¶meters, 10_000_000, 128, 2, 16).unwrap(); - - assert!(estimate >= MEASURED_PEAK); - assert!(estimate <= MEASURED_PEAK * 120 / 100); - } - - #[cfg(feature = "pipnn")] - #[test] - fn pipnn_falls_back_to_vamana_above_budget() { - let budget = MemoryBudget::try_from_gb(1.0).unwrap(); - let pq = NumPQChunks::new_with(1, 128).unwrap(); - let mut params = - DiskIndexBuildParameters::new_pipnn(budget, pq, PiPNNParameters::default()); - - let estimate = params - .use_vamana_if_pipnn_exceeds(10_000_000, 128, 2, 16) - .unwrap(); - - assert!(estimate > budget.in_bytes()); - assert!(matches!(params.build_algorithm(), BuildAlgorithm::Vamana)); - assert_eq!(params.build_quantization(), &QuantizationType::FP); - } } From 5d040be0c36490e98e6b66e15768674a7360cac2 Mon Sep 17 00:00:00 2001 From: Weiyao Luo <9347182+SeliMeli@users.noreply.github.com> Date: Fri, 31 Jul 2026 03:19:50 +0000 Subject: [PATCH 5/6] docs(disk): map PiPNN adapter lifecycle --- diskann-disk/src/build/builder/build/pipnn.rs | 31 ++++++++++++++++++- .../src/build/builder/build/pipnn/tests.rs | 13 ++++++++ 2 files changed, 43 insertions(+), 1 deletion(-) diff --git a/diskann-disk/src/build/builder/build/pipnn.rs b/diskann-disk/src/build/builder/build/pipnn.rs index a224adc02..09ad6a760 100644 --- a/diskann-disk/src/build/builder/build/pipnn.rs +++ b/diskann-disk/src/build/builder/build/pipnn.rs @@ -3,7 +3,27 @@ * Licensed under the MIT license. */ -//! PiPNN graph adapter for the common disk-build pipeline. +//! Adapter from provider-independent PiPNN adjacency to the common disk index format. +//! +//! The core crate deliberately knows nothing about dataset files, medoids, graph +//! headers, or serialization. This adapter owns that boundary: +//! +//! 1. verify on-disk dataset metadata against the requested index configuration; +//! 2. load the contiguous matrix required by batch construction; +//! 3. run PiPNN in the caller-provided Rayon pool; +//! 4. compute the production start node with the existing medoid policy; and +//! 5. serialize adjacency with the same header/layout used by Vamana. +//! +//! ```text +//! dataset file ──> metadata check ──> MatrixView ──> diskann-pipnn ──> adjacency +//! │ │ +//! └──────────────────> sampled medoid ────────────────────────────┤ +//! v +//! canonical graph writer +//! ``` +//! +//! There is no PiPNN-specific disk graph format. Keeping serialization here means +//! search and loading cannot distinguish which builder produced the graph. use diskann::{utils::VectorRepr, ANNError, ANNResult}; use diskann_pipnn::{PiPNNBuildContext, PiPNNConfig}; @@ -16,6 +36,7 @@ use diskann_utils::io::{read_bin, Metadata}; use super::{u32_try_from, DiskIndexBuilder}; use crate::data_model::GraphDataType; +/// Build PiPNN adjacency and persist it through the canonical disk graph writer. pub(super) fn build_graph( builder: &DiskIndexBuilder<'_, Data, StorageProvider>, pool: RayonThreadPoolRef<'_>, @@ -27,6 +48,8 @@ where StorageProvider: StorageReadProvider + StorageWriteProvider, { let data_path = builder.index_writer.get_dataset_file(); + // Validate metadata before allocating/loading the full matrix. A mismatch + // here otherwise turns a configuration error into a later shape failure. let (points, dimensions) = Metadata::read(&mut builder.storage_provider.open_reader(&data_path)?)?.into_dims(); if dimensions != builder.index_configuration.dim { @@ -42,6 +65,9 @@ where ))); } + // PiPNN is a batch algorithm: materialize the matrix once, while all + // partition and leaf scratch stays inside the supplied pool and is released + // before the outer disk pipeline continues. let data = read_bin::(&mut builder.storage_provider.open_reader(&data_path)?)?; let context = PiPNNBuildContext::new( @@ -52,6 +78,9 @@ where )?; let adjacency = diskann_pipnn::build_graph(data.as_view(), &context)?; + // Start-node policy belongs to the persisted index, not the core graph + // constructor. Reuse the production sampled medoid implementation so the + // serialized header has the same semantics as a Vamana-built index. let mut rng = diskann_providers::utils::create_rnd_from_optional_seed( builder.index_configuration.random_seed, ); diff --git a/diskann-disk/src/build/builder/build/pipnn/tests.rs b/diskann-disk/src/build/builder/build/pipnn/tests.rs index a53a1a6e0..8edb57f84 100644 --- a/diskann-disk/src/build/builder/build/pipnn/tests.rs +++ b/diskann-disk/src/build/builder/build/pipnn/tests.rs @@ -106,6 +106,19 @@ fn pipnn_disk_build_rejects_configuration_dataset_mismatch() { assert!(storage.exists("/index_pq_compressed.bin")); } +#[test] +fn pipnn_graph_adapter_rejects_point_count_mismatch() { + let storage = VirtualStorageProvider::new_memory(); + write_data(&storage, 2, 8); + let parameters = pipnn(); + let builder = builder(&storage, 3, 8, 1.0, 1.2, parameters.clone()); + let pool = create_thread_pool(1).unwrap(); + + let error = super::build_graph(&builder, pool.as_ref(), (¶meters).into()).unwrap_err(); + assert!(format!("{error:?}").contains("configured point count 3")); + assert!(!storage.exists(&builder.index_writer.get_mem_index_file())); +} + #[test] fn pipnn_disk_build_uses_common_pipeline() { let storage = VirtualStorageProvider::new_memory(); From f3b78cd8d4a60b779025a894a326725c0e135bf8 Mon Sep 17 00:00:00 2001 From: Weiyao Luo <9347182+SeliMeli@users.noreply.github.com> Date: Wed, 5 Aug 2026 07:31:18 +0000 Subject: [PATCH 6/6] refactor(disk): consume in-crate pipnn module Forward the disk feature to diskann/pipnn instead of depending on a separate crate. Consolidate duplicate pipeline builds and name adapter tests by the behavior they protect. --- Cargo.lock | 1 - diskann-disk/Cargo.toml | 3 +-- diskann-disk/src/build/builder/build/pipnn.rs | 6 ++--- .../src/build/builder/build/pipnn/tests.rs | 24 +++++-------------- .../build/configuration/build_algorithm.rs | 2 +- .../disk_index_build_parameter.rs | 21 +--------------- 6 files changed, 12 insertions(+), 45 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index c065d7e7f..c02628a4f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -580,7 +580,6 @@ dependencies = [ "criterion", "diskann", "diskann-linalg", - "diskann-pipnn", "diskann-providers", "diskann-quantization", "diskann-utils", diff --git a/diskann-disk/Cargo.toml b/diskann-disk/Cargo.toml index 0c53172e9..29aae9f4c 100644 --- a/diskann-disk/Cargo.toml +++ b/diskann-disk/Cargo.toml @@ -43,7 +43,6 @@ vfs = { workspace = true } # Optional dependencies opentelemetry = { workspace = true, optional = true } -diskann-pipnn = { workspace = true, optional = true } [target.'cfg(target_os = "linux")'.dependencies] io-uring = "0.6.4" @@ -84,7 +83,7 @@ proptest.workspace = true [features] default = [] perf_test = ["dep:opentelemetry"] -pipnn = ["dep:diskann-pipnn"] +pipnn = ["diskann/pipnn"] virtual_storage = ["diskann-providers/virtual_storage"] experimental_diversity_search = [ "diskann/experimental_diversity_search", diff --git a/diskann-disk/src/build/builder/build/pipnn.rs b/diskann-disk/src/build/builder/build/pipnn.rs index 09ad6a760..069bb19fc 100644 --- a/diskann-disk/src/build/builder/build/pipnn.rs +++ b/diskann-disk/src/build/builder/build/pipnn.rs @@ -15,7 +15,7 @@ //! 5. serialize adjacency with the same header/layout used by Vamana. //! //! ```text -//! dataset file ──> metadata check ──> MatrixView ──> diskann-pipnn ──> adjacency +//! dataset file ──> metadata check ──> MatrixView ──> diskann::graph::pipnn ──> adjacency //! │ │ //! └──────────────────> sampled medoid ────────────────────────────┤ //! v @@ -25,8 +25,8 @@ //! There is no PiPNN-specific disk graph format. Keeping serialization here means //! search and loading cannot distinguish which builder produced the graph. +use diskann::graph::pipnn::{PiPNNBuildContext, PiPNNConfig}; use diskann::{utils::VectorRepr, ANNError, ANNResult}; -use diskann_pipnn::{PiPNNBuildContext, PiPNNConfig}; use diskann_providers::{ storage::{save_adjacency_graph, StorageReadProvider, StorageWriteProvider}, utils::{find_medoid_with_sampling, RayonThreadPoolRef, MAX_MEDOID_SAMPLE_SIZE}, @@ -76,7 +76,7 @@ where builder.index_configuration.dist_metric, pool.as_rayon(), )?; - let adjacency = diskann_pipnn::build_graph(data.as_view(), &context)?; + let adjacency = diskann::graph::pipnn::build_graph(data.as_view(), &context)?; // Start-node policy belongs to the persisted index, not the core graph // constructor. Reuse the production sampled medoid implementation so the diff --git a/diskann-disk/src/build/builder/build/pipnn/tests.rs b/diskann-disk/src/build/builder/build/pipnn/tests.rs index 8edb57f84..28881f706 100644 --- a/diskann-disk/src/build/builder/build/pipnn/tests.rs +++ b/diskann-disk/src/build/builder/build/pipnn/tests.rs @@ -88,7 +88,7 @@ fn builder<'a>( } #[test] -fn pipnn_disk_build_rejects_configuration_dataset_mismatch() { +fn disk_build_rejects_dataset_shape_mismatch() { let storage = VirtualStorageProvider::new_memory(); write_data(&storage, 2, 8); let params = DiskIndexBuildParameters::new_pipnn( @@ -107,7 +107,7 @@ fn pipnn_disk_build_rejects_configuration_dataset_mismatch() { } #[test] -fn pipnn_graph_adapter_rejects_point_count_mismatch() { +fn graph_adapter_rejects_point_count_mismatch() { let storage = VirtualStorageProvider::new_memory(); write_data(&storage, 2, 8); let parameters = pipnn(); @@ -120,20 +120,7 @@ fn pipnn_graph_adapter_rejects_point_count_mismatch() { } #[test] -fn pipnn_disk_build_uses_common_pipeline() { - let storage = VirtualStorageProvider::new_memory(); - let (points, dimensions) = (256, 8); - write_data(&storage, points, dimensions); - let mut builder = builder(&storage, points, dimensions, 1.0, 1.2, pipnn()); - - builder.build().unwrap(); - - assert!(storage.exists(&get_disk_index_file("/index"))); - assert!(storage.exists("/index_pq_compressed.bin")); -} - -#[test] -fn pipnn_graph_adapter_writes_real_point_header() { +fn graph_adapter_writes_degree_medoid_and_frozen_count() { let storage = VirtualStorageProvider::new_memory(); let (points, dimensions) = (256, 8); write_data(&storage, points, dimensions); @@ -157,7 +144,7 @@ fn pipnn_graph_adapter_writes_real_point_header() { } #[test] -fn explicit_pipnn_selection_is_not_replaced_by_memory_budget() { +fn explicit_selection_ignores_the_vamana_memory_strategy() { let storage = VirtualStorageProvider::new_memory(); let (points, dimensions) = (256, 8); write_data(&storage, points, dimensions); @@ -176,10 +163,11 @@ fn explicit_pipnn_selection_is_not_replaced_by_memory_budget() { assert_eq!(builder.index_configuration.config.alpha(), 1.3); builder.build().unwrap(); assert!(storage.exists(&get_disk_index_file("/index"))); + assert!(storage.exists("/index_pq_compressed.bin")); } #[test] -fn pipnn_disk_build_rejects_invalid_config() { +fn builder_rejects_invalid_pipnn_config() { let storage = VirtualStorageProvider::new_memory(); let invalid = PiPNNParameters { c_max: 0, diff --git a/diskann-disk/src/build/configuration/build_algorithm.rs b/diskann-disk/src/build/configuration/build_algorithm.rs index f798138a1..109bf34c0 100644 --- a/diskann-disk/src/build/configuration/build_algorithm.rs +++ b/diskann-disk/src/build/configuration/build_algorithm.rs @@ -46,7 +46,7 @@ impl Default for PiPNNParameters { } #[cfg(feature = "pipnn")] -impl From<&PiPNNParameters> for diskann_pipnn::PiPNNConfig { +impl From<&PiPNNParameters> for diskann::graph::pipnn::PiPNNConfig { fn from(config: &PiPNNParameters) -> Self { Self { c_max: config.c_max, diff --git a/diskann-disk/src/build/configuration/disk_index_build_parameter.rs b/diskann-disk/src/build/configuration/disk_index_build_parameter.rs index f09af881d..19fef368d 100644 --- a/diskann-disk/src/build/configuration/disk_index_build_parameter.rs +++ b/diskann-disk/src/build/configuration/disk_index_build_parameter.rs @@ -196,7 +196,7 @@ impl DiskIndexBuildParameters { } #[cfg(feature = "pipnn")] - pub(crate) fn pipnn_config(&self) -> Option { + pub(crate) fn pipnn_config(&self) -> Option { match &self.build_algorithm { BuildAlgorithm::PiPNN(config) => Some(config.into()), BuildAlgorithm::Vamana => None, @@ -284,23 +284,4 @@ mod dataset_test { let chunks = NumPQChunks::new_with(64, 128).unwrap(); assert_eq!(chunks.get(), 64); } - - #[cfg(feature = "pipnn")] - #[test] - fn new_pipnn_uses_the_common_disk_pipeline_parameters() { - let pq = NumPQChunks::new_with(1, 128).unwrap(); - let parameters = PiPNNParameters::default(); - let config = diskann_pipnn::PiPNNConfig::from(¶meters); - let budget = MemoryBudget::try_from_gb(2.0).unwrap(); - let params = DiskIndexBuildParameters::new_pipnn(budget, pq, parameters); - - assert_eq!(params.pipnn_config(), Some(config)); - assert_eq!(params.build_memory_limit(), budget); - assert_eq!(params.search_pq_chunks(), pq); - assert_eq!( - params.data_compression_chunk_vector_count(), - DEFAULT_DATA_COMPRESSION_CHUNK_VECTOR_COUNT - ); - assert!(matches!(params.build_algorithm(), BuildAlgorithm::PiPNN(_))); - } }