Skip to content

Commit b8c333b

Browse files
committed
fix: enforce manifest capability checks before side effects
1 parent 1cc5983 commit b8c333b

3 files changed

Lines changed: 164 additions & 40 deletions

File tree

rust/lance-namespace-impls/src/dir/manifest.rs

Lines changed: 63 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1933,6 +1933,7 @@ impl ManifestNamespace {
19331933
/// concurrent upgrade in between is still caught.
19341934
async fn ensure_manifest_writable(&self) -> Result<()> {
19351935
let dataset_guard = self.manifest_dataset.get().await?;
1936+
ensure_can_write_manifest(dataset_guard.manifest())?;
19361937
ensure_writable(dataset_guard.metadata())
19371938
}
19381939

@@ -1953,10 +1954,11 @@ impl ManifestNamespace {
19531954

19541955
loop {
19551956
let dataset_guard = self.manifest_dataset.get_refreshed().await?;
1957+
ensure_can_write_manifest(dataset_guard.manifest())?;
19561958
let dataset = Arc::new(dataset_guard.clone());
19571959
drop(dataset_guard);
1958-
// Refuse to mutate a manifest written with a writer feature flag this
1959-
// build does not understand.
1960+
// The namespace format has its own capabilities in table metadata,
1961+
// separate from the Lance manifest capabilities checked above.
19601962
ensure_writable(dataset.metadata())?;
19611963
// Staged files, indices, the commit, and cleanup must all use the dataset's
19621964
// own object store (see `commit_manifest_overwrite`).
@@ -3859,6 +3861,7 @@ mod tests {
38593861
CreateNamespaceRequest, CreateTableRequest, DescribeTableRequest, DropTableRequest,
38603862
ListTablesRequest, TableExistsRequest,
38613863
};
3864+
use lance_table::feature_flags::FLAG_UNKNOWN;
38623865
use lance_table::format::Fragment;
38633866
use rstest::rstest;
38643867
use std::collections::{HashMap, HashSet};
@@ -4391,6 +4394,64 @@ mod tests {
43914394
);
43924395
}
43934396

4397+
#[tokio::test]
4398+
async fn test_manifest_rewrite_rejects_unknown_writer_flag_before_staging() {
4399+
let temp_dir = TempStdDir::default();
4400+
let temp_path = temp_dir.to_str().unwrap();
4401+
let manifest_ns = create_manifest_namespace(temp_path, false).await;
4402+
let data_paths_before = manifest_data_paths(&manifest_ns).await;
4403+
let original_version = {
4404+
let mut dataset = manifest_ns.manifest_dataset.get_mut().await.unwrap();
4405+
let mut manifest = dataset.manifest().clone();
4406+
manifest.writer_feature_flags |= FLAG_UNKNOWN << 1;
4407+
let version = manifest.version;
4408+
dataset.manifest = Arc::new(manifest);
4409+
version
4410+
};
4411+
4412+
let entries_before = dir_entry_names(temp_path);
4413+
let mut create_request = CreateTableRequest::new();
4414+
create_request.id = Some(vec!["new_table".to_string()]);
4415+
let error = manifest_ns
4416+
.create_table(create_request, Bytes::from(create_test_ipc_data()))
4417+
.await
4418+
.unwrap_err();
4419+
assert!(
4420+
error.to_string().to_lowercase().contains("upgrade"),
4421+
"expected an upgrade error, got: {error}"
4422+
);
4423+
assert_eq!(dir_entry_names(temp_path), entries_before);
4424+
4425+
let error = manifest_ns
4426+
.insert_into_manifest_with_metadata(
4427+
vec![ManifestEntry {
4428+
object_id: "table".to_string(),
4429+
object_type: ObjectType::Table,
4430+
location: Some("table.lance".to_string()),
4431+
metadata: None,
4432+
}],
4433+
None,
4434+
)
4435+
.await
4436+
.unwrap_err();
4437+
4438+
assert!(
4439+
error.to_string().to_lowercase().contains("upgrade"),
4440+
"expected an upgrade error, got: {error}"
4441+
);
4442+
assert_eq!(
4443+
manifest_ns
4444+
.manifest_dataset
4445+
.get()
4446+
.await
4447+
.unwrap()
4448+
.version()
4449+
.version,
4450+
original_version
4451+
);
4452+
assert_eq!(manifest_data_paths(&manifest_ns).await, data_paths_before);
4453+
}
4454+
43944455
#[tokio::test]
43954456
async fn test_manifest_noop_delete_uses_latest_snapshot() {
43964457
let temp_dir = TempStdDir::default();

rust/lance/src/dataset.rs

Lines changed: 9 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -149,7 +149,8 @@ use lance_core::box_error;
149149
use lance_index::scalar::lance_format::LanceIndexStore;
150150
use lance_namespace::models::{DeclareTableRequest, DescribeTableRequest};
151151
use lance_table::feature_flags::{
152-
apply_feature_flags, can_read_dataset, validate_paired_feature_flags,
152+
apply_feature_flags, ensure_can_read_manifest, ensure_can_write_manifest,
153+
validate_paired_feature_flags,
153154
};
154155
use lance_table::io::deletion::{DELETIONS_DIR, relative_deletion_file_path};
155156
use lance_table::rowids::{RowIdSequence, write_row_ids};
@@ -766,16 +767,7 @@ impl Dataset {
766767
read_struct(object_reader.as_ref(), offset).await
767768
}?;
768769

769-
validate_paired_feature_flags(&manifest)?;
770-
771-
if !can_read_dataset(manifest.reader_feature_flags) {
772-
let message = format!(
773-
"This dataset cannot be read by this version of Lance. \
774-
Please upgrade Lance to read this dataset.\n Flags: {}",
775-
manifest.reader_feature_flags
776-
);
777-
return Err(Error::not_supported_source(message.into()));
778-
}
770+
ensure_can_read_manifest(&manifest)?;
779771

780772
// If indices were also in the last block, we can take the opportunity to
781773
// decode them now and cache them.
@@ -848,6 +840,7 @@ impl Dataset {
848840
e_tag: manifest_location.e_tag.as_deref(),
849841
};
850842
if let Some(cached) = metadata_cache.get_with_key(&manifest_key).await {
843+
ensure_can_read_manifest(&cached)?;
851844
return Ok(cached);
852845
}
853846
let loaded =
@@ -1207,35 +1200,13 @@ impl Dataset {
12071200
.resolve_latest_location(&self.base, &self.object_store)
12081201
.await?;
12091202

1210-
// Check if manifest is in cache before reading from storage
1211-
let manifest_key = ManifestKey {
1212-
version: location.version,
1213-
e_tag: location.e_tag.as_deref(),
1214-
};
1215-
let cached_manifest = self.metadata_cache.get_with_key(&manifest_key).await;
1216-
if let Some(cached_manifest) = cached_manifest {
1217-
return Ok((cached_manifest, location));
1218-
}
1219-
12201203
if self.already_checked_out(&location, self.manifest.branch.as_deref()) {
1204+
ensure_can_read_manifest(&self.manifest)?;
12211205
return Ok((self.manifest.clone(), self.manifest_location.clone()));
12221206
}
1223-
let mut manifest = read_manifest(&self.object_store, &location.path, location.size).await?;
1224-
if manifest.schema.has_dictionary_types() {
1225-
let reader = if let Some(size) = location.size {
1226-
self.object_store
1227-
.open_with_size(&location.path, size as usize)
1228-
.await?
1229-
} else {
1230-
self.object_store.open(&location.path).await?
1231-
};
1232-
populate_manifest_schema_dictionaries(&mut manifest, reader.as_ref()).await?;
1233-
}
1234-
let manifest_arc = Arc::new(manifest);
1235-
self.metadata_cache
1236-
.insert_with_key(&manifest_key, manifest_arc.clone())
1237-
.await;
1238-
Ok((manifest_arc, location))
1207+
let manifest =
1208+
Self::get_manifest(&self.object_store, &location, &self.uri, &self.session).await?;
1209+
Ok((manifest, location))
12391210
}
12401211

12411212
/// Read the transaction file for this version of the dataset.
@@ -3307,6 +3278,7 @@ impl Dataset {
33073278

33083279
// Resolve source dataset and its manifest using checkout_version
33093280
let src_ds = self.checkout_version(version).await?;
3281+
ensure_can_write_manifest(&src_ds.manifest)?;
33103282
let src_paths = src_ds.collect_paths().await?;
33113283

33123284
// Prepare target object store and base path

rust/lance/src/dataset/tests/dataset_io.rs

Lines changed: 92 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,7 @@ use lance_file::{
4444
};
4545
use lance_io::assert_io_eq;
4646
use lance_table::feature_flags;
47-
use lance_table::format::BasePath;
47+
use lance_table::format::{BasePath, Fragment};
4848
use object_store::ObjectStoreExt;
4949

5050
use crate::index::DatasetIndexExt;
@@ -1412,6 +1412,51 @@ async fn test_restore_rejects_unknown_target_flags() {
14121412
assert!(matches!(error, Error::NotSupported { .. }), "{error}");
14131413
}
14141414

1415+
#[tokio::test]
1416+
async fn test_checkout_latest_rejects_unsupported_reader_before_caching() {
1417+
let test_uri = TempStrDir::default();
1418+
let data = gen_batch()
1419+
.col("i", array::step::<Int32Type>())
1420+
.into_reader_rows(RowCount::from(1), BatchCount::from(1));
1421+
let mut dataset = Dataset::write(data, &test_uri, None).await.unwrap();
1422+
let original_version = dataset.version().version;
1423+
1424+
let mut unsupported_manifest = dataset.manifest.as_ref().clone();
1425+
unsupported_manifest.version += 1;
1426+
unsupported_manifest.reader_feature_flags |= feature_flags::FLAG_UNKNOWN;
1427+
unsupported_manifest.writer_feature_flags |= feature_flags::FLAG_UNKNOWN;
1428+
let location = write_manifest_file(
1429+
dataset.object_store.as_ref(),
1430+
dataset.commit_handler.as_ref(),
1431+
&dataset.base,
1432+
&mut unsupported_manifest,
1433+
None,
1434+
&ManifestWriteConfig {
1435+
auto_set_feature_flags: false,
1436+
..Default::default()
1437+
},
1438+
dataset.manifest_location.naming_scheme,
1439+
None,
1440+
)
1441+
.await
1442+
.unwrap();
1443+
1444+
let error = dataset.checkout_latest().await.unwrap_err();
1445+
assert!(matches!(error, Error::NotSupported { .. }), "{error}");
1446+
assert_eq!(dataset.version().version, original_version);
1447+
assert!(
1448+
dataset
1449+
.metadata_cache
1450+
.get_with_key(&ManifestKey {
1451+
version: location.version,
1452+
e_tag: location.e_tag.as_deref(),
1453+
})
1454+
.await
1455+
.is_none(),
1456+
"unsupported manifest must not be cached"
1457+
);
1458+
}
1459+
14151460
#[tokio::test]
14161461
async fn test_rle_v2_v23_write_and_append() {
14171462
let test_uri = TempStrDir::default();
@@ -1792,6 +1837,52 @@ async fn test_deep_clone(
17921837
assert_eq!(count_files(store, &dst_root, "_deletions").await, 0);
17931838
}
17941839

1840+
#[tokio::test]
1841+
async fn test_deep_clone_rejects_unsupported_writer_before_copying() {
1842+
let test_dir = TempStdDir::default();
1843+
let source_dir = test_dir.join("source");
1844+
let target_dir = test_dir.join("target");
1845+
let mut source = Dataset::write(
1846+
gen_batch()
1847+
.col("id", array::step::<Int32Type>())
1848+
.into_reader_rows(RowCount::from(32), BatchCount::from(1)),
1849+
source_dir.to_str().unwrap(),
1850+
None,
1851+
)
1852+
.await
1853+
.unwrap();
1854+
1855+
let mut unsupported_manifest = source.manifest.as_ref().clone();
1856+
unsupported_manifest.version += 1;
1857+
unsupported_manifest.writer_feature_flags |= feature_flags::FLAG_UNKNOWN << 1;
1858+
write_manifest_file(
1859+
source.object_store.as_ref(),
1860+
source.commit_handler.as_ref(),
1861+
&source.base,
1862+
&mut unsupported_manifest,
1863+
None,
1864+
&ManifestWriteConfig {
1865+
auto_set_feature_flags: false,
1866+
..Default::default()
1867+
},
1868+
source.manifest_location.naming_scheme,
1869+
None,
1870+
)
1871+
.await
1872+
.unwrap();
1873+
1874+
let error = source
1875+
.deep_clone(
1876+
target_dir.to_str().unwrap(),
1877+
unsupported_manifest.version,
1878+
None,
1879+
)
1880+
.await
1881+
.unwrap_err();
1882+
assert!(matches!(error, Error::NotSupported { .. }));
1883+
assert!(!target_dir.exists());
1884+
}
1885+
17951886
#[tokio::test]
17961887
async fn test_deep_clone_recognizes_ambiguous_commit_as_own() {
17971888
use crate::utils::test::{AmbiguousCommitHandler, AmbiguousFailure};

0 commit comments

Comments
 (0)