diff --git a/crates/ecstore/src/core/pools_test.rs b/crates/ecstore/src/core/pools_test.rs index 9e7e9390f..76de24d8c 100644 --- a/crates/ecstore/src/core/pools_test.rs +++ b/crates/ecstore/src/core/pools_test.rs @@ -3737,7 +3737,8 @@ mod decommission_lock_order_tests { .await .expect_err("lost data-movement PUT capacity lease must not publish the object"); assert!( - crate::error::is_err_object_not_found(&object_err), + matches!(&object_err, crate::error::Error::VersionNotFound(b, o, v) + if b == RUSTFS_META_BUCKET && o == object && v == &version_id), "data-movement PUT lease loss should leave the target object absent: {object_err}" ); @@ -3849,13 +3850,17 @@ mod decommission_lock_order_tests { next_object, &ObjectOptions { versioned: true, - version_id: Some(next_version_id), + version_id: Some(next_version_id.clone()), ..Default::default() }, ) .await .expect_err("the next mutation must remain unpublished after A recovery"); - assert!(crate::error::is_err_object_not_found(&next_target_err)); + assert!( + matches!(&next_target_err, crate::error::Error::VersionNotFound(b, o, v) + if b == RUSTFS_META_BUCKET && o == next_object && v == &next_version_id), + "the next mutation must leave its exact target version absent: {next_target_err}" + ); } #[tokio::test] diff --git a/crates/ecstore/src/set_disk/metadata.rs b/crates/ecstore/src/set_disk/metadata.rs index 102743827..ec432f891 100644 --- a/crates/ecstore/src/set_disk/metadata.rs +++ b/crates/ecstore/src/set_disk/metadata.rs @@ -410,7 +410,13 @@ impl SetDisks { default_parity_count: usize, ) -> disk::error::Result<(i32, i32)> { if Self::all_not_found_metadata(errs) { - return Err(DiskError::FileNotFound); + // Preserve explicit-version absence from the disk replies before + // the object/API error boundary assigns the S3 error code. + return Err(if errs.iter().any(|err| matches!(err, Some(DiskError::FileVersionNotFound))) { + DiskError::FileVersionNotFound + } else { + DiskError::FileNotFound + }); } let expected_rquorum = if default_parity_count == 0 { diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 656835eb3..ae4f083ba 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -11527,6 +11527,76 @@ mod tests { assert_eq!(err, DiskError::FileNotFound); } + #[test] + fn test_object_quorum_from_meta_preserves_version_not_found() { + for errs in [ + vec![Some(DiskError::FileVersionNotFound); 4], + vec![ + Some(DiskError::FileVersionNotFound), + Some(DiskError::FileNotFound), + Some(DiskError::FileVersionNotFound), + Some(DiskError::FileNotFound), + ], + vec![ + Some(DiskError::FileVersionNotFound), + Some(DiskError::VolumeNotFound), + Some(DiskError::DiskNotFound), + Some(DiskError::FileVersionNotFound), + ], + ] { + let err = SetDisks::object_quorum_from_meta(&vec![FileInfo::default(); errs.len()], &errs, 2) + .expect_err("absent version metadata must remain a version miss"); + assert_eq!(err, DiskError::FileVersionNotFound, "disk replies: {errs:?}"); + } + } + + #[test] + fn test_object_quorum_from_meta_version_misses_preserve_other_failures() { + for (errs, expected) in [ + (vec![Some(DiskError::DiskNotFound); 4], DiskError::ErasureReadQuorum), + ( + vec![ + Some(DiskError::FileVersionNotFound), + Some(DiskError::FileCorrupt), + Some(DiskError::DiskNotFound), + None, + ], + DiskError::ErasureReadQuorum, + ), + ( + vec![ + Some(DiskError::FileVersionNotFound), + Some(DiskError::FileAccessDenied), + Some(DiskError::FileAccessDenied), + Some(DiskError::DiskNotFound), + ], + DiskError::FileAccessDenied, + ), + ( + vec![ + Some(DiskError::FileVersionNotFound), + Some(DiskError::VolumeNotFound), + Some(DiskError::VolumeNotFound), + None, + ], + DiskError::VolumeNotFound, + ), + ( + vec![ + Some(DiskError::FileVersionNotFound), + Some(DiskError::FileVersionNotFound), + Some(DiskError::FileCorrupt), + None, + ], + DiskError::FileVersionNotFound, + ), + ] { + let err = SetDisks::object_quorum_from_meta(&vec![FileInfo::default(); errs.len()], &errs, 2) + .expect_err("metadata failures must retain quorum reduction semantics"); + assert_eq!(err, expected, "disk replies: {errs:?}"); + } + } + #[test] fn test_object_quorum_from_meta_preserves_read_quorum_for_mixed_failures() { let errs = vec![ diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index 35517a7d0..abf875ec6 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -5591,7 +5591,8 @@ mod heal_result_report_tests { ) .await; assert!( - matches!(&resurrected, Err(Error::FileVersionNotFound) | Err(Error::ObjectNotFound(..))), + matches!(&resurrected, Err(Error::VersionNotFound(b, o, v)) + if b == bucket && o == object && v == &first_version), "a racing heal must not resurrect the deleted version: {resurrected:?}" ); diff --git a/crates/ecstore/src/set_disk/read.rs b/crates/ecstore/src/set_disk/read.rs index 0c1c8392c..4deb5bf30 100644 --- a/crates/ecstore/src/set_disk/read.rs +++ b/crates/ecstore/src/set_disk/read.rs @@ -540,7 +540,7 @@ impl SetDisks { let metadata_resolve_stage_start = get_stage_timer_if_enabled(stage_metrics_enabled); let (read_quorum, write_quorum) = match Self::object_quorum_from_meta(&parts_metadata, &errs, self.default_parity_count) - .map_err(|err| to_object_err(err.into(), vec![bucket, object])) + .map_err(|err| to_object_err(err.into(), vec![bucket, object, &vid])) { Ok(v) => v, Err(e) => { @@ -564,7 +564,7 @@ impl SetDisks { GET_STAGE_METADATA_RESOLVE, metadata_resolve_stage_start, ); - return Err(to_object_err(err.into(), vec![bucket, object])); + return Err(to_object_err(err.into(), vec![bucket, object, &vid])); } let (op_online_disks, mut fi, fileinfo_selection_quorum) = diff --git a/rustfs/src/app/object/delete.rs b/rustfs/src/app/object/delete.rs index 7d376ccde..63c8a472c 100644 --- a/rustfs/src/app/object/delete.rs +++ b/rustfs/src/app/object/delete.rs @@ -1663,11 +1663,7 @@ mod tests { .expect_err("the marker version must actually be gone"); assert_eq!( missing.code(), - &if pool_count == 1 { - S3ErrorCode::NoSuchKey - } else { - S3ErrorCode::NoSuchVersion - }, + &S3ErrorCode::NoSuchVersion, "pool_count={pool_count} suspended={suspended} batch={batch}: removed marker lookup returned {missing:?}" ); let get = GetObjectInput::builder() diff --git a/rustfs/tests/embedded_test.rs b/rustfs/tests/embedded_test.rs index a2df768c3..367729583 100644 --- a/rustfs/tests/embedded_test.rs +++ b/rustfs/tests/embedded_test.rs @@ -21,6 +21,8 @@ use aws_sdk_s3::config::{Credentials, Region}; use aws_sdk_s3::error::ProvideErrorMetadata; +use aws_sdk_s3::error::SdkError; +use aws_sdk_s3::operation::get_object::GetObjectError; use aws_sdk_s3::primitives::ByteStream; use aws_sdk_s3::types::{ BucketVersioningStatus, Delete, ObjectAttributes, ObjectIdentifier, Tag, Tagging, VersioningConfiguration, @@ -325,6 +327,320 @@ async fn assert_tagging_version( assert_eq!(tags.tag_set(), expected, "tagging contents for selector {selector:?}"); } +async fn assert_get_object_error( + client: &Client, + bucket: &str, + key: &str, + selector: Option<&str>, + status: u16, + code: &str, +) -> SdkError { + let error = client + .get_object() + .bucket(bucket) + .key(key) + .set_version_id(selector.map(str::to_owned)) + .send() + .await + .expect_err("GET must reject the selected object or version"); + assert_eq!(error.raw_response().map(|response| response.status().as_u16()), Some(status)); + assert_eq!( + error.as_service_error().and_then(ProvideErrorMetadata::code), + Some(code), + "GET {bucket}/{key}, selector {selector:?}: {error:?}" + ); + if code == "NoSuchVersion" { + let headers = error.raw_response().expect("S3 error response").headers(); + assert_eq!(headers.get("x-amz-delete-marker"), None); + assert_eq!(headers.get("x-amz-version-id"), None); + } + error +} + +#[test] +fn test_get_object_version_errors() { + if let Ok(pool_count) = std::env::var("RUSTFS_TEST_GET_VERSION_POOL_COUNT") { + let pool_count: usize = pool_count.parse().expect("test pool count"); + assert!((1..=2).contains(&pool_count)); + common::run_embedded_test(move || get_object_version_errors(pool_count)); + return; + } + + // Cache configuration and logical-drive overrides are process-wide. Each + // child exercises the real HTTP route with its own server and disk roots. + for pool_count in [1, 2] { + for mode in ["disabled", "fill_materialize_enabled"] { + let output = std::process::Command::new(std::env::current_exe().expect("test executable")) + .args(["--exact", "test_get_object_version_errors", "--nocapture"]) + .env("RUSTFS_TEST_GET_VERSION_POOL_COUNT", pool_count.to_string()) + .env("RUSTFS_UNSAFE_BYPASS_DISK_CHECK", "true") + .env("RUSTFS_OBJECT_DATA_CACHE_ENABLE", if mode == "disabled" { "false" } else { "true" }) + .env("RUSTFS_OBJECT_DATA_CACHE_MODE", mode) + .env("RUSTFS_OBJECT_DATA_CACHE_MAX_BYTES", "8388608") + .env("RUSTFS_OBJECT_DATA_CACHE_MAX_ENTRY_BYTES", "1048576") + .env("RUSTFS_OBJECT_DATA_CACHE_MIN_FREE_MEMORY_PERCENT", "0") + .env("NO_PROXY", "127.0.0.1,localhost") + .env("no_proxy", "127.0.0.1,localhost") + .output() + .expect("run GET version error child"); + let stdout = String::from_utf8_lossy(&output.stdout); + assert!( + output.status.success() && stdout.contains("test result: ok. 1 passed;"), + "GET version errors, pools={pool_count}, cache={mode}:\n{stdout}\n{}", + String::from_utf8_lossy(&output.stderr) + ); + } + } +} + +async fn get_object_version_errors(pool_count: usize) { + let root = tempfile::tempdir().expect("temporary drives"); + for pool in 0..pool_count { + for disk in 1..=4 { + std::fs::create_dir_all(root.path().join(format!("pool{pool}/disk{disk}"))).expect("create test drive"); + } + } + let server = RustFSServerBuilder::new() + .address(format!("127.0.0.1:{}", find_available_port().expect("free port"))) + .access_key("testaccesskey") + .secret_key("testsecretkey") + .volumes( + (0..pool_count) + .map(|pool| format!("{}/pool{pool}/disk{{1...4}}", root.path().display())) + .collect(), + ) + .build() + .await + .expect("start version error server"); + let client = s3_client(&server.endpoint(), server.access_key(), server.secret_key()); + let bucket = "get-version-errors"; + let key = "objects/history.bin"; + client.create_bucket().bucket(bucket).send().await.expect("create bucket"); + client + .put_bucket_versioning() + .bucket(bucket) + .versioning_configuration( + VersioningConfiguration::builder() + .status(BucketVersioningStatus::Enabled) + .build(), + ) + .send() + .await + .expect("enable versioning"); + + let mut versions = Vec::new(); + for body in [b"version one", b"version two", b"version tri"] { + let put = client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from_static(body)) + .send() + .await + .expect("write data version"); + let version = put.version_id().expect("acknowledged version ID").to_owned(); + assert_read_version(&client, bucket, key, Some(&version), Some(&version), body).await; + versions.push(version); + } + for (removed, current, body) in [ + (&versions[0], &versions[2], b"version tri"), + (&versions[2], &versions[1], b"version two"), + ] { + client + .delete_object() + .bucket(bucket) + .key(key) + .version_id(removed) + .send() + .await + .expect("delete selected data version"); + assert_get_object_error(&client, bucket, key, Some(removed), 404, "NoSuchVersion").await; + assert_read_version(&client, bucket, key, None, Some(current), body).await; + assert_read_version(&client, bucket, key, Some(&versions[1]), Some(&versions[1]), b"version two").await; + } + + let absent = uuid::Uuid::new_v4().to_string(); + for missing_key in [key, "never-created"] { + for selector in [absent.as_str(), "null"] { + assert_get_object_error(&client, bucket, missing_key, Some(selector), 404, "NoSuchVersion").await; + } + } + assert_get_object_error(&client, bucket, "never-created", None, 404, "NoSuchKey").await; + assert_get_object_error(&client, bucket, key, Some("invalid-uuid"), 400, "InvalidArgument").await; + assert_get_object_error(&client, "never-created-bucket", key, Some(&absent), 404, "NoSuchBucket").await; + let denied = reqwest::Client::builder() + .no_proxy() + .build() + .expect("anonymous HTTP client") + .get(format!("{}/{bucket}/{key}?versionId={absent}", server.endpoint())) + .send() + .await + .expect("anonymous version GET"); + assert_eq!(denied.status(), reqwest::StatusCode::FORBIDDEN); + assert!(denied.text().await.expect("denial XML").contains("AccessDenied")); + + let mut batch = Vec::new(); + for delete_in_batch in [false, true] { + let marker = client + .delete_object() + .bucket(bucket) + .key(key) + .send() + .await + .expect("create marker"); + assert_eq!(marker.delete_marker(), Some(true)); + let marker_id = marker.version_id().expect("marker ID"); + for (selector, status, code) in [(None, 404, "NoSuchKey"), (Some(marker_id), 405, "MethodNotAllowed")] { + let error = assert_get_object_error(&client, bucket, key, selector, status, code).await; + let headers = error.raw_response().expect("marker response").headers(); + assert_eq!(headers.get("x-amz-delete-marker"), Some("true")); + assert_eq!(headers.get("x-amz-version-id"), Some(marker_id)); + if selector.is_some() { + assert!(headers.get("last-modified").is_some()); + } + } + assert_read_version(&client, bucket, key, Some(&versions[1]), Some(&versions[1]), b"version two").await; + if delete_in_batch { + batch.push( + ObjectIdentifier::builder() + .key(key) + .version_id(marker_id) + .build() + .expect("marker selector"), + ); + } else { + client + .delete_object() + .bucket(bucket) + .key(key) + .version_id(marker_id) + .send() + .await + .expect("purge marker"); + assert_get_object_error(&client, bucket, key, Some(marker_id), 404, "NoSuchVersion").await; + assert_read_version(&client, bucket, key, None, Some(&versions[1]), b"version two").await; + } + } + batch.push( + ObjectIdentifier::builder() + .key(key) + .version_id(&versions[1]) + .build() + .expect("last data selector"), + ); + let only = client + .put_object() + .bucket(bucket) + .key("only-version") + .body(ByteStream::from_static(b"only data")) + .send() + .await + .expect("write only version"); + batch.push( + ObjectIdentifier::builder() + .key("only-version") + .version_id(only.version_id().expect("only version ID")) + .build() + .expect("only selector"), + ); + let deleted = client + .delete_objects() + .bucket(bucket) + .delete( + Delete::builder() + .set_objects(Some(batch.clone())) + .build() + .expect("batch delete"), + ) + .send() + .await + .expect("delete exact versions in batch"); + assert!(deleted.errors().is_empty(), "batch errors: {:?}", deleted.errors()); + assert_eq!(deleted.deleted().len(), batch.len()); + for object in batch { + assert_get_object_error(&client, bucket, object.key(), object.version_id(), 404, "NoSuchVersion").await; + } + assert_get_object_error(&client, bucket, key, None, 404, "NoSuchKey").await; + + // A pre-versioning null slot stays addressable while suspended. Removing + // it or its replacement marker must not fall back to the retained UUID. + let bucket = "get-null-version-errors"; + client + .create_bucket() + .bucket(bucket) + .send() + .await + .expect("create null bucket"); + client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from_static(b"null data")) + .send() + .await + .expect("write pre-versioning object"); + let mut retained = String::new(); + for status in [BucketVersioningStatus::Enabled, BucketVersioningStatus::Suspended] { + client + .put_bucket_versioning() + .bucket(bucket) + .versioning_configuration(VersioningConfiguration::builder().status(status.clone()).build()) + .send() + .await + .expect("set versioning state"); + assert_read_version(&client, bucket, key, Some("null"), Some("null"), b"null data").await; + if status == BucketVersioningStatus::Enabled { + let put = client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from_static(b"retained data")) + .send() + .await + .expect("write retained UUID version"); + retained = put.version_id().expect("retained version").to_owned(); + } + } + client + .delete_object() + .bucket(bucket) + .key(key) + .version_id("null") + .send() + .await + .expect("delete null slot"); + assert_get_object_error(&client, bucket, key, Some("null"), 404, "NoSuchVersion").await; + let marker = client + .delete_object() + .bucket(bucket) + .key(key) + .send() + .await + .expect("create suspended null marker"); + assert_eq!(marker.version_id(), Some("null")); + for (selector, status, code) in [(None, 404, "NoSuchKey"), (Some("null"), 405, "MethodNotAllowed")] { + let error = assert_get_object_error(&client, bucket, key, selector, status, code).await; + let headers = error.raw_response().expect("null marker response").headers(); + assert_eq!(headers.get("x-amz-delete-marker"), Some("true")); + assert_eq!(headers.get("x-amz-version-id"), Some("null")); + if selector.is_some() { + assert!(headers.get("last-modified").is_some()); + } + } + client + .delete_object() + .bucket(bucket) + .key(key) + .version_id("null") + .send() + .await + .expect("purge null marker"); + assert_get_object_error(&client, bucket, key, Some("null"), 404, "NoSuchVersion").await; + assert_read_version(&client, bucket, key, None, Some(&retained), b"retained data").await; + assert_read_version(&client, bucket, key, Some(&retained), Some(&retained), b"retained data").await; + server.shutdown().await; +} + #[test] fn test_read_version_headers_across_versioning_states() { common::run_embedded_test(test_read_version_headers_across_versioning_states_body);