diff --git a/docs/src/guide/blob.md b/docs/src/guide/blob.md index d34557e5384..68796ef7fa4 100644 --- a/docs/src/guide/blob.md +++ b/docs/src/guide/blob.md @@ -61,6 +61,29 @@ source of truth for which scheme is supported at each `data_storage_version`. | `0.1`, `2.0`, `2.1` | Supported for write/read | Not supported | | `2.2+` | Not supported for write | Supported for write/read (recommended) | +### Managed objects and client compatibility + +Current writers store out-of-line Blob v2 payloads in independently named +`_blobs/.blob` objects and publish the Managed Blob reader and writer +capability on the table. This works with file formats 2.2 and 2.3; it does not +require choosing 2.3. Clients that do not understand the capability must refuse +to open a flagged snapshot. Updating an existing table with Blob data files can +activate it, even if that batch contains only inline values. The capability +remains set across later writes and restores. + +Compaction preserves Managed payload objects and can adopt existing Packed or +Dedicated sidecars in place. It records their complete addresses, so deleting +the original data file does not require copying its sidecars. Cleanup retains +objects referenced by protected snapshots; a partially live packed object is +retained as a whole. Blob reads continue to return the same bytes, while raw +descriptor scans can now report `kind = 4` for Managed values. + +After activation, use a client that supports Managed Blobs for all table +maintenance. Older clients, including v11.0.0, can bypass capability checks when +running cleanup from an unflagged historical snapshot or a cached handle. Such +cleanup can delete adopted sidecars by treating their original data file as +their owner. The table flag does not retrofit those old maintenance paths. + ## Blob v2: Write Patterns Use `blob_field` and `blob_array` to build blob v2 columns. diff --git a/python/python/tests/test_blob.py b/python/python/tests/test_blob.py index cb31aca4060..2f81f123c3f 100644 --- a/python/python/tests/test_blob.py +++ b/python/python/tests/test_blob.py @@ -138,7 +138,7 @@ def _add_columns_blob_v2_values(tmp_path): def _assert_blob_v2_add_columns_result(dataset, column, payloads): desc = dataset.to_table(columns=[column]).column(column).chunk(0) - assert desc.field("kind").to_pylist() == [0, 1, 2, 3] + assert desc.field("kind").to_pylist() == [0, 4, 4, 3] assert desc.field("blob_id").to_pylist()[3] == 1 assert desc.field("blob_uri").to_pylist()[3] == "external_blob.bin" @@ -1411,7 +1411,7 @@ def test_blob_extension_inline_threshold_per_column(tmp_path): desc = ds.to_table(columns=["inline_blob", "packed_blob"]) assert desc.column("inline_blob").chunk(0).field("kind").to_pylist() == [0] - assert desc.column("packed_blob").chunk(0).field("kind").to_pylist() == [1] + assert desc.column("packed_blob").chunk(0).field("kind").to_pylist() == [4] def test_blob_extension_threshold_metadata_persists_after_reopen(tmp_path): @@ -1514,7 +1514,7 @@ def test_blob_extension_dedicated_threshold_precedes_inline_threshold(tmp_path): ) desc = ds.to_table(columns=["blob"]).column("blob").chunk(0) - assert desc.field("kind").to_pylist() == [2] + assert desc.field("kind").to_pylist() == [4] def test_blob_extension_write_external(tmp_path): @@ -1655,6 +1655,15 @@ def failing_reader(): ds.add_columns(failing_reader(), reader_schema=schema) assert ds.version == 1 + files_after = _dataset_file_set(dataset_path) + assert files_before <= files_after + orphans = files_after - files_before + assert orphans and all( + p.parts[0] == "_blobs" and p.suffix == ".blob" for p in orphans + ) + # Independent payloads from failed writes follow the existing orphan policy. + # No concurrent writer is running here, so immediate unverified GC is safe. + ds.cleanup_old_versions(delete_unverified=True) assert _dataset_file_set(dataset_path) == files_before assert external_blob_path.exists() @@ -1685,6 +1694,15 @@ def fail_on_second_fragment(batch): assert call_count == 2 assert ds.version == 1 + files_after = _dataset_file_set(dataset_path) + assert files_before <= files_after + orphans = files_after - files_before + assert orphans and all( + p.parts[0] == "_blobs" and p.suffix == ".blob" for p in orphans + ) + # Independent payloads from failed writes follow the existing orphan policy. + # No concurrent writer is running here, so immediate unverified GC is safe. + ds.cleanup_old_versions(delete_unverified=True) assert _dataset_file_set(dataset_path) == files_before assert external_blob_path.exists() @@ -2453,7 +2471,7 @@ def test_blob_v2_lazy_preserves_empty_and_null(tmp_path, values, has_sidecar): assert descriptions[0]["size"] == 0 if has_sidecar: assert any( - description is not None and description["kind"] == 1 + description is not None and description["kind"] == 4 for description in descriptions ) assert any(path.suffix == ".blob" for path in _dataset_file_set(dataset_path)) @@ -2712,7 +2730,7 @@ def test_write_nested_blob_v2_and_take_by_field_path(tmp_path): ) desc = dataset.to_table(columns=["info.blob"]).column("info.blob").chunk(0) - assert desc.field("kind").to_pylist()[:2] == [0, 1] + assert desc.field("kind").to_pylist()[:2] == [0, 4] blobs = dataset.take_blobs("info.blob", indices=[0, 1]) with blobs[0] as f: diff --git a/rust/lance-core/src/datatypes.rs b/rust/lance-core/src/datatypes.rs index dcad595e5aa..89cdf11a62c 100644 --- a/rust/lance-core/src/datatypes.rs +++ b/rust/lance-core/src/datatypes.rs @@ -596,6 +596,12 @@ pub enum BlobKind { /// External blobs can have a position and a size. If the position is not set, /// it defaults to 0, which points to the beginning of the blob. External = 3, + /// A Lance-owned immutable object, independently of the descriptor's data file. + /// `blob_id` is the exact manifest base ID (including zero); `blob_uri` is + /// relative to that base root. `position`/`size` select a known range, and + /// zero size is an empty value rather than a request to discover its length. + /// Tables containing this kind require the Managed Blob reader and writer feature flags. + Managed = 4, } impl TryFrom for BlobKind { @@ -607,6 +613,7 @@ impl TryFrom for BlobKind { 1 => Ok(Self::Packed), 2 => Ok(Self::Dedicated), 3 => Ok(Self::External), + 4 => Ok(Self::Managed), other => Err(Error::invalid_input_source( format!("Unknown blob kind {other:?}").into(), )), diff --git a/rust/lance-core/src/utils/blob.rs b/rust/lance-core/src/utils/blob.rs index faf97f2d3ab..2d2e8bba6e4 100644 --- a/rust/lance-core/src/utils/blob.rs +++ b/rust/lance-core/src/utils/blob.rs @@ -3,6 +3,35 @@ use object_store::path::Path; +use crate::{Error, Result}; + +/// Validate a Managed descriptor's object-relative path and known byte range. +/// +/// No base ID is reserved. Resolving the base belongs to the snapshot holding +/// the descriptor, and must fail if that snapshot has no matching binding. +pub fn validate_managed_reference(uri: &str, position: u64, size: u64) -> Result { + if uri.is_empty() + || uri.starts_with('/') + || uri.ends_with('/') + || uri.contains("://") + || uri.contains('\\') + || uri + .split('/') + .any(|part| part.is_empty() || part == "." || part == "..") + { + return Err(Error::invalid_input(format!( + "Managed blob_uri must be a non-empty canonical relative object path, got {uri:?}" + ))); + } + position.checked_add(size).ok_or_else(|| { + Error::invalid_input(format!( + "Managed blob range overflows u64: position={position}, size={size}" + )) + })?; + Path::parse(uri) + .map_err(|error| Error::invalid_input(format!("Invalid Managed blob_uri {uri:?}: {error}"))) +} + /// Format a blob sidecar path for a data file. /// /// Layout: `//.blob` @@ -17,6 +46,34 @@ pub fn blob_path(base: &Path, data_file_key: &str, blob_id: u32) -> Path { #[cfg(test)] mod tests { use super::*; + use rstest::rstest; + + #[rstest] + #[case("")] + #[case("/data/a.blob")] + #[case("data/../a.blob")] + #[case("data/./a.blob")] + #[case("data//a.blob")] + #[case("data/a.blob/")] + #[case("s3://bucket/a.blob")] + #[case("data\\a.blob")] + fn managed_paths_reject_noncanonical_references(#[case] uri: &str) { + let error = validate_managed_reference(uri, 0, 1).unwrap_err(); + assert!(matches!(error, Error::InvalidInput { .. })); + assert!(error.to_string().contains("Managed blob_uri")); + } + + #[test] + fn managed_ranges_preserve_empty_values_and_reject_overflow() { + assert!(validate_managed_reference("data/a.blob", u64::MAX, 0).is_ok()); + let error = validate_managed_reference("data/a.blob", u64::MAX, 1).unwrap_err(); + assert!(matches!(error, Error::InvalidInput { .. })); + assert!( + error + .to_string() + .contains("position=18446744073709551615, size=1") + ); + } #[test] fn test_blob_path_formatting() { diff --git a/rust/lance-encoding/src/encoder/structural.rs b/rust/lance-encoding/src/encoder/structural.rs index 3d4e49f42c8..266d19c35ff 100644 --- a/rust/lance-encoding/src/encoder/structural.rs +++ b/rust/lance-encoding/src/encoder/structural.rs @@ -72,7 +72,10 @@ impl PrimitiveFieldEncoding { } } - fn create_at( + /// Compose a primitive encoder at an already allocated physical column. + /// File grammars use this when another logical type, such as a Blob, + /// supplies the physical descriptor field instead of the logical field. + pub fn create_at( &self, field: Field, column_index: u32, diff --git a/rust/lance-encoding/src/encodings/logical/blob.rs b/rust/lance-encoding/src/encodings/logical/blob.rs index a5432fdc54c..629231f0c0e 100644 --- a/rust/lance-encoding/src/encodings/logical/blob.rs +++ b/rust/lance-encoding/src/encodings/logical/blob.rs @@ -329,67 +329,91 @@ impl FieldEncoder for BlobV2StructuralEncoder { let mut uri_builder = StringBuilder::with_capacity(row_count, row_count * 16); for i in 0..row_count { - let (kind_value, position_value, size_value, blob_id_value, uri_value) = - if struct_arr.is_null(i) || kind_col.is_null(i) { - (BlobKind::Inline as u8, 0, 0, 0, "".to_string()) - } else { - let kind_val = BlobKind::try_from(kind_col.value(i))?; - match kind_val { - BlobKind::Dedicated => ( - BlobKind::Dedicated as u8, - 0, - blob_size_col.value(i), - blob_id_col.value(i), - "".to_string(), - ), - BlobKind::External => { - let uri = uri_col.value(i).to_string(); - let position = if packed_position_col.is_null(i) { - 0 - } else { - packed_position_col.value(i) - }; - let size = if blob_size_col.is_null(i) { - 0 - } else { - blob_size_col.value(i) - }; - let external_base_id = if blob_id_col.is_null(i) { - 0 - } else { - blob_id_col.value(i) - }; - ( - BlobKind::External as u8, - position, - size, - external_base_id, - uri, - ) + let (kind_value, position_value, size_value, blob_id_value, uri_value) = if struct_arr + .is_null(i) + || kind_col.is_null(i) + { + (BlobKind::Inline as u8, 0, 0, 0, "".to_string()) + } else { + let kind_val = BlobKind::try_from(kind_col.value(i))?; + match kind_val { + BlobKind::Managed => { + if uri_col.is_null(i) + || blob_id_col.is_null(i) + || packed_position_col.is_null(i) + || blob_size_col.is_null(i) + { + return Err(Error::invalid_input(format!( + "Managed blob row {i} requires URI, base ID, position, and size" + ))); } - BlobKind::Packed => ( - BlobKind::Packed as u8, - packed_position_col.value(i), - blob_size_col.value(i), + let uri = uri_col.value(i); + let position = packed_position_col.value(i); + let size = blob_size_col.value(i); + lance_core::utils::blob::validate_managed_reference(uri, position, size)?; + ( + BlobKind::Managed as u8, + position, + size, blob_id_col.value(i), + uri.to_string(), + ) + } + BlobKind::Dedicated => ( + BlobKind::Dedicated as u8, + 0, + blob_size_col.value(i), + blob_id_col.value(i), + "".to_string(), + ), + BlobKind::External => { + let uri = uri_col.value(i).to_string(); + let position = if packed_position_col.is_null(i) { + 0 + } else { + packed_position_col.value(i) + }; + let size = if blob_size_col.is_null(i) { + 0 + } else { + blob_size_col.value(i) + }; + let external_base_id = if blob_id_col.is_null(i) { + 0 + } else { + blob_id_col.value(i) + }; + ( + BlobKind::External as u8, + position, + size, + external_base_id, + uri, + ) + } + BlobKind::Packed => ( + BlobKind::Packed as u8, + packed_position_col.value(i), + blob_size_col.value(i), + blob_id_col.value(i), + "".to_string(), + ), + BlobKind::Inline => { + let data_val = data_col.value(i); + let blob_len = data_val.len() as u64; + let position = + external_buffers.add_buffer(LanceBuffer::from(Buffer::from(data_val))); + + ( + BlobKind::Inline as u8, + position, + blob_len, + 0, "".to_string(), - ), - BlobKind::Inline => { - let data_val = data_col.value(i); - let blob_len = data_val.len() as u64; - let position = external_buffers - .add_buffer(LanceBuffer::from(Buffer::from(data_val))); - - ( - BlobKind::Inline as u8, - position, - blob_len, - 0, - "".to_string(), - ) - } + ) } - }; + } + }; kind_builder.append_value(kind_value); position_builder.append_value(position_value); diff --git a/rust/lance-file/src/concat.rs b/rust/lance-file/src/concat.rs index c67892bf412..a989759f615 100644 --- a/rust/lance-file/src/concat.rs +++ b/rust/lance-file/src/concat.rs @@ -437,6 +437,22 @@ fn validate_blob_field( )); } } + BlobKind::Managed => { + let uris = descriptors + .column_by_name("blob_uri") + .ok_or_else(|| { + Error::corrupt_file(path.clone(), "Managed descriptor has no blob_uri") + })? + .as_string::(); + lance_core::utils::blob::validate_managed_reference( + uris.value(row), + positions.value(row), + sizes.value(row), + )?; + // The ID is a snapshot base binding, not a leased sidecar + // number. Concatenation preserves the independent object + // address; the dataset caller owns the base namespace. + } BlobKind::Inline | BlobKind::External => {} } } diff --git a/rust/lance-table/src/feature_flags.rs b/rust/lance-table/src/feature_flags.rs index ca060dc056a..1418b70408c 100644 --- a/rust/lance-table/src/feature_flags.rs +++ b/rust/lance-table/src/feature_flags.rs @@ -68,7 +68,12 @@ const _: () = assert!(FLAG_MIXED_DATA_FILE_VERSIONS < FLAG_UNKNOWN); /// Bit 9 is taken by the stable-row-id FRI compatibility flag. pub const FLAG_FRAGMENT_REUSE_INDEX: u64 = 1 << 10; -pub(crate) const STICKY_PAIRED_FLAGS: u64 = FLAG_MIXED_DATA_FILE_VERSIONS; +/// Blob v2 descriptors may independently address Lance-owned objects. Readers +/// must resolve their explicit bases and writers/GC must preserve those references. +/// This capability is sticky, including across restore, and requires both words. +pub const FLAG_MANAGED_BLOBS: u64 = 1 << 11; + +pub(crate) const STICKY_PAIRED_FLAGS: u64 = FLAG_MIXED_DATA_FILE_VERSIONS | FLAG_MANAGED_BLOBS; /// Environment variable that opts a release build into reading and writing data /// overlay files before the feature is generally released. @@ -191,7 +196,7 @@ fn mark_supported(flags: &mut u64, flag: u64, feature_enabled: bool) { /// is enabled. Split out from [`supported_flags`] so the policy is testable /// without toggling the build profile or environment. fn supported_flags_when(overlay_enabled: bool) -> u64 { - let mut supported = FLAG_UNKNOWN - 1; + let mut supported = (FLAG_UNKNOWN - 1) | FLAG_MANAGED_BLOBS; mark_supported( &mut supported, FLAG_UNSTABLE_DATA_OVERLAY_FILES, @@ -257,14 +262,20 @@ pub fn has_deprecated_v2_feature_flag(writer_flags: u64) -> bool { /// commit path refuses to *produce* this, so seeing it on read means the /// manifest was written by something that did not. pub fn validate_paired_feature_flags(manifest: &Manifest) -> Result<()> { - let reader = manifest.reader_feature_flags & FLAG_MIXED_DATA_FILE_VERSIONS != 0; - let writer = manifest.writer_feature_flags & FLAG_MIXED_DATA_FILE_VERSIONS != 0; - if reader != writer { - return Err(Error::corrupt_file_named( - "manifest", - "Manifest has only one of the mixed data-file-version reader and writer feature bits set, \ - so its semantics are undefined", - )); + for (flag, name) in [ + (FLAG_MIXED_DATA_FILE_VERSIONS, "mixed data-file-version"), + (FLAG_MANAGED_BLOBS, "Managed Blob"), + ] { + let reader = manifest.reader_feature_flags & flag != 0; + let writer = manifest.writer_feature_flags & flag != 0; + if reader != writer { + return Err(Error::corrupt_file_named( + "manifest", + format!( + "Manifest has only one of the {name} reader and writer feature bits set, so its semantics are undefined" + ), + )); + } } Ok(()) } @@ -547,6 +558,38 @@ mod tests { assert!(err.to_string().contains("cannot be written"), "{err}"); } + #[rstest::rstest] + #[case::reader_only(true, false)] + #[case::writer_only(false, true)] + #[case::paired(true, true)] + fn managed_capability_is_paired_and_sticky(#[case] reader: bool, #[case] writer: bool) { + let mut source = empty_manifest(); + source.reader_feature_flags = if reader { FLAG_MANAGED_BLOBS } else { 0 }; + source.writer_feature_flags = if writer { FLAG_MANAGED_BLOBS } else { 0 }; + if reader != writer { + for error in [ + ensure_can_read_manifest(&source).unwrap_err(), + ensure_can_write_manifest(&source).unwrap_err(), + ] { + assert!(matches!(error, Error::CorruptFile { .. })); + assert!(error.to_string().contains("Managed Blob")); + } + return; + } + ensure_can_read_manifest(&source).unwrap(); + ensure_can_write_manifest(&source).unwrap(); + let mut destination = empty_manifest(); + inherit_sticky_feature_flags(&mut destination, &source).unwrap(); + apply_feature_flags(&mut destination, false, false).unwrap(); + assert!(destination.has_managed_blobs()); + assert_eq!( + destination.writer_feature_flags & FLAG_MANAGED_BLOBS, + FLAG_MANAGED_BLOBS + ); + // The released v11.0.0 client only accepts bits below 128. + assert_ne!(destination.reader_feature_flags & !(128 - 1), 0); + } + fn empty_manifest() -> Manifest { use crate::format::DataStorageFormat; use arrow_schema::{DataType, Field as ArrowField, Schema as ArrowSchema}; diff --git a/rust/lance-table/src/format/manifest.rs b/rust/lance-table/src/format/manifest.rs index 628e313a9a7..c39b0a52a20 100644 --- a/rust/lance-table/src/format/manifest.rs +++ b/rust/lance-table/src/format/manifest.rs @@ -170,6 +170,33 @@ impl From for BTreeMap { } impl Manifest { + /// Whether this table requires independently addressed Managed Blob support. + pub fn has_managed_blobs(&self) -> bool { + self.reader_feature_flags & crate::feature_flags::FLAG_MANAGED_BLOBS != 0 + } + + /// Register the writer's explicit base without rebinding encoded references. + pub fn bind_managed_base(&mut self, base: BasePath) -> Result<()> { + if !self.has_managed_blobs() { + return Ok(()); + } + if let Some(existing) = self.base_paths.get(&base.id) { + if existing.path != base.path || existing.is_dataset_root != base.is_dataset_root { + return Err(Error::invalid_input(format!( + "Managed base ID {} is bound to {:?} (is_dataset_root={}); descriptors require {:?} (is_dataset_root={})", + base.id, + existing.path, + existing.is_dataset_root, + base.path, + base.is_dataset_root + ))); + } + } else { + self.base_paths.insert(base.id, base); + } + Ok(()) + } + pub fn new( schema: Schema, fragments: Arc>, @@ -611,6 +638,25 @@ pub struct BasePath { } impl BasePath { + /// Choose an unused exact base ID without reserving zero or overflowing at + /// `u32::MAX`. The caller must publish the binding with its references and + /// reject a concurrent attempt to bind the chosen ID to another location. + pub fn unused_id(bases: impl IntoIterator) -> Result { + let mut ids = bases.into_iter().collect::>(); + ids.sort_unstable(); + ids.dedup(); + let mut candidate = 0u32; + for id in ids { + if id != candidate { + break; + } + candidate = candidate + .checked_add(1) + .ok_or_else(|| Error::invalid_input("All u32 base IDs are already registered"))?; + } + Ok(candidate) + } + /// Create a new BasePath /// /// # Arguments @@ -1152,9 +1198,40 @@ mod tests { use super::*; use arrow_schema::{Field as ArrowField, Schema as ArrowSchema}; - use lance_core::datatypes::Field; + use lance_core::datatypes::{BLOB_V2_DESC_LANCE_FIELD, Field}; use roaring::RoaringBitmap; + #[rstest::rstest] + #[case::different_path("memory://other", true)] + #[case::data_only("memory://dataset", false)] + fn managed_base_binding_preserves_address_contract( + #[case] path: &str, + #[case] is_dataset_root: bool, + ) { + let schema = Schema { + fields: vec![BLOB_V2_DESC_LANCE_FIELD.clone()], + ..Default::default() + }; + let mut manifest = Manifest::new( + schema, + Arc::new(vec![]), + DataStorageFormat::new(ConcreteFileVersion::V2_3), + HashMap::new(), + ); + manifest.reader_feature_flags |= crate::feature_flags::FLAG_MANAGED_BLOBS; + manifest.writer_feature_flags |= crate::feature_flags::FLAG_MANAGED_BLOBS; + let base = BasePath::new(7, "memory://dataset".to_string(), None, true); + manifest.bind_managed_base(base.clone()).unwrap(); + manifest.bind_managed_base(base.clone()).unwrap(); + let error = manifest + .bind_managed_base(BasePath::new(7, path.to_string(), None, is_dataset_root)) + .unwrap_err(); + assert!(matches!(error, Error::InvalidInput { .. })); + assert!(error.to_string().contains("Managed base ID 7 is bound")); + assert!(error.to_string().contains("is_dataset_root=")); + assert_eq!(manifest.base_paths[&7], base); + } + /// A shallow clone points every local file at the parent through `base_id`. /// An overlay's data file lives in the parent too, so it needs the same /// stamp; without it the clone looks for the overlay under its own root. diff --git a/rust/lance-table/src/transaction/manifest_build.rs b/rust/lance-table/src/transaction/manifest_build.rs index 515e914c0b6..a711ed3860a 100644 --- a/rust/lance-table/src/transaction/manifest_build.rs +++ b/rust/lance-table/src/transaction/manifest_build.rs @@ -11,7 +11,7 @@ //! metadata it stamps, the validation that runs before it. use crate::feature_flags::{ - FLAG_COVERED_INDEX_METADATA, FLAG_STABLE_ROW_IDS, apply_feature_flags, + FLAG_COVERED_INDEX_METADATA, FLAG_MANAGED_BLOBS, FLAG_STABLE_ROW_IDS, apply_feature_flags, ensure_can_read_manifest, ensure_can_write_manifest, inherit_sticky_feature_flags, }; use crate::format::overlay::{OverlayCoverage, TOMBSTONE_FIELD_ID}; @@ -1329,6 +1329,36 @@ impl Transaction { ) }; + // Only newly published Blob data files activate the capability. Comparing + // physical files also covers column rewrites and overlays while leaving + // metadata-only changes and deletion vectors on old tables alone. + let blob_fields: Vec<_> = manifest + .schema + .fields_pre_order() + .filter(|field| field.is_blob_v2()) + .map(|field| field.id) + .collect(); + if !blob_fields.is_empty() { + let old_files: HashSet<_> = current_manifest + .into_iter() + .flat_map(|manifest| manifest.fragments.iter()) + .flat_map(|fragment| fragment.referenced_lance_files()) + .map(|file| (file.base_id, file.path.as_str())) + .collect(); + if manifest + .fragments + .iter() + .flat_map(|fragment| fragment.referenced_lance_files()) + .any(|file| { + !old_files.contains(&(file.base_id, file.path.as_str())) + && file.fields.iter().any(|id| blob_fields.contains(id)) + }) + { + manifest.reader_feature_flags |= FLAG_MANAGED_BLOBS; + manifest.writer_feature_flags |= FLAG_MANAGED_BLOBS; + } + } + manifest.tag.clone_from(&self.tag); if config.auto_set_feature_flags { @@ -1534,13 +1564,20 @@ impl Transaction { // Assign a new ID if not already assigned let mut base_to_add = new_base.clone(); if base_to_add.id == 0 { - let next_id = manifest - .base_paths - .keys() - .max() - .map(|&id| id + 1) - .unwrap_or(1); - base_to_add.id = next_id; + base_to_add.id = crate::format::BasePath::unused_id( + manifest + .base_paths + .keys() + .copied() + .chain(std::iter::once(0)), + )?; + } else if manifest.has_managed_blobs() + && let Some(existing) = manifest.base_paths.get(&base_to_add.id) + { + return Err(Error::invalid_input(format!( + "Cannot replace base ID {} bound to {:?} with {:?}", + base_to_add.id, existing.path, base_to_add.path + ))); } manifest.base_paths.insert(base_to_add.id, base_to_add); diff --git a/rust/lance/src/blob.rs b/rust/lance/src/blob.rs index f3a6aa745cd..3c0fa1be999 100644 --- a/rust/lance/src/blob.rs +++ b/rust/lance/src/blob.rs @@ -463,6 +463,22 @@ fn validate_prepared_blob_value_array(field: &Field, array: &ArrayRef) -> Result } validate_blob_id(blob_id_col.value(row))?; } + BlobKind::Managed => { + if uri_col.is_null(row) + || blob_id_col.is_null(row) + || blob_size_col.is_null(row) + || position_col.is_null(row) + { + return Err(Error::invalid_input(format!( + "Prepared Managed blob row {row} must set `uri`, `blob_id`, `blob_size`, and `position`" + ))); + } + lance_core::utils::blob::validate_managed_reference( + uri_col.value(row), + position_col.value(row), + blob_size_col.value(row), + )?; + } BlobKind::External => { if uri_col.is_null(row) || uri_col.value(row).is_empty() { return Err(Error::invalid_input(format!( @@ -535,6 +551,14 @@ pub enum BlobDescriptor { }, /// Payload bytes stored as the full contents of a dedicated sidecar blob. Dedicated { blob_id: u32, size: u64 }, + /// A known range in an immutable Lance-owned object. The exact `base_id` + /// must be bound in the same committed snapshot as this descriptor. + Managed { + base_id: u32, + uri: String, + offset: u64, + size: u64, + }, /// Payload bytes referenced from an external object or registered base. External { base_id: u32, @@ -722,6 +746,20 @@ impl BlobDescriptorArrayBuilder { blob_size_builder.append_value(size); position_builder.append_null(); } + BlobDescriptor::Managed { + base_id, + uri, + offset, + size, + } => { + validity.append_non_null(); + kind_builder.append_value(BlobKind::Managed as u8); + data_builder.append_null(); + uri_builder.append_value(uri); + blob_id_builder.append_value(base_id); + blob_size_builder.append_value(size); + position_builder.append_value(offset); + } BlobDescriptor::External { base_id, uri, @@ -778,6 +816,12 @@ fn validate_blob_descriptor(value: &BlobDescriptor) -> Result<()> { Ok(()) } BlobDescriptor::Dedicated { blob_id, .. } => validate_blob_id(*blob_id), + BlobDescriptor::Managed { + uri, offset, size, .. + } => { + lance_core::utils::blob::validate_managed_reference(uri, *offset, *size)?; + Ok(()) + } BlobDescriptor::External { uri, offset, size, .. } => { @@ -828,6 +872,14 @@ impl PackedBlobWriter { blob_id: u32, ) -> Result { let path = sidecar_path_for_data_file(&data_file_path, blob_id)?; + Self::try_new_at(object_store, path, blob_id).await + } + + pub(crate) async fn try_new_at( + object_store: ObjectStore, + path: Path, + blob_id: u32, + ) -> Result { let writer = object_store.create(&path).await?; Ok(Self { object_store, @@ -969,6 +1021,14 @@ impl DedicatedBlobWriter { blob_id: u32, ) -> Result { let path = sidecar_path_for_data_file(&data_file_path, blob_id)?; + Self::try_new_at(object_store, path, blob_id).await + } + + pub(crate) async fn try_new_at( + object_store: ObjectStore, + path: Path, + blob_id: u32, + ) -> Result { let writer = object_store.create(&path).await?; Ok(Self { object_store, diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index ab415549d20..e278b2dcc94 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -38,6 +38,7 @@ use lance_io::utils::{ CachedFileSize, read_last_block, read_message, read_metadata_offset, read_struct, }; use lance_namespace::LanceNamespace; +use lance_table::format::BasePath; use lance_table::format::{ DataFile, DataStorageFormat, DeletionFile, Fragment, IndexMetadata, MAGIC, Manifest, ManifestBuildConfig, RowIdMeta, pb, populate_manifest_schema_dictionaries, @@ -889,7 +890,7 @@ impl Dataset { } #[allow(clippy::too_many_arguments)] - fn checkout_manifest( + pub(crate) fn checkout_manifest( object_store: Arc, base_path: Path, uri: String, @@ -2453,10 +2454,35 @@ impl Dataset { } } + pub(crate) fn managed_default_base(&self) -> Result { + if let Some(base) = self + .manifest + .base_paths + .values() + .filter(|base| base.path == self.uri && base.is_dataset_root) + .min_by_key(|base| base.id) + { + return Ok(base.clone()); + } + let id = BasePath::unused_id(self.manifest.base_paths.keys().copied())?; + Ok(BasePath::new(id, self.uri.clone(), None, true)) + } + async fn base_object_store(&self, base_id: u32) -> Result> { let base_path = self.manifest.base_paths.get(&base_id).ok_or_else(|| { Error::invalid_input(format!("Dataset base path with ID {} not found", base_id)) })?; + if base_path.path == self.uri + && !self + .base_store_params + .as_ref() + .is_some_and(|params| params.contains_key(&base_path.path)) + && self.store_params.as_ref().is_none_or(|params| { + matches!(params.scoped_to_base(Some(base_id)), Cow::Borrowed(_)) + }) + { + return Ok(self.object_store.clone()); + } let store_params = self.store_params_for_base(Some(base_path)); let cell = { diff --git a/rust/lance/src/dataset/blob.rs b/rust/lance/src/dataset/blob.rs index bed052ee315..1cf454d4489 100644 --- a/rust/lance/src/dataset/blob.rs +++ b/rust/lance/src/dataset/blob.rs @@ -1,6 +1,8 @@ // SPDX-License-Identifier: Apache-2.0 // SPDX-FileCopyrightText: Copyright The Lance Authors +pub mod clone; + use std::{ collections::{BTreeMap, HashMap, HashSet}, future::Future, @@ -359,16 +361,22 @@ impl RollingPackedBlobWriter { async fn start_new_pack( &mut self, object_store: ObjectStore, - data_dir: Path, - data_file_key: String, + data_file_path: Path, + managed_base: Option<&(u32, Path)>, blob_id_allocator: BlobIdAllocator, max_pack_size: usize, ) -> Result<()> { self.finish().await?; let blob_id = blob_id_allocator.next()?; - let data_file_path = data_dir.join(format!("{data_file_key}.lance")); - self.current = - Some(PackedBlobWriter::try_new(object_store, data_file_path, blob_id).await?); + self.current = Some(if let Some((_, root)) = managed_base { + let path = root + .clone() + .join("_blobs") + .join(format!("{}.blob", uuid::Uuid::new_v4())); + PackedBlobWriter::try_new_at(object_store, path, blob_id).await? + } else { + PackedBlobWriter::try_new(object_store, data_file_path, blob_id).await? + }); self.current_size = 0; self.current_max_pack_size = Some(max_pack_size); Ok(()) @@ -377,8 +385,8 @@ impl RollingPackedBlobWriter { async fn write( &mut self, object_store: ObjectStore, - data_dir: Path, - data_file_key: String, + data_file_path: Path, + managed_base: Option<&(u32, Path)>, blob_id_allocator: BlobIdAllocator, max_pack_size: usize, source: BlobWriteSource<'_>, @@ -394,8 +402,8 @@ impl RollingPackedBlobWriter { if needs_new_pack { self.start_new_pack( object_store, - data_dir, - data_file_key, + data_file_path, + managed_base, blob_id_allocator, max_pack_size, ) @@ -412,7 +420,7 @@ impl RollingPackedBlobWriter { } }; self.current_size += len; - Ok(value) + BlobPreprocessor::managed_descriptor(value, managed_base, writer.path()) } async fn finish(&mut self) -> Result<()> { @@ -434,9 +442,10 @@ impl RollingPackedBlobWriter { /// Preprocesses blob v2 columns on the write path so the encoder only sees lightweight descriptors: /// /// - Spills large blobs to sidecar files before encoding, reducing memory/CPU and avoiding copying huge payloads through page builders. -/// - Emits `blob_id/blob_size` tied to the data file stem, giving readers a stable path independent of temporary fragment IDs assigned during write. +/// - Emits explicit Managed addresses when a base is bound; low-level sidecar writers retain their existing descriptors. /// - Leaves small inline blobs and URI rows unchanged for compatibility. pub struct BlobPreprocessor { + managed_base: Option<(u32, Path)>, object_store: ObjectStore, data_dir: Path, data_file_key: String, @@ -649,6 +658,7 @@ impl BlobPreprocessor { .map(|field| BlobPreprocessField::new(field.as_ref())) .collect::>>()?; Ok(Self { + managed_base: None, object_store, data_dir, data_file_key, @@ -665,6 +675,40 @@ impl BlobPreprocessor { }) } + pub(super) fn with_managed_base(mut self, base_id: u32, root: Path) -> Self { + self.managed_base = Some((base_id, root)); + self + } + + fn managed_descriptor( + descriptor: BlobDescriptor, + managed_base: Option<&(u32, Path)>, + path: &Path, + ) -> Result { + let Some((base_id, root)) = managed_base else { + return Ok(descriptor); + }; + let (offset, size) = match descriptor { + BlobDescriptor::Packed { offset, size, .. } => (offset, size), + BlobDescriptor::Dedicated { size, .. } => (0, size), + other => return Ok(other), + }; + let relative = path + .prefix_match(root) + .ok_or_else(|| { + Error::internal(format!("Blob object {path} is outside writer base {root}")) + })? + .map(|part| part.as_ref().to_string()) + .collect::>() + .join("/"); + Ok(BlobDescriptor::Managed { + base_id: *base_id, + uri: relative, + offset, + size, + }) + } + pub(super) fn with_part_blob_ids(mut self, blob_ids: Range) -> Result { self.blob_id_allocator = BlobIdAllocator::from_range(blob_ids.clone())?; self.part_blob_ids = Some(blob_ids); @@ -682,18 +726,27 @@ impl BlobPreprocessor { BlobDescriptorArrayBuilder::new_with_metadata(field.name(), field.is_nullable(), metadata) } - async fn write_dedicated( - object_store: ObjectStore, - data_dir: Path, - data_file_key: String, - blob_id_allocator: BlobIdAllocator, - source: BlobWriteSource<'_>, - ) -> Result { - let blob_id = blob_id_allocator.next()?; - let data_file_path = data_dir.join(format!("{data_file_key}.lance")); - let mut writer = - crate::blob::DedicatedBlobWriter::try_new(object_store, data_file_path, blob_id) - .await?; + async fn write_dedicated(&mut self, source: BlobWriteSource<'_>) -> Result { + let blob_id = self.blob_id_allocator.next()?; + let mut writer = if let Some((_, root)) = &self.managed_base { + let path = root + .clone() + .join("_blobs") + .join(format!("{}.blob", uuid::Uuid::new_v4())); + crate::blob::DedicatedBlobWriter::try_new_at(self.object_store.clone(), path, blob_id) + .await? + } else { + let data_file_path = self + .data_dir + .clone() + .join(format!("{}.lance", self.data_file_key)); + crate::blob::DedicatedBlobWriter::try_new( + self.object_store.clone(), + data_file_path, + blob_id, + ) + .await? + }; match source { BlobWriteSource::Bytes(data) => writer.write(data).await?, BlobWriteSource::External(source) => { @@ -702,7 +755,8 @@ impl BlobPreprocessor { .await?; } } - writer.finish().await + let path = writer.path().clone(); + Self::managed_descriptor(writer.finish().await?, self.managed_base.as_ref(), &path) } async fn write_packed( @@ -716,8 +770,10 @@ impl BlobPreprocessor { self.pack_writer .write( self.object_store.clone(), - self.data_dir.clone(), - self.data_file_key.clone(), + self.data_dir + .clone() + .join(format!("{}.lance", self.data_file_key)), + self.managed_base.as_ref(), self.blob_id_allocator.clone(), max_pack_size, source, @@ -774,7 +830,7 @@ impl BlobPreprocessor { )) })?; } - BlobKind::Inline | BlobKind::External => {} + BlobKind::Inline | BlobKind::External | BlobKind::Managed => {} } } @@ -807,6 +863,14 @@ impl BlobPreprocessor { BlobKind::Dedicated => { output.push_dedicated(blob_ids.value(row), sizes.value(row))?; } + BlobKind::Managed => { + output.push(BlobDescriptor::Managed { + base_id: blob_ids.value(row), + uri: uris.value(row).to_string(), + offset: positions.value(row), + size: sizes.value(row), + })?; + } BlobKind::External => { output.push(BlobDescriptor::External { base_id: blob_ids.value(row), @@ -1202,14 +1266,9 @@ impl BlobPreprocessor { let data_len = if has_data { data_col.value(i).len() } else { 0 }; if has_data && data_len > dedicated_threshold { - let value = Self::write_dedicated( - self.object_store.clone(), - self.data_dir.clone(), - self.data_file_key.clone(), - self.blob_id_allocator.clone(), - BlobWriteSource::Bytes(data_col.value(i)), - ) - .await?; + let value = self + .write_dedicated(BlobWriteSource::Bytes(data_col.value(i))) + .await?; blob_writer.push(value)?; continue; } @@ -1247,14 +1306,9 @@ impl BlobPreprocessor { let data_len = source.size(); if data_len > dedicated_threshold as u64 { - let value = Self::write_dedicated( - self.object_store.clone(), - self.data_dir.clone(), - self.data_file_key.clone(), - self.blob_id_allocator.clone(), - BlobWriteSource::External(&source), - ) - .await?; + let value = self + .write_dedicated(BlobWriteSource::External(&source)) + .await?; blob_writer.push(value)?; continue; } @@ -1592,6 +1646,7 @@ pub struct BlobFile { /// recompute the object store, data directory, and data file key. #[derive(Clone)] struct BlobReadLocation { + base_id: Option, object_store: Arc, data_file_dir: Path, data_file_key: String, @@ -4288,6 +4343,7 @@ impl<'a> BlobV2ReadContext<'a> { BlobKind::Dedicated => self.collect_dedicated(columns, idx, row_addr).await?, BlobKind::Packed => self.collect_packed(columns, idx, row_addr).await?, BlobKind::External => self.collect_external(columns, idx).await?, + BlobKind::Managed => self.collect_managed(columns, idx).await?, }; Ok(Some(file)) @@ -4368,6 +4424,49 @@ impl<'a> BlobV2ReadContext<'a> { )) } + async fn collect_managed( + &mut self, + columns: &BlobV2DescriptorColumns<'_>, + idx: usize, + ) -> Result { + let uri = columns.blob_uris.value(idx); + let position = columns.positions.value(idx); + let size = columns.sizes.value(idx); + let relative = lance_core::utils::blob::validate_managed_reference(uri, position, size)?; + let base_id = columns.blob_ids.value(idx); + let base = self + .dataset + .manifest + .base_paths + .get(&base_id) + .ok_or_else(|| { + Error::invalid_input(format!("Managed blob references unknown base_id {base_id}")) + })?; + let root = if let Some(root) = self.external_base_path_cache.get(&base_id) { + root.clone() + } else { + let root = base.extract_path(self.dataset.session.store_registry())?; + self.external_base_path_cache.insert(base_id, root.clone()); + root + }; + let store = if let Some(store) = self.store_cache.get(&base_id) { + store.clone() + } else { + let store = self.dataset.object_store(Some(base_id)).await?; + self.store_cache.insert(base_id, store.clone()); + store + }; + let path = join_base_and_relative_path(&root, relative.as_ref())?; + let source = shared_blob_source(&mut self.source_cache, store, &path); + Ok(BlobFile::with_source( + source, + position, + size, + BlobKind::Managed, + Some(uri.to_string()), + )) + } + async fn collect_external( &mut self, columns: &BlobV2DescriptorColumns<'_>, @@ -4497,6 +4596,7 @@ async fn resolve_blob_read_location( }; let location = BlobReadLocation { + base_id: data_file.base_id, object_store, data_file_dir, data_file_key, @@ -4506,6 +4606,241 @@ async fn resolve_blob_read_location( Ok(location) } +fn field_contains_blob(field: &LanceField) -> bool { + field.is_blob_v2() || field.children.iter().any(field_contains_blob) +} + +fn visit_managed_row( + field: &LanceField, + array: &dyn Array, + row: usize, + visit: &mut impl FnMut(u32, &str, Range) -> Result<()>, +) -> Result<()> { + if array.is_null(row) { + return Ok(()); + } + if field.is_blob_v2() { + let values = array.as_struct_opt().ok_or_else(|| { + Error::internal("Expected Blob descriptor array during reference scan") + })?; + let columns = BlobV2DescriptorColumns::new(values); + if !columns.is_null_blob(row) + && BlobKind::try_from(columns.kinds.value(row))? == BlobKind::Managed + { + let uri = columns.blob_uris.value(row); + lance_core::utils::blob::validate_managed_reference( + uri, + columns.positions.value(row), + columns.sizes.value(row), + )?; + visit( + columns.blob_ids.value(row), + uri, + columns.positions.value(row) + ..columns.positions.value(row) + columns.sizes.value(row), + )?; + } + return Ok(()); + } + match array.data_type() { + ArrowDataType::Struct(_) => { + for child in field + .children + .iter() + .filter(|child| field_contains_blob(child)) + { + let values = array + .as_struct() + .column_by_name(&child.name) + .ok_or_else(|| Error::internal("Missing Blob child during reference scan"))?; + visit_managed_row(child, values.as_ref(), row, visit)?; + } + } + ArrowDataType::List(_) => { + let values = array.as_list::(); + let child = field + .children + .first() + .ok_or_else(|| Error::internal("Missing Blob list item field"))?; + for index in values.value_offsets()[row]..values.value_offsets()[row + 1] { + visit_managed_row(child, values.values().as_ref(), index as usize, visit)?; + } + } + ArrowDataType::LargeList(_) => { + let values = array.as_list::(); + let child = field + .children + .first() + .ok_or_else(|| Error::internal("Missing Blob list item field"))?; + for index in values.value_offsets()[row]..values.value_offsets()[row + 1] { + visit_managed_row(child, values.values().as_ref(), index as usize, visit)?; + } + } + _ => { + return Err(Error::not_supported(format!( + "Cannot scan Managed references in {}", + array.data_type() + ))); + } + } + Ok(()) +} + +/// Discover referenced objects from stored descriptors without reading payloads. +async fn managed_references(dataset: &Dataset) -> Result> { + let fields = dataset + .schema() + .fields + .iter() + .filter(|field| field_contains_blob(field)) + .collect::>(); + if fields.is_empty() { + return Ok(HashSet::new()); + } + let names = fields + .iter() + .map(|field| field.name.as_str()) + .collect::>(); + let mut scan = dataset.scan(); + scan.project(&names)?; + let mut stream = scan.try_into_stream().await?; + let mut references = HashSet::new(); + while let Some(batch) = stream.try_next().await? { + for field in &fields { + let values = batch + .column_by_name(&field.name) + .ok_or_else(|| Error::internal("Missing Blob column during reference scan"))?; + for row in 0..batch.num_rows() { + visit_managed_row(field, values.as_ref(), row, &mut |id, uri, _range| { + references.insert((id, uri.to_string())); + Ok(()) + })?; + } + } + } + Ok(references) +} + +/// Resolve only objects within the cleanup owner's deletion jurisdiction. +pub(super) async fn managed_paths(dataset: &Dataset, owner: &Dataset) -> Result> { + let references = managed_references(dataset).await?; + let mut paths = HashSet::new(); + let mut bases = HashMap::new(); + for (id, uri) in references { + if let std::collections::hash_map::Entry::Vacant(entry) = bases.entry(id) { + let base = dataset.manifest.base_paths.get(&id).ok_or_else(|| { + Error::invalid_input(format!("Managed reference scan found unknown base_id {id}")) + })?; + let store = dataset.object_store(Some(id)).await?; + let root = base.extract_path(dataset.session.store_registry())?; + entry.insert((store.store_prefix == owner.object_store.store_prefix, root)); + } + let (same_store, root) = bases + .get(&id) + .ok_or_else(|| Error::internal("Missing resolved Managed base"))?; + if *same_store { + let path = join_base_and_relative_path(root, &uri)?; + if path.prefix_match(&owner.base).is_some() { + paths.insert(path); + } + } + } + Ok(paths) +} + +/// Rewrite descriptors in the same snapshot base namespace, copying only Inline +/// payloads. The commit must also publish the explicit default-base binding. +pub(super) async fn preserve_managed_descriptors( + dataset: &Arc, + field_id: u32, + values: &StructArray, + row_addrs: &[u64], + field: &ArrowField, +) -> Result { + if values.len() != row_addrs.len() { + return Err(Error::internal( + "Blob descriptors and row addresses have different lengths", + )); + } + let columns = BlobV2DescriptorColumns::new(values); + let mut builder = BlobDescriptorArrayBuilder::new_with_metadata( + field.name(), + field.is_nullable(), + field.metadata().clone(), + ); + let mut context = BlobV2ReadContext::new(dataset, field_id); + for (row, row_addr) in row_addrs.iter().enumerate() { + if columns.is_null_blob(row) { + builder.push_null()?; + continue; + } + let kind = BlobKind::try_from(columns.kinds.value(row))?; + match kind { + BlobKind::Managed => { + context.collect_managed(&columns, row).await?; + builder.push(BlobDescriptor::Managed { + base_id: columns.blob_ids.value(row), + uri: columns.blob_uris.value(row).to_string(), + offset: columns.positions.value(row), + size: columns.sizes.value(row), + })?; + } + BlobKind::Packed | BlobKind::Dedicated => { + let location = context.blob_read_location(*row_addr).await?; + let base = if let Some(id) = location.base_id { + dataset + .manifest + .base_paths + .get(&id) + .cloned() + .ok_or_else(|| { + Error::invalid_input(format!( + "Blob source references unknown base_id {id}" + )) + })? + } else { + dataset.managed_default_base()? + }; + let root = base.extract_path(dataset.session.store_registry())?; + let path = blob_path( + &location.data_file_dir, + &location.data_file_key, + columns.blob_ids.value(row), + ); + let uri = path + .prefix_match(&root) + .ok_or_else(|| { + Error::internal(format!("Blob source {path} is outside base {}", base.path)) + })? + .map(|part| part.as_ref().to_string()) + .collect::>() + .join("/"); + builder.push(BlobDescriptor::Managed { + base_id: base.id, + uri, + offset: if kind == BlobKind::Dedicated { + 0 + } else { + columns.positions.value(row) + }, + size: columns.sizes.value(row), + })?; + } + BlobKind::External => builder.push(BlobDescriptor::External { + base_id: columns.blob_ids.value(row), + uri: columns.blob_uris.value(row).to_string(), + offset: columns.positions.value(row), + size: columns.sizes.value(row), + })?, + BlobKind::Inline => { + let file = context.collect_inline(&columns, row, *row_addr).await?; + builder.push_inline(file.read().await?)?; + } + } + } + Ok(builder.finish()?.into_parts().1) +} + pub(super) fn data_file_key_from_path(path: &str) -> &str { let filename = path.rsplit('/').next().unwrap_or(path); filename.strip_suffix(".lance").unwrap_or(filename) @@ -4600,6 +4935,454 @@ mod tests { expected: Vec, } + #[rstest] + #[case::inline(1024, BlobKind::Inline)] + #[case::packed(100 * 1024, BlobKind::Managed)] + #[case::dedicated(8 * 1024 * 1024, BlobKind::Managed)] + #[tokio::test] + async fn managed_write_read( + #[case] size: usize, + #[case] kind: BlobKind, + #[values(false, true)] maximum_base_id: bool, + #[values(false, true)] primary_write: bool, + #[values(LanceFileVersion::V2_2, LanceFileVersion::V2_3)] version: LanceFileVersion, + ) { + let test_dir = TempStrDir::default(); + let payload = vec![42u8; size]; + let mut blobs = BlobArrayBuilder::new(3); + blobs.push_bytes(&payload).unwrap(); + blobs.push_bytes(b"").unwrap(); + blobs.push_null().unwrap(); + let schema = Arc::new(Schema::new(vec![blob_field("blob", true)])); + let batch = RecordBatch::try_new(schema.clone(), vec![blobs.finish().unwrap()]).unwrap(); + let dataset = Arc::new( + Dataset::write( + RecordBatchIterator::new(vec![Ok(batch)], schema), + &test_dir, + Some(WriteParams { + data_storage_version: Some(version), + initial_bases: maximum_base_id.then(|| { + vec![BasePath::new( + u32::MAX, + test_dir.to_string(), + Some("payload".to_string()), + true, + )] + }), + target_bases: (maximum_base_id && !primary_write).then_some(vec![u32::MAX]), + ..Default::default() + }), + ) + .await + .unwrap(), + ); + let flag = lance_table::feature_flags::FLAG_MANAGED_BLOBS; + assert_ne!(dataset.manifest.reader_feature_flags & flag, 0); + assert_ne!(dataset.manifest.writer_feature_flags & flag, 0); + for (_, uri) in super::managed_references(&dataset).await.unwrap() { + assert!(uri.starts_with("_blobs/"), "{uri}"); + assert_eq!(uri.split('/').count(), 2); + Uuid::parse_str(uri.trim_start_matches("_blobs/").trim_end_matches(".blob")).unwrap(); + } + let id = if maximum_base_id && !primary_write { + u32::MAX + } else { + 0 + }; + assert_eq!(dataset.manifest.base_paths[&id].path, dataset.uri); + if !maximum_base_id || !primary_write { + assert_eq!(dataset.manifest.base_paths.len(), 1); + } + let values = dataset + .take_blobs_by_indices(&[0, 1, 2], "blob") + .await + .unwrap(); + assert_eq!(values[0].as_ref().unwrap().kind(), kind); + assert_eq!( + values[0].as_ref().unwrap().read().await.unwrap().as_ref(), + payload + ); + assert!(values[1].as_ref().unwrap().read().await.unwrap().is_empty()); + assert!(values[2].is_none()); + if kind == BlobKind::Managed { + let mut missing_base = dataset.as_ref().clone(); + let manifest = Arc::make_mut(&mut missing_base.manifest); + // The data file lives at the dataset root even when the writer + // selected an explicit base. Isolate the descriptor's base lookup. + for fragment in Arc::make_mut(&mut manifest.fragments) { + for file in fragment.referenced_lance_files_mut() { + file.base_id = None; + } + } + manifest.base_paths.clear(); + let error = Arc::new(missing_base) + .take_blobs_by_indices(&[0], "blob") + .await + .err() + .unwrap(); + assert!(matches!(error, Error::InvalidInput { .. })); + assert!( + error.to_string().contains(&format!("unknown base_id {id}")), + "{error}" + ); + } + } + + #[rstest] + #[tokio::test] + async fn managed_append_preserves_primary_base_root_semantics( + #[values(false, true)] is_dataset_root: bool, + ) { + let test_dir = TempStrDir::default(); + let payload = vec![7; 100 * 1024]; + let make_reader = || { + let mut blobs = BlobArrayBuilder::new(1); + blobs.push_bytes(&payload).unwrap(); + let schema = Arc::new(Schema::new(vec![blob_field("blob", false)])); + let batch = + RecordBatch::try_new(schema.clone(), vec![blobs.finish().unwrap()]).unwrap(); + RecordBatchIterator::new(vec![Ok(batch)], schema) + }; + Dataset::write( + make_reader(), + &test_dir, + Some(WriteParams { + data_storage_version: Some(LanceFileVersion::V2_2), + initial_bases: Some(vec![BasePath::new( + 7, + test_dir.to_string(), + None, + is_dataset_root, + )]), + ..Default::default() + }), + ) + .await + .unwrap(); + let dataset = Arc::new( + Dataset::write( + make_reader(), + &test_dir, + Some(WriteParams { + mode: WriteMode::Append, + data_storage_version: Some(LanceFileVersion::V2_3), + ..Default::default() + }), + ) + .await + .unwrap(), + ); + let blobs = dataset + .take_blobs_by_indices(&[0, 1], "blob") + .await + .unwrap(); + for blob in blobs { + assert_eq!(blob.unwrap().read().await.unwrap().as_ref(), payload); + } + let base = dataset.managed_default_base().unwrap(); + assert!(base.is_dataset_root); + if !is_dataset_root { + assert_eq!(base.id, 0); + } + assert_eq!( + dataset.manifest.base_paths[&7].is_dataset_root, + is_dataset_root + ); + assert_eq!(dataset.manifest.base_paths[&base.id], base); + } + + #[tokio::test] + async fn managed_empty_range_does_not_discover_object_length() { + let (store, recording) = recording_range_store(Bytes::from_static(b"not empty")); + let source = Arc::new(BlobSource::new(store, Path::from("data/a/payload.blob"))); + let blob = BlobFile::with_source(source, 4, 0, BlobKind::Managed, None); + assert!(blob.read().await.unwrap().is_empty()); + assert!(blob.read_range(0..0).await.unwrap().is_empty()); + assert_eq!(recording.head_requests.load(Ordering::Relaxed), 0); + assert!(recording.requested_ranges().is_empty()); + } + + #[rstest] + #[case::adopt_stable(true)] + #[case::reuse_managed(false)] + #[tokio::test] + async fn managed_compaction_preserves_sidecars_after_source_gc( + #[case] released: bool, + #[values(LanceFileVersion::V2_2, LanceFileVersion::V2_3)] version: LanceFileVersion, + ) { + mock_instant::thread_local::MockClock::set_system_time( + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap(), + ); + let test_dir = TempStrDir::default(); + let payload = vec![b'p'; 16]; + let dedicated = vec![b'd'; 32]; + let mut blobs = BlobArrayBuilder::new(6); + blobs.push_bytes(&payload).unwrap(); + blobs.push_bytes(b"inline").unwrap(); + blobs.push_null().unwrap(); + blobs.push_bytes(b"").unwrap(); + blobs.push_bytes(&dedicated).unwrap(); + blobs.push_bytes(b"inline").unwrap(); + let mut field = blob_field("blob", true); + let mut metadata = field.metadata().clone(); + metadata.extend(HashMap::from([ + ( + BLOB_INLINE_SIZE_THRESHOLD_META_KEY.to_string(), + "8".to_string(), + ), + ( + BLOB_DEDICATED_SIZE_THRESHOLD_META_KEY.to_string(), + "24".to_string(), + ), + ])); + field = field.with_metadata(metadata); + let schema = Arc::new(Schema::new(vec![field])); + let batch = RecordBatch::try_new(schema.clone(), vec![blobs.finish().unwrap()]).unwrap(); + let old_data = crate::utils::test::copy_test_data_to_tmp("v11.0.0/blob_sidecars").unwrap(); + let mut dataset = if released { + let mut dataset = Dataset::open(old_data.path_str().as_str()).await.unwrap(); + assert!(!dataset.manifest.has_managed_blobs()); + dataset + .update_config([("test", "metadata-only")]) + .await + .unwrap(); + assert!(!dataset.manifest.has_managed_blobs()); + dataset + } else { + Dataset::write( + RecordBatchIterator::new(vec![Ok(batch)], schema), + &test_dir, + Some(WriteParams { + max_rows_per_file: 3, + data_storage_version: Some(version), + ..Default::default() + }), + ) + .await + .unwrap() + }; + assert_eq!(dataset.manifest.fragments.len(), 2); + let before = dataset + .object_store + .read_dir_all(&dataset.base, None) + .try_collect::>() + .await + .unwrap(); + let sidecars = before + .iter() + .filter(|object| object.location.extension() == Some("blob")) + .map(|object| (object.location.clone(), object.size)) + .collect::>(); + assert!(!sidecars.is_empty()); + mock_instant::thread_local::MockClock::set_system_time( + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap(), + ); + crate::dataset::optimize::compact_files( + &mut dataset, + crate::dataset::optimize::CompactionOptions { + data_storage_version: Some(version), + ..Default::default() + }, + None, + ) + .await + .unwrap(); + assert_eq!(dataset.manifest.fragments.len(), 1); + assert!(dataset.manifest.has_managed_blobs()); + let cleanup_stats = crate::dataset::cleanup::cleanup_old_versions( + &dataset, + crate::dataset::cleanup::CleanupPolicy { + before_timestamp: Some(Utc::now() + chrono::TimeDelta::seconds(1)), + ..Default::default() + }, + ) + .await + .unwrap(); + let after = dataset + .object_store + .read_dir_all(&dataset.base, None) + .try_collect::>() + .await + .unwrap(); + let remaining_sidecars = after + .iter() + .filter(|object| object.location.extension() == Some("blob")) + .map(|object| (object.location.clone(), object.size)) + .collect::>(); + assert_eq!(sidecars, remaining_sidecars); + for source in before + .iter() + .filter(|object| object.location.extension() == Some("lance")) + { + assert!( + !after + .iter() + .any(|object| object.location == source.location), + "source={source:?}, stats={cleanup_stats:?}, fragments={:?}", + dataset.manifest.fragments + ); + } + let values = Arc::new(dataset) + .take_blobs_by_indices(&[0, 1, 2, 3, 4, 5], "blob") + .await + .unwrap(); + for index in [0, 4] { + assert_eq!(values[index].as_ref().unwrap().kind(), BlobKind::Managed); + assert_eq!( + values[index] + .as_ref() + .unwrap() + .read() + .await + .unwrap() + .as_ref(), + if index == 0 { &payload } else { &dedicated }.as_slice() + ); + } + for index in [1, 5] { + assert_eq!(values[index].as_ref().unwrap().kind(), BlobKind::Inline); + assert_eq!( + values[index] + .as_ref() + .unwrap() + .read() + .await + .unwrap() + .as_ref(), + b"inline" + ); + } + assert!(values[2].is_none()); + assert!(values[3].as_ref().unwrap().read().await.unwrap().is_empty()); + } + + #[tokio::test] + async fn managed_flag_survives_restore_of_unflagged_snapshot() { + let dir = crate::utils::test::copy_test_data_to_tmp("v11.0.0/blob_sidecars").unwrap(); + let mut dataset = Dataset::open(&dir.path_str()).await.unwrap(); + let mut old = dataset.checkout_version(1).await.unwrap(); + assert!(!old.manifest.has_managed_blobs()); + crate::dataset::optimize::compact_files(&mut dataset, Default::default(), None) + .await + .unwrap(); + assert!(dataset.manifest.has_managed_blobs()); + old.restore().await.unwrap(); + assert!(old.manifest.has_managed_blobs()); + assert_ne!( + old.manifest.writer_feature_flags & lance_table::feature_flags::FLAG_MANAGED_BLOBS, + 0 + ); + let values = Arc::new(old) + .take_blobs_by_indices(&[0, 4], "blob") + .await + .unwrap(); + assert_eq!(values[0].as_ref().unwrap().kind(), BlobKind::Packed); + assert_eq!(values[1].as_ref().unwrap().kind(), BlobKind::Dedicated); + assert_eq!( + values[0].as_ref().unwrap().read().await.unwrap().as_ref(), + b"pppppppppppppppp" + ); + } + + #[rstest] + #[tokio::test] + async fn managed_clone_preserves_existing_retention_contract( + #[values(false, true)] deep: bool, + ) { + let source_dir = TempStrDir::default(); + let clone_dir = TempStrDir::default(); + let external_dir = TempDir::default(); + let external_path = external_dir.std_path().join("external.bin"); + std::fs::write(&external_path, b"external").unwrap(); + let payload = vec![9; 100 * 1024]; + let mut blobs = BlobArrayBuilder::new(6); + for row in 0..6 { + match row { + 2 => blobs + .push_uri(format!("file://{}", external_path.display())) + .unwrap(), + 4 => blobs.push_bytes(b"inline").unwrap(), + _ => blobs.push_bytes(&payload).unwrap(), + } + } + let schema = Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int32, false), + blob_field("blob", true), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int32Array::from(vec![0, 1, 2, 3, 4, 5])), + blobs.finish().unwrap(), + ], + ) + .unwrap(); + let mut source = Dataset::write( + RecordBatchIterator::new(vec![Ok(batch)], schema), + &source_dir, + Some(WriteParams { + data_storage_version: Some(LanceFileVersion::V2_3), + max_rows_per_file: 3, + allow_external_blob_outside_bases: true, + ..Default::default() + }), + ) + .await + .unwrap(); + source.delete("id IN (1, 3, 5)").await.unwrap(); + let version = source.version_id(); + let cloned = if deep { + source.deep_clone(&clone_dir, version, None).await.unwrap() + } else { + source.tags().create("clone_source", version).await.unwrap(); + source + .shallow_clone(&clone_dir, "clone_source", None) + .await + .unwrap() + }; + source.checkout_latest().await.unwrap(); + assert_eq!(source.version_id(), version); + if deep { + std::fs::remove_dir_all(source_dir.as_ref()).unwrap(); + } else { + source.delete("true").await.unwrap(); + let version = source.version_id(); + crate::dataset::cleanup::cleanup_old_versions( + &source, + crate::dataset::cleanup::CleanupPolicy { + before_version: Some(version), + error_if_tagged_old_versions: false, + delete_unverified: true, + ..Default::default() + }, + ) + .await + .unwrap(); + source.checkout_latest().await.unwrap(); + assert_eq!(source.version_id(), version); + assert_eq!(source.count_rows(None).await.unwrap(), 0); + } + assert_eq!(cloned.count_rows(None).await.unwrap(), 3); + let blobs = Arc::new(cloned) + .take_blobs_by_indices(&[0, 1, 2], "blob") + .await + .unwrap(); + let expected = [ + (BlobKind::Managed, payload.as_slice()), + (BlobKind::External, b"external".as_slice()), + (BlobKind::Inline, b"inline".as_slice()), + ]; + for (blob, (kind, bytes)) in blobs.into_iter().zip(expected) { + let blob = blob.unwrap(); + assert_eq!(blob.kind(), kind); + assert_eq!(blob.read().await.unwrap().as_ref(), bytes); + } + } + #[test] fn test_blob_version_rejects_malformed_v2_descriptor_layout() { let descriptions = StructArray::try_new( @@ -6372,8 +7155,11 @@ mod tests { assert_eq!(data_file.column_indices.as_ref(), &[0, 1]); } + #[rstest] #[tokio::test] - async fn test_write_and_take_nested_blob_v2() { + async fn test_write_and_take_nested_blob_v2( + #[values(LanceFileVersion::V2_2, LanceFileVersion::V2_3)] version: LanceFileVersion, + ) { let test_dir = TempStrDir::default(); let packed_payload = vec![0x4A; super::INLINE_MAX + 1024]; @@ -6391,7 +7177,7 @@ mod tests { reader, &test_dir, Some(WriteParams { - data_storage_version: Some(LanceFileVersion::V2_2), + data_storage_version: Some(version), ..Default::default() }), ) @@ -6426,7 +7212,7 @@ mod tests { .unwrap() .as_primitive::() .value(1), - BlobKind::Packed as u8 + BlobKind::Managed as u8 ); let blobs = dataset @@ -6462,8 +7248,11 @@ mod tests { assert_eq!(filtered.num_rows(), 2); } + #[rstest] #[tokio::test] - async fn test_write_and_take_nested_complete_blob_v2() { + async fn test_write_and_take_nested_complete_blob_v2( + #[values(LanceFileVersion::V2_2, LanceFileVersion::V2_3)] version: LanceFileVersion, + ) { let test_dir = TempStrDir::default(); let packed_payload = vec![0x4A; super::INLINE_MAX + 1024]; @@ -6484,7 +7273,7 @@ mod tests { reader, &test_dir, Some(WriteParams { - data_storage_version: Some(LanceFileVersion::V2_2), + data_storage_version: Some(version), ..Default::default() }), ) @@ -6519,7 +7308,7 @@ mod tests { .unwrap() .as_primitive::() .value(1), - BlobKind::Packed as u8 + BlobKind::Managed as u8 ); let blobs = dataset @@ -6555,8 +7344,11 @@ mod tests { assert_eq!(filtered.num_rows(), 2); } + #[rstest] #[tokio::test] - async fn test_write_and_scan_list_blob_v2_descriptions() { + async fn test_write_and_scan_list_blob_v2_descriptions( + #[values(LanceFileVersion::V2_2, LanceFileVersion::V2_3)] version: LanceFileVersion, + ) { let test_dir = TempStrDir::default(); let packed_payload = vec![0x4B; super::INLINE_MAX + 1024]; @@ -6598,7 +7390,7 @@ mod tests { reader, &test_dir, Some(WriteParams { - data_storage_version: Some(LanceFileVersion::V2_2), + data_storage_version: Some(version), enable_stable_row_ids: true, ..Default::default() }), @@ -6638,7 +7430,7 @@ mod tests { .unwrap() .as_primitive::(); assert_eq!(kinds.value(0), BlobKind::Inline as u8); - assert_eq!(kinds.value(2), BlobKind::Packed as u8); + assert_eq!(kinds.value(2), BlobKind::Managed as u8); assert_eq!(kinds.value(3), BlobKind::Inline as u8); let filtered = dataset @@ -6712,8 +7504,11 @@ mod tests { } } + #[rstest] #[tokio::test] - async fn test_write_and_scan_struct_nested_list_blob_v2() { + async fn test_write_and_scan_struct_nested_list_blob_v2( + #[values(LanceFileVersion::V2_2, LanceFileVersion::V2_3)] version: LanceFileVersion, + ) { let test_dir = TempStrDir::default(); let mut blob_builder = BlobArrayBuilder::new(2); @@ -6758,7 +7553,7 @@ mod tests { reader, &test_dir, Some(WriteParams { - data_storage_version: Some(LanceFileVersion::V2_2), + data_storage_version: Some(version), ..Default::default() }), ) @@ -7753,8 +8548,8 @@ mod tests { #[rstest] #[case::inline(BlobKind::Inline, 1024 * 1024, 1024 * 1024)] - #[case::packed(BlobKind::Packed, 0, 1024 * 1024)] - #[case::dedicated(BlobKind::Dedicated, 0, 1)] + #[case::packed(BlobKind::Managed, 0, 1024 * 1024)] + #[case::dedicated(BlobKind::Managed, 0, 1)] #[tokio::test] async fn test_read_blob_ranges_avoids_whole_value_read_amplification_end_to_end( #[case] kind: BlobKind, @@ -7948,7 +8743,7 @@ mod tests { assert_eq!(blobs.len(), 1); let blob = blobs[0].as_ref().unwrap(); - assert_eq!(blob.kind(), BlobKind::Packed); + assert_eq!(blob.kind(), BlobKind::Managed); assert_eq!( blob.read().await.unwrap().as_ref(), fixture.expected.as_slice() @@ -7967,7 +8762,7 @@ mod tests { assert_eq!(blobs.len(), 1); let blob = blobs[0].as_ref().unwrap(); - assert_eq!(blob.kind(), BlobKind::Dedicated); + assert_eq!(blob.kind(), BlobKind::Managed); assert_eq!( blob.read().await.unwrap().as_ref(), fixture.expected.as_slice() @@ -7988,7 +8783,7 @@ mod tests { assert_eq!(blobs.len(), 1); let blob = blobs[0].as_ref().unwrap(); - assert_eq!(blob.kind(), BlobKind::Packed); + assert_eq!(blob.kind(), BlobKind::Managed); assert_eq!( blob.read().await.unwrap().as_ref(), fixture.expected.as_slice() @@ -8680,13 +9475,13 @@ mod tests { .unwrap() .as_primitive::() .value(0), - BlobKind::Packed as u8 + BlobKind::Managed as u8 ); let blobs = dataset.take_blobs_by_indices(&[0], "blob").await.unwrap(); assert_eq!(blobs.len(), 1); let blob = blobs[0].as_ref().unwrap(); - assert_eq!(blob.kind(), BlobKind::Packed); + assert_eq!(blob.kind(), BlobKind::Managed); assert_eq!(blob.read().await.unwrap().as_ref(), payload.as_slice()); } @@ -8731,7 +9526,7 @@ mod tests { let blobs = dataset.take_blobs_by_indices(&[0], "blob").await.unwrap(); assert_eq!(blobs.len(), 1); let blob = blobs[0].as_ref().unwrap(); - assert_eq!(blob.kind(), BlobKind::Packed); + assert_eq!(blob.kind(), BlobKind::Managed); assert_eq!(blob.read().await.unwrap().as_ref(), payload.as_slice()); } @@ -8780,13 +9575,13 @@ mod tests { .unwrap() .as_primitive::() .value(0), - BlobKind::Dedicated as u8 + BlobKind::Managed as u8 ); let blobs = dataset.take_blobs_by_indices(&[0], "blob").await.unwrap(); assert_eq!(blobs.len(), 1); let blob = blobs[0].as_ref().unwrap(); - assert_eq!(blob.kind(), BlobKind::Dedicated); + assert_eq!(blob.kind(), BlobKind::Managed); assert_eq!(blob.read().await.unwrap().as_ref(), payload.as_slice()); } diff --git a/rust/lance/src/dataset/blob/clone.rs b/rust/lance/src/dataset/blob/clone.rs new file mode 100644 index 00000000000..a28711c76fd --- /dev/null +++ b/rust/lance/src/dataset/blob/clone.rs @@ -0,0 +1,252 @@ +// SPDX-License-Identifier: Apache-2.0 +// SPDX-FileCopyrightText: Copyright The Lance Authors + +//! Copy Managed payloads and rebind descriptors for an independent deep clone. + +use std::collections::HashMap; +use std::sync::Arc; + +use arrow_array::cast::AsArray; +use arrow_array::{ + Array, ArrayRef, GenericListArray, RecordBatch, StringArray, StructArray, UInt32Array, +}; +use arrow_schema::DataType; +use lance_core::datatypes::{BlobHandling, BlobKind, Field}; +use lance_core::{Error, ROW_ADDR, Result}; +use lance_io::object_store::ObjectStore; +use lance_table::format::{Fragment, Manifest}; +use object_store::path::Path; +use uuid::Uuid; + +use super::{ + BlobV2DescriptorColumns, field_contains_blob, join_base_and_relative_path, managed_references, +}; +use crate::dataset::Dataset; +use crate::dataset::optimize::BlobV2BatchRewritePlan; + +type BlobCopies = HashMap<(u32, String), (u32, String)>; + +fn replace_array(copies: &BlobCopies, field: &Field, array: ArrayRef) -> Result { + if !field_contains_blob(field) { + return Ok(array); + } + if field.is_blob_v2() { + let values = array.as_struct(); + let column = |name: &str| { + values.column_by_name(name).ok_or_else(|| { + Error::internal(format!("Prepared Managed rewrite is missing {name}")) + }) + }; + let columns = BlobV2DescriptorColumns { + descriptions: values, + kinds: column("kind")?.as_primitive(), + positions: column("position")?.as_primitive(), + sizes: column("blob_size")?.as_primitive(), + blob_ids: column("blob_id")?.as_primitive(), + blob_uris: column("uri")?.as_string(), + }; + let mut ids = columns.blob_ids.values().to_vec(); + let mut uris = Vec::with_capacity(values.len()); + for (row, id) in ids.iter_mut().enumerate() { + let uri = columns.blob_uris.value(row); + if !columns.is_null_blob(row) + && columns.kinds.value(row) == BlobKind::Managed as u8 + && let Some(replacement) = copies.get(&(*id, uri.to_string())) + { + *id = replacement.0; + uris.push(Some(replacement.1.as_str())); + } else { + uris.push(columns.blob_uris.is_valid(row).then_some(uri)); + } + } + let arrays = values + .fields() + .iter() + .zip(values.columns()) + .map(|(field, array)| -> ArrayRef { + match field.name().as_str() { + "blob_id" => Arc::new(UInt32Array::new( + ids.clone().into(), + columns.blob_ids.nulls().cloned(), + )), + "uri" => Arc::new(StringArray::from(uris.clone())), + _ => array.clone(), + } + }) + .collect(); + return Ok(Arc::new(StructArray::try_new( + values.fields().clone(), + arrays, + values.nulls().cloned(), + )?)); + } + match array.data_type() { + DataType::Struct(_) => { + let values = array.as_struct(); + let arrays = field + .children + .iter() + .zip(values.columns()) + .map(|(child, array)| replace_array(copies, child, array.clone())) + .collect::>>()?; + Ok(Arc::new(StructArray::try_new( + values.fields().clone(), + arrays, + values.nulls().cloned(), + )?)) + } + DataType::List(child) => { + let values = array.as_list::(); + Ok(Arc::new(GenericListArray::::try_new( + child.clone(), + values.offsets().clone(), + replace_array(copies, &field.children[0], values.values().clone())?, + values.nulls().cloned(), + )?)) + } + DataType::LargeList(child) => { + let values = array.as_list::(); + Ok(Arc::new(GenericListArray::::try_new( + child.clone(), + values.offsets().clone(), + replace_array(copies, &field.children[0], values.values().clone())?, + values.nulls().cloned(), + )?)) + } + datatype => Err(Error::not_supported(format!( + "Managed replacement does not support {datatype}" + ))), + } +} + +/// Rewrite only top-level columns that contain selected references. Fragment +/// identity, row IDs, deletion vectors, and unrelated column files survive. +async fn rewrite_blob_columns( + copies: &BlobCopies, + source: Arc, + target: Arc, +) -> Result> { + let fields = source + .schema() + .fields + .iter() + .filter(|field| field_contains_blob(field)) + .collect::>(); + let mut updated = Vec::new(); + if copies.is_empty() || fields.is_empty() { + return Ok(updated); + } + for fragment in source.get_fragments() { + let columns = fields + .iter() + .map(|field| field.name.as_str()) + .collect::>(); + let write_schema = source.schema().project(&columns)?; + let mut projection = columns.clone(); + projection.push(ROW_ADDR); + let destination_fragment = target + .get_fragment(fragment.id()) + .ok_or_else(|| Error::internal("Managed rewrite target fragment is missing"))?; + let mut updater = destination_fragment + .updater( + Some(&projection), + Some((write_schema, target.schema().clone())), + Some(64), + Some(BlobHandling::BlobsDescriptions), + ) + .await?; + updater.allow_external_blob_outside_bases(); + while let Some(batch) = updater.next().await? { + let plan = BlobV2BatchRewritePlan::try_new( + source.schema(), + batch.schema().as_ref(), + false, + true, + )?; + let prepared = plan.transform_batch(&source, batch.clone()).await?; + let arrays = prepared + .schema() + .fields() + .iter() + .zip(prepared.columns()) + .map(|(field, array)| { + let field = source + .schema() + .field(field.name()) + .ok_or_else(|| Error::internal("Managed rewrite field is missing"))?; + replace_array(copies, field, array.clone()) + }) + .collect::>>()?; + updater + .update(RecordBatch::try_new(prepared.schema(), arrays)?) + .await?; + } + let mut fragment = updater.finish().await?; + let replaced = fragment + .files + .last() + .ok_or_else(|| Error::internal("Managed rewrite produced no data file"))? + .fields + .clone(); + for file in fragment.files.iter_mut().rev().skip(1) { + file.fields = file + .fields + .iter() + .map(|id| if replaced.contains(id) { -2 } else { *id }) + .collect::>() + .into(); + } + fragment + .files + .retain(|file| file.fields.iter().any(|id| *id != -2)); + updated.push(fragment); + } + Ok(updated) +} + +pub async fn copy_blob_columns( + source: Arc, + store: Arc, + base: Path, + uri: &str, + manifest: &mut Manifest, +) -> Result<()> { + let mut target = source.as_ref().clone(); + target.object_store = store; + target.base = base; + target.uri = uri.to_string(); + target.base_object_stores = Default::default(); + target.manifest = Arc::new(manifest.clone()); + manifest.bind_managed_base(target.managed_default_base()?)?; + target.manifest = Arc::new(manifest.clone()); + let target = Arc::new(target); + let target_base = target.managed_default_base()?; + let mut copies = HashMap::new(); + for (id, uri) in managed_references(&source).await? { + let source_base = source.manifest.base_paths.get(&id).ok_or_else(|| { + Error::invalid_input(format!("Managed clone references unknown base ID {id}")) + })?; + let path = join_base_and_relative_path( + &source_base.extract_path(source.session.store_registry())?, + &uri, + )?; + let target_uri = format!("_blobs/{}.blob", Uuid::new_v4()); + let target_path = join_base_and_relative_path(&target.base, &target_uri)?; + source + .object_store(Some(id)) + .await? + .copy_bulk(&path, &target.object_store, &target_path) + .await?; + copies.insert((id, uri), (target_base.id, target_uri)); + } + let rewritten = rewrite_blob_columns(&copies, source, target).await?; + let fragments = Arc::make_mut(&mut manifest.fragments); + for fragment in rewritten { + let existing = fragments + .iter_mut() + .find(|entry| entry.id == fragment.id) + .ok_or_else(|| Error::internal("Cloned fragment is missing"))?; + *existing = fragment; + } + Ok(()) +} diff --git a/rust/lance/src/dataset/cleanup.rs b/rust/lance/src/dataset/cleanup.rs index b9cad720efa..cdd173c02e9 100644 --- a/rust/lance/src/dataset/cleanup.rs +++ b/rust/lance/src/dataset/cleanup.rs @@ -51,6 +51,7 @@ use lance_core::{ }, }; use lance_table::{ + feature_flags::{ensure_can_read_manifest, ensure_can_write_manifest}, format::{IndexMetadata, Manifest}, io::{ commit::ManifestLocation, @@ -74,6 +75,7 @@ use tracing::{Span, debug, info, instrument, warn}; #[derive(Clone, Debug, Default)] struct ReferencedFiles { data_paths: HashSet, + managed_blob_paths: HashSet, delete_paths: HashSet, tx_paths: HashSet, index_uuids: HashSet, @@ -555,6 +557,8 @@ impl<'a> CleanupTask<'a> { let manifest_and_indexes = async { let manifest = read_manifest(&self.dataset.object_store, &location.path, location.size).await?; + ensure_can_read_manifest(&manifest)?; + ensure_can_write_manifest(&manifest)?; let indexes = read_manifest_indexes(&self.dataset.object_store, &location, &manifest).await?; Ok::<_, Error>((manifest, indexes)) @@ -589,7 +593,33 @@ impl<'a> CleanupTask<'a> { let is_latest = self.read_version <= manifest.version; let is_tagged = tagged_versions.contains(&manifest.version); let in_working_set = is_latest || !self.policy.should_clean(&manifest) || is_tagged; + let managed_paths = if manifest.has_managed_blobs() { + let snapshot = self + .dataset + .checkout_version((manifest.branch.as_deref(), Some(manifest.version))) + .await?; + match super::blob::managed_paths(&snapshot, self.dataset).await { + Ok(paths) => paths, + // A concurrent cleanup may already have removed expired data + // files. Losing deletion proof is safe; losing retained references + // is not. Other failures still stop cleanup before any deletion. + Err(error) if !in_working_set && error.is_not_found() => HashSet::new(), + Err(error) => return Err(error), + } + } else { + HashSet::new() + }; let mut inspection = inspection.lock().unwrap(); + let references = if in_working_set { + &mut inspection.referenced_files + } else { + &mut inspection.verified_files + }; + references.managed_blob_paths.extend( + managed_paths + .into_iter() + .map(|path| remove_prefix(&path, &self.dataset.base)), + ); // Track tagged old versions in case we want to return a `CleanupError` later. // Only track tagged when it is old. @@ -748,6 +778,10 @@ impl<'a> CleanupTask<'a> { build_listing_stream(self.dataset.versions_dir(), unmodified_since), build_listing_stream(self.dataset.transactions_dir(), unmodified_since), build_listing_stream(self.dataset.data_dir(), data_unmodified_since), + build_listing_stream( + self.dataset.base.clone().join("_blobs"), + data_unmodified_since, + ), // Index UUIDs from manifests being removed are proof that their files are // safe to delete. Scan every index artifact while that proof is available; // a retained-manifest cutoff can otherwise skip newer artifacts and lose @@ -974,6 +1008,33 @@ impl<'a> CleanupTask<'a> { } } Some("blob") => { + if inspection + .referenced_files + .managed_blob_paths + .contains(&relative_path) + { + return Ok(None); + } + if relative_path + .parts() + .next() + .is_some_and(|part| part.as_ref() == "_blobs") + { + let verified = inspection + .verified_files + .managed_blob_paths + .contains(&relative_path); + return if verified || !maybe_in_progress { + Ok(cleanup_file( + path, + CleanupFileKind::Data, + !verified, + size_bytes, + )) + } else { + Ok(None) + }; + } // Blob v2 sidecar files are keyed by the data file stem: // data/{data_file_key}/{obfuscated_blob_id:032b}.blob // @@ -1034,6 +1095,10 @@ impl<'a> CleanupTask<'a> { .verified_files .data_paths .contains(&parent_data_path) + || inspection + .verified_files + .managed_blob_paths + .contains(&relative_path) { Ok(cleanup_file(path, CleanupFileKind::Data, false, size_bytes)) } else { @@ -1123,26 +1188,38 @@ impl<'a> CleanupTask<'a> { let referenced_branches = &referenced_branches; async move { - let manifest_location = dataset - .commit_handler - .resolve_version_location( - &dataset.base, - *referenced_version, - &dataset.object_store.inner, - ) - .await?; + let manifest = async { + let manifest_location = dataset + .commit_handler + .resolve_version_location( + &dataset.base, + *referenced_version, + &dataset.object_store.inner, + ) + .await?; - let manifest = read_manifest( - &dataset.object_store, - &manifest_location.path, - manifest_location.size, - ) + read_manifest( + &dataset.object_store, + &manifest_location.path, + manifest_location.size, + ) + .await + } .await; - - if let Ok(manifest) = manifest - && policy.should_clean(&manifest) - { - referenced_branches.insert(branch_name.clone()); + match manifest { + Ok(manifest) => { + ensure_can_read_manifest(&manifest)?; + ensure_can_write_manifest(&manifest)?; + if policy.should_clean(&manifest) { + referenced_branches.insert(branch_name.clone()); + } + } + Err(error) if error.is_not_found() => { + // The source may be gone while descendants still use its files. + // Scan their manifests before deleting any parent data. + referenced_branches.insert(branch_name.clone()); + } + Err(error) => return Err(error), } Ok::<(), Error>(()) } @@ -1253,10 +1330,33 @@ impl<'a> CleanupTask<'a> { ) -> Result<()> { let manifest = read_manifest(&self.dataset.object_store, &location.path, location.size).await?; + ensure_can_read_manifest(&manifest)?; + ensure_can_write_manifest(&manifest)?; + let managed_paths = if manifest.has_managed_blobs() { + let snapshot = self + .dataset + .checkout_version((manifest.branch.as_deref(), Some(manifest.version))) + .await?; + super::blob::managed_paths(&snapshot, self.dataset).await? + } else { + HashSet::new() + }; let indexes = read_manifest_indexes(&self.dataset.object_store, &location, &manifest).await?; let mut inspection = inspection.lock().unwrap(); let mut is_referenced = false; + for path in managed_paths { + let relative = remove_prefix(&path, &self.dataset.base); + inspection + .verified_files + .managed_blob_paths + .remove(&relative); + inspection + .referenced_files + .managed_blob_paths + .insert(relative); + is_referenced = true; + } for fragment in manifest.fragments.iter() { for file in fragment.referenced_lance_files() { @@ -1695,7 +1795,10 @@ fn tagged_old_versions_cleanup_error( mod tests { use std::{ collections::HashMap, - sync::{Arc, Mutex}, + sync::{ + Arc, Mutex, + atomic::{AtomicUsize, Ordering}, + }, }; use super::*; @@ -1723,6 +1826,7 @@ mod tests { use lance_table::io::commit::RenameCommitHandler; use lance_testing::datagen::{BatchGenerator, IncrementingInt32, RandomVector, some_batch}; use mock_instant::thread_local::MockClock; + use object_store::ObjectStoreExt; use rstest::rstest; use uuid::Uuid; @@ -2743,6 +2847,13 @@ mod tests { fixture.create_some_data().await.unwrap(); fixture.block_commits(); assert!(fixture.append_some_data().await.is_err()); + let dataset = fixture.open().await.unwrap(); + let blob_path = dataset.base.clone().join("_blobs").join("orphan.blob"); + dataset + .object_store + .put(&blob_path, b"orphan") + .await + .unwrap(); let age = if old_files { TimeDelta::try_days(UNVERIFIED_THRESHOLD_DAYS + 1).unwrap() @@ -2768,6 +2879,11 @@ mod tests { let should_delete = override_opt.unwrap_or(false) || old_files; let after_count = fixture.count_files().await.unwrap(); + assert_eq!( + dataset.object_store.exists(&blob_path).await.unwrap(), + !should_delete, + "override={override_opt:?}, old_files={old_files}" + ); assert_eq!(removed.old_versions, 0); assert_eq!( removed.bytes_removed, @@ -2782,6 +2898,80 @@ mod tests { } } + #[rstest] + #[case::head_error(Some(0))] + #[case::read_error(Some(1))] + #[case::missing_source(None)] + #[tokio::test] + async fn cleanup_preserves_child_after_source_manifest_failure( + #[case] fail_request: Option, + ) { + MockClock::set_system_time(std::time::Duration::ZERO); + let fixture = MockDatasetFixture::try_new().unwrap(); + fixture.create_some_data().await.unwrap(); + let mut parent = fixture.open().await.unwrap(); + let source_manifest = parent.manifest_location.path.clone(); + let child = fixture + .create_branch_and_load(&mut parent, "child", (None, None)) + .await + .unwrap(); + let expected = child.scan().try_into_batch().await.unwrap(); + MockClock::set_system_time(TimeDelta::try_days(10).unwrap().to_std().unwrap()); + fixture.overwrite_some_data().await.unwrap(); + + let source_reads = Arc::new(AtomicUsize::new(0)); + if let Some(fail_request) = fail_request { + let source_reads = source_reads.clone(); + fixture.mock_store.policy.lock().unwrap().set_before_policy( + "fail_branch_source", + Arc::new(move |operation, path| { + // get_opts serves both HEAD (first) and the subsequent range read. + if operation == "get_opts" + && path == &source_manifest + && source_reads.fetch_add(1, Ordering::SeqCst) == fail_request + { + return Err(Error::internal("transient branch source read")); + } + Ok(()) + }), + ); + } else { + parent + .object_store + .inner + .delete(&source_manifest) + .await + .unwrap(); + } + let before = fixture.count_files().await.unwrap(); + let policy = CleanupPolicy { + before_version: Some(2), + error_if_tagged_old_versions: false, + ..Default::default() + }; + let result = fixture.run_cleanup_with_policy(policy.clone()).await; + if let Some(fail_request) = fail_request { + let error = result.unwrap_err(); + assert!(matches!(error, Error::IO { .. }), "{error}"); + assert_contains!(error.to_string(), "transient branch source read"); + assert_eq!(source_reads.load(Ordering::SeqCst), fail_request + 1); + let after = fixture.count_files().await.unwrap(); + assert_eq!(after.num_data_files, before.num_data_files); + assert_eq!(after.num_manifest_files, before.num_manifest_files); + fixture + .mock_store + .policy + .lock() + .unwrap() + .clear_before_policy("fail_branch_source"); + let retried = fixture.run_cleanup_with_policy(policy).await.unwrap(); + assert_eq!(retried.data_files_removed, 0); + } else { + assert_eq!(result.unwrap().data_files_removed, 0); + } + assert_eq!(child.scan().try_into_batch().await.unwrap(), expected); + } + #[tokio::test] async fn cleanup_old_index() { let fixture = MockDatasetFixture::try_new().unwrap(); diff --git a/rust/lance/src/dataset/data_file.rs b/rust/lance/src/dataset/data_file.rs index 48ce7593904..a97099dd1bc 100644 --- a/rust/lance/src/dataset/data_file.rs +++ b/rust/lance/src/dataset/data_file.rs @@ -182,8 +182,10 @@ impl DataFileTarget { Ok(()) } - /// Delete an abandoned target's staging parts, final file, and managed Blob - /// payloads, including objects left by failed writes or assembly. + /// Delete an abandoned target's staging parts, final file, and file-relative + /// Packed/Dedicated sidecars, including those left by failed writes or assembly. + /// Independent Managed objects are not owned by this target; unreferenced + /// objects follow the dataset's ordinary garbage-collection policy. /// /// The caller must stop all users and ensure no current or retained dataset /// version, checkpoint, or future commit needs this target. Lance does not diff --git a/rust/lance/src/dataset/fragment.rs b/rust/lance/src/dataset/fragment.rs index b2a2adb2ce8..05f4a2bbe2b 100644 --- a/rust/lance/src/dataset/fragment.rs +++ b/rust/lance/src/dataset/fragment.rs @@ -1801,6 +1801,7 @@ impl FileFragment { projection, stream.schema().as_ref(), false, + false, )?); Ok(stream .map(move |batch_result| { diff --git a/rust/lance/src/dataset/optimize.rs b/rust/lance/src/dataset/optimize.rs index 9632d338035..f3731450cbb 100644 --- a/rust/lance/src/dataset/optimize.rs +++ b/rust/lance/src/dataset/optimize.rs @@ -1561,7 +1561,11 @@ impl BlobV2FieldRewritePlan { } } - fn try_new(field: &LanceField, input_field: &ArrowField) -> Result { + fn try_new( + field: &LanceField, + input_field: &ArrowField, + preserve_managed: bool, + ) -> Result { if !field_contains_blob_v2(field) { return Ok(Self::passthrough(input_field)); } @@ -1582,9 +1586,14 @@ impl BlobV2FieldRewritePlan { }; let output_field = match BlobV2Layout::classify(input_children) { Some(BlobV2Layout::Logical) => Arc::new(input_field.clone()), - Some(BlobV2Layout::Descriptor) => { - transformed_arrow_field(field, BLOB_V2_LOGICAL_TYPE.clone()) - } + Some(BlobV2Layout::Descriptor) => transformed_arrow_field( + field, + if preserve_managed { + lance_core::datatypes::BLOB_V2_PREPARED_TYPE.clone() + } else { + BLOB_V2_LOGICAL_TYPE.clone() + }, + ), Some(actual) => { return Err(Error::invalid_input(format!( "Blob v2 field '{}' has {actual} input layout; expected logical or descriptor layout during rewrite", @@ -1619,7 +1628,7 @@ impl BlobV2FieldRewritePlan { .children .iter() .zip(input_children.iter()) - .map(|(child, input_child)| Self::try_new(child, input_child)) + .map(|(child, input_child)| Self::try_new(child, input_child, preserve_managed)) .collect::>>()?; let output_children = children .iter() @@ -1641,7 +1650,7 @@ impl BlobV2FieldRewritePlan { field.name )) })?; - let child = Box::new(Self::try_new(child, input_child)?); + let child = Box::new(Self::try_new(child, input_child, preserve_managed)?); Ok(Self::List { field_name: field.name.clone(), output_field: arrow_field_with_data_type( @@ -1658,7 +1667,7 @@ impl BlobV2FieldRewritePlan { field.name )) })?; - let child = Box::new(Self::try_new(child, input_child)?); + let child = Box::new(Self::try_new(child, input_child, preserve_managed)?); Ok(Self::LargeList { field_name: field.name.clone(), output_field: arrow_field_with_data_type( @@ -1697,12 +1706,15 @@ impl BlobV2FieldRewritePlan { Self::Blob { field_id, field_name, - .. + output_field, } => { let struct_arr = array.as_struct(); match BlobV2Layout::classify(struct_arr.fields()) { Some(BlobV2Layout::Logical) => Ok(array), Some(BlobV2Layout::Descriptor) => { + if matches!(output_field.data_type(), ArrowDataType::Struct(fields) if BlobV2Layout::classify(fields) == Some(BlobV2Layout::Prepared)) { + return super::blob::preserve_managed_descriptors(dataset, *field_id, struct_arr, &row_addrs, output_field).await; + } let descriptor = BlobV2Descriptor::try_from_struct(struct_arr, field_name)?; let classification = @@ -1868,6 +1880,7 @@ impl BlobV2BatchRewritePlan { schema: &lance_core::datatypes::Schema, input_schema: &ArrowSchema, keep_row_addr: bool, + preserve_managed: bool, ) -> Result { let row_addr_idx = input_schema .column_with_name(lance_core::ROW_ADDR) @@ -1889,7 +1902,7 @@ impl BlobV2BatchRewritePlan { continue; } let field_plan = if let Some(field) = schema.field(input_field.name()) { - BlobV2FieldRewritePlan::try_new(field, input_field)? + BlobV2FieldRewritePlan::try_new(field, input_field, preserve_managed)? } else { BlobV2FieldRewritePlan::passthrough(input_field) }; @@ -1949,7 +1962,8 @@ pub(crate) async fn transform_blob_v2_batch( batch: RecordBatch, keep_row_addr: bool, ) -> Result { - let plan = BlobV2BatchRewritePlan::try_new(schema, batch.schema().as_ref(), keep_row_addr)?; + let plan = + BlobV2BatchRewritePlan::try_new(schema, batch.schema().as_ref(), keep_row_addr, false)?; plan.transform_batch(dataset, batch).await } @@ -2455,6 +2469,10 @@ async fn rewrite_files( dataset.schema(), schema.as_ref(), false, + matches!( + write_version, + ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 + ), )?); let transformed_schema = rewrite_plan.output_schema.clone(); let transformed = reader_with_progress.then(move |batch_result| { @@ -9232,8 +9250,8 @@ mod tests { input_schema.fields[0].unload_blobs_recursive(); let input_field = Field::from(&input_schema.fields[0]); - let plan = - BlobV2FieldRewritePlan::try_new(&logical_schema.fields[0], &input_field).unwrap(); + let plan = BlobV2FieldRewritePlan::try_new(&logical_schema.fields[0], &input_field, false) + .unwrap(); let BlobV2FieldRewritePlan::Struct { children, .. } = &plan else { panic!("nested blob field should produce a struct rewrite plan"); }; diff --git a/rust/lance/src/dataset/schema_evolution.rs b/rust/lance/src/dataset/schema_evolution.rs index 8c1977d1605..a8c29f7fb7f 100644 --- a/rust/lance/src/dataset/schema_evolution.rs +++ b/rust/lance/src/dataset/schema_evolution.rs @@ -513,10 +513,23 @@ async fn cleanup_new_column_data_files(fragments: &[FileFragment], new_fragments }) .collect::>(); + let dataset = first_fragment.dataset(); + let cleanup_bases = match dataset.managed_default_base() { + Ok(base) => vec![super::write::TargetBaseInfo { + base_id: base.id, + object_store: dataset.object_store.clone(), + base_dir: dataset.base.clone(), + is_dataset_root: true, + }], + Err(error) => { + log::warn!("Cannot resolve primary base for failed column write cleanup: {error}"); + Vec::new() + } + }; cleanup_data_fragments( - &first_fragment.dataset().object_store, - &first_fragment.dataset().base, - None, + &dataset.object_store, + &dataset.base, + Some(&cleanup_bases), &fragments_to_cleanup, ) .await; @@ -1831,6 +1844,17 @@ mod test { baseline_files, "add_columns should clean files written by the current unfinished writer" ); + let blob_dir = StdPath::new(test_uri).join("_blobs"); + assert!(!file_paths_in(&blob_dir).is_empty()); + // Failed uploads use the existing orphan policy. There is no concurrent + // writer in this test, so explicit unverified cleanup can reclaim them. + dataset + .cleanup_with_policy(super::super::cleanup::CleanupPolicy { + delete_unverified: true, + ..Default::default() + }) + .await?; + assert!(file_paths_in(&blob_dir).is_empty()); Ok(()) } @@ -1928,15 +1952,13 @@ mod test { }) .expect("checkpoint should record the newly written data file"); let new_file_path = StdPath::new(test_uri).join("data").join(&new_file.path); - let new_blob_dir = StdPath::new(test_uri) - .join("data") - .join(StdPath::new(&new_file.path).file_stem().unwrap()); + let new_blob_dir = StdPath::new(test_uri).join("_blobs"); assert!( new_file_path.exists(), "cleanup must not delete data files after checkpoint takes ownership" ); assert!( - new_blob_dir.exists(), + !file_paths_in(&new_blob_dir).is_empty(), "cleanup must not delete blob sidecars after checkpoint takes ownership" ); @@ -2177,15 +2199,13 @@ mod test { }) .expect("checkpoint should record the newly written data file"); let new_file_path = StdPath::new(test_uri).join("data").join(&new_file.path); - let new_blob_dir = StdPath::new(test_uri) - .join("data") - .join(StdPath::new(&new_file.path).file_stem().unwrap()); + let new_blob_dir = StdPath::new(test_uri).join("_blobs"); assert!( new_file_path.exists(), "cleanup must not delete data files after checkpoint takes ownership" ); assert!( - new_blob_dir.exists(), + !file_paths_in(&new_blob_dir).is_empty(), "cleanup must not delete blob sidecars after checkpoint takes ownership" ); diff --git a/rust/lance/src/dataset/tests/data_file_part.rs b/rust/lance/src/dataset/tests/data_file_part.rs index 35ae388c7c4..90b612a0058 100644 --- a/rust/lance/src/dataset/tests/data_file_part.rs +++ b/rust/lance/src/dataset/tests/data_file_part.rs @@ -507,6 +507,7 @@ async fn abandoned_target_cleanup_includes_failed_writes_and_preserves_other_tar #[case::blob_registered(true, Some(7))] #[tokio::test] async fn restores_target_and_completed_parts_from_checkpoint( + #[values(LanceFileVersion::V2_2, LanceFileVersion::V2_3)] version: LanceFileVersion, #[case] has_blob: bool, #[case] base_id: Option, ) { @@ -538,7 +539,7 @@ async fn restores_target_and_completed_parts_from_checkpoint( RecordBatchIterator::new([Ok(original.clone())], original.schema()), dataset_uri.as_str(), Some(WriteParams { - data_storage_version: Some(LanceFileVersion::V2_2), + data_storage_version: Some(version), max_rows_per_file: 2, initial_bases: base_id.map(|id| { vec![BasePath { @@ -837,14 +838,12 @@ async fn blob_parts_write_sidecars_in_final_namespace_and_concat_descriptors() { let mut invalid = serde_json::to_value(&first).unwrap(); invalid["blob_ids"] = serde_json::json!({"start": 20, "end": 30}); let invalid: DataFilePart = serde_json::from_value(invalid).unwrap(); - let error = dataset + // Managed references use base IDs, so their identity is independent of the + // file-local sidecar lease recorded in the checkpoint. + dataset .concat_data_file_parts(&target, &[invalid]) .await - .unwrap_err(); - assert!( - error.to_string().contains("outside declared range"), - "{error}" - ); + .unwrap(); let second = write_part( &dataset, &target, diff --git a/rust/lance/src/dataset/tests/dataset_io.rs b/rust/lance/src/dataset/tests/dataset_io.rs index 1b723da363f..74cc1fd731a 100644 --- a/rust/lance/src/dataset/tests/dataset_io.rs +++ b/rust/lance/src/dataset/tests/dataset_io.rs @@ -1929,6 +1929,9 @@ async fn test_rle_v2_shallow_clone_preserves_v23_storage() { .await .unwrap(); + assert!(dataset.manifest.base_paths.is_empty()); + assert!(!dataset.manifest.has_managed_blobs()); + let clone = dataset .shallow_clone(clone_uri.as_str(), dataset.version().version, None) .await diff --git a/rust/lance/src/dataset/versions/mod.rs b/rust/lance/src/dataset/versions/mod.rs index be99a42bbe4..368e934b13e 100644 --- a/rust/lance/src/dataset/versions/mod.rs +++ b/rust/lance/src/dataset/versions/mod.rs @@ -617,6 +617,7 @@ pub async fn open_writer( }; if schema.fields_pre_order().any(Field::is_blob_v2) { write::open_current_blob_v2_writer( + version, create_file_writer, object_store, schema, @@ -651,16 +652,24 @@ pub async fn open_update_writer( } ConcreteFileVersion::V1 | ConcreteFileVersion::V2_0 | ConcreteFileVersion::V2_1 => None, }; + let mut options = WriterOptions::update( + dataset.session.store_registry(), + external_base_resolver, + allow_external_blob_outside_bases, + ); + if matches!( + version, + ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 + ) && schema.fields_pre_order().any(Field::is_blob_v2) + { + options.base_id = Some(dataset.managed_default_base()?.id); + } open_writer( version, &dataset.object_store, schema, &dataset.base, - WriterOptions::update( - dataset.session.store_registry(), - external_base_resolver, - allow_external_blob_outside_bases, - ), + options, ) .await } diff --git a/rust/lance/src/dataset/write.rs b/rust/lance/src/dataset/write.rs index cc360d02a08..2c7ee1e0c6a 100644 --- a/rust/lance/src/dataset/write.rs +++ b/rust/lance/src/dataset/write.rs @@ -118,19 +118,19 @@ impl Dataset { /// Encode one managed part and return its serializable description. /// /// Lance generates a unique staging name in the target's base. Managed Blob - /// payloads are written directly beneath the sidecar directory selected by the final - /// target using IDs from `blob_ids`; every non-empty logical Inline value is - /// spilled to Packed or Dedicated storage so final concatenation never copies - /// Blob payload bytes. + /// payloads use independent `_blobs/.blob` objects in that base. Every + /// non-empty logical Inline value is spilled so final concatenation never + /// copies Blob payload bytes. Parts still require disjoint `blob_ids` + /// reservations; Managed descriptors encode base IDs rather than these IDs. /// Every use of `target` must refer to the same dataset and resolved base; /// associating a target with that storage context is the caller's /// responsibility. /// Persist the target before writing. A failed write may leave files; after - /// stopping all users of the target, [`DataFileTarget::cleanup`] can - /// remove them without a completed part description. Retries must use fresh, - /// disjoint Blob ID ranges, including ranges from failed writes. Staging - /// `.part` files are only explicitly cleaned; ordinary dataset GC rules still - /// apply to uncommitted Blob sidecars and must be coordinated with checkpoints. + /// stopping all users of the target, [`DataFileTarget::cleanup`] removes + /// staging files and file-relative sidecars without a completed part description. + /// Retries must use fresh, disjoint Blob ID ranges, including ranges from + /// failed writes. Independent Managed objects follow ordinary dataset GC + /// rules, which must be coordinated with uncommitted writes and checkpoints. /// /// # Example /// @@ -198,6 +198,19 @@ impl Dataset { } else { None }; + if let Some(writer) = preprocessor.take() { + let base = if let Some(id) = target.base_id { + self.manifest.base_paths.get(&id).cloned().ok_or_else(|| { + Error::invalid_input(format!("Managed part target has unknown base ID {id}")) + })? + } else { + self.managed_default_base()? + }; + preprocessor = Some( + writer + .with_managed_base(base.id, base.extract_path(self.session.store_registry())?), + ); + } let file_name = format!("{}.part", generate_random_filename()); let path = target @@ -946,8 +959,32 @@ where .unwrap_or_else(|| params.store_registry()); let source_store_params = params.store_params.clone().unwrap_or_default(); - // Keep a copy so failure paths can clean up files written to target bases. - let cleanup_bases = target_bases_info.clone(); + let default_blob_base = if schema.fields_pre_order().any(|field| field.is_blob_v2()) { + Some(if let Some(dataset) = dataset { + dataset.managed_default_base()?.id + } else { + lance_table::format::BasePath::unused_id( + params.initial_bases.iter().flatten().map(|base| base.id), + )? + }) + } else { + None + }; + // Keep all physical write destinations, including the explicit primary + // alias used by Managed descriptors, available to failed-write cleanup. + let mut cleanup_bases = target_bases_info.clone().unwrap_or_default(); + if let Some(base_id) = default_blob_base { + cleanup_bases.push(TargetBaseInfo { + base_id, + object_store: object_store.clone(), + base_dir: base_dir.clone(), + is_dataset_root: true, + }); + } + let open_writer = move |object_store, schema, base_dir, mut options: WriterOptions| { + options.base_id = options.base_id.or(default_blob_base); + open_writer(object_store, schema, base_dir, options) + }; let file_writer_options = params.file_writer_options.clone().unwrap_or_default(); let writer_generator = WriterGenerator::new( object_store.clone(), @@ -1148,13 +1185,7 @@ where // Drop the writer so its in-progress file is cleaned up (LocalWriter // removes its temp file; ObjectWriter aborts the multipart upload). drop(writer.take()); - cleanup_data_fragments( - &object_store, - base_dir, - cleanup_bases.as_deref(), - &fragments, - ) - .await; + cleanup_data_fragments(&object_store, base_dir, Some(&cleanup_bases), &fragments).await; return Err(e); } @@ -1162,13 +1193,7 @@ where if let Some(mut writer) = writer.take() { if let Err(e) = flush_seed_writers(writer.as_mut(), &mut seed_writers).await { drop(writer); - cleanup_data_fragments( - &object_store, - base_dir, - cleanup_bases.as_deref(), - &fragments, - ) - .await; + cleanup_data_fragments(&object_store, base_dir, Some(&cleanup_bases), &fragments).await; return Err(e); } match writer.finish().await { @@ -1190,13 +1215,8 @@ where } Err(e) => { drop(writer); - cleanup_data_fragments( - &object_store, - base_dir, - cleanup_bases.as_deref(), - &fragments, - ) - .await; + cleanup_data_fragments(&object_store, base_dir, Some(&cleanup_bases), &fragments) + .await; return Err(e); } } @@ -2059,7 +2079,7 @@ impl GenericWriter for V2WriterAdapter { #[derive(Default)] pub(crate) struct WriterOptions { add_data_dir: bool, - base_id: Option, + pub(super) base_id: Option, external_base_resolver: Option>, allow_external_blob_outside_bases: bool, external_blob_mode: ExternalBlobMode, @@ -2151,6 +2171,7 @@ where } pub(in crate::dataset) async fn open_current_blob_v2_writer( + version: ConcreteFileVersion, create_file_writer: F, object_store: &ObjectStore, schema: &Schema, @@ -2187,7 +2208,7 @@ where base_id, file_writer_options, )?; - let preprocessor = BlobPreprocessor::new( + let mut preprocessor = BlobPreprocessor::new( object_store.clone(), data_dir, data_file_key, @@ -2199,6 +2220,14 @@ where source_store_params, blob_pack_file_size_threshold, )?; + if matches!( + version, + ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 + ) { + let base_id = base_id + .ok_or_else(|| Error::invalid_input("Managed writer requires an explicit base ID"))?; + preprocessor = preprocessor.with_managed_base(base_id, base_dir.clone()); + } Ok(Box::new(V2WriterAdapter { writer: file_writer, data_file: Some(data_file), diff --git a/rust/lance/src/dataset/write/commit.rs b/rust/lance/src/dataset/write/commit.rs index be0f4bdd5be..576aff89f19 100644 --- a/rust/lance/src/dataset/write/commit.rs +++ b/rust/lance/src/dataset/write/commit.rs @@ -468,13 +468,14 @@ impl<'a> CommitBuilder<'a> { commit_new_dataset( object_store.as_ref(), source_store.as_deref(), - commit_handler.as_ref(), + &commit_handler, &base_path, + &dest.uri(), &transaction, &manifest_config, manifest_naming_scheme, metadata_cache.as_ref(), - session.store_registry(), + session.clone(), ) .await? }; diff --git a/rust/lance/src/dataset/write/merge_insert.rs b/rust/lance/src/dataset/write/merge_insert.rs index 3b163feef44..a85c08c096a 100644 --- a/rust/lance/src/dataset/write/merge_insert.rs +++ b/rust/lance/src/dataset/write/merge_insert.rs @@ -1745,6 +1745,17 @@ impl MergeInsertJob { &write_schema, &dataset.base, super::WriterOptions { + base_id: if matches!( + write_version, + ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 + ) && write_schema + .fields_pre_order() + .any(|field| field.is_blob_v2()) + { + Some(dataset.managed_default_base()?.id) + } else { + None + }, add_data_dir: true, ..Default::default() }, diff --git a/rust/lance/src/dataset/write/merge_insert/exec/write.rs b/rust/lance/src/dataset/write/merge_insert/exec/write.rs index 86efee59943..6cff2af2b52 100644 --- a/rust/lance/src/dataset/write/merge_insert/exec/write.rs +++ b/rust/lance/src/dataset/write/merge_insert/exec/write.rs @@ -901,6 +901,7 @@ impl ExecutionPlan for FullSchemaMergeInsertExec { self.dataset.schema(), input_schema.as_ref(), true, + false, ) .map_err(|error| DataFusionError::External(Box::new(error)))?, ); diff --git a/rust/lance/src/dataset/write/update.rs b/rust/lance/src/dataset/write/update.rs index ee08ef22b74..0aaecba89c1 100644 --- a/rust/lance/src/dataset/write/update.rs +++ b/rust/lance/src/dataset/write/update.rs @@ -364,6 +364,7 @@ impl UpdateJob { self.dataset.schema(), scan_schema.as_ref(), false, + false, )?); let output_schema = rewrite_plan.output_schema().clone(); let dataset = self.dataset.clone(); diff --git a/rust/lance/src/io/commit.rs b/rust/lance/src/io/commit.rs index 64d9a73f044..1208dc76bbc 100644 --- a/rust/lance/src/io/commit.rs +++ b/rust/lance/src/io/commit.rs @@ -66,7 +66,6 @@ use futures::future::Either; use futures::{FutureExt, StreamExt, TryFutureExt, TryStreamExt}; use lance_core::{Error, Result}; use lance_index::is_system_index; -use lance_io::object_store::ObjectStoreRegistry; use log; use object_store::ObjectStoreExt; use object_store::path::Path; @@ -386,13 +385,14 @@ pub(crate) const MAX_INLINE_TRANSACTION_BYTES: usize = 64 * 1024; async fn do_commit_new_dataset( object_store: &ObjectStore, source_store: Option<&ObjectStore>, - commit_handler: &dyn CommitHandler, + commit_handler: &Arc, base_path: &Path, + uri: &str, transaction: &Transaction, write_config: &ManifestWriteConfig, manifest_naming_scheme: ManifestNamingScheme, metadata_cache: &DSMetadataCache, - store_registry: Arc, + session: Arc, ) -> Result<(Manifest, ManifestLocation)> { let pb_transaction = pb::Transaction::from(transaction); let inline_transaction = pb_transaction.encoded_len() <= MAX_INLINE_TRANSACTION_BYTES; @@ -411,7 +411,7 @@ async fn do_commit_new_dataset( // back to the destination store for same-store clones. let source_store = source_store.unwrap_or(object_store); let source_base_path = - ObjectStore::extract_path_from_uri(store_registry, ref_path.as_str())?; + ObjectStore::extract_path_from_uri(session.store_registry(), ref_path.as_str())?; let source_manifest_location = commit_handler .resolve_version_location(&source_base_path, *ref_version, &source_store.inner) .await?; @@ -419,7 +419,7 @@ async fn do_commit_new_dataset( source_store, &source_manifest_location, ref_path.as_str(), - &Session::default(), + &session, ) .await?; ensure_can_write_manifest(&source_manifest)?; @@ -446,12 +446,9 @@ async fn do_commit_new_dataset( ) = (&transaction.operation, clone_source) { if *is_shallow { - let new_base_id = source_manifest - .base_paths - .keys() - .max() - .map(|id| *id + 1) - .unwrap_or(0); + let new_base_id = lance_table::format::BasePath::unused_id( + source_manifest.base_paths.keys().copied(), + )?; let new_manifest = source_manifest.shallow_clone( ref_name.clone(), ref_path.clone(), @@ -480,7 +477,9 @@ async fn do_commit_new_dataset( } else { // Deep clone: build a manifest that references local files (no external bases) let mut new_manifest = source_manifest.clone(); - new_manifest.base_paths.clear(); + if !source_manifest.has_managed_blobs() { + new_manifest.base_paths.clear(); + } new_manifest.branch = None; new_manifest.tag = None; new_manifest.index_section = None; // will be rewritten below @@ -497,6 +496,30 @@ async fn do_commit_new_dataset( } new_manifest.fragments = Arc::new(new_frags); + if source_manifest.has_managed_blobs() { + let source = Dataset::checkout_manifest( + Arc::new(source_store.clone()), + ObjectStore::extract_path_from_uri(session.store_registry(), ref_path)?, + ref_path.clone(), + Arc::new(source_manifest.clone()), + source_manifest_location.clone(), + session.clone(), + commit_handler.clone(), + None, + None, + None, + )?; + crate::dataset::blob::clone::copy_blob_columns( + Arc::new(source), + Arc::new(object_store.clone()), + base_path.clone(), + uri, + &mut new_manifest, + ) + .boxed() + .await?; + } + // Indices: keep metadata but normalize base to local let mut updated_indices = Vec::new(); if let Some(index_section_pos) = source_manifest.index_section { @@ -525,9 +548,25 @@ async fn do_commit_new_dataset( (manifest, indices) }; + if manifest.has_managed_blobs() { + let id = lance_table::format::BasePath::unused_id(manifest.base_paths.keys().copied())?; + if manifest + .fragments + .iter() + .flat_map(|fragment| fragment.referenced_lance_files()) + .any(|file| file.base_id == Some(id)) + { + manifest.bind_managed_base(lance_table::format::BasePath::new( + id, + uri.to_string(), + None, + true, + ))?; + } + } let result = write_manifest_file( object_store, - commit_handler, + commit_handler.as_ref(), base_path, &mut manifest, if indices.is_empty() { @@ -557,7 +596,7 @@ async fn do_commit_new_dataset( // transaction file a landed manifest would reference). match verify_commit_outcome( object_store, - commit_handler, + commit_handler.as_ref(), base_path, manifest.version, transaction, @@ -595,7 +634,7 @@ async fn do_commit_new_dataset( Err(CommitError::OtherError(err)) => { match verify_commit_outcome( object_store, - commit_handler, + commit_handler.as_ref(), base_path, manifest.version, transaction, @@ -661,24 +700,26 @@ async fn record_new_dataset_commit( pub(crate) async fn commit_new_dataset( object_store: &ObjectStore, source_store: Option<&ObjectStore>, - commit_handler: &dyn CommitHandler, + commit_handler: &Arc, base_path: &Path, + uri: &str, transaction: &Transaction, write_config: &ManifestWriteConfig, manifest_naming_scheme: ManifestNamingScheme, metadata_cache: &crate::session::caches::DSMetadataCache, - store_registry: Arc, + session: Arc, ) -> Result<(Manifest, ManifestLocation)> { do_commit_new_dataset( object_store, source_store, commit_handler, base_path, + uri, transaction, write_config, manifest_naming_scheme, metadata_cache, - store_registry, + session, ) .await } @@ -1206,6 +1247,7 @@ pub(crate) async fn do_commit_detached_transaction( }; manifest.version = random_version; + manifest.bind_managed_base(dataset.managed_default_base()?)?; // recompute_stats is always false so far because detached manifests are newer than // the old stats bug. @@ -1465,6 +1507,7 @@ pub(crate) async fn commit_transaction( // The Arc is kept rather than cloned out: `load_all_indices` returns shared // cached data, so the common case is a cache hit rather than a read. let read_version_dataset = dataset.clone(); + let managed_default_base = read_version_dataset.managed_default_base()?; let read_version_indices = load_all_indices(&read_version_dataset).await?; let read_version_state = Some(crate::dataset::transaction::ReadVersionState { manifest: read_version_dataset.manifest.as_ref(), @@ -1560,6 +1603,8 @@ pub(crate) async fn commit_transaction( manifest.version = target_version; + manifest.bind_managed_base(managed_default_base.clone())?; + let previous_writer_version = &dataset.manifest.writer_version; // The versions of Lance prior to when we started writing the writer version // sometimes wrote incorrect `Fragment.physical_rows` values, so we should diff --git a/test_data/v11.0.0/blob_sidecars/_transactions/0-d1be0750-43f8-4399-a587-30f496ebf963.txn b/test_data/v11.0.0/blob_sidecars/_transactions/0-d1be0750-43f8-4399-a587-30f496ebf963.txn new file mode 100644 index 00000000000..1ff4ca23c37 Binary files /dev/null and b/test_data/v11.0.0/blob_sidecars/_transactions/0-d1be0750-43f8-4399-a587-30f496ebf963.txn differ diff --git a/test_data/v11.0.0/blob_sidecars/_versions/18446744073709551614.manifest b/test_data/v11.0.0/blob_sidecars/_versions/18446744073709551614.manifest new file mode 100644 index 00000000000..073c62b30d9 Binary files /dev/null and b/test_data/v11.0.0/blob_sidecars/_versions/18446744073709551614.manifest differ diff --git a/test_data/v11.0.0/blob_sidecars/_versions/latest_version_hint.json b/test_data/v11.0.0/blob_sidecars/_versions/latest_version_hint.json new file mode 100644 index 00000000000..491d734467a --- /dev/null +++ b/test_data/v11.0.0/blob_sidecars/_versions/latest_version_hint.json @@ -0,0 +1 @@ +{"version":1} \ No newline at end of file diff --git a/test_data/v11.0.0/blob_sidecars/data/101001100001101001100110024e524e87a2c996cccd0d0f1f.lance b/test_data/v11.0.0/blob_sidecars/data/101001100001101001100110024e524e87a2c996cccd0d0f1f.lance new file mode 100644 index 00000000000..73b8bd155d6 Binary files /dev/null and b/test_data/v11.0.0/blob_sidecars/data/101001100001101001100110024e524e87a2c996cccd0d0f1f.lance differ diff --git a/test_data/v11.0.0/blob_sidecars/data/101001100001101001100110024e524e87a2c996cccd0d0f1f/10000000000000000000000000000000.blob b/test_data/v11.0.0/blob_sidecars/data/101001100001101001100110024e524e87a2c996cccd0d0f1f/10000000000000000000000000000000.blob new file mode 100644 index 00000000000..00ea360d7cd --- /dev/null +++ b/test_data/v11.0.0/blob_sidecars/data/101001100001101001100110024e524e87a2c996cccd0d0f1f/10000000000000000000000000000000.blob @@ -0,0 +1 @@ +dddddddddddddddddddddddddddddddd \ No newline at end of file diff --git a/test_data/v11.0.0/blob_sidecars/data/101101001100001010100011cadfab455c80f603c31f0f8324.lance b/test_data/v11.0.0/blob_sidecars/data/101101001100001010100011cadfab455c80f603c31f0f8324.lance new file mode 100644 index 00000000000..9fe601b78de Binary files /dev/null and b/test_data/v11.0.0/blob_sidecars/data/101101001100001010100011cadfab455c80f603c31f0f8324.lance differ diff --git a/test_data/v11.0.0/blob_sidecars/data/101101001100001010100011cadfab455c80f603c31f0f8324/10000000000000000000000000000000.blob b/test_data/v11.0.0/blob_sidecars/data/101101001100001010100011cadfab455c80f603c31f0f8324/10000000000000000000000000000000.blob new file mode 100644 index 00000000000..357ab15b826 --- /dev/null +++ b/test_data/v11.0.0/blob_sidecars/data/101101001100001010100011cadfab455c80f603c31f0f8324/10000000000000000000000000000000.blob @@ -0,0 +1 @@ +pppppppppppppppp \ No newline at end of file diff --git a/test_data/v11.0.0/datagen.py b/test_data/v11.0.0/datagen.py new file mode 100644 index 00000000000..611087c8f5f --- /dev/null +++ b/test_data/v11.0.0/datagen.py @@ -0,0 +1,36 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright The Lance Authors + +"""Generate the released Packed/Dedicated fixture used by Managed adoption tests.""" + +import shutil +from pathlib import Path + +import lance +import pyarrow as pa + +assert lance.__version__ == "11.0.0" + +field = lance.blob_field("blob").with_metadata( + { + "lance-encoding:blob-inline-size-threshold": "8", + "lance-encoding:blob-dedicated-size-threshold": "24", + } +) +table = pa.Table.from_arrays( + [lance.blob_array([b"p" * 16, b"inline", None, b"", b"d" * 32, b"inline"])], + schema=pa.schema([field]), +) +dataset_path = Path(__file__).parent / "blob_sidecars" +shutil.rmtree(dataset_path, ignore_errors=True) +dataset = lance.write_dataset( + table, + dataset_path, + data_storage_version="2.2", + max_rows_per_file=3, + max_rows_per_group=3, +) +assert [ + row["blob"]["kind"] if row["blob"] is not None else None + for row in dataset.to_table().to_pylist() +] == [1, 0, None, 0, 2, 0]