diff --git a/Cargo.lock b/Cargo.lock index dc355d7cb..07834d903 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -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]] diff --git a/Cargo.toml b/Cargo.toml index d77dd4d78..eb8491d5d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -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" diff --git a/crates/e2e_test/src/lib.rs b/crates/e2e_test/src/lib.rs index 6ac513bcc..be3f3758f 100644 --- a/crates/e2e_test/src/lib.rs +++ b/crates/e2e_test/src/lib.rs @@ -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; diff --git a/crates/e2e_test/src/list_object_versions_metadata_extension_test.rs b/crates/e2e_test/src/list_object_versions_metadata_extension_test.rs new file mode 100644 index 000000000..150c45090 --- /dev/null +++ b/crates/e2e_test/src/list_object_versions_metadata_extension_test.rs @@ -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> { + let uri = url.parse::()?; + 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> { + 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(""), "expected at least one Version entry, got: {body}"); + assert!(body.contains(""), "expected UserMetadata extension, got: {body}"); + assert!( + body.contains("alpha"), + "expected stripped user metadata key in response, got: {body}" + ); + assert!( + body.contains("ops"), + "expected second metadata key in response, got: {body}" + ); + assert!( + body.contains("env=test&project=alpha"), + "expected UserTags extension in response, got: {body}" + ); + assert!(body.contains(""), "expected Internal extension, got: {body}"); + assert!(body.contains("1"), "expected Internal/K in response, got: {body}"); + assert!(body.contains("0"), "expected Internal/M in response, got: {body}"); + + Ok(()) + } +} diff --git a/crates/e2e_test/src/list_objects_v2_metadata_extension_test.rs b/crates/e2e_test/src/list_objects_v2_metadata_extension_test.rs new file mode 100644 index 000000000..72be61b6c --- /dev/null +++ b/crates/e2e_test/src/list_objects_v2_metadata_extension_test.rs @@ -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> { + let uri = url.parse::()?; + 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> { + 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(""), "expected at least one Contents entry, got: {body}"); + assert!(body.contains(""), "expected UserMetadata extension, got: {body}"); + assert!( + body.contains("alpha"), + "expected stripped user metadata key in response, got: {body}" + ); + assert!( + body.contains("ops"), + "expected second metadata key in response, got: {body}" + ); + assert!( + body.contains("env=test&project=alpha"), + "expected UserTags extension in response, got: {body}" + ); + assert!(body.contains(""), "expected Internal extension, got: {body}"); + assert!(body.contains("1"), "expected Internal/K in response, got: {body}"); + assert!(body.contains("0"), "expected Internal/M in response, got: {body}"); + + Ok(()) + } +} diff --git a/rustfs/src/app/bucket_usecase.rs b/rustfs/src/app/bucket_usecase.rs index bf8f041e5..323563fbe 100644 --- a/rustfs/src/app/bucket_usecase.rs +++ b/rustfs/src/app/bucket_usecase.rs @@ -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(value: &T) -> S3Result> { 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, + max_keys: i32, + encoding_type: Option, + key_marker: Option, + version_id_marker: Option, +} + +#[derive(Debug, Clone)] +struct ListObjectsV2MResponseContext { + bucket: String, + prefix: String, + delimiter: Option, + max_keys: i32, + encoding_type: Option, + continuation_token: Option, + start_after: Option, + 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::>() + .join("/") + } else { + value.to_string() + } +} + +fn build_metadata_extension_user_metadata(user_defined: &HashMap) -> Option { + let mut items = extract_user_defined_metadata(user_defined) + .into_iter() + .filter(|(key, _)| !key.is_empty()) + .map(|(key, value)| MinioMetadataEntry { key, value }) + .collect::>(); + 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( + req: &S3Request, + bucket: &str, + object: &str, + action: S3Action, +) -> S3Result { + 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( + req: &S3Request, + bucket: &str, + objects: &[ObjectInfo], +) -> S3Result> { + 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, +) -> 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::>(); + + 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::>(); + + 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, +) -> 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::>(); + + 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::>(); + + 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> { 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, + ) -> S3Result> { + 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, @@ -1574,6 +1975,60 @@ impl DefaultBucketUsecase { Ok(S3Response::new(output)) } + pub async fn execute_list_object_versions_m( + &self, + req: S3Request, + ) -> S3Result> { + 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) -> S3Result> { 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); + } } diff --git a/rustfs/src/storage/access.rs b/rustfs/src/storage/access.rs index 7cc4045df..a681dda4e 100644 --- a/rustfs/src/storage/access.rs +++ b/rustfs/src/storage/access.rs @@ -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) -> 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) -> 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. diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index c3b72d34c..28679492f 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -596,6 +596,15 @@ impl S3 for FS { usecase.execute_list_object_versions(req).await } + async fn list_object_versions_m( + &self, + req: S3Request, + ) -> S3Result> { + 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) -> S3Result> { 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) -> S3Result> { + 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) -> S3Result> { record_s3_op(S3Operation::ListParts, &req.input.bucket); diff --git a/rustfs/src/storage/ecfs_test.rs b/rustfs/src/storage/ecfs_test.rs index f6eed5d4f..f10216aca 100644 --- a/rustfs/src/storage/ecfs_test.rs +++ b/rustfs/src/storage/ecfs_test.rs @@ -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)", + "usecase.execute_list_objects_v2m(req).await", + "list_objects_v2m must delegate to DefaultBucketUsecase::execute_list_objects_v2m", + ); } #[test]