diff --git a/crates/e2e_test/src/replication_extension_test.rs b/crates/e2e_test/src/replication_extension_test.rs index e636d9db5..a4ddd0369 100644 --- a/crates/e2e_test/src/replication_extension_test.rs +++ b/crates/e2e_test/src/replication_extension_test.rs @@ -10394,13 +10394,14 @@ async fn test_get_object_tagging_proxies_unreplicated_object_to_replication_targ let target_bucket = "proxy-tag-dst"; let (target, source_env, source_client, target_client) = start_read_proxy_lab(source_bucket, target_bucket).await?; - target_client + let tagged = target_client .put_object() .bucket(target_bucket) .key("proxy-tagged") .body(ByteStream::from_static(b"tagged payload")) .send() .await?; + let tagged_version = tagged.version_id().ok_or("versioned target PUT omitted its identity")?; target_client .put_object_tagging() .bucket(target_bucket) @@ -10424,6 +10425,11 @@ async fn test_get_object_tagging_proxies_unreplicated_object_to_replication_targ assert_eq!(tags.tag_set.len(), 1, "proxied tagging read must return the target's tags"); assert_eq!(tags.tag_set[0].key.as_str(), "team"); assert_eq!(tags.tag_set[0].value.as_str(), "storage"); + assert_eq!( + tags.version_id(), + Some(tagged_version), + "proxy must preserve the resolved remote identity" + ); let record = target .requests() @@ -10437,6 +10443,33 @@ async fn test_get_object_tagging_proxies_unreplicated_object_to_replication_targ ); assert!(record.proxy_headers.replication_check.is_none()); + let empty = target_client + .put_object() + .bucket(target_bucket) + .key("proxy-tagged") + .body(ByteStream::from_static(b"new version without tags")) + .send() + .await?; + let empty_version = empty.version_id().ok_or("empty-tag version omitted its identity")?; + for selector in [None, Some(empty_version), Some(tagged_version)] { + let tags = source_client + .get_object_tagging() + .bucket(source_bucket) + .key("proxy-tagged") + .set_version_id(selector.map(str::to_owned)) + .send() + .await?; + let expected_version = selector.unwrap_or(empty_version); + assert_eq!(tags.version_id(), Some(expected_version)); + if expected_version == tagged_version { + assert_eq!(tags.tag_set().len(), 1); + assert_eq!(tags.tag_set()[0].key(), "team"); + assert_eq!(tags.tag_set()[0].value(), "storage"); + } else { + assert!(tags.tag_set().is_empty()); + } + } + drop(source_env); target.shutdown().await; Ok(()) diff --git a/rustfs/src/app/object/shared.rs b/rustfs/src/app/object/shared.rs index a34abc5cc..cca89b73c 100644 --- a/rustfs/src/app/object/shared.rs +++ b/rustfs/src/app/object/shared.rs @@ -34,7 +34,7 @@ pub(super) const LOG_SUBSYSTEM_OBJECT: &str = "object"; /// Encode the resolved local read identity. Storage distinguishes a null /// version from an unversioned object by returning a nil UUID instead of None. -pub(super) fn read_response_version_id(version_id: Option) -> Option { +pub(crate) fn read_response_version_id(version_id: Option) -> Option { version_id.map(|id| { if id.is_nil() { NULL_VERSION_ID.to_owned() diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 82038336a..0a940a178 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -198,7 +198,7 @@ impl FS { object: &str, version_id: Option, headers: &http::HeaderMap, - ) -> Option { + ) -> Option { let opts = Self::tagging_proxy_opts(bucket, object, version_id, headers).await?; let targets = get_read_proxy_targets(bucket, object, &opts).await; if targets.is_empty() { @@ -213,8 +213,8 @@ impl FS { // MinIO-aligned accounting: one total per proxy attempt, // one failed when no target served it. record_replication_proxy(bucket, "GetObjectTagging", false).await; - return Some( - remote + return Some(GetObjectTaggingOutput { + tag_set: remote .tag_set .into_iter() .map(|tag| Tag { @@ -222,7 +222,8 @@ impl FS { value: Some(tag.value), }) .collect(), - ); + version_id: remote.version_id, + }); } Err(err) if Self::proxy_sdk_error_is_not_found(&err) => { debug!(bucket, object, arn = %target.arn, "tagging proxy: target does not have the object"); @@ -1134,6 +1135,8 @@ impl S3 for FS { #[instrument(level = "debug", skip(self))] async fn get_object_tagging(&self, req: S3Request) -> S3Result> { + use crate::storage::storage_api::ecstore_bucket::versioning::VersioningApi as _; + record_s3_op(S3Operation::GetObjectTagging); let start_time = std::time::Instant::now(); let bucket = req.input.bucket.as_str(); @@ -1152,30 +1155,38 @@ impl S3 for FS { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; - let version_id = req.input.version_id.clone(); + let version_id = parse_object_version_id(req.input.version_id.clone())?.map(Into::into); + let versioning = BucketVersioningSys::get(bucket).await.map_err(ApiError::from)?; let opts = ObjectOptions { - version_id: parse_object_version_id(version_id)?.map(Into::into), + version_id, + versioned: versioning.prefix_enabled(object), + version_suspended: versioning.prefix_suspended(object), ..Default::default() }; - let tags = match store.get_object_tags(bucket, object, &opts).await { - Ok(tags) => tags, + // Tags and response identity must come from the same locked metadata snapshot. + let info = match store.get_object_info(bucket, object, &opts).await { + Ok(info) if info.delete_marker => { + return Err(S3Error::new(if opts.version_id.is_some() { + S3ErrorCode::MethodNotAllowed + } else { + S3ErrorCode::NoSuchKey + })); + } + Ok(info) => info, Err(e) => { // Replication lag window: the object may exist on a // replication target even though it is missing locally — // proxy the tagging read there (backlog#1675 P1-5). if (is_err_object_not_found(&e) || is_err_version_not_found(&e)) - && let Some(tag_set) = + && let Some(output) = Self::proxy_get_object_tagging(bucket, object, req.input.version_id.clone(), &req.headers).await { counter!("rustfs_get_object_tagging_success").increment(1); let duration = start_time.elapsed(); histogram!("rustfs_object_tagging_operation_duration_seconds", "operation" => "get") .record(duration.as_secs_f64()); - return Ok(S3Response::new(GetObjectTaggingOutput { - tag_set, - version_id: req.input.version_id.clone(), - })); + return Ok(S3Response::new(output)); } if is_err_object_not_found(&e) { debug!( @@ -1206,7 +1217,7 @@ impl S3 for FS { } }; - let tag_set = decode_tags(tags.as_str()); + let tag_set = decode_tags(info.user_tags.as_str()); debug!( component = LOG_COMPONENT_STORAGE, subsystem = LOG_SUBSYSTEM_TAGGING, @@ -1222,7 +1233,7 @@ impl S3 for FS { histogram!("rustfs_object_tagging_operation_duration_seconds", "operation" => "get").record(duration.as_secs_f64()); Ok(S3Response::new(GetObjectTaggingOutput { tag_set, - version_id: req.input.version_id.clone(), + version_id: s3_api::read_response_version_id(info.version_id), })) } diff --git a/rustfs/src/storage/s3_api/mod.rs b/rustfs/src/storage/s3_api/mod.rs index 6773d7716..f29f88e68 100644 --- a/rustfs/src/storage/s3_api/mod.rs +++ b/rustfs/src/storage/s3_api/mod.rs @@ -17,6 +17,7 @@ //! This file intentionally starts as skeleton-only. Behavior remains in place //! until each helper is moved with dedicated small refactor steps. +pub(crate) use crate::app::object::read_response_version_id; use crate::app::{ bucket_usecase::DefaultBucketUsecase, multipart_usecase::DefaultMultipartUsecase, object_usecase::DefaultObjectUsecase, }; diff --git a/rustfs/tests/embedded_test.rs b/rustfs/tests/embedded_test.rs index 654a37285..a2df768c3 100644 --- a/rustfs/tests/embedded_test.rs +++ b/rustfs/tests/embedded_test.rs @@ -20,8 +20,11 @@ #![recursion_limit = "256"] use aws_sdk_s3::config::{Credentials, Region}; +use aws_sdk_s3::error::ProvideErrorMetadata; use aws_sdk_s3::primitives::ByteStream; -use aws_sdk_s3::types::{BucketVersioningStatus, Delete, ObjectAttributes, ObjectIdentifier, VersioningConfiguration}; +use aws_sdk_s3::types::{ + BucketVersioningStatus, Delete, ObjectAttributes, ObjectIdentifier, Tag, Tagging, VersioningConfiguration, +}; use aws_sdk_s3::{Client, Config}; use rustfs::embedded::{RustFSServerBuilder, find_available_port}; @@ -294,6 +297,32 @@ async fn assert_read_version( .expect("get selected version attributes"); assert_eq!(attributes.version_id(), expected_version, "Attributes identity for selector {selector:?}"); assert_eq!(attributes.object_size(), Some(expected_body.len() as i64)); + + assert_tagging_version(client, bucket, key, selector, expected_version, None).await; +} + +async fn assert_tagging_version( + client: &Client, + bucket: &str, + key: &str, + selector: Option<&str>, + expected_version: Option<&str>, + generation: Option<&str>, +) { + let tags = client + .get_object_tagging() + .bucket(bucket) + .key(key) + .set_version_id(selector.map(str::to_owned)) + .send() + .await + .expect("read selected version tags"); + assert_eq!(tags.version_id(), expected_version, "tagging identity for selector {selector:?}"); + let expected: Vec<_> = generation + .map(|value| Tag::builder().key("generation").value(value).build().expect("valid tag")) + .into_iter() + .collect(); + assert_eq!(tags.tag_set(), expected, "tagging contents for selector {selector:?}"); } #[test] @@ -437,3 +466,222 @@ async fn test_read_version_headers_across_versioning_states_body() { server.shutdown().await; } + +#[test] +fn test_tagging_version_snapshots_single_pool() { + common::run_embedded_test(|| tagging_version_snapshots(1)); +} + +#[test] +fn test_tagging_version_snapshots_multiple_pools() { + if std::env::var("RUSTFS_UNSAFE_BYPASS_DISK_CHECK").as_deref() != Ok("true") { + // These logical drives share one temporary filesystem. Scope the + // existing local-test override to a child instead of mutating the + // environment of other embedded servers in this test process. + let output = std::process::Command::new(std::env::current_exe().expect("test executable")) + .args(["--exact", "test_tagging_version_snapshots_multiple_pools", "--nocapture"]) + .env("RUSTFS_UNSAFE_BYPASS_DISK_CHECK", "true") + .output() + .expect("run multi-pool child"); + let stdout = String::from_utf8_lossy(&output.stdout); + assert!( + output.status.success() && stdout.contains("test result: ok. 1 passed;"), + "multi-pool child must execute and pass the regression:\n{stdout}\n{}", + String::from_utf8_lossy(&output.stderr) + ); + return; + } + common::run_embedded_test(|| tagging_version_snapshots(2)); +} + +async fn tagging_version_snapshots(pool_count: usize) { + let root = tempfile::tempdir().expect("temporary drives"); + let mut builder = RustFSServerBuilder::new() + .address(format!("127.0.0.1:{}", find_available_port().expect("free port"))) + .access_key("testaccesskey") + .secret_key("testsecretkey"); + if pool_count > 1 { + for pool in 0..pool_count { + for disk in 1..=4 { + std::fs::create_dir_all(root.path().join(format!("pool{pool}/disk{disk}"))).expect("create test drive"); + } + } + builder = builder.volumes( + (0..pool_count) + .map(|pool| format!("{}/pool{pool}/disk{{1...4}}", root.path().display())) + .collect(), + ); + } + let server = builder.build().await.expect("start tagging server"); + let client = s3_client(&server.endpoint(), server.access_key(), server.secret_key()); + let bucket = "tagging-version-snapshots"; + let key = "tags/nested/中文 +%.bin"; + client.create_bucket().bucket(bucket).send().await.expect("create bucket"); + client + .put_bucket_versioning() + .bucket(bucket) + .versioning_configuration( + VersioningConfiguration::builder() + .status(BucketVersioningStatus::Enabled) + .build(), + ) + .send() + .await + .expect("enable versioning"); + + let mut history = Vec::new(); + for generation in ["oldest", "middle", "latest"] { + let put = client + .put_object() + .bucket(bucket) + .key(key) + .tagging(format!("generation={generation}")) + .body(ByteStream::from(generation.as_bytes().to_vec())) + .send() + .await + .expect("write tagged version"); + let version = put.version_id().expect("acknowledged version").to_owned(); + assert_tagging_version(&client, bucket, key, None, Some(&version), Some(generation)).await; + history.push((version, generation)); + } + for (index, value) in [(0, "updated-oldest"), (2, "updated-latest")] { + let (version, generation) = &mut history[index]; + client + .put_object_tagging() + .bucket(bucket) + .key(key) + .version_id(version.as_str()) + .tagging( + Tagging::builder() + .tag_set(Tag::builder().key("generation").value(value).build().expect("valid tag")) + .build() + .expect("valid tagging"), + ) + .send() + .await + .expect("update only the selected version's tags"); + *generation = value; + } + for (version, generation) in &history { + assert_tagging_version(&client, bucket, key, Some(version), Some(version), Some(generation)).await; + } + let current = history[2].0.as_str(); + assert_tagging_version(&client, bucket, key, None, Some(current), Some("updated-latest")).await; + client + .delete_object_tagging() + .bucket(bucket) + .key(key) + .version_id(current) + .send() + .await + .expect("clear current tags without creating a version"); + for selector in [None, Some(current)] { + assert_tagging_version(&client, bucket, key, selector, Some(current), None).await; + } + + // Start each read alongside a PUT/exact DELETE pair. A successful response + // may select either snapshot, but must never mix their identity and tags. + let barrier = tokio::sync::Barrier::new(2); + let writer = async { + let mut acknowledged = std::collections::HashMap::new(); + acknowledged.insert(current.to_owned(), None); + for round in 0..12 { + barrier.wait().await; + let generation = format!("race-{round}"); + let put = client + .put_object() + .bucket(bucket) + .key(key) + .tagging(format!("generation={generation}")) + .body(ByteStream::from(generation.as_bytes().to_vec())) + .send() + .await + .expect("concurrent tagged PUT"); + let version = put.version_id().expect("concurrent PUT identity").to_owned(); + client + .delete_object() + .bucket(bucket) + .key(key) + .version_id(&version) + .send() + .await + .expect("delete the transient current version"); + acknowledged.insert(version, Some(generation)); + barrier.wait().await; + } + acknowledged + }; + let reader = async { + let mut observed = Vec::new(); + for _ in 0..12 { + barrier.wait().await; + let tags = client + .get_object_tagging() + .bucket(bucket) + .key(key) + .send() + .await + .expect("current tagging during writes/deletes"); + let version = tags.version_id().expect("concurrent read identity").to_owned(); + assert!(tags.tag_set().len() <= 1); + let generation = tags.tag_set().first().map(|tag| { + assert_eq!(tag.key(), "generation"); + tag.value().to_owned() + }); + observed.push((version, generation)); + barrier.wait().await; + } + observed + }; + let (acknowledged, observed) = tokio::join!(writer, reader); + for (version, generation) in observed { + assert_eq!(acknowledged.get(&version), Some(&generation), "tags must belong to the returned version"); + } + + let marker = client + .delete_object() + .bucket(bucket) + .key(key) + .send() + .await + .expect("create current delete marker"); + let missing_version = uuid::Uuid::new_v4().to_string(); + let marker_version = marker.version_id().expect("delete marker identity"); + for (selector, expected_code) in [ + (None, "NoSuchKey"), + (Some(marker_version), "MethodNotAllowed"), + (Some(missing_version.as_str()), "NoSuchVersion"), + (Some("null"), "NoSuchVersion"), + ] { + let Err(err) = client + .get_object_tagging() + .bucket(bucket) + .key(key) + .set_version_id(selector.map(str::to_owned)) + .send() + .await + else { + panic!("missing versions and markers must not return tags: selector {selector:?}"); + }; + assert_eq!(err.as_service_error().and_then(ProvideErrorMetadata::code), Some(expected_code)); + } + assert_tagging_version(&client, bucket, key, Some(current), Some(current), None).await; + let listed = client + .list_object_versions() + .bucket(bucket) + .prefix(key) + .send() + .await + .expect("list retained history"); + let mut actual: Vec<_> = listed + .versions() + .iter() + .map(|version| version.version_id().expect("listed version")) + .collect(); + let mut expected: Vec<_> = history.iter().map(|(version, _)| version.as_str()).collect(); + actual.sort_unstable(); + expected.sort_unstable(); + assert_eq!(actual, expected, "tag mutations must preserve the original three versions"); + assert_eq!(listed.delete_markers().len(), 1); + server.shutdown().await; +}