diff --git a/rustfs/src/app/object/get.rs b/rustfs/src/app/object/get.rs index f9a0ca36e..85e11f35b 100644 --- a/rustfs/src/app/object/get.rs +++ b/rustfs/src/app/object/get.rs @@ -3620,7 +3620,6 @@ impl DefaultObjectUsecase { queue_status: &concurrency::IoQueueStatus, concurrent_requests: usize, part_number: Option, - versioned: bool, lifecycle: GetObjectBodyLifecycle, resume: F, ) -> S3Result @@ -3678,17 +3677,7 @@ impl DefaultObjectUsecase { let checksums = Self::build_get_object_checksums(&info, &req.headers, part_number, rs.as_ref())?; record_get_object_s3_handler_stage_duration(GET_OBJECT_STAGE_CHECKSUM_HEADERS, checksum_headers_start); - let output_version_id = if versioned { - info.version_id.map(|vid| { - if vid == Uuid::nil() { - "null".to_string() - } else { - vid.to_string() - } - }) - } else { - None - }; + let output_version_id = read_response_version_id(info.version_id); // x-amz-restore: extract from object metadata let restore = info.user_defined.get(X_AMZ_RESTORE.as_str()).and_then(|v| { @@ -4181,7 +4170,6 @@ impl DefaultObjectUsecase { &queue_status, concurrent_requests, part_number, - opts.versioned, lifecycle, |info| { Some(get_object_resume_control(GetObjectResumeContext::new( @@ -4426,17 +4414,7 @@ impl DefaultObjectUsecase { None }; - let version_id = if BucketVersioningSys::prefix_enabled(&bucket, &key).await { - info.version_id.map(|vid| { - if vid == Uuid::nil() { - "null".to_string() - } else { - vid.to_string() - } - }) - } else { - None - }; + let version_id = read_response_version_id(info.version_id); let output = GetObjectAttributesOutput { checksum, @@ -10639,7 +10617,6 @@ mod tests { &queue_status, 1, None, - false, GetObjectBodyLifecycle::disabled(), |_| panic!("a buffered output must not initialize streaming resume state"), ) diff --git a/rustfs/src/app/object/head.rs b/rustfs/src/app/object/head.rs index b649d1630..ae208e20e 100644 --- a/rustfs/src/app/object/head.rs +++ b/rustfs/src/app/object/head.rs @@ -553,7 +553,7 @@ impl DefaultObjectUsecase { last_modified, e_tag: info.etag.map(|etag| to_s3s_etag(&etag)), metadata: filter_object_metadata(&metadata_map), - version_id: info.version_id.map(|v| v.to_string()), + version_id: read_response_version_id(info.version_id), server_side_encryption, sse_customer_algorithm, sse_customer_key_md5, diff --git a/rustfs/src/app/object/shared.rs b/rustfs/src/app/object/shared.rs index 79edcfa32..7c5d09133 100644 --- a/rustfs/src/app/object/shared.rs +++ b/rustfs/src/app/object/shared.rs @@ -32,6 +32,18 @@ pub(super) const LOG_COMPONENT_APP: &str = "app"; pub(super) const LOG_SUBSYSTEM_OBJECT: &str = "object"; +/// Encode the resolved local read identity. Storage distinguishes a null +/// version from an unversioned object by returning a nil UUID instead of None. +pub(super) fn read_response_version_id(version_id: Option) -> Option { + version_id.map(|id| { + if id.is_nil() { + NULL_VERSION_ID.to_owned() + } else { + id.to_string() + } + }) +} + fn is_delete_marker_read_error(err: &S3Error, version_id: Option<&str>) -> bool { let code = if version_id.is_some() { S3ErrorCode::MethodNotAllowed diff --git a/rustfs/tests/embedded_test.rs b/rustfs/tests/embedded_test.rs index e46ccaa31..654a37285 100644 --- a/rustfs/tests/embedded_test.rs +++ b/rustfs/tests/embedded_test.rs @@ -21,7 +21,7 @@ use aws_sdk_s3::config::{Credentials, Region}; use aws_sdk_s3::primitives::ByteStream; -use aws_sdk_s3::types::{BucketVersioningStatus, Delete, ObjectIdentifier, VersioningConfiguration}; +use aws_sdk_s3::types::{BucketVersioningStatus, Delete, ObjectAttributes, ObjectIdentifier, VersioningConfiguration}; use aws_sdk_s3::{Client, Config}; use rustfs::embedded::{RustFSServerBuilder, find_available_port}; @@ -249,3 +249,191 @@ async fn test_null_version_delete_marker_round_trip_body() { server.shutdown().await; } + +async fn assert_read_version( + client: &Client, + bucket: &str, + key: &str, + selector: Option<&str>, + expected_version: Option<&str>, + expected_body: &[u8], +) { + let get = client + .get_object() + .bucket(bucket) + .key(key) + .set_version_id(selector.map(str::to_owned)) + .send() + .await + .expect("get selected version"); + assert_eq!(get.version_id(), expected_version, "GET identity for selector {selector:?}"); + assert_eq!(get.content_length(), Some(expected_body.len() as i64)); + let etag = get.e_tag().expect("GET ETag").to_owned(); + assert_eq!(get.body.collect().await.expect("read selected body").into_bytes().as_ref(), expected_body); + + let head = client + .head_object() + .bucket(bucket) + .key(key) + .set_version_id(selector.map(str::to_owned)) + .send() + .await + .expect("head selected version"); + assert_eq!(head.version_id(), expected_version, "HEAD identity for selector {selector:?}"); + assert_eq!(head.content_length(), Some(expected_body.len() as i64)); + assert_eq!(head.e_tag(), Some(etag.as_str())); + + let attributes = client + .get_object_attributes() + .bucket(bucket) + .key(key) + .set_version_id(selector.map(str::to_owned)) + .object_attributes(ObjectAttributes::ObjectSize) + .send() + .await + .expect("get selected version attributes"); + assert_eq!(attributes.version_id(), expected_version, "Attributes identity for selector {selector:?}"); + assert_eq!(attributes.object_size(), Some(expected_body.len() as i64)); +} + +#[test] +fn test_read_version_headers_across_versioning_states() { + common::run_embedded_test(test_read_version_headers_across_versioning_states_body); +} + +async fn test_read_version_headers_across_versioning_states_body() { + let port = find_available_port().expect("find free port"); + let server = RustFSServerBuilder::new() + .address(format!("127.0.0.1:{port}")) + .access_key("testaccesskey") + .secret_key("testsecretkey") + .build() + .await + .expect("start embedded server"); + let client = s3_client(&server.endpoint(), server.access_key(), server.secret_key()); + let bucket = "read-version-headers"; + let key = "history.txt"; + client.create_bucket().bucket(bucket).send().await.expect("create bucket"); + + let unversioned = client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from_static(b"before-versioning")) + .send() + .await + .expect("put unversioned object"); + assert_eq!(unversioned.version_id(), None); + for selector in [None, Some("null")] { + assert_read_version(&client, bucket, key, selector, None, b"before-versioning").await; + } + + client + .put_bucket_versioning() + .bucket(bucket) + .versioning_configuration( + VersioningConfiguration::builder() + .status(BucketVersioningStatus::Enabled) + .build(), + ) + .send() + .await + .expect("enable versioning"); + for selector in [None, Some("null")] { + assert_read_version(&client, bucket, key, selector, Some("null"), b"before-versioning").await; + } + + let mut history = Vec::new(); + for body in [b"enabled-one".as_slice(), b"enabled-second".as_slice()] { + let put = client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from(body.to_vec())) + .send() + .await + .expect("put enabled version"); + let version = put.version_id().expect("enabled PUT version").to_owned(); + assert!(!uuid::Uuid::parse_str(&version).expect("version UUID").is_nil()); + assert_read_version(&client, bucket, key, None, Some(&version), body).await; + history.push((version, body)); + } + + client + .put_bucket_versioning() + .bucket(bucket) + .versioning_configuration( + VersioningConfiguration::builder() + .status(BucketVersioningStatus::Suspended) + .build(), + ) + .send() + .await + .expect("suspend versioning"); + for (version, body) in &history { + assert_read_version(&client, bucket, key, Some(version), Some(version), body).await; + } + let (latest_version, latest_body) = history.last().expect("latest UUID version"); + assert_read_version(&client, bucket, key, None, Some(latest_version), latest_body).await; + assert_read_version(&client, bucket, key, Some("null"), Some("null"), b"before-versioning").await; + + for body in [b"null-one".as_slice(), b"null-overwritten".as_slice()] { + let put = client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from(body.to_vec())) + .send() + .await + .expect("put suspended null version"); + // S3 omits the null identity on PUT, but returns it on subsequent reads. + assert_eq!(put.version_id(), None, "suspended PUT keeps its write response contract"); + for selector in [None, Some("null")] { + assert_read_version(&client, bucket, key, selector, Some("null"), body).await; + } + } + + client + .put_bucket_versioning() + .bucket(bucket) + .versioning_configuration( + VersioningConfiguration::builder() + .status(BucketVersioningStatus::Enabled) + .build(), + ) + .send() + .await + .expect("re-enable versioning"); + let enabled_again = client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from_static(b"enabled-again")) + .send() + .await + .expect("put after re-enabling"); + let current = enabled_again.version_id().expect("re-enabled PUT version"); + assert_read_version(&client, bucket, key, None, Some(current), b"enabled-again").await; + assert_read_version(&client, bucket, key, Some("null"), Some("null"), b"null-overwritten").await; + for (version, body) in &history { + assert_read_version(&client, bucket, key, Some(version), Some(version), body).await; + } + let listed = client + .list_object_versions() + .bucket(bucket) + .prefix(key) + .send() + .await + .expect("list history"); + let mut versions: Vec<_> = listed + .versions() + .iter() + .map(|version| version.version_id().expect("listed identity")) + .collect(); + versions.sort_unstable(); + let mut expected = vec!["null", history[0].0.as_str(), history[1].0.as_str(), current]; + expected.sort_unstable(); + assert_eq!(versions, expected, "one null slot and all UUID versions must remain"); + + server.shutdown().await; +}