diff --git a/crates/e2e_test/src/object_lock/object_lock_test.rs b/crates/e2e_test/src/object_lock/object_lock_test.rs index 31a86d89a..199ffb9d5 100644 --- a/crates/e2e_test/src/object_lock/object_lock_test.rs +++ b/crates/e2e_test/src/object_lock/object_lock_test.rs @@ -25,8 +25,11 @@ //! - Default bucket retention is applied to new objects use super::common::*; +use aws_sdk_s3::Client; use aws_sdk_s3::primitives::ByteStream; -use aws_sdk_s3::types::{Delete, ObjectIdentifier, ObjectLockLegalHoldStatus, ObjectLockRetentionMode}; +use aws_sdk_s3::types::{ + CompletedMultipartUpload, CompletedPart, Delete, ObjectIdentifier, ObjectLockLegalHoldStatus, ObjectLockRetentionMode, +}; use serial_test::serial; use tracing::info; @@ -37,6 +40,45 @@ fn init_logging() { .try_init(); } +async fn put_bucket_deny_policy( + client: &Client, + bucket: &str, + sid: &str, + action: &str, +) -> Result<(), Box> { + let policy = serde_json::json!({ + "Version": "2012-10-17", + "Statement": [{ + "Sid": sid, + "Effect": "Deny", + "Principal": "*", + "Action": action, + "Resource": format!("arn:aws:s3:::{}/*", bucket) + }] + }) + .to_string(); + + client.put_bucket_policy().bucket(bucket).policy(policy).send().await?; + Ok(()) +} + +fn retention_timestamp(days: i64) -> aws_sdk_s3::primitives::DateTime { + let retain_until = future_retain_until(days).format("%Y-%m-%dT%H:%M:%SZ").to_string(); + aws_sdk_s3::primitives::DateTime::from_str(&retain_until, aws_sdk_s3::primitives::DateTimeFormat::DateTime) + .expect("retention timestamp should parse") +} + +fn assert_access_denied(result: Result, context: &str) { + let err = match result { + Ok(_) => panic!("{context}"), + Err(err) => format!("{err:?}"), + }; + assert!( + err.contains("AccessDenied") || err.to_lowercase().contains("access denied"), + "{context}: expected AccessDenied, got: {err}" + ); +} + // ============================================================================ // DeleteObject Tests // ============================================================================ @@ -149,6 +191,57 @@ async fn test_delete_object_allowed_by_governance_with_bypass() { info!("✅ Test passed: GOVERNANCE retention allows deletion with bypass"); } +#[tokio::test] +#[serial] +async fn test_delete_object_creates_delete_marker_for_retained_current_version() { + init_logging(); + info!("🧪 Test: DeleteObject creates delete marker for retained current version"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-retention-delete-marker"; + let key = "retained-object"; + let data = b"test data for retained current version"; + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + + let retain_until = future_retain_until(30); + let retained_version_id = + put_object_with_retention(&client, bucket, key, data, ObjectLockRetentionMode::Governance, retain_until) + .await + .unwrap(); + + let delete_marker_output = client.delete_object().bucket(bucket).key(key).send().await.unwrap(); + assert_eq!(delete_marker_output.delete_marker(), Some(true)); + + let delete_marker_version_id = delete_marker_output + .version_id() + .expect("delete marker should have a version id") + .to_string(); + + let protected_delete = delete_object_with_bypass(&client, bucket, key, Some(&retained_version_id), false).await; + assert!(protected_delete.is_err(), "Retained version should still reject direct deletion"); + + delete_object_with_bypass(&client, bucket, key, Some(&delete_marker_version_id), false) + .await + .unwrap(); + + let still_protected = delete_object_with_bypass(&client, bucket, key, Some(&retained_version_id), false).await; + assert!( + still_protected.is_err(), + "Retained version should remain protected after delete marker removal" + ); + + delete_object_with_bypass(&client, bucket, key, Some(&retained_version_id), true) + .await + .unwrap(); + + info!("✅ Test passed: Delete marker is allowed while retained version stays protected"); +} + #[tokio::test] #[serial] async fn test_delete_object_blocked_by_legal_hold() { @@ -182,6 +275,42 @@ async fn test_delete_object_blocked_by_legal_hold() { info!("✅ Test passed: Legal Hold blocks deletion"); } +#[tokio::test] +#[serial] +async fn test_delete_object_allowed_with_legal_hold_off() { + init_logging(); + info!("🧪 Test: DeleteObject allowed with Legal Hold OFF"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-legal-hold-off-delete"; + let key = "legal-hold-off-object"; + let data = b"test data for legal hold off"; + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + + let version_id = put_object_with_legal_hold(&client, bucket, key, data, ObjectLockLegalHoldStatus::Off) + .await + .unwrap(); + + let delete_result = delete_object_with_bypass(&client, bucket, key, Some(&version_id), false).await; + assert!(delete_result.is_ok(), "Delete should succeed when legal hold is OFF"); + + let head_result = client + .head_object() + .bucket(bucket) + .key(key) + .version_id(&version_id) + .send() + .await; + assert!(head_result.is_err(), "Object should be deleted when legal hold is OFF"); + + info!("✅ Test passed: Legal Hold OFF allows deletion"); +} + #[tokio::test] #[serial] async fn test_delete_object_after_legal_hold_removed() { @@ -216,6 +345,737 @@ async fn test_delete_object_after_legal_hold_removed() { info!("✅ Test passed: Deletion succeeds after Legal Hold removal"); } +#[tokio::test] +#[serial] +async fn test_get_object_legal_hold_returns_updated_status() { + init_logging(); + info!("🧪 Test: GetObjectLegalHold returns updated status"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-get-legal-hold"; + let key = "legal-hold-object"; + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + let version_id = put_object_with_legal_hold(&client, bucket, key, b"test data", ObjectLockLegalHoldStatus::On) + .await + .unwrap(); + + let on_hold = client + .get_object_legal_hold() + .bucket(bucket) + .key(key) + .version_id(&version_id) + .send() + .await + .unwrap(); + assert_eq!( + on_hold + .legal_hold() + .and_then(|value| value.status()) + .map(|value| value.as_str()), + Some("ON") + ); + + put_object_legal_hold(&client, bucket, key, Some(&version_id), ObjectLockLegalHoldStatus::Off) + .await + .unwrap(); + + let off_hold = client + .get_object_legal_hold() + .bucket(bucket) + .key(key) + .version_id(&version_id) + .send() + .await + .unwrap(); + assert_eq!( + off_hold + .legal_hold() + .and_then(|value| value.status()) + .map(|value| value.as_str()), + Some("OFF") + ); +} + +#[tokio::test] +#[serial] +async fn test_get_object_retention_returns_configured_values() { + init_logging(); + info!("🧪 Test: GetObjectRetention returns configured values"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-get-retention"; + let key = "retained-object"; + let retain_until = future_retain_until(30); + let retain_until_expected = retain_until.format("%Y-%m-%dT%H:%M:%SZ").to_string(); + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + let version_id = + put_object_with_retention(&client, bucket, key, b"test data", ObjectLockRetentionMode::Governance, retain_until) + .await + .unwrap(); + + let retention = client + .get_object_retention() + .bucket(bucket) + .key(key) + .version_id(&version_id) + .send() + .await + .unwrap(); + let retention = retention.retention().expect("retention should be present"); + + assert_eq!(retention.mode().map(|value| value.as_str()), Some("GOVERNANCE")); + assert_eq!( + retention + .retain_until_date() + .expect("retain_until_date should be present") + .fmt(aws_sdk_s3::primitives::DateTimeFormat::DateTime) + .unwrap(), + retain_until_expected + ); +} + +// ============================================================================ +// Put/Copy/Multipart Legal Hold Tests +// ============================================================================ + +#[tokio::test] +#[serial] +async fn test_put_object_overwrite_blocked_by_legal_hold() { + init_logging(); + info!("🧪 Test: PutObject overwrite blocked by Legal Hold"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-put-overwrite-legal-hold"; + let key = "overwrite-object"; + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + + put_object_with_legal_hold(&client, bucket, key, b"locked-body", ObjectLockLegalHoldStatus::On) + .await + .unwrap(); + + let overwrite_result = client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from(b"replacement-body".to_vec())) + .send() + .await; + + assert!(overwrite_result.is_err(), "PutObject overwrite should fail while legal hold is ON"); + + let error_str = format!("{:?}", overwrite_result.unwrap_err()); + assert!( + error_str.to_lowercase().contains("legal") || error_str.to_lowercase().contains("hold"), + "overwrite error should mention legal hold, got: {error_str}" + ); +} + +#[tokio::test] +#[serial] +async fn test_copy_object_applies_requested_legal_hold() { + init_logging(); + info!("🧪 Test: CopyObject applies requested Legal Hold"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-copy-object-legal-hold"; + let src_key = "src-object"; + let dst_key = "dst-object"; + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + client + .put_object() + .bucket(bucket) + .key(src_key) + .body(ByteStream::from(b"copy-source".to_vec())) + .send() + .await + .unwrap(); + + client + .copy_object() + .copy_source(format!("{bucket}/{src_key}")) + .bucket(bucket) + .key(dst_key) + .object_lock_legal_hold_status(ObjectLockLegalHoldStatus::On) + .send() + .await + .unwrap(); + + let legal_hold = client + .get_object_legal_hold() + .bucket(bucket) + .key(dst_key) + .send() + .await + .unwrap(); + + assert_eq!( + legal_hold + .legal_hold() + .and_then(|value| value.status()) + .map(|value| value.as_str()), + Some("ON") + ); +} + +#[tokio::test] +#[serial] +async fn test_copy_object_overwrite_blocked_by_legal_hold() { + init_logging(); + info!("🧪 Test: CopyObject overwrite blocked by Legal Hold"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-copy-overwrite-legal-hold"; + let src_key = "src-object"; + let dst_key = "dst-object"; + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + client + .put_object() + .bucket(bucket) + .key(src_key) + .body(ByteStream::from(b"copy-source".to_vec())) + .send() + .await + .unwrap(); + + put_object_with_legal_hold(&client, bucket, dst_key, b"locked-destination", ObjectLockLegalHoldStatus::On) + .await + .unwrap(); + + let copy_result = client + .copy_object() + .copy_source(format!("{bucket}/{src_key}")) + .bucket(bucket) + .key(dst_key) + .send() + .await; + + assert!( + copy_result.is_err(), + "CopyObject overwrite should fail while destination legal hold is ON" + ); + + let error_str = format!("{:?}", copy_result.unwrap_err()); + assert!( + error_str.to_lowercase().contains("legal") || error_str.to_lowercase().contains("hold"), + "copy overwrite error should mention legal hold, got: {error_str}" + ); +} + +#[tokio::test] +#[serial] +async fn test_create_multipart_upload_applies_requested_legal_hold() { + init_logging(); + info!("🧪 Test: CreateMultipartUpload applies requested Legal Hold"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-multipart-legal-hold"; + let key = "multipart-object"; + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + let create_output = client + .create_multipart_upload() + .bucket(bucket) + .key(key) + .object_lock_legal_hold_status(ObjectLockLegalHoldStatus::On) + .send() + .await + .unwrap(); + + let upload_id = create_output.upload_id().unwrap(); + let upload_part_output = client + .upload_part() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .part_number(1) + .body(ByteStream::from(b"multipart-body".to_vec())) + .send() + .await + .unwrap(); + + let completed_upload = CompletedMultipartUpload::builder() + .parts( + CompletedPart::builder() + .part_number(1) + .e_tag(upload_part_output.e_tag().unwrap_or_default()) + .build(), + ) + .build(); + + client + .complete_multipart_upload() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .multipart_upload(completed_upload) + .send() + .await + .unwrap(); + + let legal_hold = client.get_object_legal_hold().bucket(bucket).key(key).send().await.unwrap(); + + assert_eq!( + legal_hold + .legal_hold() + .and_then(|value| value.status()) + .map(|value| value.as_str()), + Some("ON") + ); +} + +#[tokio::test] +#[serial] +async fn test_create_multipart_upload_blocked_by_compliance_retention() { + init_logging(); + info!("🧪 Test: CreateMultipartUpload blocked by COMPLIANCE retention"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-multipart-create-compliance"; + let key = "protected-object"; + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + put_object_with_retention( + &client, + bucket, + key, + b"locked-destination", + ObjectLockRetentionMode::Compliance, + future_retain_until(30), + ) + .await + .unwrap(); + + let create_result = client.create_multipart_upload().bucket(bucket).key(key).send().await; + + assert!( + create_result.is_err(), + "CreateMultipartUpload should fail while destination is under active COMPLIANCE retention" + ); + + let error_str = format!("{:?}", create_result.unwrap_err()); + assert!( + error_str.to_lowercase().contains("retention") || error_str.to_lowercase().contains("compliance"), + "multipart create error should mention retention, got: {error_str}" + ); +} + +#[tokio::test] +#[serial] +async fn test_delete_completed_multipart_object_blocked_by_legal_hold() { + init_logging(); + info!("🧪 Test: Delete completed multipart object blocked by Legal Hold"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-multipart-delete-legal-hold"; + let key = "multipart-object"; + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + let create_output = client + .create_multipart_upload() + .bucket(bucket) + .key(key) + .object_lock_legal_hold_status(ObjectLockLegalHoldStatus::On) + .send() + .await + .unwrap(); + + let upload_id = create_output.upload_id().unwrap(); + let upload_part_output = client + .upload_part() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .part_number(1) + .body(ByteStream::from(b"multipart-body".to_vec())) + .send() + .await + .unwrap(); + + let completed_upload = CompletedMultipartUpload::builder() + .parts( + CompletedPart::builder() + .part_number(1) + .e_tag(upload_part_output.e_tag().unwrap_or_default()) + .build(), + ) + .build(); + + let complete_output = client + .complete_multipart_upload() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .multipart_upload(completed_upload) + .send() + .await + .unwrap(); + + let version_id = complete_output.version_id().expect("multipart object should be versioned"); + let delete_result = delete_object_with_bypass(&client, bucket, key, Some(version_id), false).await; + assert!(delete_result.is_err(), "Delete should fail for multipart object protected by legal hold"); +} + +#[tokio::test] +#[serial] +async fn test_delete_completed_multipart_object_blocked_by_retention() { + init_logging(); + info!("🧪 Test: Delete completed multipart object blocked by retention"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-multipart-delete-retention"; + let key = "multipart-object"; + let retain_until = retention_timestamp(30); + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + let create_output = client + .create_multipart_upload() + .bucket(bucket) + .key(key) + .object_lock_mode(aws_sdk_s3::types::ObjectLockMode::Compliance) + .object_lock_retain_until_date(retain_until) + .send() + .await + .unwrap(); + + let upload_id = create_output.upload_id().unwrap(); + let upload_part_output = client + .upload_part() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .part_number(1) + .body(ByteStream::from(b"multipart-body".to_vec())) + .send() + .await + .unwrap(); + + let completed_upload = CompletedMultipartUpload::builder() + .parts( + CompletedPart::builder() + .part_number(1) + .e_tag(upload_part_output.e_tag().unwrap_or_default()) + .build(), + ) + .build(); + + let complete_output = client + .complete_multipart_upload() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .multipart_upload(completed_upload) + .send() + .await + .unwrap(); + + let version_id = complete_output.version_id().expect("multipart object should be versioned"); + let delete_result = delete_object_with_bypass(&client, bucket, key, Some(version_id), false).await; + assert!(delete_result.is_err(), "Delete should fail for multipart object protected by retention"); +} + +#[tokio::test] +#[serial] +async fn test_complete_multipart_upload_blocked_when_legal_hold_added_after_create() { + init_logging(); + info!("🧪 Test: CompleteMultipartUpload blocked when Legal Hold appears after MPU creation"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-complete-multipart-legal-hold"; + let key = "multipart-race-object"; + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + let create_output = client.create_multipart_upload().bucket(bucket).key(key).send().await.unwrap(); + + let upload_id = create_output.upload_id().unwrap(); + let upload_part_output = client + .upload_part() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .part_number(1) + .body(ByteStream::from(b"multipart-body".to_vec())) + .send() + .await + .unwrap(); + + put_object_with_legal_hold(&client, bucket, key, b"locked-current-version", ObjectLockLegalHoldStatus::On) + .await + .unwrap(); + + let completed_upload = CompletedMultipartUpload::builder() + .parts( + CompletedPart::builder() + .part_number(1) + .e_tag(upload_part_output.e_tag().unwrap_or_default()) + .build(), + ) + .build(); + + let complete_result = client + .complete_multipart_upload() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .multipart_upload(completed_upload) + .send() + .await; + + assert!(complete_result.is_err(), "CompleteMultipartUpload should fail once legal hold is enabled"); + + let error_str = format!("{:?}", complete_result.unwrap_err()); + assert!( + error_str.to_lowercase().contains("legal") || error_str.to_lowercase().contains("hold"), + "complete error should mention legal hold, got: {error_str}" + ); +} + +#[tokio::test] +#[serial] +async fn test_complete_multipart_upload_blocked_when_compliance_retention_added_after_create() { + init_logging(); + info!("🧪 Test: CompleteMultipartUpload blocked when COMPLIANCE retention appears after MPU creation"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-complete-multipart-compliance"; + let key = "multipart-race-object"; + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + let create_output = client.create_multipart_upload().bucket(bucket).key(key).send().await.unwrap(); + + let upload_id = create_output.upload_id().unwrap(); + let upload_part_output = client + .upload_part() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .part_number(1) + .body(ByteStream::from(b"multipart-body".to_vec())) + .send() + .await + .unwrap(); + + put_object_with_retention( + &client, + bucket, + key, + b"locked-current-version", + ObjectLockRetentionMode::Compliance, + future_retain_until(30), + ) + .await + .unwrap(); + + let completed_upload = CompletedMultipartUpload::builder() + .parts( + CompletedPart::builder() + .part_number(1) + .e_tag(upload_part_output.e_tag().unwrap_or_default()) + .build(), + ) + .build(); + + let complete_result = client + .complete_multipart_upload() + .bucket(bucket) + .key(key) + .upload_id(upload_id) + .multipart_upload(completed_upload) + .send() + .await; + + assert!( + complete_result.is_err(), + "CompleteMultipartUpload should fail once COMPLIANCE retention is enabled" + ); + + let error_str = format!("{:?}", complete_result.unwrap_err()); + assert!( + error_str.to_lowercase().contains("retention") || error_str.to_lowercase().contains("compliance"), + "complete error should mention retention, got: {error_str}" + ); +} + +#[tokio::test] +#[serial] +async fn test_write_paths_require_put_object_legal_hold_permission() { + init_logging(); + info!("🧪 Test: write paths require PutObjectLegalHold permission"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-legal-hold-permissions"; + let src_key = "src-object"; + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + client + .put_object() + .bucket(bucket) + .key(src_key) + .body(ByteStream::from(b"copy-source".to_vec())) + .send() + .await + .unwrap(); + + put_bucket_deny_policy(&client, bucket, "DenyPutObjectLegalHold", "s3:PutObjectLegalHold") + .await + .unwrap(); + + assert_access_denied( + client + .put_object() + .bucket(bucket) + .key("put-target") + .body(ByteStream::from(b"put-body".to_vec())) + .object_lock_legal_hold_status(ObjectLockLegalHoldStatus::On) + .send() + .await, + "PutObject with legal hold should require s3:PutObjectLegalHold", + ); + + assert_access_denied( + client + .copy_object() + .copy_source(format!("{bucket}/{src_key}")) + .bucket(bucket) + .key("copy-target") + .object_lock_legal_hold_status(ObjectLockLegalHoldStatus::On) + .send() + .await, + "CopyObject with legal hold should require s3:PutObjectLegalHold", + ); + + assert_access_denied( + client + .create_multipart_upload() + .bucket(bucket) + .key("multipart-target") + .object_lock_legal_hold_status(ObjectLockLegalHoldStatus::On) + .send() + .await, + "CreateMultipartUpload with legal hold should require s3:PutObjectLegalHold", + ); +} + +#[tokio::test] +#[serial] +async fn test_write_paths_require_put_object_retention_permission() { + init_logging(); + info!("🧪 Test: write paths require PutObjectRetention permission"); + + let mut env = ObjectLockTestEnvironment::new().await.unwrap(); + env.start_rustfs().await.unwrap(); + + let bucket = "test-retention-permissions"; + let src_key = "src-object"; + let retain_until = retention_timestamp(30); + + env.create_object_lock_bucket(bucket).await.unwrap(); + + let client = env.s3_client(); + client + .put_object() + .bucket(bucket) + .key(src_key) + .body(ByteStream::from(b"copy-source".to_vec())) + .send() + .await + .unwrap(); + + put_bucket_deny_policy(&client, bucket, "DenyPutObjectRetention", "s3:PutObjectRetention") + .await + .unwrap(); + + assert_access_denied( + client + .put_object() + .bucket(bucket) + .key("put-target") + .body(ByteStream::from(b"put-body".to_vec())) + .object_lock_mode(aws_sdk_s3::types::ObjectLockMode::Governance) + .object_lock_retain_until_date(retain_until) + .send() + .await, + "PutObject with retention should require s3:PutObjectRetention", + ); + + assert_access_denied( + client + .copy_object() + .copy_source(format!("{bucket}/{src_key}")) + .bucket(bucket) + .key("copy-target") + .object_lock_mode(aws_sdk_s3::types::ObjectLockMode::Governance) + .object_lock_retain_until_date(retain_until) + .send() + .await, + "CopyObject with retention should require s3:PutObjectRetention", + ); + + assert_access_denied( + client + .create_multipart_upload() + .bucket(bucket) + .key("multipart-target") + .object_lock_mode(aws_sdk_s3::types::ObjectLockMode::Governance) + .object_lock_retain_until_date(retain_until) + .send() + .await, + "CreateMultipartUpload with retention should require s3:PutObjectRetention", + ); +} + // ============================================================================ // DeleteObjects (Batch Delete) Tests // ============================================================================ diff --git a/rustfs/src/app/multipart_usecase.rs b/rustfs/src/app/multipart_usecase.rs index 1143e6cfc..6079c5892 100644 --- a/rustfs/src/app/multipart_usecase.rs +++ b/rustfs/src/app/multipart_usecase.rs @@ -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) -> S3Result 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() diff --git a/rustfs/src/app/object_usecase.rs b/rustfs/src/app/object_usecase.rs index daaf01b2f..f72b6c897 100644 --- a/rustfs/src/app/object_usecase.rs +++ b/rustfs/src/app/object_usecase.rs @@ -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, 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, + object_lock_mode: Option, + object_lock_retain_until_date: Option, +) -> S3Result>> { 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 = 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) diff --git a/rustfs/src/storage/access.rs b/rustfs/src/storage/access.rs index 341697336..6a3a9def9 100644 --- a/rustfs/src/storage/access.rs +++ b/rustfs/src/storage/access.rs @@ -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] { diff --git a/scripts/s3-tests/excluded_tests.txt b/scripts/s3-tests/excluded_tests.txt index 7616b1340..295147e3c 100644 --- a/scripts/s3-tests/excluded_tests.txt +++ b/scripts/s3-tests/excluded_tests.txt @@ -240,20 +240,8 @@ test_object_header_acl_grants test_object_lock_changing_mode_from_compliance test_object_lock_changing_mode_from_governance_with_bypass test_object_lock_changing_mode_from_governance_without_bypass -test_object_lock_delete_multipart_object_with_legal_hold_on -test_object_lock_delete_multipart_object_with_retention -test_object_lock_delete_object_with_legal_hold_off -test_object_lock_delete_object_with_legal_hold_on -test_object_lock_delete_object_with_retention -test_object_lock_delete_object_with_retention_and_marker -test_object_lock_get_legal_hold test_object_lock_get_obj_lock test_object_lock_get_obj_metadata -test_object_lock_get_obj_retention -test_object_lock_get_obj_retention_iso8601 -test_object_lock_multi_delete_object_with_retention -test_object_lock_put_legal_hold -test_object_lock_put_legal_hold_invalid_status test_object_lock_put_obj_lock test_object_lock_put_obj_lock_invalid_days test_object_lock_put_obj_lock_invalid_mode diff --git a/scripts/s3-tests/implemented_tests.txt b/scripts/s3-tests/implemented_tests.txt index 0ce1c6190..e24a2691c 100644 --- a/scripts/s3-tests/implemented_tests.txt +++ b/scripts/s3-tests/implemented_tests.txt @@ -391,6 +391,18 @@ test_atomic_multipart_upload_write # Object Lock tests test_object_lock_put_obj_lock_enable_after_create +test_object_lock_get_legal_hold +test_object_lock_get_obj_retention +test_object_lock_get_obj_retention_iso8601 +test_object_lock_put_legal_hold +test_object_lock_put_legal_hold_invalid_status +test_object_lock_delete_object_with_legal_hold_off +test_object_lock_delete_object_with_legal_hold_on +test_object_lock_delete_object_with_retention +test_object_lock_delete_object_with_retention_and_marker +test_object_lock_delete_multipart_object_with_legal_hold_on +test_object_lock_delete_multipart_object_with_retention +test_object_lock_multi_delete_object_with_retention # Checksum validation tests test_object_checksum_sha256