diff --git a/.github/workflows/disk-benchmarks.yml b/.github/workflows/disk-benchmarks.yml index 29c0a7af86..b7960cf9ea 100644 --- a/.github/workflows/disk-benchmarks.yml +++ b/.github/workflows/disk-benchmarks.yml @@ -116,7 +116,7 @@ jobs: working-directory: baseline run: | cargo run -p diskann-benchmark --features disk-index --release -- \ - run --input-file ../diskann_rust/${{ env.PERF_INPUTS }}/${{ matrix.config }} \ + run --input-file ${{ env.PERF_INPUTS }}/${{ matrix.config }} \ --output-file target/tmp/${{ matrix.dataset }}_baseline.json - name: Run current branch benchmark @@ -144,4 +144,4 @@ jobs: path: | diskann_rust/target/tmp/${{ matrix.dataset }}_target.json baseline/target/tmp/${{ matrix.dataset }}_baseline.json - retention-days: 30 \ No newline at end of file + retention-days: 30 diff --git a/diskann-benchmark/example/disk-index-determinant-diversity.json b/diskann-benchmark/example/disk-index-determinant-diversity.json index e32e7c07eb..cab0f93fd3 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": { + "mode": "graph-diverse", + "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 a3f35ca91a..167b4646dc 100644 --- a/diskann-benchmark/example/disk-index-filter.json +++ b/diskann-benchmark/example/disk-index-filter.json @@ -27,9 +27,15 @@ "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": { + "mode": "graph-inline-filter", + "vector_filters_file": "disk_index_10pts_idx_uint32_range_res_r_100000.bin", + "adaptive_l": { + "sample_count": 10, + "scale_factor": 16.0 + } + }, + "distance": "squared_l2" } } }, @@ -57,9 +63,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": { + "mode": "flat", + "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 46db8e118a..0b8534951a 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": { "mode": "graph" }, + "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": { "mode": "flat" }, + "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 6b3e3b42d5..2e535d4f4c 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": { "mode": "graph" }, + "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 59c439017d..3593ad8d62 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": { "mode": "graph" }, + "distance": "inner_product" } } } diff --git a/diskann-benchmark/src/disk_index/search.rs b/diskann-benchmark/src/disk_index/search.rs index 8cd0f5e606..a9e63e469e 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::{ @@ -18,7 +19,7 @@ use diskann_disk::{ disk_provider::DiskIndexSearcher, disk_vertex_provider_factory::DiskVertexProviderFactory, }, - search_mode::SearchMode, + search_mode::{SearchMode, SearchPredicate}, }, storage::disk_index_reader::DiskIndexReader, utils::{instrumentation::PerfLogger, statistics, QueryStatistics}, @@ -32,27 +33,83 @@ use diskann_providers::{ }; use diskann_tools::utils::{search_index_utils, KRecallAtN}; use diskann_utils::views::Matrix; +use scopeguard::defer; use serde::{Deserialize, Serialize}; use crate::{ disk_index::json_spancollector::JsonSpanCollector, - inputs::disk::{DiskIndexLoad, DiskSearchPhase}, + inputs::disk::{DiskIndexLoad, DiskSearchMode, DiskSearchPhase, DiskSearchStrategy}, + inputs::post_processor::TopkPostProcessor, utils::{datafiles, SimilarityMeasure}, }; -#[derive(Serialize, Deserialize, Debug)] +#[derive(Serialize, Debug)] pub(super) struct DiskSearchStats { pub(super) num_threads: usize, pub(super) beam_width: usize, pub(super) recall_at: u32, - pub(crate) is_flat_search: bool, + pub(crate) search_strategy: DiskSearchStrategy, pub(crate) distance: SimilarityMeasure, - pub(crate) uses_vector_filters: bool, pub(super) num_nodes_to_cache: Option, pub(super) search_results_per_l: Vec, span_metrics: serde_json::Value, } +#[derive(Deserialize)] +struct RawDiskSearchStats { + num_threads: usize, + beam_width: usize, + recall_at: u32, + #[serde(default)] + search_strategy: Option, + #[serde(default)] + is_flat_search: Option, + distance: SimilarityMeasure, + #[serde(default)] + uses_vector_filters: Option, + num_nodes_to_cache: Option, + search_results_per_l: Vec, + span_metrics: serde_json::Value, +} + +impl<'de> Deserialize<'de> for DiskSearchStats { + fn deserialize(deserializer: D) -> Result + where + D: serde::Deserializer<'de>, + { + let raw = RawDiskSearchStats::deserialize(deserializer)?; + let search_strategy = if let Some(strategy) = raw.search_strategy { + strategy + } else if let Some(is_flat) = raw.is_flat_search { + let uses_vector_filters = raw.uses_vector_filters.unwrap_or(false); + if is_flat { + DiskSearchStrategy::Flat { + uses_vector_filters, + } + } else { + DiskSearchStrategy::Graph { + uses_vector_filters, + } + } + } else { + return Err(serde::de::Error::custom( + "missing search_strategy or legacy is_flat_search", + )); + }; + + Ok(Self { + num_threads: raw.num_threads, + beam_width: raw.beam_width, + recall_at: raw.recall_at, + search_strategy, + distance: raw.distance, + num_nodes_to_cache: raw.num_nodes_to_cache, + search_results_per_l: raw.search_results_per_l, + span_metrics: raw.span_metrics, + }) + } +} + #[derive(Serialize, Deserialize, Debug)] pub(super) struct DiskSearchResult { pub(super) search_l: u32, @@ -158,6 +215,59 @@ impl DiskSearchResult { } } +fn build_adaptive_l(mode: &DiskSearchMode) -> anyhow::Result> { + let DiskSearchMode::GraphInlineFilter { adaptive_l, .. } = mode else { + return Ok(None); + }; + + adaptive_l + .as_ref() + .map(|adaptive_l| { + graph::search::AdaptiveL::new(adaptive_l.sample_count.into(), adaptive_l.scale_factor) + .map_err(Into::into) + }) + .transpose() +} + +/// Construct the backend [`SearchMode`] from the JSON-configured strategy, +/// the per-query vector filter, and the pre-validated adaptive-L settings. +fn build_search_mode<'a>( + mode: &DiskSearchMode, + vector_filter: Option<&'a HashSet>, + adaptive_l: Option<&graph::search::AdaptiveL>, +) -> anyhow::Result> { + if adaptive_l.is_some() && !matches!(mode, DiskSearchMode::GraphInlineFilter { .. }) { + anyhow::bail!("adaptive-L is only valid for inline-filter search"); + } + + let filter = vector_filter.map(|vector_filter| { + Box::new(move |vid: &u32| vector_filter.contains(vid)) as SearchPredicate<'a> + }); + + let mode = match mode { + DiskSearchMode::Flat { .. } => SearchMode::FlatScan { filter }, + DiskSearchMode::Graph { .. } => SearchMode::Graph { filter }, + DiskSearchMode::GraphInlineFilter { .. } => { + let vector_filter = vector_filter.ok_or_else(|| { + anyhow::anyhow!("inline-filter search requires a vector filter for every query") + })?; + SearchMode::inline_filter( + move |vid: &u32| vector_filter.contains(vid), + adaptive_l.cloned(), + ) + } + DiskSearchMode::GraphDiverse { post_processor, .. } => { + let TopkPostProcessor::DeterminantDiversity(params) = post_processor; + SearchMode::DiverseGraph { + filter, + params: *params, + } + } + }; + + Ok(mode) +} + pub(super) fn search_disk_index( index_load: &DiskIndexLoad, search_params: &DiskSearchPhase, @@ -176,6 +286,9 @@ where global::set_tracer_provider(provider.clone()); Some((collector, provider)) }; + defer! { + global::set_tracer_provider(previous_tracer_provider); + } // Use PerfLogger for consistent checkpoint logging let mut logger = PerfLogger::new("search_disk_index", true); @@ -185,21 +298,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, @@ -236,6 +355,7 @@ where logger.log_checkpoint("index_loaded"); + let adaptive_l = build_adaptive_l(&search_params.search_mode)?; let pool = create_thread_pool(search_params.num_threads)?; let mut search_results_per_l = Vec::with_capacity(search_params.search_list.len()); let has_any_search_failed = AtomicBool::new(false); @@ -259,24 +379,23 @@ 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()) .zip(result_counts.par_iter_mut()); - zipped.for_each_in_pool( + zipped.try_for_each_in_pool( pool.as_ref(), - |(((((q, vf), 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(), - ); + |(((((query_index, q), id_chunk), dist_chunk), stats), rc)| { + let vector_filter = vector_filters + .as_ref() + .and_then(|filters| filters.get(query_index)); + let mode = build_search_mode( + &search_params.search_mode, + vector_filter, + adaptive_l.as_ref(), + )?; match searcher.search( q, @@ -310,8 +429,10 @@ where has_any_search_failed.store(true, std::sync::atomic::Ordering::Release); } } + + Ok::<(), anyhow::Error>(()) }, - ); + )?; let total_time = start.elapsed(); if has_any_search_failed.load(std::sync::atomic::Ordering::Acquire) { @@ -343,15 +464,12 @@ where serde_json::json!({ "span_data": [] }) }; - global::set_tracer_provider(previous_tracer_provider); - Ok(DiskSearchStats { num_threads: search_params.num_threads, beam_width: search_params.beam_width, recall_at: search_params.recall_at, - is_flat_search: search_params.search_mode.is_flat_search, + search_strategy: search_params.search_mode.strategy(), distance: search_params.distance, - uses_vector_filters: search_params.vector_filters_file.is_some(), num_nodes_to_cache: search_params.num_nodes_to_cache, search_results_per_l, span_metrics, @@ -433,9 +551,8 @@ impl fmt::Display for DiskSearchStats { writeln!(f, "Threads, : {}", self.num_threads)?; writeln!(f, "Beam width, : {}", self.beam_width)?; writeln!(f, "Recall at, : {}", self.recall_at)?; - writeln!(f, "Flat search, : {}", self.is_flat_search)?; + writeln!(f, "Search strategy, : {}", self.search_strategy)?; writeln!(f, "Distance, : {}", self.distance)?; - writeln!(f, "Vector filters, : {}", self.uses_vector_filters)?; writeln!( f, "Nodes to cache, : {}", @@ -482,3 +599,73 @@ impl fmt::Display for DiskSearchStats { Ok(()) } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn legacy_search_stats_deserialize_to_strategy() { + let stats: DiskSearchStats = serde_json::from_value(serde_json::json!({ + "num_threads": 1, + "beam_width": 4, + "recall_at": 10, + "is_flat_search": false, + "distance": "squared_l2", + "num_nodes_to_cache": null, + "search_results_per_l": [], + "span_metrics": {} + })) + .unwrap(); + + assert!(matches!( + stats.search_strategy, + DiskSearchStrategy::Graph { + uses_vector_filters: false + } + )); + } + + #[test] + fn new_search_stats_serialize_only_complete_strategy() { + let stats: DiskSearchStats = serde_json::from_value(serde_json::json!({ + "num_threads": 1, + "beam_width": 4, + "recall_at": 10, + "search_strategy": { + "mode": "graph-inline-filter", + "adaptive_l": { + "sample_count": 10, + "scale_factor": 16.0 + } + }, + "is_flat_search": false, + "uses_vector_filters": true, + "distance": "squared_l2", + "num_nodes_to_cache": null, + "search_results_per_l": [], + "span_metrics": {} + })) + .unwrap(); + + let value = serde_json::to_value(stats).unwrap(); + assert_eq!(value["search_strategy"]["mode"], "graph-inline-filter"); + assert_eq!(value["search_strategy"]["adaptive_l"]["sample_count"], 10); + assert!(value.get("is_flat_search").is_none()); + assert!(value.get("uses_vector_filters").is_none()); + } + + #[test] + fn inline_search_mode_requires_query_filter() { + let mode: DiskSearchMode = serde_json::from_value(serde_json::json!({ + "mode": "graph-inline-filter", + "vector_filters_file": "filters.bin" + })) + .unwrap(); + + let error = build_search_mode(&mode, None, None) + .err() + .expect("inline-filter mode must reject a missing query filter"); + assert!(error.to_string().contains("requires a vector filter")); + } +} diff --git a/diskann-benchmark/src/inputs/disk.rs b/diskann-benchmark/src/inputs/disk.rs index 7ed521a874..5a33f6b811 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,80 +63,267 @@ 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, +/// Disk search mode, modeled after the four backend search strategies. +/// Strategy-specific settings live only on the variants that use them. +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(tag = "mode", rename_all = "kebab-case")] +pub(crate) enum DiskSearchMode { + /// Brute-force flat scan, optionally restricted by a per-query vector filter. + Flat { + #[serde(default)] + vector_filters_file: Option, + }, + /// Greedy graph search, optionally post-filtered by a per-query vector filter. + Graph { + #[serde(default)] + vector_filters_file: Option, + }, + /// Graph search that checks a required vector filter during traversal. + GraphInlineFilter { + vector_filters_file: InputFile, + #[serde(default)] + adaptive_l: Option, + }, + /// Graph search followed by determinant-diversity selection. + GraphDiverse { + #[serde(default)] + vector_filters_file: Option, + post_processor: TopkPostProcessor, + }, +} + +/// Path-independent search metadata written to benchmark results. +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(tag = "mode", rename_all = "kebab-case")] +pub(crate) enum DiskSearchStrategy { + Flat { + uses_vector_filters: bool, + }, + Graph { + uses_vector_filters: bool, + }, + GraphInlineFilter { + #[serde(default)] + adaptive_l: Option, + }, + GraphDiverse { + uses_vector_filters: bool, + post_processor: TopkPostProcessor, + }, +} + +impl Default for DiskSearchMode { + fn default() -> Self { + Self::Graph { + vector_filters_file: None, + } + } } -#[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)) + pub(crate) fn vector_filters_file(&self) -> Option<&InputFile> { + match self { + Self::Flat { + vector_filters_file, } - (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)) + | Self::Graph { + vector_filters_file, } + | Self::GraphDiverse { + vector_filters_file, + .. + } => vector_filters_file.as_ref(), + Self::GraphInlineFilter { + vector_filters_file, + .. + } => Some(vector_filters_file), + } + } + + pub(crate) fn post_processor(&self) -> Option<&TopkPostProcessor> { + match self { + Self::GraphDiverse { post_processor, .. } => Some(post_processor), + _ => None, + } + } + + pub(crate) fn strategy(&self) -> DiskSearchStrategy { + match self { + Self::Flat { + vector_filters_file, + } => DiskSearchStrategy::Flat { + uses_vector_filters: vector_filters_file.is_some(), + }, + Self::Graph { + vector_filters_file, + } => DiskSearchStrategy::Graph { + uses_vector_filters: vector_filters_file.is_some(), + }, + Self::GraphInlineFilter { adaptive_l, .. } => DiskSearchStrategy::GraphInlineFilter { + adaptive_l: adaptive_l.clone(), + }, + Self::GraphDiverse { + vector_filters_file, + post_processor, + } => DiskSearchStrategy::GraphDiverse { + uses_vector_filters: vector_filters_file.is_some(), + post_processor: post_processor.clone(), + }, } } pub(crate) fn validate(&mut self, checker: &mut Checker) -> Result<(), anyhow::Error> { - if let Some(adaptive_l) = self.adaptive_l.as_mut() { - adaptive_l.validate(checker)?; + match self { + Self::Flat { + vector_filters_file, + } + | Self::Graph { + vector_filters_file, + } => { + if let Some(vf) = vector_filters_file.as_mut() { + vf.resolve(checker).context("invalid vector_filters_file")?; + } + } + Self::GraphInlineFilter { + vector_filters_file, + adaptive_l, + } => { + vector_filters_file + .resolve(checker) + .context("invalid vector_filters_file")?; + if let Some(adaptive_l) = adaptive_l { + adaptive_l.validate(checker)?; + } + } + Self::GraphDiverse { + vector_filters_file, + post_processor, + } => { + if let Some(vf) = vector_filters_file.as_mut() { + vf.resolve(checker).context("invalid vector_filters_file")?; + } + post_processor + .validate(checker) + .context("invalid disk search post processor")?; + } } Ok(()) } } -#[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" }; - if self.adaptive_l.is_some() { - write!(f, "{} + adaptive-l", base) - } else { - write!(f, "{}", base) + self.strategy().fmt(f) + } +} + +impl fmt::Display for DiskSearchStrategy { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Flat { + uses_vector_filters: false, + } => write!(f, "flat"), + Self::Flat { + uses_vector_filters: true, + } => write!(f, "flat + vector-filter"), + Self::Graph { + uses_vector_filters: false, + } => write!(f, "graph"), + Self::Graph { + uses_vector_filters: true, + } => write!(f, "graph + vector-filter"), + Self::GraphInlineFilter { adaptive_l: None } => write!(f, "graph inline-filter"), + Self::GraphInlineFilter { + adaptive_l: Some(_), + } => write!(f, "graph inline-filter + adaptive-l"), + Self::GraphDiverse { + uses_vector_filters: false, + .. + } => write!(f, "graph diverse"), + Self::GraphDiverse { + uses_vector_filters: true, + .. + } => write!(f, "graph diverse + vector-filter"), } } } +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn disk_search_modes_deserialize() { + let flat: DiskSearchMode = + serde_json::from_str(r#"{ "mode": "flat" }"#).expect("flat mode must deserialize"); + assert!(matches!(flat, DiskSearchMode::Flat { .. })); + + let graph: DiskSearchMode = + serde_json::from_str(r#"{ "mode": "graph", "vector_filters_file": "filters.bin" }"#) + .expect("graph mode must deserialize"); + assert!(matches!(graph, DiskSearchMode::Graph { .. })); + + let inline: DiskSearchMode = serde_json::from_str( + r#"{ + "mode": "graph-inline-filter", + "vector_filters_file": "filters.bin", + "adaptive_l": { "sample_count": 1, "scale_factor": 2.0 } + }"#, + ) + .expect("inline-filter mode must deserialize"); + assert!(matches!( + inline, + DiskSearchMode::GraphInlineFilter { + adaptive_l: Some(_), + .. + } + )); + + let diverse: DiskSearchMode = serde_json::from_str( + r#"{ + "mode": "graph-diverse", + "post_processor": { + "type": "determinant-diversity", + "power": 2.0, + "eta": 1.0 + } + }"#, + ) + .expect("diverse mode must deserialize"); + assert!(matches!(diverse, DiskSearchMode::GraphDiverse { .. })); + } + + #[test] + fn strategy_specific_fields_are_required() { + let inline_error = + serde_json::from_str::(r#"{ "mode": "graph-inline-filter" }"#) + .expect_err("inline-filter mode must require a vector filter"); + assert!(inline_error.to_string().contains("vector_filters_file")); + + let diverse_error = + serde_json::from_str::(r#"{ "mode": "graph-diverse" }"#) + .expect_err("diverse mode must require a post-processor"); + assert!(diverse_error.to_string().contains("post_processor")); + } + + #[test] + fn search_strategy_omits_filter_file_paths() { + let mode: DiskSearchMode = serde_json::from_str( + r#"{ + "mode": "graph-inline-filter", + "vector_filters_file": "private/filters.bin", + "adaptive_l": { "sample_count": 1, "scale_factor": 2.0 } + }"#, + ) + .unwrap(); + + let value = serde_json::to_value(mode.strategy()).unwrap(); + assert_eq!(value["mode"], "graph-inline-filter"); + assert!(value.get("vector_filters_file").is_none()); + assert_eq!(value["adaptive_l"]["sample_count"], 1); + } +} + /// Search phase configuration #[derive(Debug, Deserialize, Serialize)] pub(crate) struct DiskSearchPhase { @@ -155,21 +333,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 +438,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 +474,6 @@ impl DiskSearchPhase { } } - if let Some(pp) = self.post_processor.as_mut() { - pp.validate(checker) - .context("invalid disk search post processor")?; - } - Ok(()) } } @@ -353,20 +507,12 @@ 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, + search_mode: DiskSearchMode::Graph { + vector_filters_file: 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 +624,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 +638,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/inputs/graph_index.rs b/diskann-benchmark/src/inputs/graph_index.rs index 2137c6e576..21904cb587 100644 --- a/diskann-benchmark/src/inputs/graph_index.rs +++ b/diskann-benchmark/src/inputs/graph_index.rs @@ -244,7 +244,7 @@ impl MultihopFilterSearchPhase { } } -#[derive(Debug, Serialize, Deserialize)] +#[derive(Debug, Clone, Serialize, Deserialize)] pub(crate) struct AdaptiveL { pub(crate) sample_count: NonZeroUsize, pub(crate) scale_factor: f64, diff --git a/diskann-benchmark/src/main.rs b/diskann-benchmark/src/main.rs index cf9a0ea522..73530b80ef 100644 --- a/diskann-benchmark/src/main.rs +++ b/diskann-benchmark/src/main.rs @@ -238,6 +238,114 @@ mod tests { } } + /// Redirect disk-index build artifacts into a temporary directory. + /// + /// Only existing `save_path` fields are replaced. Jobs without a `save_path` + /// are left unchanged. + fn redirect_save_paths(raw: &mut serde_json::Value, directory: &std::path::Path) { + let Some(jobs) = raw.get_mut("jobs").and_then(Value::as_array_mut) else { + return; + }; + + for (index, job) in jobs.iter_mut().enumerate() { + let Some(save_path) = job + .get_mut("content") + .and_then(|content| content.get_mut("source")) + .and_then(Value::as_object_mut) + .and_then(|source| source.get_mut("save_path")) + else { + continue; + }; + + *save_path = Value::String( + directory + .join(format!("disk_index_job_{index}")) + .to_string_lossy() + .into_owned(), + ); + } + } + + /// Unit-test only the in-memory JSON rewrite; no index is built or loaded. + #[test] + fn redirect_save_paths_updates_builds_without_modifying_loads() { + let mut raw = serde_json::json!({ + "jobs": [ + { + "content": { + "source": { + "disk-index-source": "Build", + "save_path": "build-only-index" + } + } + }, + { + "content": { + "source": { + "disk-index-source": "Load", + "load_path": "existing-index" + } + } + }, + { + "content": { + "source": { + "disk-index-source": "Build", + "save_path": "unrelated-build" + } + } + }, + { + "content": { + "source": { + "disk-index-source": "Load", + "load_path": "unrelated-load" + } + } + } + ] + }); + let directory = std::path::Path::new("temporary-output"); + + redirect_save_paths(&mut raw, directory); + + assert_eq!( + raw["jobs"][0]["content"]["source"]["save_path"], + directory + .join("disk_index_job_0") + .to_string_lossy() + .as_ref() + ); + assert_eq!( + raw["jobs"][1]["content"]["source"]["load_path"], + "existing-index" + ); + assert_eq!( + raw["jobs"][2]["content"]["source"]["save_path"], + directory + .join("disk_index_job_2") + .to_string_lossy() + .as_ref() + ); + assert_eq!( + raw["jobs"][3]["content"]["source"]["load_path"], + "unrelated-load" + ); + for index in [1, 3] { + assert!(raw["jobs"][index]["content"]["source"] + .get("save_path") + .is_none()); + } + } + + /// Unit-test that `redirect_save_paths` safely ignores malformed JSON input. + #[test] + fn redirect_save_paths_ignores_malformed_input() { + for mut raw in [serde_json::json!(null), serde_json::json!({ "jobs": null })] { + redirect_save_paths(&mut raw, std::path::Path::new("temporary-output")); + } + } + // Retrieve the number of jobs in the raw input JSON. // // The format is @@ -264,13 +372,18 @@ mod tests { } } - fn run_integration_test(mut raw: serde_json::Value) { + fn run_integration_test(raw: serde_json::Value) { + run_integration_test_with_results(raw); + } + + fn run_integration_test_with_results(mut raw: serde_json::Value) -> Vec { // First, parse and modify the input file to establish paths relative to the // directory building the dispatcher. // let mut raw = serde_json::from_str(json_string).unwrap(); prefix_search_directories(&mut raw, &root_directory()); let tempdir = tempfile::tempdir().unwrap(); + redirect_save_paths(&mut raw, tempdir.path()); let input_path = tempdir.path().join("input.json"); save_to_file(&input_path, &raw); @@ -299,6 +412,7 @@ mod tests { let results: Vec = load_from_file(&output_path); assert_eq!(results.len(), num_jobs(&raw)); + results } //////////////////////////////// @@ -701,34 +815,20 @@ mod tests { } /// Filtered disk search end-to-end: drives the disk-index backend through - /// `disk-index-filter.json` + /// `disk-index-filter.json`. #[test] #[cfg(feature = "disk-index")] fn disk_index_filter_integration() { - let mut raw = value_from_file(&example_directory().join("disk-index-filter.json")); - prefix_search_directories(&mut raw, &root_directory()); - - let tempdir = tempfile::tempdir().unwrap(); - let input_path = tempdir.path().join("disk-index-filter.json"); - save_to_file(&input_path, &raw); - let output_path = tempdir.path().join("output.json"); - - let command = Commands::Run { - input_file: input_path.to_owned(), - output_file: output_path.to_owned(), - dry_run: false, - allow_debug: true, - }; - let cli = Cli::from_commands(command, true); - let mut output = Memory::new(); - let result = cli.run(&mut output); - let output_str = String::from_utf8(output.into_inner()).unwrap(); - println!("output = {}", output_str); - result.expect("disk-index-filter run failed"); - - assert!(output_path.exists()); - let results: Vec = load_from_file(&output_path); - assert_eq!(results.len(), num_jobs(&raw)); + let raw = value_from_file(&example_directory().join("disk-index-filter.json")); + let results = run_integration_test_with_results(raw); + + let strategy = &results[0]["results"]["search"]["search_strategy"]; + assert_eq!(strategy["mode"], "graph-inline-filter"); + assert_eq!(strategy["adaptive_l"]["sample_count"], 10); + assert_eq!( + results[1]["results"]["search"]["search_strategy"]["mode"], + "flat" + ); } #[test]