From 1e7fb8438a9ab94b20d53a9161384edab218300e Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Wed, 15 Jul 2026 17:33:42 +0800 Subject: [PATCH] fix(table-catalog): accept inherited snapshot manifests (#4833) Co-authored-by: Henry Guo --- rustfs/src/admin/handlers/table_catalog.rs | 634 +++++++++++++++++++-- 1 file changed, 580 insertions(+), 54 deletions(-) diff --git a/rustfs/src/admin/handlers/table_catalog.rs b/rustfs/src/admin/handlers/table_catalog.rs index ba749c760..c7981f13b 100644 --- a/rustfs/src/admin/handlers/table_catalog.rs +++ b/rustfs/src/admin/handlers/table_catalog.rs @@ -2794,13 +2794,26 @@ fn snapshot_has_manifest_references(snapshot: &serde_json::Value) -> bool { #[derive(Default)] struct SnapshotLiveFiles { - data_files: BTreeSet, - delete_files: BTreeSet, + data_files: BTreeMap, + delete_files: BTreeMap, + manifest_files: BTreeMap, } 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, + added_snapshot_id: Option, +} + +#[derive(PartialEq, Eq)] +struct SnapshotFileIdentity { + sequence_number: Option, + file_sequence_number: Option, +} + async fn validate_table_snapshot_commit_conflicts( metadata_backend: &B, bucket: &str, @@ -2870,7 +2895,8 @@ where table, entry, snapshot, - SnapshotChangeIdentity { + SnapshotChangeContext { + current_live_files: ¤t_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( table: &crate::table_catalog::IdentifierSegment, entry: &crate::table_catalog::TableEntry, snapshot: &serde_json::Value, - change_identity: SnapshotChangeIdentity, + context: SnapshotChangeContext<'_>, ) -> S3Result 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, +} + async fn read_snapshot_manifest_references( metadata_backend: &B, bucket: &str, @@ -3043,12 +3152,12 @@ async fn read_snapshot_manifest_references( table: &crate::table_catalog::IdentifierSegment, entry: &crate::table_catalog::TableEntry, snapshot: &serde_json::Value, -) -> S3Result> +) -> S3Result> 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 { + let manifests = manifest_paths + .iter() + .map(|manifest_path| (*manifest_path, sequence_number, snapshot_id)) + .collect::>(); + test_manifest_list_avro_entries(&manifests) + } + + fn test_manifest_list_avro_entries(manifests: &[(&str, i64, i64)]) -> Vec { 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"); }