diff --git a/diskann-benchmark/example/disk-index-determinant-diversity.json b/diskann-benchmark/example/disk-index-determinant-diversity.json index 2962c1d97..a8c236a99 100644 --- a/diskann-benchmark/example/disk-index-determinant-diversity.json +++ b/diskann-benchmark/example/disk-index-determinant-diversity.json @@ -27,14 +27,15 @@ "beam_width": 4, "recall_at": 10, "num_threads": 1, - "is_flat_search": false, - "distance": "squared_l2", - "vector_filters_file": null, - "post_processor": { - "type": "determinant-diversity", - "power": 2.0, - "eta": 1.0 - } + "search_mode": { + "is_flat_search": false, + "post_processor": { + "type": "determinant-diversity", + "power": 2.0, + "eta": 1.0 + } + }, + "distance": "squared_l2" } } } diff --git a/diskann-benchmark/example/disk-index-filter.json b/diskann-benchmark/example/disk-index-filter.json index a3f35ca91..cdb9842d7 100644 --- a/diskann-benchmark/example/disk-index-filter.json +++ b/diskann-benchmark/example/disk-index-filter.json @@ -27,9 +27,11 @@ "beam_width": 4, "recall_at": 10, "num_threads": 1, - "is_flat_search": false, - "distance": "squared_l2", - "vector_filters_file": "disk_index_10pts_idx_uint32_range_res_r_100000.bin" + "search_mode": { + "is_flat_search": false, + "vector_filters_file": "disk_index_10pts_idx_uint32_range_res_r_100000.bin" + }, + "distance": "squared_l2" } } }, @@ -57,9 +59,11 @@ "beam_width": 4, "recall_at": 10, "num_threads": 1, - "is_flat_search": true, - "distance": "squared_l2", - "vector_filters_file": "disk_index_10pts_idx_uint32_range_res_r_100000.bin" + "search_mode": { + "is_flat_search": true, + "vector_filters_file": "disk_index_10pts_idx_uint32_range_res_r_100000.bin" + }, + "distance": "squared_l2" } } } diff --git a/diskann-benchmark/example/disk-index.json b/diskann-benchmark/example/disk-index.json index 4d60fdb8e..c8a0c7b7d 100644 --- a/diskann-benchmark/example/disk-index.json +++ b/diskann-benchmark/example/disk-index.json @@ -27,9 +27,8 @@ "beam_width": 4, "recall_at": 10, "num_threads": 1, - "is_flat_search": false, - "distance": "squared_l2", - "vector_filters_file": null + "search_mode": { "is_flat_search": false }, + "distance": "squared_l2" } } }, @@ -48,9 +47,8 @@ "beam_width": 4, "recall_at": 10, "num_threads": 1, - "is_flat_search": true, - "distance": "squared_l2", - "vector_filters_file": null + "search_mode": { "is_flat_search": true }, + "distance": "squared_l2" } } } diff --git a/diskann-benchmark/perf_test_inputs/openai-100K-disk-index.json b/diskann-benchmark/perf_test_inputs/openai-100K-disk-index.json index 6b3e3b42d..241782899 100644 --- a/diskann-benchmark/perf_test_inputs/openai-100K-disk-index.json +++ b/diskann-benchmark/perf_test_inputs/openai-100K-disk-index.json @@ -29,9 +29,8 @@ "beam_width": 4, "recall_at": 100, "num_threads": 4, - "is_flat_search": false, - "distance": "squared_l2", - "vector_filters_file": null + "search_mode": { "is_flat_search": false }, + "distance": "squared_l2" } } } diff --git a/diskann-benchmark/perf_test_inputs/wikipedia-100K-disk-index.json b/diskann-benchmark/perf_test_inputs/wikipedia-100K-disk-index.json index 59c439017..aa00ab4ac 100644 --- a/diskann-benchmark/perf_test_inputs/wikipedia-100K-disk-index.json +++ b/diskann-benchmark/perf_test_inputs/wikipedia-100K-disk-index.json @@ -29,9 +29,8 @@ "beam_width": 4, "recall_at": 100, "num_threads": 4, - "is_flat_search": false, - "distance": "inner_product", - "vector_filters_file": null + "search_mode": { "is_flat_search": false }, + "distance": "inner_product" } } } diff --git a/diskann-benchmark/src/disk_index/search.rs b/diskann-benchmark/src/disk_index/search.rs index 8cd0f5e60..5ae45c2da 100644 --- a/diskann-benchmark/src/disk_index/search.rs +++ b/diskann-benchmark/src/disk_index/search.rs @@ -9,6 +9,7 @@ use std::{collections::HashSet, fmt, sync::atomic::AtomicBool, time::Instant}; use opentelemetry::{global, trace::Span, trace::Tracer}; use opentelemetry_sdk::trace::SdkTracerProvider; +use diskann::graph; use diskann::utils::VectorRepr; use diskann_benchmark_runner::{files::InputFile, utils::MicroSeconds}; use diskann_disk::{ @@ -36,7 +37,8 @@ use serde::{Deserialize, Serialize}; use crate::{ disk_index::json_spancollector::JsonSpanCollector, - inputs::disk::{DiskIndexLoad, DiskSearchPhase}, + inputs::disk::{DiskIndexLoad, DiskSearchMode, DiskSearchPhase}, + inputs::post_processor::TopkPostProcessor, utils::{datafiles, SimilarityMeasure}, }; @@ -158,6 +160,51 @@ impl DiskSearchResult { } } +/// Construct the disk [`SearchMode`] from the JSON-driven [`DiskSearchMode`] +/// config plus the per-query filter and post-processor supplied at search time. +fn build_search_mode<'a>( + mode: &'a DiskSearchMode, + vector_filter: Option<&'a HashSet>, + post_processor: Option<&TopkPostProcessor>, +) -> SearchMode<'a> { + let adaptive_l = mode.adaptive_l.as_ref().map(|adaptive_l| { + graph::search::AdaptiveL::new(adaptive_l.sample_count.into(), adaptive_l.scale_factor) + .expect("validated adaptive L must construct") + }); + + match ( + mode.is_flat_search, + vector_filter, + post_processor, + adaptive_l, + ) { + (true, None, _, _) => SearchMode::flat(), + (true, Some(vector_filter), _, _) => { + SearchMode::flat_filtered(move |vid: &u32| vector_filter.contains(vid)) + } + (false, None, Some(TopkPostProcessor::DeterminantDiversity(params)), _) => { + SearchMode::diverse_graph(*params) + } + (false, Some(vector_filter), Some(TopkPostProcessor::DeterminantDiversity(params)), _) => { + SearchMode::diverse_graph_filtered( + move |vid: &u32| vector_filter.contains(vid), + *params, + ) + } + (false, None, None, Some(adaptive_l)) => { + SearchMode::inline_filter(|_| true, Some(adaptive_l)) + } + (false, Some(vector_filter), None, Some(adaptive_l)) => SearchMode::inline_filter( + move |vid: &u32| vector_filter.contains(vid), + Some(adaptive_l), + ), + (false, None, None, None) => SearchMode::graph(), + (false, Some(vector_filter), None, None) => { + SearchMode::graph_filtered(move |vid: &u32| vector_filter.contains(vid)) + } + } +} + pub(super) fn search_disk_index( index_load: &DiskIndexLoad, search_params: &DiskSearchPhase, @@ -185,21 +232,27 @@ where let num_queries = queries.nrows(); // Load the vector filters - let vector_filters = match &search_params.vector_filters_file { + let vector_filters = match &search_params.search_mode.vector_filters_file { Some(vector_filters_file) => { let vector_filters_file = vector_filters_file.to_string_lossy().to_string(); - search_index_utils::load_vector_filters(storage_provider, &vector_filters_file)? + Some(search_index_utils::load_vector_filters( + storage_provider, + &vector_filters_file, + )?) } - None => vec![HashSet::::new(); num_queries], + None => None, }; - if vector_filters.len() != num_queries { + if vector_filters + .as_ref() + .is_some_and(|filters| filters.len() != num_queries) + { anyhow::bail!("Mismatch in query and vector filter sizes"); } // Prepare ground truth context let gt_context = prepare_ground_truth_context( - search_params.vector_filters_file.is_some(), + search_params.search_mode.vector_filters_file.is_some(), &search_params.groundtruth, search_params.recall_at, storage_provider, @@ -259,7 +312,7 @@ where let zipped = queries .par_row_iter() - .zip(vector_filters.par_iter()) + .enumerate() .zip(result_ids.par_chunks_mut(search_params.recall_at as usize)) .zip(result_dists.par_chunks_mut(search_params.recall_at as usize)) .zip(statistics_vec.par_iter_mut()) @@ -267,15 +320,17 @@ where zipped.for_each_in_pool( pool.as_ref(), - |(((((q, vf), id_chunk), dist_chunk), stats), rc)| { + |(((((query_index, q), id_chunk), dist_chunk), stats), rc)| { // Construct the SearchMode from the JSON-driven // `adaptive_l` is now encapsulated in `DiskSearchMode`, so the // benchmark only supplies the per-query filter and post-processor. - let has_filter = search_params.vector_filters_file.is_some(); - let mode: SearchMode<'_> = search_params.search_mode.search_mode( - has_filter, - vf, - search_params.post_processor.as_ref(), + let vector_filter = vector_filters + .as_ref() + .and_then(|filters| filters.get(query_index)); + let mode: SearchMode<'_> = build_search_mode( + &search_params.search_mode, + vector_filter, + search_params.search_mode.post_processor.as_ref(), ); match searcher.search( @@ -351,7 +406,7 @@ where recall_at: search_params.recall_at, is_flat_search: search_params.search_mode.is_flat_search, distance: search_params.distance, - uses_vector_filters: search_params.vector_filters_file.is_some(), + uses_vector_filters: search_params.search_mode.vector_filters_file.is_some(), num_nodes_to_cache: search_params.num_nodes_to_cache, search_results_per_l, span_metrics, diff --git a/diskann-benchmark/src/inputs/disk.rs b/diskann-benchmark/src/inputs/disk.rs index 7ed521a87..b6744a29c 100644 --- a/diskann-benchmark/src/inputs/disk.rs +++ b/diskann-benchmark/src/inputs/disk.rs @@ -6,27 +6,18 @@ use std::{fmt, num::NonZeroUsize, path::Path}; use anyhow::Context; -#[cfg(feature = "disk-index")] -use std::collections::HashSet; -#[cfg(feature = "disk-index")] -use diskann::graph; use diskann_benchmark_runner::{files::InputFile, utils::datatype::DataType, Checker}; #[cfg(feature = "disk-index")] -use diskann_disk::search::search_mode::SearchMode; -#[cfg(feature = "disk-index")] use diskann_disk::QuantizationType; use diskann_providers::storage::{get_compressed_pq_file, get_disk_index_file, get_pq_pivot_file}; use serde::{Deserialize, Serialize}; use crate::{ - inputs::{as_input, post_processor::TopkPostProcessor, Example}, + inputs::{as_input, graph_index::AdaptiveL, post_processor::TopkPostProcessor, Example}, utils::SimilarityMeasure, }; -#[cfg(feature = "disk-index")] -use crate::inputs::graph_index::AdaptiveL; - ////////////// // Registry // ////////////// @@ -72,69 +63,50 @@ pub(crate) struct DiskIndexBuild { pub(crate) save_path: String, } -#[cfg(feature = "disk-index")] #[derive(Debug, Serialize, Deserialize, Default)] pub(crate) struct DiskSearchMode { pub(crate) is_flat_search: bool, #[serde(default)] pub(crate) adaptive_l: Option, + #[serde(default)] + pub(crate) vector_filters_file: Option, + #[serde(default)] + pub(crate) post_processor: Option, } -#[cfg(feature = "disk-index")] impl DiskSearchMode { - pub(crate) fn search_mode<'a>( - &'a self, - has_vector_filters: bool, - vector_filter: &'a HashSet, - post_processor: Option<&TopkPostProcessor>, - ) -> SearchMode<'a> { - let adaptive_l = self.adaptive_l.as_ref().map(|adaptive_l| { - graph::search::AdaptiveL::new(adaptive_l.sample_count.into(), adaptive_l.scale_factor) - .expect("validated adaptive L must construct") - }); - - match ( - self.is_flat_search, - has_vector_filters, - post_processor, - adaptive_l, - ) { - (true, false, _, _) => SearchMode::flat(), - (true, true, _, _) => { - SearchMode::flat_filtered(move |vid: &u32| vector_filter.contains(vid)) - } - (false, false, Some(TopkPostProcessor::DeterminantDiversity(params)), _) => { - SearchMode::diverse_graph(*params) - } - (false, true, Some(TopkPostProcessor::DeterminantDiversity(params)), _) => { - SearchMode::diverse_graph_filtered( - move |vid: &u32| vector_filter.contains(vid), - *params, - ) - } - (false, false, None, Some(adaptive_l)) => { - SearchMode::inline_filter(|_| true, Some(adaptive_l)) - } - (false, true, None, Some(adaptive_l)) => SearchMode::inline_filter( - move |vid: &u32| vector_filter.contains(vid), - Some(adaptive_l), - ), - (false, false, None, None) => SearchMode::graph(), - (false, true, None, None) => { - SearchMode::graph_filtered(move |vid: &u32| vector_filter.contains(vid)) - } - } - } - pub(crate) fn validate(&mut self, checker: &mut Checker) -> Result<(), anyhow::Error> { + self.validate_compatibility()?; + if let Some(adaptive_l) = self.adaptive_l.as_mut() { adaptive_l.validate(checker)?; } + if let Some(vf) = self.vector_filters_file.as_mut() { + vf.resolve(checker).context("invalid vector_filters_file")?; + } + if let Some(pp) = self.post_processor.as_mut() { + pp.validate(checker) + .context("invalid disk search post processor")?; + } Ok(()) } + + fn validate_compatibility(&self) -> Result<(), anyhow::Error> { + if !self.is_flat_search { + return Ok(()); + } + + match (self.adaptive_l.is_some(), self.post_processor.is_some()) { + (false, false) => Ok(()), + (true, false) => anyhow::bail!("flat disk search does not support adaptive_l"), + (false, true) => anyhow::bail!("flat disk search does not support post_processor"), + (true, true) => { + anyhow::bail!("flat disk search does not support adaptive_l or post_processor") + } + } + } } -#[cfg(feature = "disk-index")] impl fmt::Display for DiskSearchMode { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { let base = if self.is_flat_search { "flat" } else { "graph" }; @@ -146,8 +118,73 @@ impl fmt::Display for DiskSearchMode { } } +#[cfg(test)] +mod tests { + use std::num::NonZeroUsize; + + use super::*; + + #[test] + fn flat_disk_search_rejects_adaptive_l() { + let mode = DiskSearchMode { + is_flat_search: true, + adaptive_l: Some(AdaptiveL { + sample_count: NonZeroUsize::MIN, + scale_factor: 1.0, + }), + vector_filters_file: None, + post_processor: None, + }; + + let err = mode + .validate_compatibility() + .expect_err("flat search with adaptive_l must be invalid"); + assert!(err.to_string().contains("does not support adaptive_l")); + } + + #[test] + fn flat_disk_search_rejects_post_processor() { + let mode: DiskSearchMode = serde_json::from_str( + r#"{ + "is_flat_search": true, + "post_processor": { + "type": "determinant-diversity", + "power": 1.0, + "eta": 0.0 + } + }"#, + ) + .expect("test post-processor configuration must deserialize"); + + let err = mode + .validate_compatibility() + .expect_err("flat search with a post_processor must be invalid"); + assert!(err.to_string().contains("does not support post_processor")); + } + + #[test] + fn disk_search_phase_rejects_legacy_phase_level_search_mode_fields() { + let error = serde_json::from_str::( + r#"{ + "queries": "queries.fbin", + "groundtruth": "groundtruth.bin", + "num_threads": 1, + "beam_width": 1, + "search_list": [1], + "recall_at": 1, + "distance": "squared_l2", + "is_flat_search": true + }"#, + ) + .expect_err("legacy phase-level search settings must be rejected"); + + assert!(error.to_string().contains("is_flat_search")); + } +} + /// Search phase configuration #[derive(Debug, Deserialize, Serialize)] +#[serde(deny_unknown_fields)] pub(crate) struct DiskSearchPhase { pub(crate) queries: InputFile, pub(crate) groundtruth: InputFile, @@ -155,21 +192,11 @@ pub(crate) struct DiskSearchPhase { pub(crate) beam_width: usize, pub(crate) search_list: Vec, pub(crate) recall_at: u32, - #[cfg(feature = "disk-index")] #[serde(default)] pub(crate) search_mode: DiskSearchMode, - // Backward compatibility for older benchmark inputs that used - // `is_flat_search` directly at the search-phase level. - #[cfg(feature = "disk-index")] - #[serde(default, skip_serializing)] - pub(crate) is_flat_search: Option, - #[cfg(not(feature = "disk-index"))] - pub(crate) is_flat_search: bool, pub(crate) distance: SimilarityMeasure, - pub(crate) vector_filters_file: Option, pub(crate) num_nodes_to_cache: Option, pub(crate) search_io_limit: Option, - pub(crate) post_processor: Option, } ///////// @@ -270,16 +297,7 @@ impl DiskSearchPhase { self.groundtruth .resolve(checker) .context("invalid groundtruth file")?; - if let Some(vf) = self.vector_filters_file.as_mut() { - vf.resolve(checker).context("invalid vector_filters_file")?; - } - #[cfg(feature = "disk-index")] - if let Some(is_flat_search) = self.is_flat_search { - self.search_mode.is_flat_search = is_flat_search; - } - - #[cfg(feature = "disk-index")] self.search_mode .validate(checker) .context("invalid disk search mode")?; @@ -315,11 +333,6 @@ impl DiskSearchPhase { } } - if let Some(pp) = self.post_processor.as_mut() { - pp.validate(checker) - .context("invalid disk search post processor")?; - } - Ok(()) } } @@ -353,20 +366,15 @@ impl Example for DiskIndexOperation { beam_width: 16, recall_at: 10, num_threads: 8, - #[cfg(feature = "disk-index")] search_mode: DiskSearchMode { is_flat_search: false, adaptive_l: None, + vector_filters_file: None, + post_processor: None, }, - #[cfg(feature = "disk-index")] - is_flat_search: None, - #[cfg(not(feature = "disk-index"))] - is_flat_search: false, distance: SimilarityMeasure::SquaredL2, - vector_filters_file: None, num_nodes_to_cache: None, search_io_limit: None, - post_processor: None, }; Self { @@ -478,12 +486,9 @@ impl DiskSearchPhase { write_field!(f, "Beam Width", self.beam_width)?; write_field!(f, "Recall@", self.recall_at)?; write_field!(f, "Threads", self.num_threads)?; - #[cfg(feature = "disk-index")] write_field!(f, "Search Mode", self.search_mode)?; - #[cfg(not(feature = "disk-index"))] - write_field!(f, "Flat Search", self.is_flat_search)?; write_field!(f, "Distance", self.distance)?; - match &self.vector_filters_file { + match &self.search_mode.vector_filters_file { Some(vf) => write_field!(f, "Vector Filters File", vf.display())?, None => write_field!(f, "Vector Filters File", "none")?, } @@ -495,7 +500,7 @@ impl DiskSearchPhase { Some(lim) => write_field!(f, "Search IO Limit", format!("{lim}"))?, None => write_field!(f, "Search IO Limit", "none (defaults to `usize::MAX`)")?, } - match &self.post_processor { + match &self.search_mode.post_processor { Some(pp) => write_field!(f, "Post Processor", pp)?, None => write_field!(f, "Post Processor", "none")?, } diff --git a/diskann-benchmark/src/main.rs b/diskann-benchmark/src/main.rs index 55f6bd019..7707fb12b 100644 --- a/diskann-benchmark/src/main.rs +++ b/diskann-benchmark/src/main.rs @@ -709,6 +709,18 @@ mod tests { prefix_search_directories(&mut raw, &root_directory()); let tempdir = tempfile::tempdir().unwrap(); + + // Redirect each build job's `save_path` into the tempdir so the disk index + // artifacts are not written relative to the process cwd (the repo tree). + let jobs = raw["jobs"] + .as_array_mut() + .expect("\"jobs\" should be an array"); + for (i, job) in jobs.iter_mut().enumerate() { + let save_path = tempdir.path().join(format!("disk_index_filter_job_{i}")); + job["content"]["source"]["save_path"] = + serde_json::Value::String(save_path.to_str().unwrap().to_string()); + } + let input_path = tempdir.path().join("disk-index-filter.json"); save_to_file(&input_path, &raw); let output_path = tempdir.path().join("output.json");