fix(table-catalog): accept inherited snapshot manifests (#4833)

Co-authored-by: Henry Guo <marshawcoco@users.noreply.github.com>
This commit is contained in:
Henry Guo
2026-07-15 17:33:42 +08:00
committed by GitHub
parent aa11ee4341
commit 1e7fb8438a
+580 -54
View File
@@ -2794,13 +2794,26 @@ fn snapshot_has_manifest_references(snapshot: &serde_json::Value) -> bool {
#[derive(Default)]
struct SnapshotLiveFiles {
data_files: BTreeSet<String>,
delete_files: BTreeSet<String>,
data_files: BTreeMap<String, SnapshotFileIdentity>,
delete_files: BTreeMap<String, SnapshotFileIdentity>,
manifest_files: BTreeMap<String, SnapshotManifestIdentity>,
}
impl SnapshotLiveFiles {
fn contains(&self, location: &str) -> bool {
self.data_files.contains(location) || self.delete_files.contains(location)
self.data_files.contains_key(location) || self.delete_files.contains_key(location)
}
fn identity(
&self,
location: &str,
object_kind: &crate::table_catalog::TableMetadataMaintenanceObjectKind,
) -> Option<&SnapshotFileIdentity> {
match object_kind {
crate::table_catalog::TableMetadataMaintenanceObjectKind::DataFile => self.data_files.get(location),
crate::table_catalog::TableMetadataMaintenanceObjectKind::DeleteFile => self.delete_files.get(location),
_ => None,
}
}
}
@@ -2826,12 +2839,24 @@ impl SnapshotFileChanges {
}
}
#[derive(Clone, Copy)]
struct SnapshotChangeIdentity {
struct SnapshotChangeContext<'a> {
current_live_files: &'a SnapshotLiveFiles,
snapshot_id: i64,
sequence_number: i64,
}
#[derive(PartialEq, Eq)]
struct SnapshotManifestIdentity {
sequence_number: Option<i64>,
added_snapshot_id: Option<i64>,
}
#[derive(PartialEq, Eq)]
struct SnapshotFileIdentity {
sequence_number: Option<i64>,
file_sequence_number: Option<i64>,
}
async fn validate_table_snapshot_commit_conflicts<B>(
metadata_backend: &B,
bucket: &str,
@@ -2870,7 +2895,8 @@ where
table,
entry,
snapshot,
SnapshotChangeIdentity {
SnapshotChangeContext {
current_live_files: &current_live_files,
snapshot_id,
sequence_number,
},
@@ -2968,22 +2994,42 @@ where
.ok_or_else(|| s3_error!(InvalidRequest, "current snapshot metadata is missing"))?;
let mut live_files = SnapshotLiveFiles::default();
for reference in read_snapshot_manifest_references(metadata_backend, bucket, namespace, table, entry, snapshot).await? {
let status = reference
.entry_status
.ok_or_else(|| s3_error!(InvalidRequest, "manifest entry status is required"))?;
match status {
0 | 1 => match reference.object_kind {
crate::table_catalog::TableMetadataMaintenanceObjectKind::DataFile => {
live_files.data_files.insert(reference.location);
}
crate::table_catalog::TableMetadataMaintenanceObjectKind::DeleteFile => {
live_files.delete_files.insert(reference.location);
}
_ => {}
for manifest in read_snapshot_manifest_references(metadata_backend, bucket, namespace, table, entry, snapshot).await? {
let SnapshotManifestLocation {
manifest_path,
sequence_number,
added_snapshot_id,
} = manifest.location;
live_files.manifest_files.insert(
manifest_path,
SnapshotManifestIdentity {
sequence_number,
added_snapshot_id,
},
2 => {}
_ => return Err(s3_error!(InvalidRequest, "manifest entry status is unsupported")),
);
for reference in manifest.references {
let status = reference
.entry_status
.ok_or_else(|| s3_error!(InvalidRequest, "manifest entry status is required"))?;
match status {
0 | 1 => {
let identity = SnapshotFileIdentity {
sequence_number: reference.sequence_number,
file_sequence_number: reference.file_sequence_number,
};
match reference.object_kind {
crate::table_catalog::TableMetadataMaintenanceObjectKind::DataFile => {
live_files.data_files.insert(reference.location, identity);
}
crate::table_catalog::TableMetadataMaintenanceObjectKind::DeleteFile => {
live_files.delete_files.insert(reference.location, identity);
}
_ => {}
}
}
2 => {}
_ => return Err(s3_error!(InvalidRequest, "manifest entry status is unsupported")),
}
}
}
Ok(live_files)
@@ -2996,46 +3042,109 @@ async fn load_snapshot_file_changes<B>(
table: &crate::table_catalog::IdentifierSegment,
entry: &crate::table_catalog::TableEntry,
snapshot: &serde_json::Value,
change_identity: SnapshotChangeIdentity,
context: SnapshotChangeContext<'_>,
) -> S3Result<SnapshotFileChanges>
where
B: crate::table_catalog::TableCatalogObjectBackend,
{
let mut changes = SnapshotFileChanges::default();
for reference in read_snapshot_manifest_references(metadata_backend, bucket, namespace, table, entry, snapshot).await? {
let status = reference
.entry_status
.ok_or_else(|| s3_error!(InvalidRequest, "manifest entry status is required"))?;
if matches!(status, 1 | 2)
&& (reference.snapshot_id != Some(change_identity.snapshot_id)
|| reference.sequence_number != Some(change_identity.sequence_number))
for manifest in read_snapshot_manifest_references(metadata_backend, bucket, namespace, table, entry, snapshot).await? {
let inherited_identity = context
.current_live_files
.manifest_files
.get(&manifest.location.manifest_path);
if let Some(inherited_identity) = inherited_identity {
let candidate_identity = SnapshotManifestIdentity {
sequence_number: manifest.location.sequence_number,
added_snapshot_id: manifest.location.added_snapshot_id,
};
if inherited_identity
.sequence_number
.is_some_and(|sequence_number| candidate_identity.sequence_number != Some(sequence_number))
|| inherited_identity
.added_snapshot_id
.is_some_and(|snapshot_id| candidate_identity.added_snapshot_id != Some(snapshot_id))
{
return Err(s3_error!(InvalidRequest, "inherited manifest must preserve its manifest-list identity"));
}
continue;
}
if manifest
.location
.added_snapshot_id
.is_some_and(|added_snapshot_id| added_snapshot_id != context.snapshot_id)
{
return Err(s3_error!(
InvalidRequest,
"manifest changed entries must belong to the committed snapshot"
));
return Err(s3_error!(InvalidRequest, "new manifest must belong to the committed snapshot"));
}
if manifest
.location
.sequence_number
.is_some_and(|sequence_number| sequence_number != context.sequence_number)
{
return Err(s3_error!(InvalidRequest, "new manifest sequence must match the committed snapshot"));
}
match (status, reference.object_kind) {
(0, _) => {}
(1, crate::table_catalog::TableMetadataMaintenanceObjectKind::DataFile) => {
changes.added_data_files.insert(reference.location);
for reference in manifest.references {
let status = reference
.entry_status
.ok_or_else(|| s3_error!(InvalidRequest, "manifest entry status is required"))?;
if matches!(status, 1 | 2) && reference.snapshot_id != Some(context.snapshot_id) {
return Err(s3_error!(
InvalidRequest,
"manifest changed entries must belong to the committed snapshot"
));
}
(1, crate::table_catalog::TableMetadataMaintenanceObjectKind::DeleteFile) => {
changes.added_delete_files.insert(reference.location);
if status == 1
&& manifest.location.sequence_number.is_some_and(|sequence_number| {
reference.sequence_number != Some(sequence_number) || reference.file_sequence_number != Some(sequence_number)
})
{
return Err(s3_error!(InvalidRequest, "added manifest entry sequence must match its manifest"));
}
(2, crate::table_catalog::TableMetadataMaintenanceObjectKind::DataFile) => {
changes.deleted_data_files.insert(reference.location);
if status == 2
&& context
.current_live_files
.identity(&reference.location, &reference.object_kind)
.is_some_and(|current_identity| {
current_identity
!= &SnapshotFileIdentity {
sequence_number: reference.sequence_number,
file_sequence_number: reference.file_sequence_number,
}
})
{
return Err(s3_error!(
InvalidRequest,
"deleted manifest entry must preserve the current file sequence"
));
}
(2, crate::table_catalog::TableMetadataMaintenanceObjectKind::DeleteFile) => {
changes.deleted_delete_files.insert(reference.location);
match (status, reference.object_kind) {
(0, _) => {}
(1, crate::table_catalog::TableMetadataMaintenanceObjectKind::DataFile) => {
changes.added_data_files.insert(reference.location);
}
(1, crate::table_catalog::TableMetadataMaintenanceObjectKind::DeleteFile) => {
changes.added_delete_files.insert(reference.location);
}
(2, crate::table_catalog::TableMetadataMaintenanceObjectKind::DataFile) => {
changes.deleted_data_files.insert(reference.location);
}
(2, crate::table_catalog::TableMetadataMaintenanceObjectKind::DeleteFile) => {
changes.deleted_delete_files.insert(reference.location);
}
_ => return Err(s3_error!(InvalidRequest, "manifest entry status is unsupported")),
}
_ => return Err(s3_error!(InvalidRequest, "manifest entry status is unsupported")),
}
}
Ok(changes)
}
struct SnapshotManifestReferences {
location: SnapshotManifestLocation,
references: Vec<crate::table_catalog::ManifestDataFileReference>,
}
async fn read_snapshot_manifest_references<B>(
metadata_backend: &B,
bucket: &str,
@@ -3043,12 +3152,12 @@ async fn read_snapshot_manifest_references<B>(
table: &crate::table_catalog::IdentifierSegment,
entry: &crate::table_catalog::TableEntry,
snapshot: &serde_json::Value,
) -> S3Result<Vec<crate::table_catalog::ManifestDataFileReference>>
) -> S3Result<Vec<SnapshotManifestReferences>>
where
B: crate::table_catalog::TableCatalogObjectBackend,
{
let manifest_locations = snapshot_manifest_locations(metadata_backend, bucket, namespace, table, entry, snapshot).await?;
let mut references = Vec::new();
let mut manifests = Vec::new();
for manifest_location in manifest_locations {
let manifest_key = table_commit_object_key(
bucket,
@@ -3065,6 +3174,7 @@ where
.ok_or_else(|| s3_error!(InvalidRequest, "snapshot manifest object is missing"))?;
let file_references =
crate::table_catalog::data_file_references_from_manifest_avro(&manifest_object.data).map_err(catalog_store_error)?;
let mut references = Vec::with_capacity(file_references.len());
for mut reference in file_references {
if reference.snapshot_id.is_none() {
reference.snapshot_id = manifest_location.added_snapshot_id;
@@ -3078,8 +3188,12 @@ where
validate_manifest_data_file_reference(metadata_backend, bucket, namespace, table, entry, &reference).await?;
references.push(reference);
}
manifests.push(SnapshotManifestReferences {
location: manifest_location,
references,
});
}
Ok(references)
Ok(manifests)
}
#[derive(Debug)]
@@ -8206,7 +8320,7 @@ mod tests {
&[(&old_data_file, 0, 2, 11, 2), (&replacement_data_file, 0, 1, 11, 2)],
)
.await;
let overwrite_request: RestCommitTableRequest = serde_json::from_value(serde_json::json!({
let overwrite_request_json = serde_json::json!({
"requirements": [
{
"type": "assert-current-snapshot-id",
@@ -8234,8 +8348,26 @@ mod tests {
"type": "branch"
}
]
}))
.expect("overwrite request should parse");
});
let stale_overwrite_request: RestCommitTableRequest =
serde_json::from_value(overwrite_request_json.clone()).expect("stale overwrite request should parse");
let error = commit_table_response(&store, &metadata_backend, "warehouse", &namespace, "events", stale_overwrite_request)
.await
.expect_err("deleted file sequence must not change");
assert_eq!(error.code(), &s3s::S3ErrorCode::InvalidRequest);
seed_test_snapshot_manifest(
&metadata_backend,
"warehouse",
&overwrite_manifest_list,
11,
2,
&[(&old_data_file, 0, 2, 11, 1), (&replacement_data_file, 0, 1, 11, 2)],
)
.await;
let overwrite_request: RestCommitTableRequest =
serde_json::from_value(overwrite_request_json).expect("overwrite request should parse");
let commit = commit_table_response(&store, &metadata_backend, "warehouse", &namespace, "events", overwrite_request)
.await
@@ -8266,7 +8398,7 @@ mod tests {
"sequence-number": 1,
"timestamp-ms": 1234,
"manifests": [
manifest
manifest.clone()
],
"summary": {
"operation": "append"
@@ -8289,6 +8421,56 @@ mod tests {
assert_eq!(commit.metadata["current-snapshot-id"], 10);
assert_eq!(commit.metadata["last-sequence-number"], 1);
let second_manifest_list = format!("{table_location}/metadata/snap-11.avro");
let second_manifest = format!("{table_location}/metadata/manifest-11.avro");
let second_data_file = format!("{table_location}/data/part-11.parquet");
let second_manifest_list_key = test_snapshot_object_key("warehouse", &second_manifest_list);
metadata_backend
.put_bytes(
"warehouse",
&second_manifest_list_key,
test_manifest_list_avro_entries(&[(&manifest, 1, 10), (&second_manifest, 2, 11)]),
)
.await;
seed_test_manifest(&metadata_backend, "warehouse", &second_manifest, &[(&second_data_file, 0, 1, 11, 2)]).await;
let second_append: RestCommitTableRequest = serde_json::from_value(serde_json::json!({
"requirements": [
{
"type": "assert-current-snapshot-id",
"snapshot-id": 10
}
],
"updates": [
{
"action": "add-snapshot",
"snapshot": {
"snapshot-id": 11,
"parent-snapshot-id": 10,
"sequence-number": 2,
"timestamp-ms": 2234,
"manifest-list": second_manifest_list,
"summary": {
"operation": "append"
}
}
},
{
"action": "set-snapshot-ref",
"ref-name": "main",
"snapshot-id": 11,
"type": "branch"
}
]
}))
.expect("second append request should parse");
let upgraded = commit_table_response(&store, &metadata_backend, "warehouse", &namespace, "events", second_append)
.await
.expect("manifest-list snapshot should inherit a legacy manifest with unknown provenance");
assert_eq!(upgraded.metadata["current-snapshot-id"], 11);
assert_eq!(upgraded.metadata["last-sequence-number"], 2);
}
#[tokio::test]
@@ -8338,6 +8520,342 @@ mod tests {
assert_eq!(commit.metadata["last-sequence-number"], 1);
}
#[tokio::test]
async fn row_level_conflict_allows_inherited_manifests_on_append() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
.as_str()
.expect("created metadata should have table location");
let first_manifest_list = format!("{table_location}/metadata/snap-10.avro");
let first_manifest = format!("{table_location}/metadata/manifest-10.avro");
let first_data_file = format!("{table_location}/data/part-10.parquet");
seed_test_manifest_list(&metadata_backend, "warehouse", &first_manifest_list, &[&first_manifest], 1, 10).await;
seed_test_manifest(&metadata_backend, "warehouse", &first_manifest, &[(&first_data_file, 0, 1, 10, 1)]).await;
let first_append: RestCommitTableRequest = serde_json::from_value(serde_json::json!({
"updates": [
{
"action": "add-snapshot",
"snapshot": {
"snapshot-id": 10,
"sequence-number": 1,
"timestamp-ms": 1234,
"manifest-list": first_manifest_list,
"summary": {
"operation": "append"
}
}
},
{
"action": "set-snapshot-ref",
"ref-name": "main",
"snapshot-id": 10,
"type": "branch"
}
]
}))
.expect("first append request should parse");
commit_table_response(&store, &metadata_backend, "warehouse", &namespace, "events", first_append)
.await
.expect("first append should commit");
let second_manifest_list = format!("{table_location}/metadata/snap-11.avro");
let second_manifest = format!("{table_location}/metadata/manifest-11.avro");
let second_data_file = format!("{table_location}/data/part-11.parquet");
let second_manifest_list_key = test_snapshot_object_key("warehouse", &second_manifest_list);
metadata_backend
.put_bytes(
"warehouse",
&second_manifest_list_key,
test_manifest_list_avro_entries(&[(&first_manifest, 1, 10), (&second_manifest, 2, 11)]),
)
.await;
seed_test_manifest(&metadata_backend, "warehouse", &second_manifest, &[(&second_data_file, 0, 1, 11, 2)]).await;
let second_append: RestCommitTableRequest = serde_json::from_value(serde_json::json!({
"requirements": [
{
"type": "assert-current-snapshot-id",
"snapshot-id": 10
}
],
"updates": [
{
"action": "add-snapshot",
"snapshot": {
"snapshot-id": 11,
"parent-snapshot-id": 10,
"sequence-number": 2,
"timestamp-ms": 2234,
"manifest-list": second_manifest_list,
"summary": {
"operation": "append"
}
}
},
{
"action": "set-snapshot-ref",
"ref-name": "main",
"snapshot-id": 11,
"type": "branch"
}
]
}))
.expect("second append request should parse");
let commit = commit_table_response(&store, &metadata_backend, "warehouse", &namespace, "events", second_append)
.await
.expect("append should preserve inherited manifests");
assert_eq!(commit.metadata["current-snapshot-id"], 11);
assert_eq!(commit.metadata["last-sequence-number"], 2);
}
#[tokio::test]
async fn row_level_conflict_rejects_changed_inherited_manifest_identity() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
.as_str()
.expect("created metadata should have table location");
let first_manifest_list = format!("{table_location}/metadata/snap-10.avro");
let first_manifest = format!("{table_location}/metadata/manifest-10.avro");
let first_data_file = format!("{table_location}/data/part-10.parquet");
seed_test_manifest_list(&metadata_backend, "warehouse", &first_manifest_list, &[&first_manifest], 1, 10).await;
seed_test_manifest(&metadata_backend, "warehouse", &first_manifest, &[(&first_data_file, 0, 1, 10, 1)]).await;
let first_append: RestCommitTableRequest = serde_json::from_value(serde_json::json!({
"updates": [
{
"action": "add-snapshot",
"snapshot": {
"snapshot-id": 10,
"sequence-number": 1,
"timestamp-ms": 1234,
"manifest-list": first_manifest_list,
"summary": {
"operation": "append"
}
}
},
{
"action": "set-snapshot-ref",
"ref-name": "main",
"snapshot-id": 10,
"type": "branch"
}
]
}))
.expect("first append request should parse");
commit_table_response(&store, &metadata_backend, "warehouse", &namespace, "events", first_append)
.await
.expect("first append should commit");
let current = store
.load_table("warehouse", "analytics", "events")
.await
.expect("table lookup should succeed")
.expect("table should exist");
let second_manifest_list = format!("{table_location}/metadata/snap-11.avro");
seed_test_manifest_list(&metadata_backend, "warehouse", &second_manifest_list, &[&first_manifest], 2, 10).await;
let second_append: RestCommitTableRequest = serde_json::from_value(serde_json::json!({
"requirements": [
{
"type": "assert-current-snapshot-id",
"snapshot-id": 10
}
],
"updates": [
{
"action": "add-snapshot",
"snapshot": {
"snapshot-id": 11,
"parent-snapshot-id": 10,
"sequence-number": 2,
"timestamp-ms": 2234,
"manifest-list": second_manifest_list,
"summary": {
"operation": "append"
}
}
}
]
}))
.expect("second append request should parse");
let error = commit_table_response(&store, &metadata_backend, "warehouse", &namespace, "events", second_append)
.await
.expect_err("inherited manifest identity must not change");
assert_eq!(error.code(), &s3s::S3ErrorCode::InvalidRequest);
let unchanged = store
.load_table("warehouse", "analytics", "events")
.await
.expect("table lookup should succeed")
.expect("table should still exist");
assert_eq!(unchanged.metadata_location, current.metadata_location);
assert_eq!(unchanged.version_token, current.version_token);
assert_eq!(unchanged.generation, current.generation);
}
#[tokio::test]
async fn row_level_conflict_rejects_stale_new_manifest_sequence() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
.as_str()
.expect("created metadata should have table location");
let current = store
.load_table("warehouse", "analytics", "events")
.await
.expect("table lookup should succeed")
.expect("table should exist");
let manifest_list = format!("{table_location}/metadata/snap-11.avro");
let manifest = format!("{table_location}/metadata/manifest-11.avro");
let data_file = format!("{table_location}/data/part-11.parquet");
seed_test_manifest_list(&metadata_backend, "warehouse", &manifest_list, &[&manifest], 1, 11).await;
seed_test_manifest(&metadata_backend, "warehouse", &manifest, &[(&data_file, 0, 1, 11, 1)]).await;
let append_request: RestCommitTableRequest = serde_json::from_value(serde_json::json!({
"updates": [
{
"action": "add-snapshot",
"snapshot": {
"snapshot-id": 11,
"sequence-number": 2,
"timestamp-ms": 2234,
"manifest-list": manifest_list,
"summary": {
"operation": "append"
}
}
}
]
}))
.expect("append request should parse");
let error = commit_table_response(&store, &metadata_backend, "warehouse", &namespace, "events", append_request)
.await
.expect_err("new manifest sequence must match the committed snapshot");
assert_eq!(error.code(), &s3s::S3ErrorCode::InvalidRequest);
let unchanged = store
.load_table("warehouse", "analytics", "events")
.await
.expect("table lookup should succeed")
.expect("table should still exist");
assert_eq!(unchanged.metadata_location, current.metadata_location);
assert_eq!(unchanged.version_token, current.version_token);
assert_eq!(unchanged.generation, current.generation);
}
#[tokio::test]
async fn row_level_conflict_rejects_stale_added_entry_sequence() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
.as_str()
.expect("created metadata should have table location");
let current = store
.load_table("warehouse", "analytics", "events")
.await
.expect("table lookup should succeed")
.expect("table should exist");
let manifest_list = format!("{table_location}/metadata/snap-11.avro");
let manifest = format!("{table_location}/metadata/manifest-11.avro");
let data_file = format!("{table_location}/data/part-11.parquet");
seed_test_manifest_list(&metadata_backend, "warehouse", &manifest_list, &[&manifest], 2, 11).await;
seed_test_manifest(&metadata_backend, "warehouse", &manifest, &[(&data_file, 0, 1, 11, 1)]).await;
let append_request: RestCommitTableRequest = serde_json::from_value(serde_json::json!({
"updates": [
{
"action": "add-snapshot",
"snapshot": {
"snapshot-id": 11,
"sequence-number": 2,
"timestamp-ms": 2234,
"manifest-list": manifest_list,
"summary": {
"operation": "append"
}
}
}
]
}))
.expect("append request should parse");
let error = commit_table_response(&store, &metadata_backend, "warehouse", &namespace, "events", append_request)
.await
.expect_err("added file sequence must match the new manifest");
assert_eq!(error.code(), &s3s::S3ErrorCode::InvalidRequest);
let unchanged = store
.load_table("warehouse", "analytics", "events")
.await
.expect("table lookup should succeed")
.expect("table should still exist");
assert_eq!(unchanged.metadata_location, current.metadata_location);
assert_eq!(unchanged.version_token, current.version_token);
assert_eq!(unchanged.generation, current.generation);
}
#[tokio::test]
async fn row_level_conflict_rejects_historical_change_in_new_manifest() {
let store = TestTableCatalogStore::default();
let metadata_backend = TestTableCatalogObjectBackend::default();
let namespace = crate::table_catalog::Namespace::parse("analytics").expect("namespace should parse");
let created = create_standard_events_table(&store, &metadata_backend, &namespace).await;
let table_location = created.metadata["location"]
.as_str()
.expect("created metadata should have table location");
let current = store
.load_table("warehouse", "analytics", "events")
.await
.expect("table lookup should succeed")
.expect("table should exist");
let manifest_list = format!("{table_location}/metadata/snap-11.avro");
let manifest = format!("{table_location}/metadata/manifest-11.avro");
let data_file = format!("{table_location}/data/part-10.parquet");
seed_test_manifest_list(&metadata_backend, "warehouse", &manifest_list, &[&manifest], 2, 11).await;
seed_test_manifest(&metadata_backend, "warehouse", &manifest, &[(&data_file, 0, 1, 10, 1)]).await;
let append_request: RestCommitTableRequest = serde_json::from_value(serde_json::json!({
"updates": [
{
"action": "add-snapshot",
"snapshot": {
"snapshot-id": 11,
"sequence-number": 2,
"timestamp-ms": 2234,
"manifest-list": manifest_list,
"summary": {
"operation": "append"
}
}
}
]
}))
.expect("append request should parse");
let error = commit_table_response(&store, &metadata_backend, "warehouse", &namespace, "events", append_request)
.await
.expect_err("new manifest must not claim a historical changed entry");
assert_eq!(error.code(), &s3s::S3ErrorCode::InvalidRequest);
let unchanged = store
.load_table("warehouse", "analytics", "events")
.await
.expect("table lookup should succeed")
.expect("table should still exist");
assert_eq!(unchanged.metadata_location, current.metadata_location);
assert_eq!(unchanged.version_token, current.version_token);
assert_eq!(unchanged.generation, current.generation);
}
#[tokio::test]
async fn row_level_conflict_allows_add_only_overwrite_snapshot() {
let store = TestTableCatalogStore::default();
@@ -9630,6 +10148,14 @@ mod tests {
}
fn test_manifest_list_avro_bytes(manifest_paths: &[&str], sequence_number: i64, snapshot_id: i64) -> Vec<u8> {
let manifests = manifest_paths
.iter()
.map(|manifest_path| (*manifest_path, sequence_number, snapshot_id))
.collect::<Vec<_>>();
test_manifest_list_avro_entries(&manifests)
}
fn test_manifest_list_avro_entries(manifests: &[(&str, i64, i64)]) -> Vec<u8> {
let schema = apache_avro::Schema::parse_str(
r#"
{
@@ -9645,15 +10171,15 @@ mod tests {
)
.expect("manifest list avro schema should parse");
let mut writer = apache_avro::Writer::new(&schema, Vec::new());
for manifest_path in manifest_paths {
for (manifest_path, sequence_number, snapshot_id) in manifests {
writer
.append(apache_avro::types::Value::Record(vec![
(
"manifest_path".to_string(),
apache_avro::types::Value::String((*manifest_path).to_string()),
),
("sequence_number".to_string(), apache_avro::types::Value::Long(sequence_number)),
("added_snapshot_id".to_string(), apache_avro::types::Value::Long(snapshot_id)),
("sequence_number".to_string(), apache_avro::types::Value::Long(*sequence_number)),
("added_snapshot_id".to_string(), apache_avro::types::Value::Long(*snapshot_id)),
]))
.expect("manifest list record should append");
}