mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-03 20:20:29 +00:00
fix(s3): return resolved object tagging version identities (#7757)
This commit is contained in:
@@ -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(())
|
||||
|
||||
@@ -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<Uuid>) -> Option<String> {
|
||||
pub(crate) fn read_response_version_id(version_id: Option<Uuid>) -> Option<String> {
|
||||
version_id.map(|id| {
|
||||
if id.is_nil() {
|
||||
NULL_VERSION_ID.to_owned()
|
||||
|
||||
+26
-15
@@ -198,7 +198,7 @@ impl FS {
|
||||
object: &str,
|
||||
version_id: Option<String>,
|
||||
headers: &http::HeaderMap,
|
||||
) -> Option<TagSet> {
|
||||
) -> Option<GetObjectTaggingOutput> {
|
||||
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<GetObjectTaggingInput>) -> S3Result<S3Response<GetObjectTaggingOutput>> {
|
||||
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),
|
||||
}))
|
||||
}
|
||||
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user