diff --git a/Cargo.lock b/Cargo.lock index 8d219f3..3135cef 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2451,8 +2451,8 @@ checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c" [[package]] name = "fsst" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow-array", "rand 0.9.2", @@ -3654,8 +3654,8 @@ checksum = "e037a2e1d8d5fdbd49b16a4ea09d5d6401c1f29eca5ff29d03d3824dba16256a" [[package]] name = "lance" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arc-swap", "arrow", @@ -3726,8 +3726,8 @@ dependencies = [ [[package]] name = "lance-arrow" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow-array", "arrow-buffer", @@ -3749,7 +3749,7 @@ dependencies = [ [[package]] name = "lance-arrow-scalar" version = "58.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow-array", "arrow-buffer", @@ -3763,7 +3763,7 @@ dependencies = [ [[package]] name = "lance-arrow-stats" version = "58.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow-array", "arrow-schema", @@ -3772,8 +3772,8 @@ dependencies = [ [[package]] name = "lance-bitpacking" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrayref", "crunchy", @@ -3814,8 +3814,8 @@ dependencies = [ [[package]] name = "lance-core" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow-array", "arrow-buffer", @@ -3852,8 +3852,8 @@ dependencies = [ [[package]] name = "lance-datafusion" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow", "arrow-array", @@ -3884,8 +3884,8 @@ dependencies = [ [[package]] name = "lance-datagen" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow", "arrow-array", @@ -3902,8 +3902,8 @@ dependencies = [ [[package]] name = "lance-derive" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "proc-macro2", "quote", @@ -3912,8 +3912,8 @@ dependencies = [ [[package]] name = "lance-encoding" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow-arith", "arrow-array", @@ -3946,8 +3946,8 @@ dependencies = [ [[package]] name = "lance-file" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow-arith", "arrow-array", @@ -3978,8 +3978,8 @@ dependencies = [ [[package]] name = "lance-geo" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "datafusion", "geo-traits", @@ -3993,8 +3993,8 @@ dependencies = [ [[package]] name = "lance-index" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arc-swap", "arrow", @@ -4061,8 +4061,8 @@ dependencies = [ [[package]] name = "lance-index-core" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow-array", "arrow-schema", @@ -4084,8 +4084,8 @@ dependencies = [ [[package]] name = "lance-io" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow", "arrow-array", @@ -4124,8 +4124,8 @@ dependencies = [ [[package]] name = "lance-linalg" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow-array", "arrow-schema", @@ -4139,8 +4139,8 @@ dependencies = [ [[package]] name = "lance-namespace" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow", "async-trait", @@ -4166,8 +4166,8 @@ dependencies = [ [[package]] name = "lance-select" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow-array", "arrow-buffer", @@ -4181,8 +4181,8 @@ dependencies = [ [[package]] name = "lance-table" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "arrow", "arrow-array", @@ -4220,8 +4220,8 @@ dependencies = [ [[package]] name = "lance-tokenizer" -version = "11.0.0" -source = "git+https://github.com/lance-format/lance.git?rev=ab6b5bbe#ab6b5bbe46009ed78746b444df8db59a8bc5d842" +version = "11.0.1-beta.0" +source = "git+https://github.com/lance-format/lance.git?rev=356acb0d333c96e970f6f84b97314fc5bc4193f7#356acb0d333c96e970f6f84b97314fc5bc4193f7" dependencies = [ "frostem", "icu_segmenter", diff --git a/Cargo.toml b/Cargo.toml index 34147fe..330bfd4 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -18,14 +18,14 @@ rust-version = "1.91.0" crate-type = ["cdylib", "staticlib", "rlib"] [dependencies] -lance = { git = "https://github.com/lance-format/lance.git", rev = "ab6b5bbe", features = ["substrait"] } -lance-core = { git = "https://github.com/lance-format/lance.git", rev = "ab6b5bbe" } -lance-file = { git = "https://github.com/lance-format/lance.git", rev = "ab6b5bbe" } -lance-index = { git = "https://github.com/lance-format/lance.git", rev = "ab6b5bbe" } -lance-io = { git = "https://github.com/lance-format/lance.git", rev = "ab6b5bbe" } -lance-linalg = { git = "https://github.com/lance-format/lance.git", rev = "ab6b5bbe" } -lance-table = { git = "https://github.com/lance-format/lance.git", rev = "ab6b5bbe" } -lance-datafusion = { git = "https://github.com/lance-format/lance.git", rev = "ab6b5bbe", features = ["substrait"] } +lance = { git = "https://github.com/lance-format/lance.git", rev = "356acb0d333c96e970f6f84b97314fc5bc4193f7", features = ["substrait"] } +lance-core = { git = "https://github.com/lance-format/lance.git", rev = "356acb0d333c96e970f6f84b97314fc5bc4193f7" } +lance-file = { git = "https://github.com/lance-format/lance.git", rev = "356acb0d333c96e970f6f84b97314fc5bc4193f7" } +lance-index = { git = "https://github.com/lance-format/lance.git", rev = "356acb0d333c96e970f6f84b97314fc5bc4193f7" } +lance-io = { git = "https://github.com/lance-format/lance.git", rev = "356acb0d333c96e970f6f84b97314fc5bc4193f7" } +lance-linalg = { git = "https://github.com/lance-format/lance.git", rev = "356acb0d333c96e970f6f84b97314fc5bc4193f7" } +lance-table = { git = "https://github.com/lance-format/lance.git", rev = "356acb0d333c96e970f6f84b97314fc5bc4193f7" } +lance-datafusion = { git = "https://github.com/lance-format/lance.git", rev = "356acb0d333c96e970f6f84b97314fc5bc4193f7", features = ["substrait"] } datafusion = { version = "54.0.0", default-features = false } arrow = { version = "58.0.0", features = ["prettyprint", "ffi"] } arrow-array = "58.0.0" @@ -47,9 +47,9 @@ snafu = "0.9" uuid = { version = "1", features = ["v4"] } [dev-dependencies] -lance = { git = "https://github.com/lance-format/lance.git", rev = "ab6b5bbe", features = ["substrait"] } -lance-datagen = { git = "https://github.com/lance-format/lance.git", rev = "ab6b5bbe" } -lance-file = { git = "https://github.com/lance-format/lance.git", rev = "ab6b5bbe" } +lance = { git = "https://github.com/lance-format/lance.git", rev = "356acb0d333c96e970f6f84b97314fc5bc4193f7", features = ["substrait"] } +lance-datagen = { git = "https://github.com/lance-format/lance.git", rev = "356acb0d333c96e970f6f84b97314fc5bc4193f7" } +lance-file = { git = "https://github.com/lance-format/lance.git", rev = "356acb0d333c96e970f6f84b97314fc5bc4193f7" } tokio = { version = "1", features = ["rt-multi-thread", "macros"] } arrow-array = "58.0.0" arrow-schema = "58.0.0" diff --git a/include/lance/lance.h b/include/lance/lance.h index 9aa66d0..e23c184 100644 --- a/include/lance/lance.h +++ b/include/lance/lance.h @@ -1710,6 +1710,58 @@ int32_t lance_index_segment_metadata_fragment_ids( /** Free parsed segment metadata. NULL-safe. */ void lance_index_segment_metadata_free(LanceIndexSegmentMetadata* metadata); +/** + * Commit previously built uncommitted index segments as one logical index. + * + * `segment_metadata_bytes[i]` must point to + * `segment_metadata_lens[i]` bytes of protobuf-encoded IndexMetadata produced + * by lance_index_segment_builder_execute_uncommitted() (typically built on + * distributed workers). All segments are registered under `index_name` on + * `column` in a single commit, so the dataset version increases by exactly + * one on success. + * + * The segment set is validated by the Lance core and rejected with + * LANCE_ERR_INVALID_ARGUMENT when it is empty, contains duplicate segment + * UUIDs, or has overlapping fragment coverage. All segments must share one + * index type, and the commit fails if `column` does not exist. Every segment + * must declare `column` as its keyed field — that is, have been built for + * `column` — or the commit fails with LANCE_ERR_INVALID_ARGUMENT. + * Vector segments that will coexist (incoming segments and retained existing + * segments) must have compatible distance metrics, dimensions, sub-index + * types, and quantizer kinds. Independently trained IVF centroids and PQ + * codebooks may differ. Incompatible segments are rejected with + * LANCE_ERR_INVALID_ARGUMENT without changing the dataset version or index. + * + * Replacement is automatic and coverage-driven — there is no replace flag: + * existing same-name segments of the same index type whose fragment coverage + * is fully covered by the incoming set are replaced, while existing segments + * covering disjoint fragments are retained as additional deltas of the + * logical index. A commit that would orphan fragments from an existing + * segment (partial overlap) is rejected. A commit whose index type differs + * from the existing same-name index replaces that index entirely, and + * therefore requires the incoming segments to cover every current fragment; + * a partial-coverage type change is rejected with LANCE_ERR_INVALID_ARGUMENT. + * Vector compatibility is checked after selecting replacements, so a complete + * replacement may change the metric without conflicting with removed segments. + * + * @param dataset Open dataset (mutated; same handle remains valid). + * @param index_name Logical index name; must not be NULL or empty. + * @param column Indexed column; must not be NULL or empty. + * @param segment_metadata_bytes Array of pointers to encoded IndexMetadata. + * @param segment_metadata_lens Array of byte lengths, parallel to + * segment_metadata_bytes. + * @param segment_count Number of segments; must be > 0. + * @return 0 on success, -1 on error. + */ +int32_t lance_dataset_commit_index_segments( + LanceDataset* dataset, + const char* index_name, + const char* column, + const uint8_t* const* segment_metadata_bytes, + const size_t* segment_metadata_lens, + size_t segment_count +); + /** Drop an index by name. Returns -1 (NOT_FOUND) if no such index. */ int32_t lance_dataset_drop_index(LanceDataset* dataset, const char* name); diff --git a/include/lance/lance.hpp b/include/lance/lance.hpp index 027e6d5..e17bbe9 100644 --- a/include/lance/lance.hpp +++ b/include/lance/lance.hpp @@ -936,6 +936,41 @@ class Dataset { return out; } + /// Commit previously built uncommitted index segments as one logical + /// index under `index_name` on `column`. Each entry of + /// `segment_metadata` is the protobuf-encoded IndexMetadata produced by + /// `IndexSegmentBuilder::execute_uncommitted()`. The commit is a single + /// dataset version bump. Every segment must have been built for `column`. + /// Coexisting vector segments, including retained existing segments, must + /// have compatible metrics, dimensions, sub-index types, and quantizer + /// kinds; independently trained IVF centroids and PQ codebooks may differ. + /// Replacement of existing same-name segments is automatic and + /// coverage-driven: fully covered segments are replaced, disjoint ones + /// are retained as deltas, and partial overlap is rejected. A commit + /// whose index type differs from the existing same-name index replaces + /// that index entirely, so it must cover every current fragment; a + /// partial-coverage type change is rejected. + /// Fully replaced vector segments do not constrain the new metric. + /// Throws lance::Error on validation failures (empty set, duplicate + /// segment UUIDs, overlapping fragment coverage, unknown or mismatched + /// column, incompatible vector segments), leaving the version and index + /// unchanged. + void commit_index_segments( + const std::string& index_name, + const std::string& column, + const std::vector>& segment_metadata) { + std::vector bytes(segment_metadata.size()); + std::vector lens(segment_metadata.size()); + for (size_t i = 0; i < segment_metadata.size(); ++i) { + bytes[i] = segment_metadata[i].data(); + lens[i] = segment_metadata[i].size(); + } + if (lance_dataset_commit_index_segments( + handle_.get(), index_name.c_str(), column.c_str(), bytes.data(), + lens.data(), segment_metadata.size()) != 0) + check_error(); + } + /// Access the underlying C handle (does not transfer ownership). const LanceDataset* c_handle() const { return handle_.get(); } diff --git a/src/index_segment.rs b/src/index_segment.rs index d4e143c..26db56c 100644 --- a/src/index_segment.rs +++ b/src/index_segment.rs @@ -1048,35 +1048,115 @@ pub unsafe extern "C" fn lance_index_segment_builder_free(builder: *mut LanceInd } } -/// Parse a protobuf-encoded Lance `IndexMetadata` value. +/// Commit previously built uncommitted index segments as one logical index. +/// +/// All segments are committed in a single dataset version. Validation of the +/// segment set (distinct UUIDs, disjoint fragment coverage, consistent index +/// details) is performed by the Lance core; replacement of existing same-name +/// segments is automatic and coverage-driven. #[unsafe(no_mangle)] -pub unsafe extern "C" fn lance_index_segment_metadata_parse( - bytes: *const u8, - len: usize, - out_metadata: *mut *mut LanceIndexSegmentMetadata, +pub unsafe extern "C" fn lance_dataset_commit_index_segments( + dataset: *mut LanceDataset, + index_name: *const c_char, + column: *const c_char, + segment_metadata_bytes: *const *const u8, + segment_metadata_lens: *const usize, + segment_count: usize, ) -> i32 { ffi_try!( - unsafe { parse_metadata_inner(bytes, len, out_metadata) }, + unsafe { + commit_index_segments_inner( + dataset, + index_name, + column, + segment_metadata_bytes, + segment_metadata_lens, + segment_count, + ) + }, neg ) } -unsafe fn parse_metadata_inner( - bytes: *const u8, - len: usize, - out_metadata: *mut *mut LanceIndexSegmentMetadata, +unsafe fn commit_index_segments_inner( + dataset: *mut LanceDataset, + index_name: *const c_char, + column: *const c_char, + segment_metadata_bytes: *const *const u8, + segment_metadata_lens: *const usize, + segment_count: usize, ) -> Result { - if bytes.is_null() || len == 0 || out_metadata.is_null() { + if dataset.is_null() || index_name.is_null() || column.is_null() { + return Err(invalid_input( + "dataset, index_name, and column must not be NULL", + )); + } + let index_name = unsafe { helpers::parse_c_string(index_name)? } + .filter(|value| !value.is_empty()) + .ok_or_else(|| invalid_input("index_name must not be NULL or empty"))?; + let column = unsafe { helpers::parse_c_string(column)? } + .filter(|value| !value.is_empty()) + .ok_or_else(|| invalid_input("column must not be NULL or empty"))?; + if segment_count == 0 { + return Err(invalid_input( + "segment_count must be > 0; at least one index segment is required to commit an index", + )); + } + if segment_metadata_bytes.is_null() || segment_metadata_lens.is_null() { return Err(invalid_input(format!( - "bytes must be non-NULL, len must be > 0, and out_metadata must be non-NULL; bytes={bytes:p}, len={len}, out_metadata={out_metadata:p}" + "segment_metadata_bytes and segment_metadata_lens must not be NULL when segment_count is {segment_count}" ))); } - if len > isize::MAX as usize { + if segment_count > isize::MAX as usize / std::mem::size_of::<*const u8>() { return Err(invalid_input(format!( - "len {len} exceeds the maximum addressable byte slice length" + "segment_count {segment_count} exceeds the maximum addressable pointer slice length" ))); } - let proto = pb::IndexMetadata::decode(unsafe { slice::from_raw_parts(bytes, len) }) + let bytes_array = unsafe { slice::from_raw_parts(segment_metadata_bytes, segment_count) }; + let lens_array = unsafe { slice::from_raw_parts(segment_metadata_lens, segment_count) }; + let mut segments = Vec::with_capacity(segment_count); + for (position, (&bytes, &len)) in bytes_array.iter().zip(lens_array.iter()).enumerate() { + if bytes.is_null() || len == 0 { + return Err(invalid_input(format!( + "segment_metadata_bytes[{position}] must be non-NULL and segment_metadata_lens[{position}] must be > 0; bytes={bytes:p}, len={len}" + ))); + } + if len > isize::MAX as usize { + return Err(invalid_input(format!( + "segment_metadata_lens[{position}]={len} exceeds the maximum addressable byte slice length" + ))); + } + let metadata = decode_segment_metadata(unsafe { slice::from_raw_parts(bytes, len) }) + .map_err(|error| { + invalid_input(format!("segment_metadata_bytes[{position}]: {error}")) + })?; + segments.push(metadata); + } + + let ds = unsafe { &*dataset }; + ds.with_mut(|dataset| { + block_on(dataset.commit_existing_index_segments(index_name, column, segments)) + })?; + Ok(0) +} + +/// Parse a protobuf-encoded Lance `IndexMetadata` value. +#[unsafe(no_mangle)] +pub unsafe extern "C" fn lance_index_segment_metadata_parse( + bytes: *const u8, + len: usize, + out_metadata: *mut *mut LanceIndexSegmentMetadata, +) -> i32 { + ffi_try!( + unsafe { parse_metadata_inner(bytes, len, out_metadata) }, + neg + ) +} + +/// Decode protobuf-encoded `IndexMetadata` bytes into the table-format type, +/// rejecting values whose ranges cannot be represented safely. +fn decode_segment_metadata(bytes: &[u8]) -> Result { + let proto = pb::IndexMetadata::decode(bytes) .map_err(|error| invalid_input(format!("invalid IndexMetadata protobuf: {error}")))?; if let Some(created_at) = proto.created_at { let created_at = i64::try_from(created_at).map_err(|_| { @@ -1106,7 +1186,25 @@ unsafe fn parse_metadata_inner( "IndexMetadata fields[{position}] must be >= 0, got {field_id}" ))); } - let metadata = IndexMetadata::try_from(proto)?; + IndexMetadata::try_from(proto) +} + +unsafe fn parse_metadata_inner( + bytes: *const u8, + len: usize, + out_metadata: *mut *mut LanceIndexSegmentMetadata, +) -> Result { + if bytes.is_null() || len == 0 || out_metadata.is_null() { + return Err(invalid_input(format!( + "bytes must be non-NULL, len must be > 0, and out_metadata must be non-NULL; bytes={bytes:p}, len={len}, out_metadata={out_metadata:p}" + ))); + } + if len > isize::MAX as usize { + return Err(invalid_input(format!( + "len {len} exceeds the maximum addressable byte slice length" + ))); + } + let metadata = decode_segment_metadata(unsafe { slice::from_raw_parts(bytes, len) })?; let name = CString::new(metadata.name.as_str()) .map_err(|_| invalid_input("index metadata name contains an embedded NUL byte"))?; let index_details_type_url = metadata diff --git a/tests/c_api_test.rs b/tests/c_api_test.rs index 0424c97..36c4392 100644 --- a/tests/c_api_test.rs +++ b/tests/c_api_test.rs @@ -4733,6 +4733,1157 @@ fn test_index_segment_options_reject_invalid_fragment_and_train_combinations() { unsafe { lance_dataset_close(dataset) }; } +/// Scalar (bitmap) segment builds reserve tens of MB from the shared +/// datafusion spill pool; serialize them so parallel commit tests cannot +/// exhaust the pool. +static SCALAR_SEGMENT_BUILD_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(()); + +/// Build one uncommitted scalar segment on the `id` column and return the +/// malloc-owned protobuf metadata bytes (free with `lance_free_bytes`). +fn build_scalar_segment_bytes( + dataset: *mut LanceDataset, + index_name: &CString, + index_type: LanceScalarIndexType, + fragment_ids: Option<&[u32]>, +) -> (*mut u8, usize) { + let _build_guard = SCALAR_SEGMENT_BUILD_LOCK + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + let column = c_str("id"); + let options = LanceIndexSegmentBuildOptions { + fragment_ids: fragment_ids.map_or(ptr::null(), |ids| ids.as_ptr()), + fragment_count: fragment_ids.map_or(0, |ids| ids.len()), + index_uuid: ptr::null(), + ivf_centroids: ptr::null_mut(), + ivf_centroids_schema: ptr::null(), + pq_codebook: ptr::null_mut(), + pq_codebook_schema: ptr::null(), + mode: LanceIndexSegmentBuildMode::Auto as i32, + }; + let builder = unsafe { + lance_index_segment_builder_new_scalar( + dataset, + column.as_ptr(), + index_name.as_ptr(), + index_type as i32, + ptr::null(), + &options, + ) + }; + assert!(!builder.is_null()); + let mut bytes = ptr::null_mut(); + let mut len = 0_usize; + assert_eq!( + unsafe { lance_index_segment_builder_execute_uncommitted(builder, &mut bytes, &mut len) }, + 0, + "{}", + unsafe { std::ffi::CStr::from_ptr(lance_last_error_message()).to_string_lossy() } + ); + unsafe { lance_index_segment_builder_free(builder) }; + (bytes, len) +} + +/// Read the UUID of an encoded segment without freeing the bytes. +fn segment_uuid(bytes: *const u8, len: usize) -> [u8; 16] { + let mut metadata = ptr::null_mut(); + assert_eq!( + unsafe { lance_index_segment_metadata_parse(bytes, len, &mut metadata) }, + 0 + ); + let mut uuid = [0_u8; 16]; + assert_eq!( + unsafe { lance_index_segment_metadata_uuid(metadata, uuid.as_mut_ptr()) }, + 0 + ); + unsafe { lance_index_segment_metadata_free(metadata) }; + uuid +} + +fn build_vector_segment_bytes( + dataset: *mut LanceDataset, + metric: LanceMetricType, + fragment_ids: &[u32], +) -> Vec { + let column = c_str("embedding"); + let params = LanceVectorIndexSegmentParams { + index_type: LanceVectorIndexType::IvfFlat as i32, + metric: metric as i32, + num_partitions: 2, + num_sub_vectors: 0, + num_bits: 0, + max_iterations: 2, + hnsw_m: 0, + hnsw_ef_construction: 0, + sample_rate: 16, + }; + let options = LanceIndexSegmentBuildOptions { + fragment_ids: fragment_ids.as_ptr(), + fragment_count: fragment_ids.len(), + index_uuid: ptr::null(), + ivf_centroids: ptr::null_mut(), + ivf_centroids_schema: ptr::null(), + pq_codebook: ptr::null_mut(), + pq_codebook_schema: ptr::null(), + mode: LanceIndexSegmentBuildMode::Auto as i32, + }; + let builder = unsafe { + lance_index_segment_builder_new_vector( + dataset, + column.as_ptr(), + c_str("worker_idx").as_ptr(), + ¶ms, + &options, + ) + }; + assert!(!builder.is_null(), "{}", take_last_error_message()); + let mut bytes = ptr::null_mut(); + let mut len = 0; + assert_eq!( + unsafe { lance_index_segment_builder_execute_uncommitted(builder, &mut bytes, &mut len) }, + 0, + "{}", + take_last_error_message() + ); + let metadata = unsafe { std::slice::from_raw_parts(bytes, len) }.to_vec(); + unsafe { + lance_free_bytes(bytes); + lance_index_segment_builder_free(builder); + } + metadata +} + +fn commit_vector_segments(dataset: *mut LanceDataset, segments: &[&[u8]]) -> i32 { + let bytes = segments + .iter() + .map(|segment| segment.as_ptr()) + .collect::>(); + let lengths = segments + .iter() + .map(|segment| segment.len()) + .collect::>(); + unsafe { + lance_dataset_commit_index_segments( + dataset, + c_str("embedding_idx").as_ptr(), + c_str("embedding").as_ptr(), + bytes.as_ptr(), + lengths.as_ptr(), + segments.len(), + ) + } +} + +fn vector_segment_query_ids(dataset: *mut LanceDataset, use_index: bool) -> Vec { + let scanner = unsafe { lance_scanner_new(dataset, ptr::null(), ptr::null()) }; + assert!(!scanner.is_null()); + // Offset from row 5 to avoid tied distances in the top three. + let query: [f32; 8] = std::array::from_fn(|component| 5.25 + component as f32 / 8.0); + assert_eq!( + unsafe { + lance_scanner_nearest( + scanner, + c_str("embedding").as_ptr(), + query.as_ptr().cast(), + query.len(), + LanceDataType::Float32 as i32, + 3, + ) + }, + 0 + ); + assert_eq!( + unsafe { lance_scanner_set_metric(scanner, LanceMetricType::L2 as i32) }, + 0 + ); + assert_eq!( + unsafe { lance_scanner_set_use_index(scanner, use_index) }, + 0 + ); + // Probe every partition so the assertion does not depend on ANN recall. + assert_eq!(unsafe { lance_scanner_set_nprobes(scanner, 2) }, 0); + let mut stream = FFI_ArrowArrayStream::empty(); + assert_eq!( + unsafe { lance_scanner_to_arrow_stream(scanner, &mut stream) }, + 0 + ); + let reader = unsafe { ArrowArrayStreamReader::from_raw(&mut stream) }.unwrap(); + let ids = reader + .flat_map(|batch| { + batch + .unwrap() + .column_by_name("id") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap() + .values() + .to_vec() + }) + .collect(); + unsafe { lance_scanner_close(scanner) }; + ids +} + +fn assert_mixed_vector_metrics_rejected(retain_existing: bool) { + let (_tmp, uri) = create_multi_fragment_vector_dataset(2, 64, 8, false); + let dataset = unsafe { lance_dataset_open(c_str(&uri).as_ptr(), ptr::null(), 0) }; + assert!(!dataset.is_null()); + let l2 = build_vector_segment_bytes(dataset, LanceMetricType::L2, &[0]); + let cosine = build_vector_segment_bytes(dataset, LanceMetricType::Cosine, &[1]); + if retain_existing { + assert_eq!(commit_vector_segments(dataset, &[&l2]), 0); + } + let version_before = unsafe { lance_dataset_version(dataset) }; + let incoming: Vec<&[u8]> = if retain_existing { + vec![&cosine] + } else { + vec![&l2, &cosine] + }; + assert_eq!( + commit_vector_segments(dataset, &incoming), + -1, + "incompatible vector metrics must be rejected before committing" + ); + assert_eq!(lance_last_error_code(), LanceErrorCode::InvalidArgument); + let message = take_last_error_message(); + assert!(message.to_lowercase().contains("metric"), "{message}"); + assert_eq!(unsafe { lance_dataset_version(dataset) }, version_before); + + // Check both the caller's handle and a fresh reader of the persisted manifest. + let reopened = unsafe { lance_dataset_open(c_str(&uri).as_ptr(), ptr::null(), 0) }; + assert!(!reopened.is_null()); + for handle in [dataset, reopened] { + assert_eq!(unsafe { lance_dataset_version(handle) }, version_before); + assert_eq!( + unsafe { lance_dataset_index_count(handle) }, + if retain_existing { 1 } else { 0 } + ); + if retain_existing { + let mut uuid = [0; 16]; + let mut count = 0; + assert_eq!( + unsafe { + lance_dataset_index_segments( + handle, + c_str("embedding_idx").as_ptr(), + uuid.as_mut_ptr(), + 1, + &mut count, + ) + }, + 0 + ); + assert_eq!(count, 1); + assert_eq!(uuid, segment_uuid(l2.as_ptr(), l2.len())); + assert_eq!( + vector_segment_query_ids(handle, true), + vector_segment_query_ids(handle, false) + ); + } + } + unsafe { + lance_dataset_close(reopened); + lance_dataset_close(dataset); + } +} + +#[test] +fn test_commit_index_segments_rejects_mixed_vector_metrics() { + assert_mixed_vector_metrics_rejected(false); +} + +#[test] +fn test_commit_index_segments_rejects_metric_mismatch_with_retained_segment() { + assert_mixed_vector_metrics_rejected(true); +} + +#[test] +fn test_commit_index_segments_vector_delta_and_complete_metric_replacement() { + let (_tmp, uri) = create_multi_fragment_vector_dataset(2, 64, 8, false); + let dataset = unsafe { lance_dataset_open(c_str(&uri).as_ptr(), ptr::null(), 0) }; + assert!(!dataset.is_null()); + let version_before = unsafe { lance_dataset_version(dataset) }; + // Each worker independently trains its IVF model on different data. + let first = build_vector_segment_bytes(dataset, LanceMetricType::L2, &[0]); + assert_eq!(commit_vector_segments(dataset, &[&first]), 0); + let second = build_vector_segment_bytes(dataset, LanceMetricType::L2, &[1]); + assert_eq!(commit_vector_segments(dataset, &[&second]), 0); + assert_eq!( + unsafe { lance_dataset_version(dataset) }, + version_before + 2 + ); + assert_eq!(unsafe { lance_dataset_index_count(dataset) }, 2); + let indexed_ids = vector_segment_query_ids(dataset, true); + assert_eq!(indexed_ids, [5, 6, 4]); + assert_eq!(indexed_ids, vector_segment_query_ids(dataset, false)); + + // A new metric is valid when no old segment will remain in the index. + let replacement = build_vector_segment_bytes(dataset, LanceMetricType::Cosine, &[0, 1]); + assert_eq!( + commit_vector_segments(dataset, &[&replacement]), + 0, + "{}", + take_last_error_message() + ); + assert_eq!( + unsafe { lance_dataset_version(dataset) }, + version_before + 3 + ); + assert_eq!(unsafe { lance_dataset_index_count(dataset) }, 1); + let mut uuid = [0; 16]; + let mut count = 0; + assert_eq!( + unsafe { + lance_dataset_index_segments( + dataset, + c_str("embedding_idx").as_ptr(), + uuid.as_mut_ptr(), + 1, + &mut count, + ) + }, + 0 + ); + assert_eq!(count, 1); + assert_eq!(uuid, segment_uuid(replacement.as_ptr(), replacement.len())); + unsafe { lance_dataset_close(dataset) }; +} + +#[test] +fn test_commit_index_segments_happy_path_multi_segment_vector_index() { + let (_tmp, uri) = create_multi_fragment_vector_dataset(2, 64, 8, false); + let uri_c = c_str(&uri); + let dataset = unsafe { lance_dataset_open(uri_c.as_ptr(), ptr::null(), 0) }; + assert!(!dataset.is_null()); + let mut fragment_ids = [0_u64; 2]; + assert_eq!( + unsafe { lance_dataset_fragment_ids(dataset, fragment_ids.as_mut_ptr()) }, + 0 + ); + + let column = c_str("embedding"); + let index_name = c_str("embedding_distributed_idx"); + let params = LanceVectorIndexSegmentParams { + index_type: LanceVectorIndexType::IvfFlat as i32, + metric: LanceMetricType::L2 as i32, + num_partitions: 2, + num_sub_vectors: 0, + num_bits: 0, + max_iterations: 2, + hnsw_m: 0, + hnsw_ef_construction: 0, + sample_rate: 16, + }; + + // Build one uncommitted segment per fragment (the distributed workers). + let mut segment_bytes = [ptr::null_mut(); 2]; + let mut segment_lens = [0_usize; 2]; + let mut expected_uuids = Vec::new(); + for (worker, fragment_id) in fragment_ids.iter().enumerate() { + let fragment = *fragment_id as u32; + let options = LanceIndexSegmentBuildOptions { + fragment_ids: &fragment, + fragment_count: 1, + index_uuid: ptr::null(), + ivf_centroids: ptr::null_mut(), + ivf_centroids_schema: ptr::null(), + pq_codebook: ptr::null_mut(), + pq_codebook_schema: ptr::null(), + mode: LanceIndexSegmentBuildMode::Auto as i32, + }; + let builder = unsafe { + lance_index_segment_builder_new_vector( + dataset, + column.as_ptr(), + index_name.as_ptr(), + ¶ms, + &options, + ) + }; + assert!(!builder.is_null()); + assert_eq!( + unsafe { + lance_index_segment_builder_execute_uncommitted( + builder, + &mut segment_bytes[worker], + &mut segment_lens[worker], + ) + }, + 0, + "{}", + unsafe { std::ffi::CStr::from_ptr(lance_last_error_message()).to_string_lossy() } + ); + unsafe { lance_index_segment_builder_free(builder) }; + expected_uuids.push(segment_uuid(segment_bytes[worker], segment_lens[worker])); + } + + let version_before = unsafe { lance_dataset_version(dataset) }; + assert_eq!( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr().cast::<*const u8>(), + segment_lens.as_ptr(), + segment_lens.len(), + ) + }, + 0, + "{}", + unsafe { std::ffi::CStr::from_ptr(lance_last_error_message()).to_string_lossy() } + ); + + // One commit for the whole segment set: exactly one version bump. + assert_eq!( + unsafe { lance_dataset_version(dataset) }, + version_before + 1 + ); + // index_count counts physical segments; both segments share one logical + // index name, which index_segment_count/index_segments resolve below. + assert_eq!(unsafe { lance_dataset_index_count(dataset) }, 2); + assert_eq!( + unsafe { lance_dataset_index_segment_count(dataset, index_name.as_ptr()) }, + 2 + ); + let mut committed_uuids = [0_u8; 32]; + let mut committed_count = 0_u64; + assert_eq!( + unsafe { + lance_dataset_index_segments( + dataset, + index_name.as_ptr(), + committed_uuids.as_mut_ptr(), + 2, + &mut committed_count, + ) + }, + 0 + ); + assert_eq!(committed_count, 2); + for (worker, expected_uuid) in expected_uuids.iter().enumerate() { + assert_eq!( + &committed_uuids[worker * 16..(worker + 1) * 16], + expected_uuid + ); + } + + // A k-NN query resolves through the committed multi-segment index. + let scanner = unsafe { lance_scanner_new(dataset, ptr::null(), ptr::null()) }; + assert!(!scanner.is_null()); + // Row 5's vector: component i is 5 + i/8, so the nearest neighbor is row 5. + let query: Vec = (0..8).map(|i| 5.0 + i as f32 / 8.0).collect(); + assert_eq!( + unsafe { + lance_scanner_nearest( + scanner, + column.as_ptr(), + query.as_ptr() as *const c_void, + 8, + LanceDataType::Float32 as i32, + 3, + ) + }, + 0 + ); + let mut stream = FFI_ArrowArrayStream::empty(); + assert_eq!( + unsafe { lance_scanner_to_arrow_stream(scanner, &mut stream) }, + 0, + "{}", + unsafe { std::ffi::CStr::from_ptr(lance_last_error_message()).to_string_lossy() } + ); + let reader = unsafe { ArrowArrayStreamReader::from_raw(&mut stream) }.unwrap(); + let ids = reader + .flat_map(|batch| { + let batch = batch.unwrap(); + batch + .column_by_name("id") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap() + .values() + .to_vec() + }) + .collect::>(); + assert_eq!(ids.len(), 3); + assert_eq!( + ids[0], 5, + "nearest neighbor of row 5's vector must be row 5" + ); + + unsafe { + lance_scanner_close(scanner); + for bytes in segment_bytes { + lance_free_bytes(bytes); + } + lance_dataset_close(dataset); + } +} + +#[test] +fn test_commit_index_segments_rejects_duplicate_segment_uuids() { + let (_tmp, uri) = create_many_small_fragments(2); + let uri_c = c_str(&uri); + let dataset = unsafe { lance_dataset_open(uri_c.as_ptr(), ptr::null(), 0) }; + let index_name = c_str("id_idx"); + let fragment = 0_u32; + let (bytes, len) = build_scalar_segment_bytes( + dataset, + &index_name, + LanceScalarIndexType::Bitmap, + Some(&[fragment]), + ); + + let column = c_str("id"); + let version_before = unsafe { lance_dataset_version(dataset) }; + // The same encoded segment (hence the same UUID) appears twice in the set. + let segment_bytes = [bytes as *const u8, bytes as *const u8]; + let segment_lens = [len, len]; + assert_eq!( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 2, + ) + }, + -1 + ); + assert_eq!(lance_last_error_code(), LanceErrorCode::InvalidArgument); + assert_eq!(unsafe { lance_dataset_version(dataset) }, version_before); + assert_eq!(unsafe { lance_dataset_index_count(dataset) }, 0); + + unsafe { + lance_free_bytes(bytes); + lance_dataset_close(dataset); + } +} + +#[test] +fn test_commit_index_segments_rejects_overlapping_fragment_coverage() { + let (_tmp, uri) = create_many_small_fragments(2); + let uri_c = c_str(&uri); + let dataset = unsafe { lance_dataset_open(uri_c.as_ptr(), ptr::null(), 0) }; + let index_name = c_str("id_idx"); + let fragment = 0_u32; + let (bytes_a, len_a) = build_scalar_segment_bytes( + dataset, + &index_name, + LanceScalarIndexType::Bitmap, + Some(&[fragment]), + ); + let (bytes_b, len_b) = build_scalar_segment_bytes( + dataset, + &index_name, + LanceScalarIndexType::Bitmap, + Some(&[fragment]), + ); + assert_ne!(segment_uuid(bytes_a, len_a), segment_uuid(bytes_b, len_b)); + + let column = c_str("id"); + let version_before = unsafe { lance_dataset_version(dataset) }; + let segment_bytes = [bytes_a as *const u8, bytes_b as *const u8]; + let segment_lens = [len_a, len_b]; + assert_eq!( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 2, + ) + }, + -1 + ); + assert_eq!(lance_last_error_code(), LanceErrorCode::InvalidArgument); + assert_eq!(unsafe { lance_dataset_version(dataset) }, version_before); + assert_eq!(unsafe { lance_dataset_index_count(dataset) }, 0); + + unsafe { + lance_free_bytes(bytes_a); + lance_free_bytes(bytes_b); + lance_dataset_close(dataset); + } +} + +#[test] +fn test_commit_index_segments_rejects_malformed_metadata() { + let (_tmp, uri) = create_many_small_fragments(2); + let uri_c = c_str(&uri); + let dataset = unsafe { lance_dataset_open(uri_c.as_ptr(), ptr::null(), 0) }; + let index_name = c_str("id_idx"); + let column = c_str("id"); + let fragment = 0_u32; + let (valid_bytes, valid_len) = build_scalar_segment_bytes( + dataset, + &index_name, + LanceScalarIndexType::Bitmap, + Some(&[fragment]), + ); + + // Garbage that is not a protobuf message at all. + let garbage = [0xab_u8, 0xcd, 0xef]; + let segment_bytes = [garbage.as_ptr()]; + let segment_lens = [garbage.len()]; + assert_eq!( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 1, + ) + }, + -1 + ); + assert_eq!(lance_last_error_code(), LanceErrorCode::InvalidArgument); + + // A valid message truncated mid-record. + let truncated_len = valid_len / 2; + let segment_bytes = [valid_bytes as *const u8]; + let segment_lens = [truncated_len]; + assert_eq!( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 1, + ) + }, + -1 + ); + assert_eq!(lance_last_error_code(), LanceErrorCode::InvalidArgument); + assert_eq!(unsafe { lance_dataset_index_count(dataset) }, 0); + + unsafe { + lance_free_bytes(valid_bytes); + lance_dataset_close(dataset); + } +} + +#[test] +fn test_commit_index_segments_validates_null_and_empty_inputs() { + let (_tmp, uri) = create_many_small_fragments(2); + let uri_c = c_str(&uri); + let dataset = unsafe { lance_dataset_open(uri_c.as_ptr(), ptr::null(), 0) }; + let index_name = c_str("id_idx"); + let column = c_str("id"); + let empty_name = c_str(""); + // Every case below is rejected at the FFI boundary before the metadata + // bytes are decoded, so a placeholder buffer is sufficient — no real + // segment build is needed. + let placeholder = [0x01_u8, 0x02, 0x03]; + let segment_bytes = [placeholder.as_ptr()]; + let segment_lens = [placeholder.len()]; + let version_before = unsafe { lance_dataset_version(dataset) }; + + let expect_invalid = |rc: i32, case: &str| { + assert_eq!(rc, -1, "{case}"); + assert_eq!( + lance_last_error_code(), + LanceErrorCode::InvalidArgument, + "{case}" + ); + }; + + expect_invalid( + unsafe { + lance_dataset_commit_index_segments( + ptr::null_mut(), + index_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 1, + ) + }, + "NULL dataset", + ); + expect_invalid( + unsafe { + lance_dataset_commit_index_segments( + dataset, + ptr::null(), + column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 1, + ) + }, + "NULL index_name", + ); + expect_invalid( + unsafe { + lance_dataset_commit_index_segments( + dataset, + empty_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 1, + ) + }, + "empty index_name", + ); + expect_invalid( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + ptr::null(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 1, + ) + }, + "NULL column", + ); + expect_invalid( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 0, + ) + }, + "segment_count 0", + ); + expect_invalid( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + ptr::null(), + segment_lens.as_ptr(), + 1, + ) + }, + "NULL segment_metadata_bytes", + ); + expect_invalid( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr(), + ptr::null(), + 1, + ) + }, + "NULL segment_metadata_lens", + ); + let null_element = [ptr::null()]; + expect_invalid( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + null_element.as_ptr(), + segment_lens.as_ptr(), + 1, + ) + }, + "NULL segment element", + ); + let zero_len = [0_usize]; + expect_invalid( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr(), + zero_len.as_ptr(), + 1, + ) + }, + "zero-length segment element", + ); + + // None of the rejected calls touched the dataset. + assert_eq!(unsafe { lance_dataset_version(dataset) }, version_before); + assert_eq!(unsafe { lance_dataset_index_count(dataset) }, 0); + + unsafe { lance_dataset_close(dataset) }; +} + +#[test] +fn test_commit_index_segments_rejects_unknown_column() { + let (_tmp, uri) = create_many_small_fragments(2); + let uri_c = c_str(&uri); + let dataset = unsafe { lance_dataset_open(uri_c.as_ptr(), ptr::null(), 0) }; + let index_name = c_str("id_idx"); + let missing_column = c_str("no_such_column"); + let fragment = 0_u32; + let (bytes, len) = build_scalar_segment_bytes( + dataset, + &index_name, + LanceScalarIndexType::Bitmap, + Some(&[fragment]), + ); + + let segment_bytes = [bytes as *const u8]; + let segment_lens = [len]; + assert_eq!( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + missing_column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 1, + ) + }, + -1 + ); + assert_eq!(unsafe { lance_dataset_index_count(dataset) }, 0); + + unsafe { + lance_free_bytes(bytes); + lance_dataset_close(dataset); + } +} + +#[test] +fn test_commit_index_segments_replaces_fully_covered_segments() { + let (_tmp, uri) = create_many_small_fragments(2); + let uri_c = c_str(&uri); + let dataset = unsafe { lance_dataset_open(uri_c.as_ptr(), ptr::null(), 0) }; + let mut fragment_ids = [0_u64; 2]; + assert_eq!( + unsafe { lance_dataset_fragment_ids(dataset, fragment_ids.as_mut_ptr()) }, + 0 + ); + let all_fragments = [fragment_ids[0] as u32, fragment_ids[1] as u32]; + let index_name = c_str("id_idx"); + let column = c_str("id"); + + // Commit one segment covering every fragment. + let (bytes_a, len_a) = build_scalar_segment_bytes( + dataset, + &index_name, + LanceScalarIndexType::Bitmap, + Some(&all_fragments), + ); + let uuid_a = segment_uuid(bytes_a, len_a); + let segment_bytes = [bytes_a as *const u8]; + let segment_lens = [len_a]; + let version_before = unsafe { lance_dataset_version(dataset) }; + assert_eq!( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 1, + ) + }, + 0, + "{}", + unsafe { std::ffi::CStr::from_ptr(lance_last_error_message()).to_string_lossy() } + ); + assert_eq!( + unsafe { lance_dataset_version(dataset) }, + version_before + 1 + ); + assert_eq!( + unsafe { lance_dataset_index_segment_count(dataset, index_name.as_ptr()) }, + 1 + ); + + // Rebuild the same coverage under a fresh UUID and commit again: the old + // segment is replaced automatically (no replace flag). The uncommitted + // builder refuses to reuse a name that is already committed, so the + // rebuild happens under a scratch name; the commit registers it under + // `index_name` regardless of the name the segment was built with. + let rebuild_name = c_str("id_idx_rebuild"); + let (bytes_b, len_b) = build_scalar_segment_bytes( + dataset, + &rebuild_name, + LanceScalarIndexType::Bitmap, + Some(&all_fragments), + ); + let uuid_b = segment_uuid(bytes_b, len_b); + assert_ne!(uuid_a, uuid_b); + let segment_bytes = [bytes_b as *const u8]; + let segment_lens = [len_b]; + assert_eq!( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 1, + ) + }, + 0, + "{}", + unsafe { std::ffi::CStr::from_ptr(lance_last_error_message()).to_string_lossy() } + ); + assert_eq!( + unsafe { lance_dataset_version(dataset) }, + version_before + 2 + ); + assert_eq!( + unsafe { lance_dataset_index_segment_count(dataset, index_name.as_ptr()) }, + 1 + ); + let mut committed_uuid = [0_u8; 16]; + let mut committed_count = 0_u64; + assert_eq!( + unsafe { + lance_dataset_index_segments( + dataset, + index_name.as_ptr(), + committed_uuid.as_mut_ptr(), + 1, + &mut committed_count, + ) + }, + 0 + ); + assert_eq!(committed_count, 1); + assert_eq!(committed_uuid, uuid_b); + + // A later commit covering only a strict subset of the live coverage + // would orphan the remaining fragment, so it is rejected. + let first_fragment = [all_fragments[0]]; + let delta_name = c_str("id_idx_delta"); + let (bytes_c, len_c) = build_scalar_segment_bytes( + dataset, + &delta_name, + LanceScalarIndexType::Bitmap, + Some(&first_fragment), + ); + let segment_bytes = [bytes_c as *const u8]; + let segment_lens = [len_c]; + assert_eq!( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 1, + ) + }, + -1, + "partial overlap must be rejected instead of orphaning fragments" + ); + + unsafe { + lance_free_bytes(bytes_a); + lance_free_bytes(bytes_b); + lance_free_bytes(bytes_c); + lance_dataset_close(dataset); + } +} + +#[test] +fn test_commit_index_segments_rejects_wrong_column() { + // A segment built for one column cannot be committed under another + // existing column: the core rejects segments whose keyed field does not + // match the commit-time column's field id. + let (_tmp, uri) = create_multi_fragment_vector_dataset(2, 16, 8, false); + let uri_c = c_str(&uri); + let dataset = unsafe { lance_dataset_open(uri_c.as_ptr(), ptr::null(), 0) }; + let index_name = c_str("id_idx"); + let (bytes, len) = build_scalar_segment_bytes( + dataset, + &index_name, + LanceScalarIndexType::Bitmap, + Some(&[0]), + ); + + let wrong_column = c_str("embedding"); + let segment_bytes = [bytes as *const u8]; + let segment_lens = [len]; + let version_before = unsafe { lance_dataset_version(dataset) }; + assert_eq!( + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + wrong_column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 1, + ) + }, + -1 + ); + assert_eq!(lance_last_error_code(), LanceErrorCode::InvalidArgument); + let message = unsafe { + std::ffi::CStr::from_ptr(lance_last_error_message()) + .to_string_lossy() + .into_owned() + }; + assert!(message.contains("keyed field"), "{message}"); + assert_eq!(unsafe { lance_dataset_version(dataset) }, version_before); + assert_eq!(unsafe { lance_dataset_index_count(dataset) }, 0); + + unsafe { + lance_free_bytes(bytes); + lance_dataset_close(dataset); + } +} + +#[test] +fn test_commit_index_segments_type_change() { + let (_tmp, uri) = create_many_small_fragments(2); + let uri_c = c_str(&uri); + let dataset = unsafe { lance_dataset_open(uri_c.as_ptr(), ptr::null(), 0) }; + let mut fragment_ids = [0_u64; 2]; + assert_eq!( + unsafe { lance_dataset_fragment_ids(dataset, fragment_ids.as_mut_ptr()) }, + 0 + ); + let all_fragments = [fragment_ids[0] as u32, fragment_ids[1] as u32]; + let index_name = c_str("id_idx"); + let column = c_str("id"); + + let commit = |bytes: *const u8, len: usize| -> i32 { + let segment_bytes = [bytes]; + let segment_lens = [len]; + unsafe { + lance_dataset_commit_index_segments( + dataset, + index_name.as_ptr(), + column.as_ptr(), + segment_bytes.as_ptr(), + segment_lens.as_ptr(), + 1, + ) + } + }; + + // Commit a BTree index covering every fragment. + let (bytes_a, len_a) = build_scalar_segment_bytes( + dataset, + &index_name, + LanceScalarIndexType::BTree, + Some(&all_fragments), + ); + let uuid_a = segment_uuid(bytes_a, len_a); + let version_before = unsafe { lance_dataset_version(dataset) }; + assert_eq!(commit(bytes_a, len_a), 0, "{}", unsafe { + std::ffi::CStr::from_ptr(lance_last_error_message()).to_string_lossy() + }); + assert_eq!( + unsafe { lance_dataset_version(dataset) }, + version_before + 1 + ); + assert_eq!( + unsafe { lance_dataset_index_segment_count(dataset, index_name.as_ptr()) }, + 1 + ); + + // A full-coverage commit of a different index type replaces the existing + // index entirely. The builder refuses to reuse a committed index name, + // so the Bitmap rebuild happens under a scratch name. + let rebuild_name = c_str("id_idx_bitmap"); + let (bytes_b, len_b) = build_scalar_segment_bytes( + dataset, + &rebuild_name, + LanceScalarIndexType::Bitmap, + Some(&all_fragments), + ); + let uuid_b = segment_uuid(bytes_b, len_b); + assert_ne!(uuid_a, uuid_b); + assert_eq!(commit(bytes_b, len_b), 0, "{}", unsafe { + std::ffi::CStr::from_ptr(lance_last_error_message()).to_string_lossy() + }); + assert_eq!( + unsafe { lance_dataset_version(dataset) }, + version_before + 2 + ); + assert_eq!( + unsafe { lance_dataset_index_segment_count(dataset, index_name.as_ptr()) }, + 1 + ); + let mut committed_uuid = [0_u8; 16]; + let mut committed_count = 0_u64; + assert_eq!( + unsafe { + lance_dataset_index_segments( + dataset, + index_name.as_ptr(), + committed_uuid.as_mut_ptr(), + 1, + &mut committed_count, + ) + }, + 0 + ); + assert_eq!(committed_count, 1); + assert_eq!( + committed_uuid, uuid_b, + "type change must replace the old segment" + ); + + // A type change with partial coverage is rejected: it would orphan the + // uncovered fragments of the existing index. + let first_fragment = [all_fragments[0]]; + let partial_name = c_str("id_idx_partial"); + let (bytes_c, len_c) = build_scalar_segment_bytes( + dataset, + &partial_name, + LanceScalarIndexType::BTree, + Some(&first_fragment), + ); + assert_eq!( + commit(bytes_c, len_c), + -1, + "partial-coverage type change must be rejected" + ); + assert_eq!(lance_last_error_code(), LanceErrorCode::InvalidArgument); + let message = unsafe { + std::ffi::CStr::from_ptr(lance_last_error_message()) + .to_string_lossy() + .into_owned() + }; + assert!(message.contains("partial fragment coverage"), "{message}"); + assert_eq!( + unsafe { lance_dataset_version(dataset) }, + version_before + 2 + ); + assert_eq!( + unsafe { lance_dataset_index_segment_count(dataset, index_name.as_ptr()) }, + 1 + ); + + unsafe { + lance_free_bytes(bytes_a); + lance_free_bytes(bytes_b); + lance_free_bytes(bytes_c); + lance_dataset_close(dataset); + } +} + #[test] fn test_vector_model_rejects_malformed_arrow_inputs_without_panicking() { let (_tmp, uri) = create_multi_fragment_vector_dataset(2, 32, 8, false); diff --git a/tests/cpp/test_c_api.c b/tests/cpp/test_c_api.c index c49ecfa..7ac24f1 100644 --- a/tests/cpp/test_c_api.c +++ b/tests/cpp/test_c_api.c @@ -902,6 +902,92 @@ static void test_vector_models_and_reusable_segments(const char *uri) { printf("OK\n"); } +/* Builds one uncommitted vector segment per fragment and commits them as a + * single logical multi-segment index from a real C caller. */ +static void test_commit_index_segments(const char *uri) { + printf(" test_commit_index_segments... "); + LanceDataset *ds = lance_dataset_open(uri, NULL, 0); + ASSERT(ds != NULL, "open failed"); + uint64_t all_ids[2] = {0, 0}; + ASSERT(lance_dataset_fragment_ids(ds, all_ids) == 0, + "fragment enumeration failed"); + uint32_t fragment_ids[2] = {(uint32_t)all_ids[0], (uint32_t)all_ids[1]}; + + LanceVectorIndexSegmentParams params = { + LANCE_INDEX_IVF_FLAT, LANCE_METRIC_L2, 2, 0, 0, 2, 0, 0, 16, + }; + uint8_t *segment_bytes[2] = {NULL, NULL}; + size_t segment_lens[2] = {0, 0}; + uint8_t expected_uuids[2][16]; + memset(expected_uuids, 0, sizeof(expected_uuids)); + for (size_t i = 0; i < 2; i++) { + LanceIndexSegmentBuildOptions options = {0}; + options.fragment_ids = &fragment_ids[i]; + options.fragment_count = 1; + options.mode = LANCE_INDEX_SEGMENT_BUILD_AUTO; + LanceIndexSegmentBuilder *builder = + lance_index_segment_builder_new_vector( + ds, "embedding", "c_distributed_idx", ¶ms, &options); + ASSERT(builder != NULL, "vector segment builder failed"); + ASSERT(lance_index_segment_builder_execute_uncommitted( + builder, &segment_bytes[i], &segment_lens[i]) == 0, + "vector segment execution failed"); + LanceIndexSegmentMetadata *metadata = NULL; + ASSERT(lance_index_segment_metadata_parse( + segment_bytes[i], segment_lens[i], &metadata) == 0, + "metadata parse failed"); + ASSERT(lance_index_segment_metadata_uuid(metadata, + expected_uuids[i]) == 0, + "metadata UUID read failed"); + lance_index_segment_metadata_free(metadata); + lance_index_segment_builder_free(builder); + } + + /* One commit registers both segments as a single logical index. */ + uint64_t version_before = lance_dataset_version(ds); + const uint8_t *const_bytes[2] = {segment_bytes[0], segment_bytes[1]}; + int32_t rc = lance_dataset_commit_index_segments( + ds, "c_distributed_idx", "embedding", const_bytes, segment_lens, 2); + ASSERT(rc == 0, "commit_index_segments failed"); + ASSERT(lance_dataset_version(ds) == version_before + 1, + "commit must bump the dataset version exactly once"); + ASSERT(lance_dataset_index_segment_count(ds, "c_distributed_idx") == 2, + "committed index must have two segments"); + uint8_t committed_uuids[32] = {0}; + uint64_t committed_count = 0; + ASSERT(lance_dataset_index_segments(ds, "c_distributed_idx", + committed_uuids, 2, + &committed_count) == 0, + "segment enumeration failed"); + ASSERT(committed_count == 2, "committed segment count mismatch"); + ASSERT(memcmp(committed_uuids, expected_uuids[0], 16) == 0 && + memcmp(committed_uuids + 16, expected_uuids[1], 16) == 0, + "committed segment UUIDs mismatch"); + + /* Duplicate segment UUIDs in the commit set are rejected. */ + const uint8_t *dup_bytes[2] = {segment_bytes[0], segment_bytes[0]}; + size_t dup_lens[2] = {segment_lens[0], segment_lens[0]}; + rc = lance_dataset_commit_index_segments(ds, "c_dup_idx", "embedding", + dup_bytes, dup_lens, 2); + ASSERT(rc == -1, "duplicate segment UUIDs must fail"); + ASSERT(lance_last_error_code() == LANCE_ERR_INVALID_ARGUMENT, + "expected INVALID_ARGUMENT"); + + /* An empty commit set is rejected. */ + rc = lance_dataset_commit_index_segments(ds, "c_empty_idx", "embedding", + const_bytes, segment_lens, 0); + ASSERT(rc == -1, "empty commit set must fail"); + ASSERT(lance_last_error_code() == LANCE_ERR_INVALID_ARGUMENT, + "expected INVALID_ARGUMENT"); + ASSERT(lance_dataset_version(ds) == version_before + 1, + "rejected commits must not bump the version"); + + lance_free_bytes(segment_bytes[0]); + lance_free_bytes(segment_bytes[1]); + lance_dataset_close(ds); + printf("OK\n"); +} + /* Re-opens the dataset just written by `test_dataset_write_roundtrip` and * exercises `lance_dataset_compact_files`. The smoke fixture is a single * fragment, so the default planner has nothing to compact — we expect @@ -980,6 +1066,7 @@ int main(int argc, char **argv) { test_error_handling(); test_index_segment_builder(uri); test_vector_models_and_reusable_segments(uri); + test_commit_index_segments(uri); test_dataset_write_roundtrip(uri, write_uri); test_data_statistics(write_uri); test_update(write_uri); diff --git a/tests/cpp/test_cpp_api.cpp b/tests/cpp/test_cpp_api.cpp index ff5ed1e..7611106 100644 --- a/tests/cpp/test_cpp_api.cpp +++ b/tests/cpp/test_cpp_api.cpp @@ -533,6 +533,68 @@ static void test_vector_models_and_reusable_segments(const std::string& uri) { PASS(); } +static void test_commit_index_segments(const std::string& uri) { + TEST(test_commit_index_segments); + + auto ds = lance::Dataset::open(uri); + auto all_ids = ds.fragment_ids(); + assert(all_ids.size() >= 2); + + LanceVectorIndexParams params = { + LANCE_INDEX_IVF_FLAT, LANCE_METRIC_L2, 2, 0, 0, 2, 0, 0, 16, + }; + + // Build one uncommitted segment per fragment (the distributed workers). + std::vector> segments; + std::vector> expected_uuids; + for (size_t i = 0; i < 2; ++i) { + uint32_t fragment_id = static_cast(all_ids[i]); + LanceIndexSegmentBuildOptions options = {}; + options.fragment_ids = &fragment_id; + options.fragment_count = 1; + options.mode = LANCE_INDEX_SEGMENT_BUILD_AUTO; + auto builder = ds.new_vector_index_segment_builder( + "embedding", params, "cpp_distributed_idx", &options); + segments.push_back(builder.execute_uncommitted()); + auto metadata = lance::IndexSegmentMetadata::parse(segments.back()); + expected_uuids.push_back(metadata.uuid()); + } + + // One commit registers both segments as a single logical index. + uint64_t version_before = ds.version(); + ds.commit_index_segments("cpp_distributed_idx", "embedding", segments); + assert(ds.version() == version_before + 1); + assert(ds.index_segment_count("cpp_distributed_idx") == 2); + auto committed = ds.index_segments("cpp_distributed_idx"); + assert(committed.size() == 2); + for (size_t i = 0; i < 2; ++i) assert(committed[i] == expected_uuids[i]); + + // Duplicate segment UUIDs in the commit set are rejected. + bool caught = false; + try { + ds.commit_index_segments( + "cpp_dup_idx", "embedding", {segments[0], segments[0]}); + } catch (const lance::Error& e) { + caught = true; + assert(e.code == LANCE_ERR_INVALID_ARGUMENT); + } + assert(caught); + + // An empty commit set is rejected. + caught = false; + try { + ds.commit_index_segments( + "cpp_empty_idx", "embedding", {}); + } catch (const lance::Error& e) { + caught = true; + assert(e.code == LANCE_ERR_INVALID_ARGUMENT); + } + assert(caught); + assert(ds.version() == version_before + 1); + + PASS(); +} + static void test_fts_smoke(const std::string& uri) { TEST(test_fts_smoke); @@ -966,6 +1028,7 @@ int main(int argc, char** argv) { test_index_segments_smoke(uri); test_index_segment_builder(uri); test_vector_models_and_reusable_segments(uri); + test_commit_index_segments(uri); test_fts_smoke(uri); test_dataset_write_roundtrip(uri, write_uri); test_data_statistics(write_uri);