fix(s3): preserve version identities on suspended reads (#7751)

This commit is contained in:
cxymds
2026-09-13 21:44:01 +08:00
committed by GitHub
parent 41770983d6
commit 13093c5afc
4 changed files with 204 additions and 27 deletions
+2 -25
View File
@@ -3620,7 +3620,6 @@ impl DefaultObjectUsecase {
queue_status: &concurrency::IoQueueStatus,
concurrent_requests: usize,
part_number: Option<usize>,
versioned: bool,
lifecycle: GetObjectBodyLifecycle,
resume: F,
) -> S3Result<GetObjectOutputContext>
@@ -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"),
)
+1 -1
View File
@@ -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,
+12
View File
@@ -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<Uuid>) -> Option<String> {
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
+189 -1
View File
@@ -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;
}