diff --git a/Cargo.lock b/Cargo.lock index 1aae5eaf8d7..11f95bc38ef 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -548,9 +548,9 @@ dependencies = [ [[package]] name = "aws-lc-rs" -version = "1.17.0" +version = "1.18.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5ec2f1fc3ec205783a5da9a7e6c1509cc69dedf09a1949e412c1e18469326d00" +checksum = "b281d307588d634de920874890732659e2e7672f72b5e10e81badc1a8a83621e" dependencies = [ "aws-lc-sys", "zeroize", @@ -558,14 +558,15 @@ dependencies = [ [[package]] name = "aws-lc-sys" -version = "0.41.0" +version = "0.45.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1a2f9779ce85b93ab6170dd940ad0169b5766ff848247aff13bb788b832fe3f4" +checksum = "9bff6c3b54fad79a2e60b8102caf565819711497c1f5f092f49508e2f5c31b27" dependencies = [ "cc", "cmake", "dunce", "fs_extra", + "pkg-config", ] [[package]] @@ -7870,9 +7871,9 @@ dependencies = [ [[package]] name = "rustls" -version = "0.23.40" +version = "0.23.45" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ef86cd5876211988985292b91c96a8f2d298df24e75989a43a3c73f2d4d8168b" +checksum = "0d41d731c7d2f962d1ccc364cec258de3c0e93b38c2fb3ba97ac74513048d634" dependencies = [ "aws-lc-rs", "log", @@ -7935,9 +7936,9 @@ checksum = "f87165f0995f63a9fbeea62b64d10b4d9d8e78ec6d7d51fb2125fda7bb36788f" [[package]] name = "rustls-webpki" -version = "0.103.13" +version = "0.103.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "61c429a8649f110dddef65e2a5ad240f747e85f7758a6bccc7e5777bd33f756e" +checksum = "f3c3cf1d8b1e7d4927e2d154c3fcb02979afb9939629c62cd9048d4f07b60ac2" dependencies = [ "aws-lc-rs", "ring", diff --git a/java/lance-jni/Cargo.lock b/java/lance-jni/Cargo.lock index 7edc85e4922..5ca99f0c411 100644 --- a/java/lance-jni/Cargo.lock +++ b/java/lance-jni/Cargo.lock @@ -499,9 +499,9 @@ dependencies = [ [[package]] name = "aws-lc-rs" -version = "1.17.0" +version = "1.18.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5ec2f1fc3ec205783a5da9a7e6c1509cc69dedf09a1949e412c1e18469326d00" +checksum = "b281d307588d634de920874890732659e2e7672f72b5e10e81badc1a8a83621e" dependencies = [ "aws-lc-sys", "zeroize", @@ -509,14 +509,15 @@ dependencies = [ [[package]] name = "aws-lc-sys" -version = "0.41.0" +version = "0.45.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1a2f9779ce85b93ab6170dd940ad0169b5766ff848247aff13bb788b832fe3f4" +checksum = "9bff6c3b54fad79a2e60b8102caf565819711497c1f5f092f49508e2f5c31b27" dependencies = [ "cc", "cmake", "dunce", "fs_extra", + "pkg-config", ] [[package]] @@ -6243,9 +6244,9 @@ dependencies = [ [[package]] name = "rustls" -version = "0.23.40" +version = "0.23.45" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ef86cd5876211988985292b91c96a8f2d298df24e75989a43a3c73f2d4d8168b" +checksum = "0d41d731c7d2f962d1ccc364cec258de3c0e93b38c2fb3ba97ac74513048d634" dependencies = [ "aws-lc-rs", "once_cell", @@ -6307,9 +6308,9 @@ checksum = "f87165f0995f63a9fbeea62b64d10b4d9d8e78ec6d7d51fb2125fda7bb36788f" [[package]] name = "rustls-webpki" -version = "0.103.13" +version = "0.103.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "61c429a8649f110dddef65e2a5ad240f747e85f7758a6bccc7e5777bd33f756e" +checksum = "f3c3cf1d8b1e7d4927e2d154c3fcb02979afb9939629c62cd9048d4f07b60ac2" dependencies = [ "aws-lc-rs", "ring", diff --git a/python/Cargo.lock b/python/Cargo.lock index ad025fff804..8d9b3cf4e79 100644 --- a/python/Cargo.lock +++ b/python/Cargo.lock @@ -533,9 +533,9 @@ dependencies = [ [[package]] name = "aws-lc-rs" -version = "1.17.0" +version = "1.18.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5ec2f1fc3ec205783a5da9a7e6c1509cc69dedf09a1949e412c1e18469326d00" +checksum = "b281d307588d634de920874890732659e2e7672f72b5e10e81badc1a8a83621e" dependencies = [ "aws-lc-sys", "zeroize", @@ -543,14 +543,15 @@ dependencies = [ [[package]] name = "aws-lc-sys" -version = "0.41.0" +version = "0.45.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1a2f9779ce85b93ab6170dd940ad0169b5766ff848247aff13bb788b832fe3f4" +checksum = "9bff6c3b54fad79a2e60b8102caf565819711497c1f5f092f49508e2f5c31b27" dependencies = [ "cc", "cmake", "dunce", "fs_extra", + "pkg-config", ] [[package]] @@ -6966,9 +6967,9 @@ dependencies = [ [[package]] name = "rustls" -version = "0.23.40" +version = "0.23.45" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ef86cd5876211988985292b91c96a8f2d298df24e75989a43a3c73f2d4d8168b" +checksum = "0d41d731c7d2f962d1ccc364cec258de3c0e93b38c2fb3ba97ac74513048d634" dependencies = [ "aws-lc-rs", "once_cell", @@ -7030,9 +7031,9 @@ checksum = "f87165f0995f63a9fbeea62b64d10b4d9d8e78ec6d7d51fb2125fda7bb36788f" [[package]] name = "rustls-webpki" -version = "0.103.13" +version = "0.103.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "61c429a8649f110dddef65e2a5ad240f747e85f7758a6bccc7e5777bd33f756e" +checksum = "f3c3cf1d8b1e7d4927e2d154c3fcb02979afb9939629c62cd9048d4f07b60ac2" dependencies = [ "aws-lc-rs", "ring", diff --git a/python/python/tests/compat/test_venv_manager.py b/python/python/tests/compat/test_venv_manager.py index 57f8b94910a..18a576bf088 100644 --- a/python/python/tests/compat/test_venv_manager.py +++ b/python/python/tests/compat/test_venv_manager.py @@ -18,6 +18,9 @@ ("6.0.0", "lance-namespace>=0.7.2,<0.8"), ("7.2.0b5", "lance-namespace>=0.8.0,<0.9"), ("7.2.0", "lance-namespace>=0.8.0,<0.9"), + ("12.0.0b5", "lance-namespace>=0.8.0,<0.9"), + ("12.0.0b6", "lance-namespace>=0.11.1,<0.12"), + ("12.0.0", "lance-namespace>=0.11.1,<0.12"), ], ) def test_lance_namespace_dependency(version: str, expected: str): diff --git a/python/python/tests/compat/venv_manager.py b/python/python/tests/compat/venv_manager.py index 04e00d7af00..6b1e3b85e2d 100644 --- a/python/python/tests/compat/venv_manager.py +++ b/python/python/tests/compat/venv_manager.py @@ -76,9 +76,14 @@ def _pip_install(python: Union[str, Path], args: list[str]) -> None: NAMESPACE_0_6_DEPENDENCY = "lance-namespace<0.7" NAMESPACE_0_7_DEPENDENCY = "lance-namespace>=0.7.2,<0.8" NAMESPACE_0_8_DEPENDENCY = "lance-namespace>=0.8.0,<0.9" +NAMESPACE_0_11_DEPENDENCY = "lance-namespace>=0.11.1,<0.12" def _lance_namespace_dependency(pylance_version: str) -> str: + # 12.0.0b5 is the last release published while pylance still pinned + # lance-namespace <0.9; releases cut after that carry the 0.11 range. + if Version(pylance_version) > Version("12.0.0b5"): + return NAMESPACE_0_11_DEPENDENCY if Version(pylance_version) >= Version("7.2.0b5"): return NAMESPACE_0_8_DEPENDENCY if Version(pylance_version) >= Version("6.0.0b0"): diff --git a/rust/lance/src/index.rs b/rust/lance/src/index.rs index d10ec624d77..c869755f603 100644 --- a/rust/lance/src/index.rs +++ b/rust/lance/src/index.rs @@ -2099,13 +2099,14 @@ impl DatasetIndexExt for Dataset { } let is_index_type_change = existing_different_type_url.is_some(); - let removed_indices = existing_named_indices - .into_iter() - .map(|idx| -> Result> { + let mut removed_indices = Vec::new(); + let mut retained_indices = Vec::new(); + for idx in existing_named_indices { + let removed = (|| -> Result { // A logical index cannot combine segment types. Full current-fragment // coverage was verified above, so a type change replaces every segment. if is_index_type_change { - return Ok(Some(idx)); + return Ok(true); } let Some(existing_fragments) = idx.effective_fragment_bitmap(&dataset_fragments) @@ -2116,18 +2117,18 @@ impl DatasetIndexExt for Dataset { idx.uuid, index_name ))); } - return Ok(Some(idx)); + return Ok(true); }; // A zero-fragment segment can be used to create an index while // deferring the actual build. Such a segment is disjoint from every // other segment but should still be removed. if existing_fragments.is_empty() { - return Ok(Some(idx)); + return Ok(true); } if existing_fragments.is_disjoint(&incoming_fragments) { - return Ok(None); + return Ok(false); } let uncovered = existing_fragments - &incoming_fragments; @@ -2140,12 +2141,53 @@ impl DatasetIndexExt for Dataset { ))); } - Ok(Some(idx)) - }) - .collect::>>()? - .into_iter() - .flatten() - .collect::>(); + Ok(true) + })()?; + if removed { + removed_indices.push(idx); + } else { + retained_indices.push(idx); + } + } + + // Segments of a logical vector index are planned and ranked together, + // so coexisting segments — the incoming set plus every retained + // existing segment — must share one query contract (metric, dimension, + // sub-index type, quantizer kind). Reject incompatible combinations + // here; otherwise scan planning would derive the metric from the + // first segment and silently rank the rest under the wrong metric. + // Replacement is already selected, so a complete replacement may + // change these settings without conflicting with removed segments. + let coexisting_indices = new_indices.iter().chain(retained_indices.iter()); + let vector_segment_count = coexisting_indices + .clone() + .filter(|segment| segment_has_vector_details(segment)) + .count(); + if vector_segment_count > 1 { + if vector_segment_count != new_indices.len() + retained_indices.len() { + return Err(Error::invalid_input(format!( + "CreateIndex: segment set for index '{index_name}' mixes vector and non-vector segments" + ))); + } + let mut vector_indices = Vec::with_capacity(vector_segment_count); + for segment in new_indices.iter().chain(retained_indices.iter()) { + let index = self + .open_vector_index_from_metadata(column, segment, &NoOpMetricsCollector) + .await + .map_err(|error| { + Error::invalid_input(format!( + "CreateIndex: cannot open vector segment {} of index '{index_name}' for compatibility validation: {error}", + segment.uuid + )) + })?; + vector_indices.push(index); + } + vector::ivf::validate_vector_query_compatibility( + &vector_indices, + &format!("CreateIndex: index '{index_name}'"), + ) + .map_err(|error| Error::invalid_input(error.to_string()))?; + } let transaction = Transaction::new( self.manifest.version, @@ -2711,102 +2753,21 @@ pub trait DatasetIndexInternalExt: DatasetIndexExt { async fn initialize_indices(&mut self, source_dataset: &Dataset) -> Result<()>; } -#[async_trait] -impl DatasetIndexInternalExt for Dataset { - async fn open_generic_index( - &self, - column: &str, - uuid: &Uuid, - metrics: &dyn MetricsCollector, - ) -> Result> { - // Checking for cache existence is cheap so we just check the vector caches. - // Scalar indices cache themselves inside `open_scalar_index` (the cache - // key is a plugin detail), so there is no cheap scalar check here. - let frag_reuse_uuid = self.frag_reuse_index_uuid().await; - - // Check sized cache for IvfIndexState (v2+ indices). - let state_key = IvfIndexStateCacheKey::new(uuid, frag_reuse_uuid.as_ref()); - if self.index_cache.get_with_key(&state_key).await.is_some() { - // Reconstruct via open_vector_index which will hit the same sized key. - let index = self.open_vector_index(column, uuid, metrics).await?; - return Ok(index.as_index()); - } - - // Fallback: in-memory cache for legacy indices. - let vector_cache_key = LegacyVectorIndexCacheKey::new(uuid, frag_reuse_uuid.as_ref()); - if let Some(cached) = self.index_cache.get_with_key(&vector_cache_key).await { - return Ok(cached.0.clone().as_index()); - } - - let frag_reuse_cache_key = FragReuseIndexCacheKey::new(uuid, frag_reuse_uuid.as_ref()); - if let Some(index) = self.index_cache.get_with_key(&frag_reuse_cache_key).await { - return Ok(Arc::new(FragReuseIndexHandle(index)).as_index()); - } - - // Sometimes we want to open an index and we don't care if it is a scalar or vector index. - // For example, we might want to get statistics for an index, regardless of type. - // - // We determine if this is a vector index by checking if INDEX_FILE_NAME exists in the - // file list (available since file sizes tracking was added). If the file list is not - // available (older indices), we fall back to checking file existence via HEAD request. - let index_meta = self - .load_index(uuid) - .await? - .ok_or_else(|| Error::index(format!("Index with id {} does not exist", uuid)))?; - - // Check if this is a vector index by looking at the files list - let is_vector_index = if let Some(files) = &index_meta.files { - // If we have file metadata, check if INDEX_FILE_NAME is in the list - files.iter().any(|f| f.path == INDEX_FILE_NAME) - } else { - // Fall back to file existence check for older indices without file metadata - let index_dir = self.indice_files_dir(&index_meta)?; - let index_file = index_dir - .clone() - .join(uuid.to_string()) - .join(INDEX_FILE_NAME); - let object_store = self.object_store_for_index(&index_meta).await?; - object_store.exists(&index_file).await? - }; - - if is_vector_index { - let index = self.open_vector_index(column, uuid, metrics).await?; - Ok(index.as_index()) - } else { - let index = self.open_scalar_index(column, uuid, metrics).await?; - Ok(index.as_index()) - } - } - - #[instrument(level = "debug", skip_all)] - async fn open_scalar_index( - &self, - column: &str, - uuid: &Uuid, - metrics: &dyn MetricsCollector, - ) -> Result> { - // Caching (including the choice of in-memory vs. serializable state) is - // a plugin implementation detail handled inside `scalar::open_scalar_index`. - let index_meta = self - .load_index(uuid) - .await? - .ok_or_else(|| Error::index(format!("Index with id {} does not exist", uuid)))?; - - scalar::open_scalar_index(self, column, &index_meta, metrics).await - } - - async fn open_vector_index( +impl Dataset { + /// Opens a vector index from its manifest metadata. + /// + /// Unlike [`DatasetIndexInternalExt::open_vector_index`], this does not + /// look the segment up in the manifest, so it can also open segments that + /// have been built but not committed yet. + pub(crate) async fn open_vector_index_from_metadata( &self, column: &str, - uuid: &Uuid, + index_meta: &IndexMetadata, metrics: &dyn MetricsCollector, ) -> Result> { + let uuid = &index_meta.uuid; let frag_reuse_uuid = self.frag_reuse_index_uuid().await; - let index_meta = self - .load_index(uuid) - .await? - .ok_or_else(|| Error::index(format!("Index with id {} does not exist", uuid)))?; - let object_store = self.object_store_for_index(&index_meta).await?; + let object_store = self.object_store_for_index(index_meta).await?; // Check sized cache first (v2+ indices with serializable state). let state_key = IvfIndexStateCacheKey::new(uuid, frag_reuse_uuid.as_ref()); @@ -2832,7 +2793,7 @@ impl DatasetIndexInternalExt for Dataset { } let frag_reuse_index = self.open_frag_reuse_index(metrics).await?; - let index_dir = self.indice_files_dir(&index_meta)?; + let index_dir = self.indice_files_dir(index_meta)?; let index_file = index_dir .clone() .join(uuid.to_string()) @@ -2936,7 +2897,7 @@ impl DatasetIndexInternalExt for Dataset { serde_json::from_str(index_metadata)?; // Resolve the column name and field - let (field_path, field) = resolve_index_column(self.schema(), &index_meta, column)?; + let (field_path, field) = resolve_index_column(self.schema(), index_meta, column)?; let (_, element_type) = get_vector_type(self.schema(), &field_path)?; @@ -3108,7 +3069,105 @@ impl DatasetIndexInternalExt for Dataset { } Ok(index) } +} + +#[async_trait] +impl DatasetIndexInternalExt for Dataset { + async fn open_generic_index( + &self, + column: &str, + uuid: &Uuid, + metrics: &dyn MetricsCollector, + ) -> Result> { + // Checking for cache existence is cheap so we just check the vector caches. + // Scalar indices cache themselves inside `open_scalar_index` (the cache + // key is a plugin detail), so there is no cheap scalar check here. + let frag_reuse_uuid = self.frag_reuse_index_uuid().await; + + // Check sized cache for IvfIndexState (v2+ indices). + let state_key = IvfIndexStateCacheKey::new(uuid, frag_reuse_uuid.as_ref()); + if self.index_cache.get_with_key(&state_key).await.is_some() { + // Reconstruct via open_vector_index which will hit the same sized key. + let index = self.open_vector_index(column, uuid, metrics).await?; + return Ok(index.as_index()); + } + + // Fallback: in-memory cache for legacy indices. + let vector_cache_key = LegacyVectorIndexCacheKey::new(uuid, frag_reuse_uuid.as_ref()); + if let Some(cached) = self.index_cache.get_with_key(&vector_cache_key).await { + return Ok(cached.0.clone().as_index()); + } + + let frag_reuse_cache_key = FragReuseIndexCacheKey::new(uuid, frag_reuse_uuid.as_ref()); + if let Some(index) = self.index_cache.get_with_key(&frag_reuse_cache_key).await { + return Ok(Arc::new(FragReuseIndexHandle(index)).as_index()); + } + + // Sometimes we want to open an index and we don't care if it is a scalar or vector index. + // For example, we might want to get statistics for an index, regardless of type. + // + // We determine if this is a vector index by checking if INDEX_FILE_NAME exists in the + // file list (available since file sizes tracking was added). If the file list is not + // available (older indices), we fall back to checking file existence via HEAD request. + let index_meta = self + .load_index(uuid) + .await? + .ok_or_else(|| Error::index(format!("Index with id {} does not exist", uuid)))?; + + // Check if this is a vector index by looking at the files list + let is_vector_index = if let Some(files) = &index_meta.files { + // If we have file metadata, check if INDEX_FILE_NAME is in the list + files.iter().any(|f| f.path == INDEX_FILE_NAME) + } else { + // Fall back to file existence check for older indices without file metadata + let index_dir = self.indice_files_dir(&index_meta)?; + let index_file = index_dir + .clone() + .join(uuid.to_string()) + .join(INDEX_FILE_NAME); + let object_store = self.object_store_for_index(&index_meta).await?; + object_store.exists(&index_file).await? + }; + + if is_vector_index { + let index = self.open_vector_index(column, uuid, metrics).await?; + Ok(index.as_index()) + } else { + let index = self.open_scalar_index(column, uuid, metrics).await?; + Ok(index.as_index()) + } + } + #[instrument(level = "debug", skip_all)] + async fn open_scalar_index( + &self, + column: &str, + uuid: &Uuid, + metrics: &dyn MetricsCollector, + ) -> Result> { + // Caching (including the choice of in-memory vs. serializable state) is + // a plugin implementation detail handled inside `scalar::open_scalar_index`. + let index_meta = self + .load_index(uuid) + .await? + .ok_or_else(|| Error::index(format!("Index with id {} does not exist", uuid)))?; + + scalar::open_scalar_index(self, column, &index_meta, metrics).await + } + + async fn open_vector_index( + &self, + column: &str, + uuid: &Uuid, + metrics: &dyn MetricsCollector, + ) -> Result> { + let index_meta = self + .load_index(uuid) + .await? + .ok_or_else(|| Error::index(format!("Index with id {} does not exist", uuid)))?; + self.open_vector_index_from_metadata(column, &index_meta, metrics) + .await + } async fn open_logical_vector_index( &self, column: &str, @@ -8369,32 +8428,32 @@ mod tests { .await .unwrap(); - let field_id = dataset.schema().field("vector").unwrap().id; - let seg0 = write_vector_segment_metadata( - &dataset, - "vector_idx", - field_id, - Uuid::new_v4(), - [0_u32], - b"seg0", - ) - .await; - let seg1 = write_vector_segment_metadata( - &dataset, - "vector_idx", - field_id, - Uuid::new_v4(), - [1_u32], - b"seg1", - ) - .await; + // Compatibility validation opens the coexisting segments, so this + // test builds real IVF segments instead of metadata-only fakes. + let params = crate::index::vector::VectorIndexParams::ivf_flat( + 2, + lance_linalg::distance::DistanceType::L2, + ); + let mut segments = Vec::new(); + for fragment in dataset.get_fragments().into_iter().take(2) { + segments.push( + dataset + .create_index_builder(&["vector"], IndexType::Vector, ¶ms) + .name("worker_idx".to_string()) + .fragments(vec![fragment.id() as u32]) + .execute_uncommitted() + .await + .unwrap(), + ); + } + let seg_uuids = segments + .iter() + .map(|segment| segment.uuid) + .collect::>(); + assert_eq!(segments.len(), 2); dataset - .commit_existing_index_segments( - "vector_idx", - "vector", - vec![segment_from_metadata(&seg0), segment_from_metadata(&seg1)], - ) + .commit_existing_index_segments("vector_idx", "vector", segments) .await .unwrap(); @@ -8403,7 +8462,7 @@ mod tests { let committed_uuids = committed.iter().map(|idx| idx.uuid).collect::>(); assert_eq!( committed_uuids, - HashSet::from([seg0.uuid, seg1.uuid]), + seg_uuids.into_iter().collect::>(), "all committed segment uuids should be preserved" ); assert_eq!( diff --git a/rust/lance/src/index/append.rs b/rust/lance/src/index/append.rs index 5403a14fd88..4e3c5d0f5a9 100644 --- a/rust/lance/src/index/append.rs +++ b/rust/lance/src/index/append.rs @@ -1837,7 +1837,7 @@ mod tests { "has type" )] #[tokio::test] - async fn test_vector_append_validates_logical_query_compatibility( + async fn test_commit_index_segments_rejects_logical_query_incompatibility( #[case] first_params: VectorIndexParams, #[case] second_params: VectorIndexParams, #[case] expected_error: &str, @@ -1892,38 +1892,8 @@ mod tests { .unwrap(), ); } - dataset - .commit_existing_index_segments(INDEX_NAME, "vector", segments) - .await - .unwrap(); - - let (appended_batch, _) = clustered_vector_batch( - schema.clone(), - (2 * ROWS_PER_FRAGMENT) as i32, - ROWS_PER_FRAGMENT, - DIMENSION, - 200.0, - ); - dataset - .append( - RecordBatchIterator::new(vec![Ok(appended_batch)], schema), - None, - ) - .await - .unwrap(); - - let version_before = dataset.version().version; - let segments_before = dataset.load_indices_by_name(INDEX_NAME).await.unwrap(); - let object_store = dataset.object_store.clone(); - let directories_before = object_store - .read_dir(dataset.indices_dir()) - .await - .unwrap() - .into_iter() - .collect::>(); - let error = dataset - .optimize_indices(&OptimizeOptions::append()) + .commit_existing_index_segments(INDEX_NAME, "vector", segments) .await .unwrap_err(); assert!( @@ -1937,27 +1907,16 @@ mod tests { .unwrap(); assert_eq!( latest.version().version, - version_before, - "incompatible logical segments must fail before committing" - ); - let mut segments_after = latest.load_indices_by_name(INDEX_NAME).await.unwrap(); - let mut segments_before = segments_before; - for segment in segments_after.iter_mut().chain(segments_before.iter_mut()) { - segment.created_at = None; - } - assert_eq!( - segments_after, segments_before, - "incompatible logical segments must remain unchanged" + dataset.version().version, + "a rejected segment commit must not create a new dataset version" ); - let directories_after = object_store - .read_dir(latest.indices_dir()) - .await - .unwrap() - .into_iter() - .collect::>(); - assert_eq!( - directories_after, directories_before, - "query compatibility validation must run before staging a new segment" + assert!( + latest + .load_indices_by_name(INDEX_NAME) + .await + .unwrap() + .is_empty(), + "a rejected segment commit must not persist any segment" ); } diff --git a/rust/lance/src/index/vector/ivf.rs b/rust/lance/src/index/vector/ivf.rs index 6fcbfc634a3..1d554c50e96 100644 --- a/rust/lance/src/index/vector/ivf.rs +++ b/rust/lance/src/index/vector/ivf.rs @@ -424,7 +424,13 @@ fn vector_index_dimension(index: &dyn VectorIndex) -> usize { } } -fn validate_vector_query_compatibility( +/// Reject vector index segments that cannot serve one logical index query. +/// +/// All segments of a logical vector index are planned and ranked together, so +/// they must agree on the distance metric, vector dimension, sub-index type, +/// and quantizer kind. Independently trained IVF centroids and PQ codebooks +/// may differ; those only affect recall, not the query contract. +pub(crate) fn validate_vector_query_compatibility( indices: &[Arc], operation: &str, ) -> Result<()> { diff --git a/rust/lance/src/index/vector/ivf/v2.rs b/rust/lance/src/index/vector/ivf/v2.rs index 6ceeac8d070..44d5cbd6b0d 100644 --- a/rust/lance/src/index/vector/ivf/v2.rs +++ b/rust/lance/src/index/vector/ivf/v2.rs @@ -3739,6 +3739,150 @@ mod tests { assert!(err.to_string().contains("overlapping fragment coverage")); } + async fn build_ivf_flat_segment( + dataset: &mut Dataset, + metric: DistanceType, + fragment_ids: Vec, + ) -> IndexMetadata { + // Each segment trains its own IVF model, as distributed workers do. + // The build name differs from the committed index name on purpose: + // builders may not reuse the name of an existing index. + let params = VectorIndexParams::ivf_flat(TWO_FRAG_NUM_PARTITIONS, metric); + dataset + .create_index_builder(&["vector"], IndexType::Vector, ¶ms) + .name("worker_idx".to_string()) + .fragments(fragment_ids) + .execute_uncommitted() + .await + .unwrap() + } + + #[tokio::test] + async fn test_commit_index_segments_rejects_mixed_vector_metrics() { + let test_dir = TempStrDir::default(); + let (schema, batches) = make_two_fragment_batches(); + let dataset_uri = format!("{}/mixed_metric_segments", test_dir.as_str()); + let mut dataset = write_dataset_from_batches(&dataset_uri, schema, batches).await; + + let fragments = dataset.get_fragments(); + assert!(fragments.len() >= 2); + let l2_segment = build_ivf_flat_segment( + &mut dataset, + DistanceType::L2, + vec![fragments[0].id() as u32], + ) + .await; + let cosine_segment = build_ivf_flat_segment( + &mut dataset, + DistanceType::Cosine, + vec![fragments[1].id() as u32], + ) + .await; + + let version_before = dataset.manifest.version; + let err = dataset + .commit_existing_index_segments( + "vector_idx", + "vector", + vec![l2_segment, cosine_segment], + ) + .await + .unwrap_err(); + assert!( + matches!(err, lance_core::Error::InvalidInput { .. }), + "expected InvalidInput, got {err:?}" + ); + assert!( + err.to_string().contains("metric"), + "error should name the metric mismatch: {err}" + ); + assert_eq!(dataset.manifest.version, version_before); + assert!( + dataset + .load_indices_by_name("vector_idx") + .await + .unwrap() + .is_empty() + ); + } + + #[tokio::test] + async fn test_commit_index_segments_rejects_metric_mismatch_with_retained_segment() { + let test_dir = TempStrDir::default(); + let (schema, batches) = make_two_fragment_batches(); + let dataset_uri = format!("{}/retained_mixed_metric_segments", test_dir.as_str()); + let mut dataset = write_dataset_from_batches(&dataset_uri, schema, batches).await; + + let fragments = dataset.get_fragments(); + assert!(fragments.len() >= 2); + let l2_segment = build_ivf_flat_segment( + &mut dataset, + DistanceType::L2, + vec![fragments[0].id() as u32], + ) + .await; + dataset + .commit_existing_index_segments("vector_idx", "vector", vec![l2_segment]) + .await + .unwrap(); + + let version_before = dataset.manifest.version; + let cosine_segment = build_ivf_flat_segment( + &mut dataset, + DistanceType::Cosine, + vec![fragments[1].id() as u32], + ) + .await; + let err = dataset + .commit_existing_index_segments("vector_idx", "vector", vec![cosine_segment]) + .await + .unwrap_err(); + assert!( + matches!(err, lance_core::Error::InvalidInput { .. }), + "expected InvalidInput, got {err:?}" + ); + assert!( + err.to_string().contains("metric"), + "error should name the metric mismatch: {err}" + ); + assert_eq!(dataset.manifest.version, version_before); + // The retained L2 segment still serves queries, unmodified. + let indices = dataset.load_indices_by_name("vector_idx").await.unwrap(); + assert_eq!(indices.len(), 1); + } + + #[tokio::test] + async fn test_commit_index_segments_allows_metric_change_on_full_replacement() { + let test_dir = TempStrDir::default(); + let (schema, batches) = make_two_fragment_batches(); + let dataset_uri = format!("{}/full_replacement_metric_change", test_dir.as_str()); + let mut dataset = write_dataset_from_batches(&dataset_uri, schema, batches).await; + + let fragments = dataset.get_fragments(); + assert!(fragments.len() >= 2); + let all_fragment_ids = fragments.iter().map(|f| f.id() as u32).collect::>(); + let l2_segment = build_ivf_flat_segment( + &mut dataset, + DistanceType::L2, + vec![fragments[0].id() as u32], + ) + .await; + dataset + .commit_existing_index_segments("vector_idx", "vector", vec![l2_segment]) + .await + .unwrap(); + + // A new metric is valid when no old segment remains in the index. + let cosine_segment = + build_ivf_flat_segment(&mut dataset, DistanceType::Cosine, all_fragment_ids).await; + dataset + .commit_existing_index_segments("vector_idx", "vector", vec![cosine_segment]) + .await + .unwrap(); + let indices = dataset.load_indices_by_name("vector_idx").await.unwrap(); + assert_eq!(indices.len(), 1); + } + #[tokio::test] async fn test_distributed_vector_build_supports_hnsw_variants() { let test_dir = TempStrDir::default();