mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-08 22:33:22 +00:00
feat(object-lock): complete legal hold enforcement (#2293)
This commit is contained in:
@@ -15,13 +15,15 @@
|
||||
//! Multipart application use-case contracts.
|
||||
|
||||
use crate::app::context::{AppContext, get_global_app_context};
|
||||
use crate::app::object_usecase::{build_put_like_object_lock_metadata, validate_existing_object_lock_for_write};
|
||||
use crate::error::ApiError;
|
||||
use crate::storage::access::has_bypass_governance_header;
|
||||
use crate::storage::concurrency::get_concurrency_manager;
|
||||
use crate::storage::entity;
|
||||
use crate::storage::helper::OperationHelper;
|
||||
use crate::storage::options::{
|
||||
copy_src_opts, extract_metadata, get_complete_multipart_upload_opts, get_content_sha256_with_query, parse_copy_source_range,
|
||||
put_opts,
|
||||
copy_src_opts, extract_metadata, get_complete_multipart_upload_opts, get_content_sha256_with_query, get_opts,
|
||||
parse_copy_source_range, put_opts,
|
||||
};
|
||||
use crate::storage::s3_api::multipart::build_list_parts_output;
|
||||
use crate::storage::*;
|
||||
@@ -53,6 +55,7 @@ use rustfs_utils::http::{
|
||||
headers::{AMZ_DECODED_CONTENT_LENGTH, AMZ_OBJECT_TAGGING},
|
||||
};
|
||||
use s3s::dto::*;
|
||||
use s3s::header::{X_AMZ_OBJECT_LOCK_LEGAL_HOLD, X_AMZ_OBJECT_LOCK_MODE, X_AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE};
|
||||
use s3s::{S3Error, S3ErrorCode, S3Request, S3Response, S3Result, s3_error};
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::str::FromStr;
|
||||
@@ -102,6 +105,13 @@ fn normalize_complete_multipart_parts(parts: Vec<CompletePart>) -> S3Result<Vec<
|
||||
Ok(deduped_reversed)
|
||||
}
|
||||
|
||||
fn has_complete_multipart_object_lock_headers(headers: &HeaderMap) -> bool {
|
||||
headers.contains_key(X_AMZ_OBJECT_LOCK_MODE)
|
||||
|| headers.contains_key(X_AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE)
|
||||
|| headers.contains_key(X_AMZ_OBJECT_LOCK_LEGAL_HOLD)
|
||||
|| has_bypass_governance_header(headers)
|
||||
}
|
||||
|
||||
fn encode_s3_path(path: &str) -> String {
|
||||
path.split('/')
|
||||
.map(|part| encode(part).to_string())
|
||||
@@ -285,12 +295,29 @@ impl DefaultMultipartUsecase {
|
||||
|
||||
let uploaded_parts = normalize_complete_multipart_parts(uploaded_parts_vec)?;
|
||||
|
||||
// TODO: check object lock
|
||||
if has_complete_multipart_object_lock_headers(&req.headers) {
|
||||
return Err(S3Error::with_message(
|
||||
S3ErrorCode::InvalidRequest,
|
||||
"CompleteMultipartUpload does not accept object lock or governance bypass headers.".to_string(),
|
||||
));
|
||||
}
|
||||
|
||||
let Some(store) = new_object_layer_fn() else {
|
||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
||||
};
|
||||
|
||||
let current_opts = get_opts(&bucket, &key, None, None, &req.headers)
|
||||
.await
|
||||
.map_err(ApiError::from)?;
|
||||
match store.get_object_info(&bucket, &key, ¤t_opts).await {
|
||||
Ok(existing_obj_info) => validate_existing_object_lock_for_write(&existing_obj_info)?,
|
||||
Err(err) => {
|
||||
if !is_err_object_not_found(&err) && !is_err_version_not_found(&err) {
|
||||
return Err(ApiError::from(err).into());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TDD: Get multipart info to extract encryption configuration before completing
|
||||
info!(
|
||||
"TDD: Attempting to get multipart info for bucket={}, key={}, upload_id={}",
|
||||
@@ -501,6 +528,9 @@ impl DefaultMultipartUsecase {
|
||||
sse_customer_algorithm,
|
||||
sse_customer_key_md5,
|
||||
ssekms_key_id,
|
||||
object_lock_legal_hold_status,
|
||||
object_lock_mode,
|
||||
object_lock_retain_until_date,
|
||||
..
|
||||
} = req.input.clone();
|
||||
|
||||
@@ -529,6 +559,17 @@ impl DefaultMultipartUsecase {
|
||||
metadata.insert(AMZ_OBJECT_TAGGING.to_owned(), tags);
|
||||
}
|
||||
|
||||
if let Some(object_lock_metadata) = build_put_like_object_lock_metadata(
|
||||
&bucket,
|
||||
object_lock_legal_hold_status,
|
||||
object_lock_mode,
|
||||
object_lock_retain_until_date,
|
||||
)
|
||||
.await?
|
||||
{
|
||||
metadata.extend(object_lock_metadata);
|
||||
}
|
||||
|
||||
let encryption_request = PrepareEncryptionRequest {
|
||||
bucket: &bucket,
|
||||
key: &key,
|
||||
@@ -562,6 +603,18 @@ impl DefaultMultipartUsecase {
|
||||
.await
|
||||
.map_err(ApiError::from)?;
|
||||
|
||||
let current_opts: ObjectOptions = get_opts(&bucket, &key, opts.version_id.clone(), None, &req.headers)
|
||||
.await
|
||||
.map_err(ApiError::from)?;
|
||||
match store.get_object_info(&bucket, &key, ¤t_opts).await {
|
||||
Ok(existing_obj_info) => validate_existing_object_lock_for_write(&existing_obj_info)?,
|
||||
Err(err) => {
|
||||
if !is_err_object_not_found(&err) && !is_err_version_not_found(&err) {
|
||||
return Err(ApiError::from(err).into());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let checksum_type = rustfs_rio::ChecksumType::from_header(&req.headers);
|
||||
if checksum_type.is(rustfs_rio::ChecksumType::INVALID) {
|
||||
return Err(s3_error!(InvalidArgument, "Invalid checksum type"));
|
||||
@@ -1342,6 +1395,36 @@ mod tests {
|
||||
assert_eq!(normalized[0].etag.as_deref(), Some("new"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn execute_complete_multipart_upload_rejects_object_lock_headers() {
|
||||
let multipart_upload = CompletedMultipartUpload {
|
||||
parts: Some(vec![CompletedPart {
|
||||
part_number: Some(1),
|
||||
..Default::default()
|
||||
}]),
|
||||
};
|
||||
|
||||
for (header_name, header_value) in [
|
||||
("x-amz-object-lock-mode", "GOVERNANCE"),
|
||||
("x-amz-object-lock-retain-until-date", "2030-01-01T00:00:00Z"),
|
||||
("x-amz-object-lock-legal-hold", "ON"),
|
||||
("x-amz-bypass-governance-retention", "true"),
|
||||
] {
|
||||
let input = CompleteMultipartUploadInput::builder()
|
||||
.bucket("bucket".to_string())
|
||||
.key("object".to_string())
|
||||
.upload_id("upload-id".to_string())
|
||||
.multipart_upload(Some(multipart_upload.clone()))
|
||||
.build()
|
||||
.unwrap();
|
||||
let mut req = build_request(input, Method::POST);
|
||||
req.headers.insert(header_name, HeaderValue::from_str(header_value).unwrap());
|
||||
|
||||
let err = make_usecase().execute_complete_multipart_upload(req).await.unwrap_err();
|
||||
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest, "header {header_name} should be rejected");
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn execute_list_multipart_uploads_returns_internal_error_when_store_uninitialized() {
|
||||
let input = ListMultipartUploadsInput::builder()
|
||||
|
||||
@@ -51,7 +51,12 @@ use rustfs_ecstore::bucket::{
|
||||
},
|
||||
metadata::{BUCKET_VERSIONING_CONFIG, OBJECT_LOCK_CONFIG},
|
||||
metadata_sys,
|
||||
object_lock::objectlock_sys::{BucketObjectLockSys, check_object_lock_for_deletion, check_retention_for_modification},
|
||||
object_lock::{
|
||||
objectlock::{get_object_legalhold_meta, get_object_retention_meta},
|
||||
objectlock_sys::{
|
||||
BucketObjectLockSys, check_object_lock_for_deletion, check_retention_for_modification, is_retention_active,
|
||||
},
|
||||
},
|
||||
quota::QuotaOperation,
|
||||
replication::{
|
||||
DeletedObjectReplicationInfo, check_replicate_delete, get_must_replicate_options, must_replicate, schedule_replication,
|
||||
@@ -499,8 +504,28 @@ async fn apply_put_request_object_lock_opts(
|
||||
object_lock_retain_until_date: Option<Timestamp>,
|
||||
opts: &mut ObjectOptions,
|
||||
) -> S3Result<()> {
|
||||
if let Some(eval_metadata) = build_put_like_object_lock_metadata(
|
||||
bucket,
|
||||
object_lock_legal_hold_status,
|
||||
object_lock_mode,
|
||||
object_lock_retain_until_date,
|
||||
)
|
||||
.await?
|
||||
{
|
||||
opts.eval_metadata = Some(eval_metadata);
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn build_put_like_object_lock_metadata(
|
||||
bucket: &str,
|
||||
object_lock_legal_hold_status: Option<ObjectLockLegalHoldStatus>,
|
||||
object_lock_mode: Option<ObjectLockMode>,
|
||||
object_lock_retain_until_date: Option<Timestamp>,
|
||||
) -> S3Result<Option<HashMap<String, String>>> {
|
||||
if object_lock_legal_hold_status.is_none() && object_lock_mode.is_none() && object_lock_retain_until_date.is_none() {
|
||||
return Ok(());
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
validate_bucket_object_lock_enabled(bucket).await?;
|
||||
@@ -522,13 +547,44 @@ async fn apply_put_request_object_lock_opts(
|
||||
object_lock_legal_hold_status.map(|status| ObjectLockLegalHold { status: Some(status) }),
|
||||
)?);
|
||||
|
||||
if !eval_metadata.is_empty() {
|
||||
opts.eval_metadata = Some(eval_metadata);
|
||||
if eval_metadata.is_empty() {
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
Ok(Some(eval_metadata))
|
||||
}
|
||||
|
||||
pub(crate) fn validate_existing_object_lock_for_write(existing_obj_info: &ObjectInfo) -> S3Result<()> {
|
||||
let legal_hold = get_object_legalhold_meta(&existing_obj_info.user_defined);
|
||||
if legal_hold
|
||||
.status
|
||||
.as_ref()
|
||||
.is_some_and(|status| status.as_str() == ObjectLockLegalHoldStatus::ON)
|
||||
{
|
||||
return Err(S3Error::with_message(
|
||||
S3ErrorCode::AccessDenied,
|
||||
"Object has a legal hold and cannot be overwritten. Remove the legal hold first.".to_string(),
|
||||
));
|
||||
}
|
||||
|
||||
let retention = get_object_retention_meta(&existing_obj_info.user_defined);
|
||||
if let Some(mode) = retention.mode.as_ref()
|
||||
&& mode.as_str() == ObjectLockRetentionMode::COMPLIANCE
|
||||
&& is_retention_active(mode.as_str(), retention.retain_until_date.as_ref())
|
||||
{
|
||||
return Err(S3Error::with_message(
|
||||
S3ErrorCode::AccessDenied,
|
||||
"Object is under COMPLIANCE retention and cannot be overwritten.".to_string(),
|
||||
));
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn delete_creates_delete_marker(opts: &ObjectOptions) -> bool {
|
||||
opts.version_id.is_none() && opts.versioned && !opts.version_suspended
|
||||
}
|
||||
|
||||
fn resolve_put_object_extract_options(headers: &HeaderMap) -> PutObjectExtractOptions {
|
||||
let prefix = snowball_meta_value_by_suffix(headers, AMZ_SNOWBALL_PREFIX_INTERNAL, SNOWBALL_PREFIX_SUFFIX_LOWER)
|
||||
.and_then(|value| normalize_snowball_prefix(&value));
|
||||
@@ -811,6 +867,18 @@ impl DefaultObjectUsecase {
|
||||
)
|
||||
.await?;
|
||||
|
||||
let current_opts: ObjectOptions = get_opts(&bucket, &key, version_id.clone(), None, &req.headers)
|
||||
.await
|
||||
.map_err(ApiError::from)?;
|
||||
match store.get_object_info(&bucket, &key, ¤t_opts).await {
|
||||
Ok(existing_obj_info) => validate_existing_object_lock_for_write(&existing_obj_info)?,
|
||||
Err(err) => {
|
||||
if !is_err_object_not_found(&err) && !is_err_version_not_found(&err) {
|
||||
return Err(ApiError::from(err).into());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let mut reader: Box<dyn Reader> = Box::new(WarpReader::new(body));
|
||||
|
||||
let actual_size = size;
|
||||
@@ -2504,6 +2572,7 @@ impl DefaultObjectUsecase {
|
||||
copy_source,
|
||||
bucket,
|
||||
key,
|
||||
version_id: dest_version_id,
|
||||
server_side_encryption: requested_sse,
|
||||
ssekms_key_id: requested_kms_key_id,
|
||||
sse_customer_algorithm,
|
||||
@@ -2516,6 +2585,9 @@ impl DefaultObjectUsecase {
|
||||
copy_source_if_match,
|
||||
copy_source_if_none_match,
|
||||
content_type,
|
||||
object_lock_legal_hold_status,
|
||||
object_lock_mode,
|
||||
object_lock_retain_until_date,
|
||||
..
|
||||
} = req.input.clone();
|
||||
let (src_bucket, src_key, version_id) = match copy_source {
|
||||
@@ -2551,7 +2623,7 @@ impl DefaultObjectUsecase {
|
||||
|
||||
src_opts.version_id = version_id.clone();
|
||||
|
||||
let mut get_opts = ObjectOptions {
|
||||
let mut src_get_opts = ObjectOptions {
|
||||
version_id: src_opts.version_id.clone(),
|
||||
versioned: src_opts.versioned,
|
||||
version_suspended: src_opts.version_suspended,
|
||||
@@ -2565,13 +2637,25 @@ impl DefaultObjectUsecase {
|
||||
let cp_src_dst_same = path_join_buf(&[&src_bucket, &src_key]) == path_join_buf(&[&bucket, &key]);
|
||||
|
||||
if cp_src_dst_same {
|
||||
get_opts.no_lock = true;
|
||||
src_get_opts.no_lock = true;
|
||||
}
|
||||
|
||||
let Some(store) = new_object_layer_fn() else {
|
||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
||||
};
|
||||
|
||||
let current_opts: ObjectOptions = get_opts(&bucket, &key, dest_version_id.clone(), None, &req.headers)
|
||||
.await
|
||||
.map_err(ApiError::from)?;
|
||||
match store.get_object_info(&bucket, &key, ¤t_opts).await {
|
||||
Ok(existing_obj_info) => validate_existing_object_lock_for_write(&existing_obj_info)?,
|
||||
Err(err) => {
|
||||
if !is_err_object_not_found(&err) && !is_err_version_not_found(&err) {
|
||||
return Err(ApiError::from(err).into());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
let bucket_sse_config = metadata_sys::get_sse_config(&bucket).await.ok();
|
||||
let mut effective_sse = requested_sse.or_else(|| {
|
||||
bucket_sse_config.as_ref().and_then(|(config, _)| {
|
||||
@@ -2599,7 +2683,7 @@ impl DefaultObjectUsecase {
|
||||
let h = HeaderMap::new();
|
||||
|
||||
let gr = store
|
||||
.get_object_reader(&src_bucket, &src_key, None, h, &get_opts)
|
||||
.get_object_reader(&src_bucket, &src_key, None, h, &src_get_opts)
|
||||
.await
|
||||
.map_err(ApiError::from)?;
|
||||
|
||||
@@ -2691,6 +2775,17 @@ impl DefaultObjectUsecase {
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(object_lock_metadata) = build_put_like_object_lock_metadata(
|
||||
&bucket,
|
||||
object_lock_legal_hold_status,
|
||||
object_lock_mode,
|
||||
object_lock_retain_until_date,
|
||||
)
|
||||
.await?
|
||||
{
|
||||
src_info.user_defined.extend(object_lock_metadata);
|
||||
}
|
||||
|
||||
let mut reader = HashReader::new(reader, length, actual_size, None, None, false).map_err(ApiError::from)?;
|
||||
|
||||
let encryption_request = EncryptionRequest {
|
||||
@@ -2887,6 +2982,7 @@ impl DefaultObjectUsecase {
|
||||
};
|
||||
|
||||
if gerr.is_none()
|
||||
&& !delete_creates_delete_marker(&opts)
|
||||
&& let Some(block_reason) = check_object_lock_for_deletion(&bucket, &goi, bypass_governance).await
|
||||
{
|
||||
delete_results[idx].error = Some(Error {
|
||||
@@ -3176,7 +3272,9 @@ impl DefaultObjectUsecase {
|
||||
// Check for bypass governance retention header (permission already verified in access.rs)
|
||||
let bypass_governance = has_bypass_governance_header(&req.headers);
|
||||
|
||||
if let Some(block_reason) = check_object_lock_for_deletion(&bucket, &obj_info, bypass_governance).await {
|
||||
if !delete_creates_delete_marker(&opts)
|
||||
&& let Some(block_reason) = check_object_lock_for_deletion(&bucket, &obj_info, bypass_governance).await
|
||||
{
|
||||
return Err(S3Error::with_message(S3ErrorCode::AccessDenied, block_reason.error_message()));
|
||||
}
|
||||
Some(obj_info)
|
||||
|
||||
@@ -355,6 +355,17 @@ pub fn has_bypass_governance_header(headers: &http::HeaderMap) -> bool {
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
fn legal_hold_write_requested(object_lock_legal_hold_status: Option<&ObjectLockLegalHoldStatus>) -> bool {
|
||||
object_lock_legal_hold_status.is_some()
|
||||
}
|
||||
|
||||
fn retention_write_requested(
|
||||
object_lock_mode: Option<&ObjectLockMode>,
|
||||
object_lock_retain_until_date: Option<&Timestamp>,
|
||||
) -> bool {
|
||||
object_lock_mode.is_some() || object_lock_retain_until_date.is_some()
|
||||
}
|
||||
|
||||
fn get_bucket_policy_authorize_action() -> Action {
|
||||
Action::S3Action(S3Action::GetBucketPolicyAction)
|
||||
}
|
||||
@@ -521,7 +532,17 @@ impl S3Access for FS {
|
||||
req_info.object = Some(req.input.key.clone());
|
||||
req_info.version_id = req.input.version_id.clone();
|
||||
|
||||
authorize_request(req, Action::S3Action(S3Action::PutObjectAction)).await
|
||||
authorize_request(req, Action::S3Action(S3Action::PutObjectAction)).await?;
|
||||
|
||||
if legal_hold_write_requested(req.input.object_lock_legal_hold_status.as_ref()) {
|
||||
authorize_request(req, Action::S3Action(S3Action::PutObjectLegalHoldAction)).await?;
|
||||
}
|
||||
|
||||
if retention_write_requested(req.input.object_lock_mode.as_ref(), req.input.object_lock_retain_until_date.as_ref()) {
|
||||
authorize_request(req, Action::S3Action(S3Action::PutObjectRetentionAction)).await?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Checks whether the CreateMultipartUpload request has accesses to the resources.
|
||||
@@ -530,7 +551,17 @@ impl S3Access for FS {
|
||||
req_info.bucket = Some(req.input.bucket.clone());
|
||||
req_info.object = Some(req.input.key.clone());
|
||||
|
||||
authorize_request(req, Action::S3Action(S3Action::PutObjectAction)).await
|
||||
authorize_request(req, Action::S3Action(S3Action::PutObjectAction)).await?;
|
||||
|
||||
if legal_hold_write_requested(req.input.object_lock_legal_hold_status.as_ref()) {
|
||||
authorize_request(req, Action::S3Action(S3Action::PutObjectLegalHoldAction)).await?;
|
||||
}
|
||||
|
||||
if retention_write_requested(req.input.object_lock_mode.as_ref(), req.input.object_lock_retain_until_date.as_ref()) {
|
||||
authorize_request(req, Action::S3Action(S3Action::PutObjectRetentionAction)).await?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Checks whether the DeleteBucket request has accesses to the resources.
|
||||
@@ -1405,7 +1436,17 @@ impl S3Access for FS {
|
||||
req_info.object = Some(req.input.key.clone());
|
||||
req_info.version_id = req.input.version_id.clone();
|
||||
|
||||
authorize_request(req, Action::S3Action(S3Action::PutObjectAction)).await
|
||||
authorize_request(req, Action::S3Action(S3Action::PutObjectAction)).await?;
|
||||
|
||||
if legal_hold_write_requested(req.input.object_lock_legal_hold_status.as_ref()) {
|
||||
authorize_request(req, Action::S3Action(S3Action::PutObjectLegalHoldAction)).await?;
|
||||
}
|
||||
|
||||
if retention_write_requested(req.input.object_lock_mode.as_ref(), req.input.object_lock_retain_until_date.as_ref()) {
|
||||
authorize_request(req, Action::S3Action(S3Action::PutObjectRetentionAction)).await?;
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Checks whether the PutObjectAcl request has accesses to the resources.
|
||||
@@ -1567,6 +1608,7 @@ mod tests {
|
||||
use super::*;
|
||||
use http::{HeaderMap, Method, Uri};
|
||||
use std::collections::HashMap;
|
||||
use time::OffsetDateTime;
|
||||
|
||||
#[test]
|
||||
fn get_bucket_policy_uses_get_bucket_policy_action() {
|
||||
@@ -1593,6 +1635,26 @@ mod tests {
|
||||
assert_eq!(list_parts_authorize_action(), Action::S3Action(S3Action::ListMultipartUploadPartsAction));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn legal_hold_write_requested_is_true_when_status_present() {
|
||||
assert!(legal_hold_write_requested(Some(&ObjectLockLegalHoldStatus::from_static(
|
||||
ObjectLockLegalHoldStatus::ON
|
||||
))));
|
||||
assert!(!legal_hold_write_requested(None));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn retention_write_requested_is_true_when_mode_or_date_present() {
|
||||
let retain_until = OffsetDateTime::now_utc().into();
|
||||
|
||||
assert!(retention_write_requested(
|
||||
Some(&ObjectLockMode::from_static(ObjectLockMode::GOVERNANCE)),
|
||||
None
|
||||
));
|
||||
assert!(retention_write_requested(None, Some(&retain_until)));
|
||||
assert!(!retention_write_requested(None, None));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn validate_post_object_success_controls_accepts_supported_status_codes() {
|
||||
for status in [200, 201, 204] {
|
||||
|
||||
Reference in New Issue
Block a user