feat(s3): support metadata extensions for bucket listings (#2344)

This commit is contained in:
weisd
2026-03-30 22:03:13 +08:00
committed by GitHub
parent 172086ff42
commit 0fb070e912
9 changed files with 963 additions and 13 deletions
Generated
+10 -10
View File
@@ -3136,7 +3136,7 @@ dependencies = [
"libc",
"option-ext",
"redox_users",
"windows-sys 0.61.2",
"windows-sys 0.59.0",
]
[[package]]
@@ -3412,7 +3412,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb"
dependencies = [
"libc",
"windows-sys 0.61.2",
"windows-sys 0.52.0",
]
[[package]]
@@ -4667,7 +4667,7 @@ checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46"
dependencies = [
"hermit-abi",
"libc",
"windows-sys 0.61.2",
"windows-sys 0.52.0",
]
[[package]]
@@ -4744,7 +4744,7 @@ dependencies = [
"portable-atomic",
"portable-atomic-util",
"serde_core",
"windows-sys 0.61.2",
"windows-sys 0.52.0",
]
[[package]]
@@ -5620,7 +5620,7 @@ version = "0.50.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5"
dependencies = [
"windows-sys 0.61.2",
"windows-sys 0.59.0",
]
[[package]]
@@ -8567,7 +8567,7 @@ dependencies = [
"errno",
"libc",
"linux-raw-sys 0.12.1",
"windows-sys 0.61.2",
"windows-sys 0.52.0",
]
[[package]]
@@ -8635,7 +8635,7 @@ dependencies = [
"security-framework",
"security-framework-sys",
"webpki-root-certs",
"windows-sys 0.61.2",
"windows-sys 0.52.0",
]
[[package]]
@@ -8682,7 +8682,7 @@ checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f"
[[package]]
name = "s3s"
version = "0.14.0-dev"
source = "git+https://github.com/rustfs/s3s?rev=b296762bc9e7fa608f1bc44f5cd625d606e0dd31#b296762bc9e7fa608f1bc44f5cd625d606e0dd31"
source = "git+https://github.com/rustfs/s3s?rev=f1815ced732e180f71935feee6ae5ef44fe39b22#f1815ced732e180f71935feee6ae5ef44fe39b22"
dependencies = [
"arc-swap",
"arrayvec",
@@ -9640,7 +9640,7 @@ dependencies = [
"getrandom 0.4.2",
"once_cell",
"rustix 1.1.4",
"windows-sys 0.61.2",
"windows-sys 0.52.0",
]
[[package]]
@@ -10599,7 +10599,7 @@ version = "0.1.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22"
dependencies = [
"windows-sys 0.61.2",
"windows-sys 0.52.0",
]
[[package]]
+1 -1
View File
@@ -250,7 +250,7 @@ rumqttc = { version = "0.25.1" }
rustix = { version = "1.1.4", features = ["fs"] }
rust-embed = { version = "8.11.0" }
rustc-hash = { version = "2.1.2" }
s3s = { git = "https://github.com/rustfs/s3s", rev = "b296762bc9e7fa608f1bc44f5cd625d606e0dd31", features = ["minio"] }
s3s = { git = "https://github.com/rustfs/s3s", rev = "f1815ced732e180f71935feee6ae5ef44fe39b22", features = ["minio"] }
serial_test = "3.4.0"
shadow-rs = { version = "1.7.1", default-features = false }
siphasher = "1.0.2"
+8
View File
@@ -75,6 +75,14 @@ mod delete_objects_versioning_test;
#[cfg(test)]
mod list_object_versions_regression_test;
// versions&metadata=true extension regression test
#[cfg(test)]
mod list_object_versions_metadata_extension_test;
// list-type=2&metadata=true extension regression test
#[cfg(test)]
mod list_objects_v2_metadata_extension_test;
#[cfg(test)]
mod protocols;
@@ -0,0 +1,122 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! End-to-end regression test for the versions metadata extension:
//! `GET /{bucket}?versions&metadata=true`
#[cfg(test)]
mod tests {
use crate::common::{RustFSTestEnvironment, init_logging, local_http_client};
use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration};
use http::header::HOST;
use reqwest::StatusCode;
use rustfs_signer::constants::UNSIGNED_PAYLOAD;
use rustfs_signer::sign_v4;
use s3s::Body;
use serial_test::serial;
use std::error::Error;
use tracing::info;
async fn signed_get(
url: &str,
access_key: &str,
secret_key: &str,
) -> Result<reqwest::Response, Box<dyn Error + Send + Sync>> {
let uri = url.parse::<http::Uri>()?;
let authority = uri.authority().ok_or("request URL missing authority")?.to_string();
let request = http::Request::builder()
.method(http::Method::GET)
.uri(uri)
.header(HOST, authority)
.header("x-amz-content-sha256", UNSIGNED_PAYLOAD)
.body(Body::empty())?;
let signed = sign_v4(request, 0, access_key, secret_key, "", "us-east-1");
let client = local_http_client();
let mut request_builder = client.get(url);
for (name, value) in signed.headers() {
request_builder = request_builder.header(name, value);
}
Ok(request_builder.send().await?)
}
#[tokio::test]
#[serial]
async fn test_list_object_versions_metadata_extension_returns_metadata_tags_and_internal()
-> Result<(), Box<dyn Error + Send + Sync>> {
init_logging();
info!("🧪 TEST: versions&metadata=true returns versions metadata extensions");
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let client = env.create_s3_client();
let bucket = "test-list-versions-metadata-ext";
let key = "versions/metadata-object.txt";
client.create_bucket().bucket(bucket).send().await?;
client
.put_bucket_versioning()
.bucket(bucket)
.versioning_configuration(
VersioningConfiguration::builder()
.status(BucketVersioningStatus::Enabled)
.build(),
)
.send()
.await?;
client
.put_object()
.bucket(bucket)
.key(key)
.metadata("project", "alpha")
.metadata("owner", "ops")
.tagging("env=test&project=alpha")
.body(aws_sdk_s3::primitives::ByteStream::from_static(b"metadata extension body"))
.send()
.await?;
let url = format!("{}/{}?versions&metadata=true&prefix={}", env.url, bucket, urlencoding::encode(key));
let response = signed_get(&url, &env.access_key, &env.secret_key).await?;
assert_eq!(response.status(), StatusCode::OK, "versions&metadata=true should succeed");
let body = response.text().await?;
info!("versions&metadata=true response body: {}", body);
assert!(body.contains("<ListVersionsResult"), "expected ListVersionsResult XML, got: {body}");
assert!(body.contains("<Version>"), "expected at least one Version entry, got: {body}");
assert!(body.contains("<UserMetadata>"), "expected UserMetadata extension, got: {body}");
assert!(
body.contains("<project>alpha</project>"),
"expected stripped user metadata key in response, got: {body}"
);
assert!(
body.contains("<owner>ops</owner>"),
"expected second metadata key in response, got: {body}"
);
assert!(
body.contains("<UserTags>env=test&amp;project=alpha</UserTags>"),
"expected UserTags extension in response, got: {body}"
);
assert!(body.contains("<Internal>"), "expected Internal extension, got: {body}");
assert!(body.contains("<K>1</K>"), "expected Internal/K in response, got: {body}");
assert!(body.contains("<M>0</M>"), "expected Internal/M in response, got: {body}");
Ok(())
}
}
@@ -0,0 +1,110 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! End-to-end regression test for the objects metadata extension:
//! `GET /{bucket}?list-type=2&metadata=true`
#[cfg(test)]
mod tests {
use crate::common::{RustFSTestEnvironment, init_logging, local_http_client};
use http::header::HOST;
use reqwest::StatusCode;
use rustfs_signer::constants::UNSIGNED_PAYLOAD;
use rustfs_signer::sign_v4;
use s3s::Body;
use serial_test::serial;
use std::error::Error;
use tracing::info;
async fn signed_get(
url: &str,
access_key: &str,
secret_key: &str,
) -> Result<reqwest::Response, Box<dyn Error + Send + Sync>> {
let uri = url.parse::<http::Uri>()?;
let authority = uri.authority().ok_or("request URL missing authority")?.to_string();
let request = http::Request::builder()
.method(http::Method::GET)
.uri(uri)
.header(HOST, authority)
.header("x-amz-content-sha256", UNSIGNED_PAYLOAD)
.body(Body::empty())?;
let signed = sign_v4(request, 0, access_key, secret_key, "", "us-east-1");
let client = local_http_client();
let mut request_builder = client.get(url);
for (name, value) in signed.headers() {
request_builder = request_builder.header(name, value);
}
Ok(request_builder.send().await?)
}
#[tokio::test]
#[serial]
async fn test_list_objects_v2_metadata_extension_returns_metadata_tags_and_internal()
-> Result<(), Box<dyn Error + Send + Sync>> {
init_logging();
info!("🧪 TEST: list-type=2&metadata=true returns object metadata extensions");
let mut env = RustFSTestEnvironment::new().await?;
env.start_rustfs_server(vec![]).await?;
let client = env.create_s3_client();
let bucket = "test-list-objects-v2-metadata-ext";
let key = "objects/metadata-object.txt";
client.create_bucket().bucket(bucket).send().await?;
client
.put_object()
.bucket(bucket)
.key(key)
.metadata("project", "alpha")
.metadata("owner", "ops")
.tagging("env=test&project=alpha")
.body(aws_sdk_s3::primitives::ByteStream::from_static(b"metadata extension body"))
.send()
.await?;
let url = format!("{}/{}?list-type=2&metadata=true&prefix={}", env.url, bucket, urlencoding::encode(key));
let response = signed_get(&url, &env.access_key, &env.secret_key).await?;
assert_eq!(response.status(), StatusCode::OK, "list-type=2&metadata=true should succeed");
let body = response.text().await?;
info!("list-type=2&metadata=true response body: {}", body);
assert!(body.contains("<ListBucketResult"), "expected ListBucketResult XML, got: {body}");
assert!(body.contains("<Contents>"), "expected at least one Contents entry, got: {body}");
assert!(body.contains("<UserMetadata>"), "expected UserMetadata extension, got: {body}");
assert!(
body.contains("<project>alpha</project>"),
"expected stripped user metadata key in response, got: {body}"
);
assert!(
body.contains("<owner>ops</owner>"),
"expected second metadata key in response, got: {body}"
);
assert!(
body.contains("<UserTags>env=test&amp;project=alpha</UserTags>"),
"expected UserTags extension in response, got: {body}"
);
assert!(body.contains("<Internal>"), "expected Internal extension, got: {body}");
assert!(body.contains("<K>1</K>"), "expected Internal/K in response, got: {body}");
assert!(body.contains("<M>0</M>"), "expected Internal/M in response, got: {body}");
Ok(())
}
}
+676 -2
View File
@@ -21,6 +21,7 @@ use crate::server::RemoteAddr;
use crate::storage::access::{ReqInfo, authorize_request, req_info_ref};
use crate::storage::helper::OperationHelper;
use crate::storage::s3_api::bucket::{build_list_buckets_output, build_list_objects_v2_output};
use crate::storage::s3_api::common::rustfs_owner;
use crate::storage::s3_api::{acl, encryption, replication, tagging};
use crate::storage::*;
use futures::StreamExt;
@@ -44,7 +45,10 @@ use rustfs_ecstore::bucket::{
use rustfs_ecstore::client::object_api_utils::to_s3s_etag;
use rustfs_ecstore::error::StorageError;
use rustfs_ecstore::new_object_layer_fn;
use rustfs_ecstore::store_api::{BucketOperations, BucketOptions, DeleteBucketOptions, ListOperations, MakeBucketOptions};
use rustfs_ecstore::store_api::{
BucketOperations, BucketOptions, DeleteBucketOptions, ListObjectVersionsInfo, ListObjectsV2Info, ListOperations,
MakeBucketOptions, ObjectInfo,
};
use rustfs_policy::policy::{
action::{Action, S3Action},
{BucketPolicy, BucketPolicyArgs, Effect, Validator},
@@ -55,13 +59,19 @@ use rustfs_targets::{
arn::{ARN, TargetIDError},
};
use rustfs_utils::http::{SUFFIX_FORCE_DELETE, get_header};
use rustfs_utils::obj::extract_user_defined_metadata;
use rustfs_utils::string::parse_bool;
use s3s::dto::*;
use s3s::region::Region;
use s3s::xml;
use s3s::{S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error};
use std::{collections::HashSet, fmt::Display, sync::Arc};
use std::{
collections::{HashMap, HashSet},
fmt::Display,
sync::Arc,
};
use tracing::{debug, error, info, instrument, warn};
use urlencoding::encode;
fn serialize_config<T: xml::Serialize>(value: &T) -> S3Result<Vec<u8>> {
serialize(value).map_err(to_internal_error)
@@ -97,6 +107,316 @@ async fn validate_bucket_versioning_update(bucket: &str, config: &VersioningConf
Ok(())
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
struct ObjectMetadataPermissions {
metadata_allowed: bool,
tags_allowed: bool,
}
#[derive(Debug, Clone)]
struct ListObjectVersionsMResponseContext {
bucket: String,
prefix: String,
delimiter: Option<String>,
max_keys: i32,
encoding_type: Option<EncodingType>,
key_marker: Option<String>,
version_id_marker: Option<String>,
}
#[derive(Debug, Clone)]
struct ListObjectsV2MResponseContext {
bucket: String,
prefix: String,
delimiter: Option<String>,
max_keys: i32,
encoding_type: Option<EncodingType>,
continuation_token: Option<String>,
start_after: Option<String>,
fetch_owner: bool,
}
fn encode_list_versions_value(value: &str, encoding_type: Option<&EncodingType>) -> String {
if encoding_type.is_some_and(|encoding| encoding.as_str() == EncodingType::URL) {
encode(value).into_owned()
} else {
value.to_string()
}
}
fn encode_list_objects_v2_value(value: &str, encoding_type: Option<&EncodingType>) -> String {
if encoding_type.is_some_and(|encoding| encoding.as_str() == EncodingType::URL) {
value
.split('/')
.map(|part| encode(part).into_owned())
.collect::<Vec<_>>()
.join("/")
} else {
value.to_string()
}
}
fn build_metadata_extension_user_metadata(user_defined: &HashMap<String, String>) -> Option<MinioUserMetadata> {
let mut items = extract_user_defined_metadata(user_defined)
.into_iter()
.filter(|(key, _)| !key.is_empty())
.map(|(key, value)| MinioMetadataEntry { key, value })
.collect::<Vec<_>>();
items.sort_by(|left, right| left.key.cmp(&right.key));
if items.is_empty() {
None
} else {
Some(MinioUserMetadata { items })
}
}
async fn is_list_objects_metadata_action_allowed<T>(
req: &S3Request<T>,
bucket: &str,
object: &str,
action: S3Action,
) -> S3Result<bool> {
let mut auth_req = S3Request {
input: (),
method: req.method.clone(),
uri: req.uri.clone(),
headers: req.headers.clone(),
extensions: req.extensions.clone(),
credentials: req.credentials.clone(),
region: req.region.clone(),
service: req.service.clone(),
trailing_headers: req.trailing_headers.clone(),
};
let mut req_info = req_info_ref(req)?.clone();
req_info.bucket = Some(bucket.to_string());
req_info.object = Some(object.to_string());
req_info.version_id = None;
auth_req.extensions.insert(req_info);
match authorize_request(&mut auth_req, Action::S3Action(action)).await {
Ok(()) => Ok(true),
Err(err) if err.code() == &S3ErrorCode::AccessDenied => Ok(false),
Err(err) => Err(err),
}
}
async fn collect_list_objects_metadata_permissions<T>(
req: &S3Request<T>,
bucket: &str,
objects: &[ObjectInfo],
) -> S3Result<HashMap<String, ObjectMetadataPermissions>> {
let mut permissions = HashMap::new();
for object in objects {
if object.name.is_empty() || permissions.contains_key(&object.name) {
continue;
}
let metadata_allowed =
is_list_objects_metadata_action_allowed(req, bucket, &object.name, S3Action::GetObjectAction).await?;
let tags_allowed =
is_list_objects_metadata_action_allowed(req, bucket, &object.name, S3Action::GetObjectTaggingAction).await?;
permissions.insert(
object.name.clone(),
ObjectMetadataPermissions {
metadata_allowed,
tags_allowed,
},
);
}
Ok(permissions)
}
fn build_list_object_versions_m_output(
object_infos: ListObjectVersionsInfo,
context: &ListObjectVersionsMResponseContext,
permissions: &HashMap<String, ObjectMetadataPermissions>,
) -> ListObjectVersionsMOutput {
let owner = rustfs_owner();
let common_prefixes = object_infos
.prefixes
.into_iter()
.map(|prefix_value| CommonPrefix {
prefix: Some(encode_list_versions_value(&prefix_value, context.encoding_type.as_ref())),
})
.collect::<Vec<_>>();
let entries = object_infos
.objects
.into_iter()
.filter(|object| !object.name.is_empty())
.map(|object| {
let object_name = encode_list_versions_value(&object.name, context.encoding_type.as_ref());
let version_id = object
.version_id
.map(|version| version.to_string())
.unwrap_or_else(|| "null".to_string());
let permission = permissions.get(&object.name).copied().unwrap_or_default();
let user_metadata = if permission.metadata_allowed {
build_metadata_extension_user_metadata(&object.user_defined)
} else {
None
};
let user_tags = if permission.tags_allowed && !object.user_tags.is_empty() {
Some(object.user_tags.clone())
} else {
None
};
let internal = if permission.metadata_allowed && (object.data_blocks > 0 || object.parity_blocks > 0) {
Some(ObjectInternalInfo {
k: object.data_blocks as i32,
m: object.parity_blocks as i32,
})
} else {
None
};
if object.delete_marker {
ListObjectVersionMEntry::DeleteMarker(DeleteMarkerM {
key: Some(object_name),
last_modified: object.mod_time.map(Timestamp::from),
owner: Some(owner.clone()),
version_id: Some(version_id),
is_latest: Some(object.is_latest),
user_metadata,
user_tags,
internal,
})
} else {
ListObjectVersionMEntry::Version(ObjectVersionM {
key: Some(object_name),
last_modified: object.mod_time.map(Timestamp::from),
size: Some(object.size),
version_id: Some(version_id),
is_latest: Some(object.is_latest),
e_tag: object.etag.clone().map(|etag| to_s3s_etag(&etag)),
storage_class: Some(ObjectVersionStorageClass::from(
object
.storage_class
.unwrap_or_else(|| ObjectVersionStorageClass::STANDARD.to_string()),
)),
owner: Some(owner.clone()),
user_metadata,
user_tags,
internal,
})
}
})
.collect::<Vec<_>>();
let next_key_marker = object_infos
.next_marker
.filter(|marker| !marker.is_empty())
.map(|marker| encode_list_versions_value(&marker, context.encoding_type.as_ref()));
ListObjectVersionsMOutput {
common_prefixes: Some(common_prefixes),
delimiter: context
.delimiter
.clone()
.map(|value| encode_list_versions_value(&value, context.encoding_type.as_ref())),
encoding_type: context.encoding_type.clone(),
is_truncated: Some(object_infos.is_truncated),
key_marker: Some(encode_list_versions_value(
context.key_marker.as_deref().unwrap_or_default(),
context.encoding_type.as_ref(),
)),
max_keys: Some(context.max_keys),
name: Some(context.bucket.clone()),
next_key_marker,
next_version_id_marker: Some(object_infos.next_version_idmarker.unwrap_or_default()),
prefix: Some(encode_list_versions_value(&context.prefix, context.encoding_type.as_ref())),
request_charged: None,
version_id_marker: Some(context.version_id_marker.clone().unwrap_or_default()),
entries,
}
}
fn build_list_objects_v2m_output(
object_infos: ListObjectsV2Info,
context: &ListObjectsV2MResponseContext,
permissions: &HashMap<String, ObjectMetadataPermissions>,
) -> ListObjectsV2MOutput {
let owner = rustfs_owner();
let contents = object_infos
.objects
.iter()
.filter(|object| !object.name.is_empty())
.map(|object| {
let permission = permissions.get(&object.name).copied().unwrap_or_default();
let user_metadata = if permission.metadata_allowed {
build_metadata_extension_user_metadata(&object.user_defined)
} else {
None
};
let user_tags = if permission.tags_allowed && !object.user_tags.is_empty() {
Some(object.user_tags.clone())
} else {
None
};
let internal = if permission.metadata_allowed && (object.data_blocks > 0 || object.parity_blocks > 0) {
Some(ObjectInternalInfo {
k: object.data_blocks as i32,
m: object.parity_blocks as i32,
})
} else {
None
};
ObjectM {
key: Some(encode_list_objects_v2_value(&object.name, context.encoding_type.as_ref())),
last_modified: object.mod_time.map(Timestamp::from),
size: Some(object.get_actual_size().unwrap_or_default()),
e_tag: object.etag.clone().map(|etag| to_s3s_etag(&etag)),
storage_class: Some(ObjectStorageClass::from(
object
.storage_class
.clone()
.unwrap_or_else(|| ObjectStorageClass::STANDARD.to_string()),
)),
owner: context.fetch_owner.then_some(owner.clone()),
user_metadata,
user_tags,
internal,
}
})
.collect::<Vec<_>>();
let common_prefixes = object_infos
.prefixes
.into_iter()
.map(|prefix| CommonPrefix {
prefix: Some(encode_list_objects_v2_value(&prefix, context.encoding_type.as_ref())),
})
.collect::<Vec<_>>();
let key_count = (contents.len() + common_prefixes.len()) as i32;
let next_continuation_token = object_infos
.next_continuation_token
.map(|token| base64_simd::STANDARD.encode_to_string(token.as_bytes()));
ListObjectsV2MOutput {
name: Some(context.bucket.clone()),
prefix: Some(context.prefix.clone()),
max_keys: Some(context.max_keys),
key_count: Some(key_count),
continuation_token: context.continuation_token.clone(),
is_truncated: Some(object_infos.is_truncated),
next_continuation_token,
contents: Some(contents),
common_prefixes: Some(common_prefixes),
delimiter: context.delimiter.clone(),
encoding_type: context.encoding_type.clone(),
start_after: context.start_after.clone(),
..Default::default()
}
}
fn create_bucket_exists_response(is_owner: bool) -> S3Result<S3Response<CreateBucketOutput>> {
if is_owner {
return Ok(S3Response::new(CreateBucketOutput::default()));
@@ -1487,6 +1807,87 @@ impl DefaultBucketUsecase {
Ok(S3Response::new(output))
}
pub async fn execute_list_objects_v2m(
&self,
req: S3Request<ListObjectsV2Input>,
) -> S3Result<S3Response<ListObjectsV2MOutput>> {
if let Some(context) = &self.context {
let _ = context.object_store();
}
let input = req.input.clone();
let ListObjectsV2Input {
bucket,
continuation_token,
delimiter,
encoding_type,
fetch_owner,
max_keys,
prefix,
start_after,
..
} = input;
let prefix = prefix.unwrap_or_default();
let max_keys = max_keys.unwrap_or(1000);
if max_keys < 0 {
return Err(S3Error::with_message(S3ErrorCode::InvalidArgument, "Invalid max keys".to_string()));
}
let delimiter = delimiter.filter(|value| !value.is_empty());
validate_list_object_unordered_with_delimiter(delimiter.as_ref(), req.uri.query())?;
let response_start_after = start_after.clone();
let start_after_for_query = start_after.filter(|value| !value.is_empty());
let response_continuation_token = continuation_token.clone();
let continuation_token_for_query = continuation_token.filter(|value| !value.is_empty());
let decoded_continuation_token = continuation_token_for_query
.map(|token| {
base64_simd::STANDARD
.decode_to_vec(token.as_bytes())
.map_err(|_| s3_error!(InvalidArgument, "Invalid continuation token"))
.and_then(|bytes| {
String::from_utf8(bytes).map_err(|_| s3_error!(InvalidArgument, "Invalid continuation token"))
})
})
.transpose()?;
let store = get_validated_store(&bucket).await?;
let incl_deleted = rustfs_utils::http::get_header(&req.headers, rustfs_utils::http::SUFFIX_INCLUDE_DELETED)
.map(|value| value.as_ref() == "true")
.unwrap_or_default();
let object_infos = store
.list_objects_v2(
&bucket,
&prefix,
decoded_continuation_token,
delimiter.clone(),
max_keys,
fetch_owner.unwrap_or_default(),
start_after_for_query,
incl_deleted,
)
.await
.map_err(ApiError::from)?;
let permissions = collect_list_objects_metadata_permissions(&req, &bucket, &object_infos.objects).await?;
let context = ListObjectsV2MResponseContext {
bucket,
prefix,
delimiter,
max_keys,
encoding_type,
continuation_token: response_continuation_token,
start_after: response_start_after,
fetch_owner: fetch_owner.unwrap_or_default(),
};
let output = build_list_objects_v2m_output(object_infos, &context, &permissions);
Ok(S3Response::new(output))
}
pub async fn execute_list_object_versions(
&self,
req: S3Request<ListObjectVersionsInput>,
@@ -1574,6 +1975,60 @@ impl DefaultBucketUsecase {
Ok(S3Response::new(output))
}
pub async fn execute_list_object_versions_m(
&self,
req: S3Request<ListObjectVersionsInput>,
) -> S3Result<S3Response<ListObjectVersionsMOutput>> {
if let Some(context) = &self.context {
let _ = context.object_store();
}
let input = req.input.clone();
let ListObjectVersionsInput {
bucket,
delimiter,
encoding_type,
key_marker,
version_id_marker,
max_keys,
prefix,
..
} = input;
let prefix = prefix.unwrap_or_default();
let max_keys = max_keys.unwrap_or(1000);
let key_marker = key_marker.filter(|value| !value.is_empty());
let version_id_marker = version_id_marker.filter(|value| !value.is_empty());
let delimiter = delimiter.filter(|value| !value.is_empty());
let store = get_validated_store(&bucket).await?;
let object_infos = store
.list_object_versions(
&bucket,
&prefix,
key_marker.clone(),
version_id_marker.clone(),
delimiter.clone(),
max_keys,
)
.await
.map_err(ApiError::from)?;
let permissions = collect_list_objects_metadata_permissions(&req, &bucket, &object_infos.objects).await?;
let context = ListObjectVersionsMResponseContext {
bucket,
prefix,
delimiter,
max_keys,
encoding_type,
key_marker,
version_id_marker,
};
let output = build_list_object_versions_m_output(object_infos, &context, &permissions);
Ok(S3Response::new(output))
}
#[instrument(level = "debug", skip(self, req))]
pub async fn execute_list_objects(&self, req: S3Request<ListObjectsInput>) -> S3Result<S3Response<ListObjectsOutput>> {
if let Some(context) = &self.context {
@@ -2071,6 +2526,126 @@ mod tests {
assert_eq!(err.code(), &S3ErrorCode::InternalError);
}
#[tokio::test]
async fn execute_list_object_versions_m_returns_internal_error_when_store_uninitialized() {
let input = ListObjectVersionsInput::builder()
.bucket("test-bucket".to_string())
.build()
.unwrap();
let req = build_request(input, Method::GET);
let usecase = DefaultBucketUsecase::without_context();
let err = usecase.execute_list_object_versions_m(req).await.unwrap_err();
assert_eq!(err.code(), &S3ErrorCode::InternalError);
}
#[test]
fn build_list_object_versions_m_output_maps_metadata_and_preserves_entry_order() {
use time::macros::datetime;
use uuid::Uuid;
let object_infos = ListObjectVersionsInfo {
is_truncated: true,
next_marker: Some("obj-z".to_string()),
next_version_idmarker: Some("null".to_string()),
prefixes: vec!["logs/".to_string()],
objects: vec![
ObjectInfo {
bucket: "demo-bucket".to_string(),
name: "obj-a".to_string(),
mod_time: Some(datetime!(2025-01-01 00:00 UTC)),
size: 11,
user_defined: HashMap::from([("project".to_string(), "alpha".to_string())]),
parity_blocks: 2,
data_blocks: 4,
version_id: Some(Uuid::nil()),
user_tags: "env=prod".to_string(),
is_latest: true,
etag: Some("0123456789abcdef0123456789abcdef".to_string()),
..Default::default()
},
ObjectInfo {
bucket: "demo-bucket".to_string(),
name: "obj-b".to_string(),
mod_time: Some(datetime!(2025-01-02 00:00 UTC)),
delete_marker: true,
user_defined: HashMap::from([("marker".to_string(), "true".to_string())]),
version_id: None,
..Default::default()
},
],
};
let permissions = HashMap::from([
(
"obj-a".to_string(),
ObjectMetadataPermissions {
metadata_allowed: true,
tags_allowed: true,
},
),
(
"obj-b".to_string(),
ObjectMetadataPermissions {
metadata_allowed: true,
tags_allowed: false,
},
),
]);
let context = ListObjectVersionsMResponseContext {
bucket: "demo-bucket".to_string(),
prefix: "pre".to_string(),
delimiter: Some("/".to_string()),
max_keys: 1000,
encoding_type: Some(EncodingType::from_static(EncodingType::URL)),
key_marker: Some("start marker".to_string()),
version_id_marker: Some("vid-1".to_string()),
};
let output = build_list_object_versions_m_output(object_infos, &context, &permissions);
assert_eq!(output.name.as_deref(), Some("demo-bucket"));
assert_eq!(output.prefix.as_deref(), Some("pre"));
assert_eq!(output.key_marker.as_deref(), Some("start%20marker"));
assert_eq!(output.next_key_marker.as_deref(), Some("obj-z"));
assert_eq!(output.next_version_id_marker.as_deref(), Some("null"));
assert_eq!(output.entries.len(), 2);
match &output.entries[0] {
ListObjectVersionMEntry::Version(version) => {
assert_eq!(version.key.as_deref(), Some("obj-a"));
assert_eq!(version.version_id.as_deref(), Some(Uuid::nil().to_string().as_str()));
assert_eq!(version.user_tags.as_deref(), Some("env=prod"));
assert_eq!(version.internal, Some(ObjectInternalInfo { k: 4, m: 2 }));
assert_eq!(
version.user_metadata.as_ref().map(|metadata| metadata.items.clone()),
Some(vec![MinioMetadataEntry {
key: "project".to_string(),
value: "alpha".to_string(),
}])
);
}
other => panic!("expected version entry, got {other:?}"),
}
match &output.entries[1] {
ListObjectVersionMEntry::DeleteMarker(marker) => {
assert_eq!(marker.key.as_deref(), Some("obj-b"));
assert_eq!(marker.version_id.as_deref(), Some("null"));
assert!(marker.user_tags.is_none());
assert_eq!(
marker.user_metadata.as_ref().map(|metadata| metadata.items.clone()),
Some(vec![MinioMetadataEntry {
key: "marker".to_string(),
value: "true".to_string(),
}])
);
}
other => panic!("expected delete marker entry, got {other:?}"),
}
}
#[tokio::test]
async fn execute_list_objects_returns_internal_error_when_store_uninitialized() {
let input = ListObjectsInput::builder().bucket("test-bucket".to_string()).build().unwrap();
@@ -2082,6 +2657,90 @@ mod tests {
assert_eq!(err.code(), &S3ErrorCode::InternalError);
}
#[tokio::test]
async fn execute_list_objects_v2m_returns_internal_error_when_store_uninitialized() {
let input = ListObjectsV2Input::builder()
.bucket("test-bucket".to_string())
.build()
.unwrap();
let req = build_request(input, Method::GET);
let usecase = DefaultBucketUsecase::without_context();
let err = usecase.execute_list_objects_v2m(req).await.unwrap_err();
assert_eq!(err.code(), &S3ErrorCode::InternalError);
}
#[test]
fn build_list_objects_v2m_output_maps_metadata_and_key_count() {
use time::macros::datetime;
let object_infos = ListObjectsV2Info {
is_truncated: true,
next_continuation_token: Some("next-token".to_string()),
objects: vec![ObjectInfo {
bucket: "demo-bucket".to_string(),
name: "logs/obj a.txt".to_string(),
mod_time: Some(datetime!(2025-01-03 00:00 UTC)),
size: 11,
user_defined: HashMap::from([("project".to_string(), "alpha".to_string())]),
parity_blocks: 2,
data_blocks: 4,
user_tags: "env=prod".to_string(),
etag: Some("0123456789abcdef0123456789abcdef".to_string()),
..Default::default()
}],
prefixes: vec!["logs/archive/".to_string()],
..Default::default()
};
let permissions = HashMap::from([(
"logs/obj a.txt".to_string(),
ObjectMetadataPermissions {
metadata_allowed: true,
tags_allowed: true,
},
)]);
let context = ListObjectsV2MResponseContext {
bucket: "demo-bucket".to_string(),
prefix: "logs/".to_string(),
delimiter: Some("/".to_string()),
max_keys: 1000,
encoding_type: Some(EncodingType::from_static(EncodingType::URL)),
continuation_token: Some("start token".to_string()),
start_after: Some("logs/start after".to_string()),
fetch_owner: true,
};
let output = build_list_objects_v2m_output(object_infos, &context, &permissions);
assert_eq!(output.name.as_deref(), Some("demo-bucket"));
assert_eq!(output.prefix.as_deref(), Some("logs/"));
assert_eq!(output.continuation_token.as_deref(), Some("start token"));
assert_eq!(output.start_after.as_deref(), Some("logs/start after"));
assert_eq!(output.next_continuation_token.as_deref(), Some("bmV4dC10b2tlbg=="));
assert_eq!(output.key_count, Some(2));
assert_eq!(output.contents.as_ref().map(Vec::len), Some(1));
assert_eq!(output.common_prefixes.as_ref().map(Vec::len), Some(1));
let object = output.contents.as_ref().unwrap().first().unwrap();
assert_eq!(object.key.as_deref(), Some("logs/obj%20a.txt"));
assert_eq!(object.user_tags.as_deref(), Some("env=prod"));
assert_eq!(object.internal, Some(ObjectInternalInfo { k: 4, m: 2 }));
assert!(object.owner.is_some());
assert_eq!(
object.user_metadata.as_ref().map(|metadata| metadata.items.clone()),
Some(vec![MinioMetadataEntry {
key: "project".to_string(),
value: "alpha".to_string(),
}])
);
let prefix = output.common_prefixes.as_ref().unwrap().first().unwrap();
assert_eq!(prefix.prefix.as_deref(), Some("logs/archive/"));
}
#[tokio::test]
async fn execute_put_bucket_lifecycle_configuration_rejects_missing_configuration() {
let input = PutBucketLifecycleConfigurationInput::builder()
@@ -2187,4 +2846,19 @@ mod tests {
let err = usecase.execute_list_objects_v2(req).await.unwrap_err();
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
}
#[tokio::test]
async fn execute_list_objects_v2m_rejects_negative_max_keys() {
let input = ListObjectsV2Input::builder()
.bucket("test-bucket".to_string())
.max_keys(Some(-1))
.build()
.unwrap();
let req = build_request(input, Method::GET);
let usecase = DefaultBucketUsecase::without_context();
let err = usecase.execute_list_objects_v2m(req).await.unwrap_err();
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument);
}
}
+13
View File
@@ -1449,6 +1449,12 @@ impl S3Access for FS {
authorize_request(req, Action::S3Action(S3Action::ListBucketVersionsAction)).await
}
async fn list_object_versions_m(&self, req: &mut S3Request<ListObjectVersionsInput>) -> S3Result<()> {
let req_info = ext_req_info_mut(&mut req.extensions)?;
req_info.bucket = Some(req.input.bucket.clone());
authorize_request(req, Action::S3Action(S3Action::ListBucketVersionsAction)).await
}
/// Checks whether the ListObjects request has accesses to the resources.
///
/// This method returns `Ok(())` by default.
@@ -1469,6 +1475,13 @@ impl S3Access for FS {
authorize_request(req, Action::S3Action(S3Action::ListBucketAction)).await
}
async fn list_objects_v2m(&self, req: &mut S3Request<ListObjectsV2Input>) -> S3Result<()> {
let req_info = ext_req_info_mut(&mut req.extensions)?;
req_info.bucket = Some(req.input.bucket.clone());
authorize_request(req, Action::S3Action(S3Action::ListBucketAction)).await
}
/// Checks whether the ListParts request has accesses to the resources.
///
/// This method returns `Ok(())` by default.
+16
View File
@@ -596,6 +596,15 @@ impl S3 for FS {
usecase.execute_list_object_versions(req).await
}
async fn list_object_versions_m(
&self,
req: S3Request<ListObjectVersionsInput>,
) -> S3Result<S3Response<ListObjectVersionsMOutput>> {
record_s3_op(S3Operation::ListObjectVersions, &req.input.bucket);
let usecase = DefaultBucketUsecase::from_global();
usecase.execute_list_object_versions_m(req).await
}
#[instrument(level = "debug", skip(self, req))]
async fn list_objects(&self, req: S3Request<ListObjectsInput>) -> S3Result<S3Response<ListObjectsOutput>> {
record_s3_op(S3Operation::ListObjects, &req.input.bucket);
@@ -610,6 +619,13 @@ impl S3 for FS {
usecase.execute_list_objects_v2(req).await
}
#[instrument(level = "debug", skip(self, req))]
async fn list_objects_v2m(&self, req: S3Request<ListObjectsV2Input>) -> S3Result<S3Response<ListObjectsV2MOutput>> {
record_s3_op(S3Operation::ListObjectsV2, &req.input.bucket);
let usecase = DefaultBucketUsecase::from_global();
usecase.execute_list_objects_v2m(req).await
}
#[instrument(level = "debug", skip(self, req))]
async fn list_parts(&self, req: S3Request<ListPartsInput>) -> S3Result<S3Response<ListPartsOutput>> {
record_s3_op(S3Operation::ListParts, &req.input.bucket);
+7
View File
@@ -400,6 +400,13 @@ mod tests {
"usecase.execute_list_objects_v2(req).await",
"list_objects_v2 must delegate to DefaultBucketUsecase::execute_list_objects_v2",
);
assert_delegates_within_method(
src,
"async fn list_objects_v2m(&self, req: S3Request<ListObjectsV2Input>)",
"usecase.execute_list_objects_v2m(req).await",
"list_objects_v2m must delegate to DefaultBucketUsecase::execute_list_objects_v2m",
);
}
#[test]