From 123967e729c74903a7f3d6dcc0a388736acd4f6c Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sat, 5 Sep 2026 13:02:23 +0800 Subject: [PATCH 1/4] fix(ecstore): fail closed on an unreadable bucket-targets blob (#7172) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(ecstore): correct sealed-credential test helper parameter type The helper took a HashMap that nothing imports, so the ecstore test target did not compile. * fix(ecstore): fail closed on an unreadable bucket-targets blob An undecodable bucket-targets.json was replaced by an empty BucketTargets, so every replication target of that bucket disappeared, replication stopped, and no caller saw an error. A missing secretKey alone triggers it, because Credentials has no struct-level serde(default). parse_all_configs now retains the failure instead: the raw bytes stay and the typed field stays None, which BucketMetadata::bucket_targets_unreadable reads as "exists but cannot be read" — the same distinction the fabricated marker draws for bucket metadata as a whole. One corrupt sub-config still never fails the metadata load, so an unreadable bucket cannot take down its neighbours or the node. BucketTargetSys records such buckets and answers every targets query with the new BucketRemoteTargetsUnreadable, leaving any snapshot from an earlier readable load in place so in-flight replication is not torn down. The replication heal queue reports Missed rather than scheduling against an empty target set, and the admin listing surfaces the fault instead of an empty list. Refs: rustfs/backlog#2282 * fix(ecstore): report corrupt permissive bucket configs as invalid Audit of the remaining parse_all_configs branches. Policy, versioning, object lock and replication already fail closed at their accessors; encryption, public access block and quota did not, and for those three "absent" is exactly the state that grants something — plaintext storage, anonymous access, unbounded capacity. They now report a stored-but-undecodable payload as invalid rather than as ConfigNotFound, matching the guard the versioning and object-lock accessors already use. The quota enforcement path already refused such a payload; only the metadata read path was misreporting it. The branches left degrading, and the concrete reason each is safe, are recorded in the table on parse_all_configs. Refs: rustfs/backlog#2282 --- .../ecstore/src/bucket/bucket_target_sys.rs | 109 ++++++++--- crates/ecstore/src/bucket/metadata.rs | 170 +++++++++++++++++- crates/ecstore/src/bucket/metadata_sys.rs | 165 ++++++++++++++++- .../bucket/replication/replication_pool.rs | 19 +- .../replication_target_boundary.rs | 3 +- rustfs/src/admin/handlers/replication.rs | 12 +- 6 files changed, 438 insertions(+), 40 deletions(-) diff --git a/crates/ecstore/src/bucket/bucket_target_sys.rs b/crates/ecstore/src/bucket/bucket_target_sys.rs index dca41cf72..db4e80f29 100644 --- a/crates/ecstore/src/bucket/bucket_target_sys.rs +++ b/crates/ecstore/src/bucket/bucket_target_sys.rs @@ -59,7 +59,7 @@ use rustfs_utils::http::{ insert_header, }; use serde::{Deserialize, Serialize}; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::error::Error; use std::fmt; use std::str::FromStr as _; @@ -376,6 +376,11 @@ pub struct BucketTargetSys { /// [`SsecPassthroughCapability`]; reset alongside `arn_remotes_map`. ssec_passthrough_map: Arc>>, pub targets_map: Arc>>>, + /// Buckets whose persisted `bucket-targets.json` exists but cannot be + /// decoded (rustfs/backlog#2282). Written under the bucket's update mutex + /// alongside `targets_map`, and read before it so an unreadable + /// configuration surfaces as a typed error instead of an empty target set. + unreadable_targets: Arc>>, pub h_mutex: Arc>>, target_h_mutex: Arc>>, pub hc_client: Arc, @@ -419,6 +424,7 @@ impl BucketTargetSys { arn_remotes_map: Arc::new(RwLock::new(HashMap::new())), ssec_passthrough_map: Arc::new(RwLock::new(HashMap::new())), targets_map: Arc::new(RwLock::new(HashMap::new())), + unreadable_targets: Arc::new(RwLock::new(HashSet::new())), h_mutex: Arc::new(RwLock::new(HashMap::new())), target_h_mutex: Arc::new(RwLock::new(HashMap::new())), hc_client: Arc::new(build_health_check_client()), @@ -628,30 +634,40 @@ impl BucketTargetSys { health_map.clone() } - pub async fn list_targets(&self, bucket: &str, arn_type: &str) -> Vec { + /// Targets of one bucket, or of every bucket when `bucket` is empty. + /// + /// A bucket that simply has no targets yields an empty list; a bucket + /// whose persisted configuration cannot be decoded is an error, so an + /// admin listing reports the fault instead of an empty list that reads as + /// "replication is not configured" (rustfs/backlog#2282). + pub async fn list_targets(&self, bucket: &str, arn_type: &str) -> Result, BucketTargetError> { let health_stats = self.target_health_stats().await; let mut targets = Vec::new(); if !bucket.is_empty() { - if let Ok(bucket_targets) = self.list_bucket_targets(bucket).await { - for mut target in bucket_targets.targets { - if arn_type.is_empty() || target.target_type.to_string() == arn_type { - if let Some(health) = health_stats.get(&target.arn) { - target.total_downtime = health.offline_duration; - target.online = health.online; - target.last_online = health.last_online; - target.latency = target::LatencyStat { - curr: health.latency.curr, - avg: health.latency.avg, - max: health.latency.peak, - }; - target.offline_count = health.offline_count; + match self.list_bucket_targets(bucket).await { + Ok(bucket_targets) => { + for mut target in bucket_targets.targets { + if arn_type.is_empty() || target.target_type.to_string() == arn_type { + if let Some(health) = health_stats.get(&target.arn) { + target.total_downtime = health.offline_duration; + target.online = health.online; + target.last_online = health.last_online; + target.latency = target::LatencyStat { + curr: health.latency.curr, + avg: health.latency.avg, + max: health.latency.peak, + }; + target.offline_count = health.offline_count; + } + targets.push(target); } - targets.push(target); } } + Err(BucketTargetError::BucketRemoteTargetNotFound { .. }) => {} + Err(err) => return Err(err), } - return targets; + return Ok(targets); } let targets_map = self.targets_map.read().await; @@ -674,10 +690,16 @@ impl BucketTargetSys { } } - targets + Ok(targets) } pub async fn list_bucket_targets(&self, bucket: &str) -> Result { + if self.unreadable_targets.read().await.contains(bucket) { + return Err(BucketTargetError::BucketRemoteTargetsUnreadable { + bucket: bucket.to_string(), + }); + } + let targets_map = self.targets_map.read().await; if let Some(targets) = targets_map.get(bucket) { Ok(BucketTargets { @@ -690,13 +712,30 @@ impl BucketTargetSys { } } + /// Record that this bucket's persisted targets configuration exists but + /// cannot be decoded (rustfs/backlog#2282). + /// + /// Any snapshot published from an earlier readable load is deliberately + /// left in place: withdrawing it would produce exactly the silent "no + /// targets configured" state this marker exists to prevent. The marker is + /// cleared by the next successful publish, which is what makes a repaired + /// configuration take effect without a restart. + pub async fn mark_targets_unreadable(&self, bucket: &str) { + let update_mutex = self.target_update_mutex(bucket).await; + let _update_guard = update_mutex.lock().await; + + self.unreadable_targets.write().await.insert(bucket.to_string()); + } + pub async fn delete(&self, bucket: &str) { let update_mutex = self.target_update_mutex(bucket).await; let _update_guard = update_mutex.lock().await; - // Lock order: targets_map, then arn_remotes_map, then target_h_mutex, - // then ssec_passthrough_map (always last; also taken standalone by the - // capability accessors). + // Lock order: unreadable_targets, then targets_map, then + // arn_remotes_map, then target_h_mutex, then ssec_passthrough_map + // (always last; also taken standalone by the capability accessors). + self.unreadable_targets.write().await.remove(bucket); + let mut targets_map = self.targets_map.write().await; let mut arn_remotes_map = self.arn_remotes_map.write().await; let mut health_map = self.target_h_mutex.write().await; @@ -1093,6 +1132,11 @@ impl BucketTargetSys { /// Keeping persisted-config reads under the same mutex prevents a stale /// reload from overwriting a concurrent credential rotation. async fn update_all_targets_locked(&self, bucket: &str, targets: Option<&BucketTargets>) { + // Reaching here means the persisted configuration decoded, so the + // unreadable marker (if any) is stale. Cleared before the maps below + // so `unreadable_targets` stays the outermost of this module's locks. + self.unreadable_targets.write().await.remove(bucket); + let mut clients = Vec::new(); if let Some(new_targets) = targets { for target in &new_targets.targets { @@ -1100,9 +1144,9 @@ impl BucketTargetSys { } } - // Lock order: targets_map, then arn_remotes_map, then target_h_mutex, - // then ssec_passthrough_map (always last; also taken standalone by the - // capability accessors). + // Lock order: unreadable_targets (above), then targets_map, then + // arn_remotes_map, then target_h_mutex, then ssec_passthrough_map + // (always last; also taken standalone by the capability accessors). let mut targets_map = self.targets_map.write().await; let mut arn_remotes_map = self.arn_remotes_map.write().await; let mut health_map = self.target_h_mutex.write().await; @@ -1161,6 +1205,11 @@ impl BucketTargetSys { } pub async fn set(&self, bucket: &str, meta: &BucketMetadata) { + if meta.bucket_targets_unreadable() { + self.mark_targets_unreadable(bucket).await; + return; + } + let Some(config) = &meta.bucket_target_config else { return; }; @@ -2276,6 +2325,13 @@ pub enum BucketTargetError { BucketRemoteTargetNotFound { bucket: String, }, + /// The bucket's persisted targets configuration exists but cannot be + /// decoded. Distinct from `BucketRemoteTargetNotFound`, which means the + /// bucket genuinely has no targets: callers must not degrade this one to + /// an empty target set (rustfs/backlog#2282). + BucketRemoteTargetsUnreadable { + bucket: String, + }, BucketRemoteArnTypeInvalid { bucket: String, }, @@ -2309,6 +2365,9 @@ impl fmt::Display for BucketTargetError { BucketTargetError::BucketRemoteTargetNotFound { bucket } => { write!(f, "Remote target not found for bucket: {bucket}") } + BucketTargetError::BucketRemoteTargetsUnreadable { bucket } => { + write!(f, "Persisted replication target configuration is unreadable for bucket: {bucket}") + } BucketTargetError::BucketRemoteArnTypeInvalid { bucket } => { write!(f, "Invalid ARN type for bucket: {bucket}") } @@ -3256,7 +3315,7 @@ mod tests { }], ); - let targets = sys.list_targets("", "").await; + let targets = sys.list_targets("", "").await.expect("listing every bucket's targets"); assert_eq!(targets.len(), 1); assert!(!targets[0].online); diff --git a/crates/ecstore/src/bucket/metadata.rs b/crates/ecstore/src/bucket/metadata.rs index bd10f0f63..41b2afdc5 100644 --- a/crates/ecstore/src/bucket/metadata.rs +++ b/crates/ecstore/src/bucket/metadata.rs @@ -477,6 +477,18 @@ impl BucketMetadata { !self.table_bucket_config_json.is_empty() } + /// `bucket-targets.json` is stored for this bucket but this build cannot + /// decode it. + /// + /// Keeps "no replication targets configured" and "the target + /// configuration cannot be read" apart, the same distinction the + /// `fabricated` marker draws for the bucket metadata as a whole. Only + /// meaningful after [`Self::parse_all_configs`] has run; readers must fail + /// closed on `true` instead of serving an empty target set. + pub fn bucket_targets_unreadable(&self) -> bool { + !self.bucket_targets_config_json.is_empty() && self.bucket_target_config.is_none() + } + /// Parsed per-bucket durability override, if a valid one is stored. /// /// Absent/empty/unparsable payloads all mean "no override" (the bucket @@ -964,7 +976,32 @@ impl BucketMetadata { Ok(()) } - fn parse_all_configs(&mut self) -> Result<()> { + /// Decode every stored sub-configuration into its typed field. + /// + /// A decode failure never fails the whole load: this runs on every bucket + /// metadata read, including startup and peer reload, so one bucket's + /// corrupt sub-configuration must not make the bucket — or the node — + /// unloadable. Instead the failure is *retained*: the raw bytes stay + /// untouched and the typed field stays `None`, so `!raw.is_empty() && + /// typed.is_none()` is the durable "exists but cannot be read" signal that + /// each accessor keys off. Which accessors must fail closed on it: + /// + /// | Config | Verdict | + /// |---|---| + /// | policy | Fails closed: `get_bucket_policy` re-parses the raw JSON and propagates the error; `get_bucket_policy_raw` returns the stored bytes. | + /// | object lock | Fails closed in `object_lock_config_state_from_authoritative_metadata`; a retention decision may never be taken on a guess. | + /// | versioning | Fails closed in `get_versioning_config`; guessing Unversioned would make delete markers and version ids diverge from what is on disk. | + /// | replication | Fails closed in `get_replication_config`. | + /// | bucket targets | Fails closed in `get_bucket_targets_config`, and `sync_bucket_target_sys` marks the bucket unreadable in `BucketTargetSys` instead of publishing an empty target set (rustfs/backlog#2282). | + /// | encryption | Fails closed in `get_sse_config`: degrading to "no default encryption" stores plaintext objects the operator required to be encrypted. | + /// | public access block | Fails closed in `get_public_access_block_config`: degrading grants the anonymous access the operator asked to block. | + /// | quota | Fails closed in `get_quota_config`; the enforcement path in `quota::checker` already re-parses the raw JSON and refuses on error. | + /// | lifecycle | Safe to degrade: no rules means no expiration and no transition, so nothing is deleted or moved on the strength of an unreadable rule set. The bucket keeps serving reads and writes. | + /// | notification | Safe to degrade: events are an outbound side channel; no consumer draws a durability or authorization conclusion from their absence. | + /// | tagging | Safe to degrade: bucket tags are cost-allocation labels here; object-level tag conditions come from object metadata, not this blob. | + /// | CORS | Safe to degrade: an absent CORS configuration rejects cross-origin browser requests, which is already the restrictive direction. | + /// | logging, website, accelerate, request payment, bucket ACL | Safe to degrade: each only shapes an optional response or an optional side channel, and none of them authorizes an action or decides whether data is retained. | + pub(super) fn parse_all_configs(&mut self) -> Result<()> { if let Err(e) = self.parse_policy_config() { tracing::warn!( event = "bucket_metadata_parse_failed", @@ -1088,20 +1125,26 @@ impl BucketMetadata { "Failed to parse bucket metadata config" ); } + // A stored targets blob that cannot be decoded must not collapse into + // the empty target set: that is indistinguishable from "no replication + // configured", so replication stops and no caller ever sees an error + // (rustfs/backlog#2282). Leaving the typed field `None` while the raw + // bytes stay non-empty is the retained parse failure every targets + // reader keys off; the bytes are preserved so the configuration is + // still recoverable. + self.bucket_target_config = None; if !self.bucket_targets_config_json.is_empty() { - if let Err(e) = serde_json::from_slice::(&self.bucket_targets_config_json) - .map(|t| self.bucket_target_config = Some(t)) - { - tracing::warn!( + match serde_json::from_slice::(&self.bucket_targets_config_json) { + Ok(targets) => self.bucket_target_config = Some(targets), + Err(e) => tracing::error!( event = "bucket_metadata_parse_failed", component = "ecstore", subsystem = "bucket_metadata", bucket = %self.name, config = "bucket_targets", error = %e, - "Failed to parse bucket metadata config" - ); - self.bucket_target_config = Some(BucketTargets::default()); + "Bucket replication targets are unreadable; replication for this bucket fails closed" + ), } } else { self.bucket_target_config = Some(BucketTargets::default()); @@ -1535,6 +1578,117 @@ mod test { assert_eq!(bucket_targets.targets[0].target_bucket, "target-bucket"); } + /// rustfs/backlog#2282: a stored targets blob this build cannot decode + /// must not become the empty target set, and must stay distinguishable + /// from a bucket that never configured a target. + #[test] + fn unreadable_bucket_targets_never_degrade_to_an_empty_target_set() { + let truncated = br#"{"targets":[{"endpoint":"s3.example.com","#.to_vec(); + let mut corrupt = BucketMetadata::new("corrupt-targets"); + corrupt.bucket_targets_config_json = truncated.clone(); + + corrupt + .parse_all_configs() + .expect("one unreadable sub-config must not fail the whole metadata load"); + + assert!( + corrupt.bucket_target_config.is_none(), + "an undecodable targets blob must not produce a target set at all" + ); + assert!(corrupt.bucket_targets_unreadable()); + assert_eq!( + corrupt.bucket_targets_config_json, truncated, + "the raw bytes must survive so the configuration stays recoverable" + ); + + // The genuinely-absent case is unchanged, and the two now diverge. + let mut absent = BucketMetadata::new("no-targets"); + absent.parse_all_configs().expect("absent targets parse"); + assert!( + absent.bucket_target_config.as_ref().is_some_and(BucketTargets::is_empty), + "a bucket that configured no target still reads as an empty target set" + ); + assert!(!absent.bucket_targets_unreadable()); + } + + /// `Credentials` carries no struct-level `serde(default)`, so one target + /// missing `secretKey` is a hard parse error for the whole document. That + /// must surface as "unreadable", never as "no targets configured". + #[test] + fn bucket_targets_missing_secret_key_are_unreadable_not_empty() { + let mut bm = BucketMetadata::new("missing-secret-key"); + bm.bucket_targets_config_json = br#"{"targets":[{"endpoint":"s3.example.com","targetbucket":"remote","arn":"arn:rustfs:replication:us-east-1:src:1","credentials":{"accessKey":"AKIAEXAMPLE"}}]}"#.to_vec(); + + bm.parse_all_configs() + .expect("a rejected targets document must not fail the whole metadata load"); + + assert!( + bm.bucket_targets_unreadable(), + "a targets document rejected for a missing secretKey is unreadable, not empty" + ); + assert!(bm.bucket_target_config.is_none()); + } + + /// The invariant every branch of `parse_all_configs` shares: a stored but + /// undecodable payload keeps its raw bytes and leaves the typed field + /// `None`, so no branch fabricates a value. What a reader may then do with + /// that state is decided per config; see the table on `parse_all_configs`. + #[test] + fn every_config_branch_retains_its_parse_failure_instead_of_defaulting() { + let malformed_xml = b">) { } async fn sync_bucket_target_sys(bucket: &str, bm: &BucketMetadata) { + if bm.bucket_targets_unreadable() { + // "The configuration cannot be read" is not "no targets configured". + // Publishing an empty snapshot here is what silently stopped + // replication (rustfs/backlog#2282): mark the bucket instead, so every + // targets reader gets a typed error, and leave any snapshot from an + // earlier readable load in place rather than withdrawing it. + BucketTargetSys::get().mark_targets_unreadable(bucket).await; + return; + } + BucketTargetSys::get() .update_all_targets(bucket, bm.bucket_target_config.as_ref()) .await; @@ -2118,7 +2128,9 @@ impl BucketMetadataSys { pub async fn get_public_access_block_config(&self, bucket: &str) -> Result<(PublicAccessBlockConfiguration, OffsetDateTime)> { let (bm, _) = self.get_config(bucket).await?; - if let Some(config) = &bm.public_access_block_config { + if !bm.public_access_block_config_xml.is_empty() && bm.public_access_block_config.is_none() { + Err(Error::other("persisted bucket public access block configuration is invalid")) + } else if let Some(config) = &bm.public_access_block_config { Ok((config.clone(), bm.public_access_block_config_updated_at)) } else { Err(Error::ConfigNotFound) @@ -2429,7 +2441,9 @@ impl BucketMetadataSys { pub async fn get_sse_config(&self, bucket: &str) -> Result<(ServerSideEncryptionConfiguration, OffsetDateTime)> { let (bm, _) = self.get_config(bucket).await?; - if let Some(config) = &bm.sse_config { + if !bm.encryption_config_xml.is_empty() && bm.sse_config.is_none() { + Err(Error::other("persisted bucket encryption configuration is invalid")) + } else if let Some(config) = &bm.sse_config { Ok((config.clone(), bm.encryption_config_updated_at)) } else { Err(Error::ConfigNotFound) @@ -2500,7 +2514,9 @@ impl BucketMetadataSys { pub async fn get_quota_config(&self, bucket: &str) -> Result<(BucketQuota, OffsetDateTime)> { let (bm, _) = self.get_config(bucket).await?; - if let Some(config) = &bm.quota_config { + if !bm.quota_config_json.is_empty() && bm.quota_config.is_none() { + Err(Error::other("persisted bucket quota configuration is invalid")) + } else if let Some(config) = &bm.quota_config { Ok((config.clone(), bm.quota_config_updated_at)) } else { Err(Error::ConfigNotFound) @@ -2522,7 +2538,9 @@ impl BucketMetadataSys { pub async fn get_bucket_targets_config(&self, bucket: &str) -> Result { let (bm, _) = self.get_config(bucket).await?; - if let Some(config) = &bm.bucket_target_config { + if bm.bucket_targets_unreadable() { + Err(Error::other("persisted bucket replication target configuration is invalid")) + } else if let Some(config) = &bm.bucket_target_config { Ok(config.clone()) } else { Err(Error::ConfigNotFound) @@ -2593,6 +2611,7 @@ pub(crate) mod test_support { mod tests { use super::test_support::isolated_store_over_temp_disks; use super::*; + use crate::bucket::bucket_target_sys::BucketTargetError; use crate::bucket::metadata::{ BUCKET_ACCELERATE_CONFIG, BUCKET_CORS_CONFIG, BUCKET_LIFECYCLE_CONFIG, BUCKET_LOGGING_CONFIG, BUCKET_NOTIFICATION_CONFIG, BUCKET_POLICY_CONFIG, BUCKET_PUBLIC_ACCESS_BLOCK_CONFIG, BUCKET_REPLICATION_CONFIG, BUCKET_REQUEST_PAYMENT_CONFIG, @@ -2788,6 +2807,36 @@ mod tests { ); } + /// The `parse_all_configs` audit (rustfs/backlog#2282): every accessor + /// whose configuration grants something — plaintext storage, anonymous + /// access, capacity, replication targets — reports a corrupt payload as + /// invalid rather than as absent, because "absent" is what grants it. + #[tokio::test] + async fn malformed_permissive_configs_are_not_reported_as_absent() { + let (_dirs, ecstore) = isolated_store_over_temp_disks().await; + let sys = BucketMetadataSys::new(ecstore); + let bucket = "malformed-permissive-config"; + let mut metadata = BucketMetadata::new(bucket); + metadata.encryption_config_xml = b" Some(targets), + // A bucket whose persisted target configuration cannot be decoded has + // an unknown target set, not an empty one: scheduling against `None` + // here would drop every heal for it without a trace + // (rustfs/backlog#2282). Report it missed so the object is retried + // once the configuration is readable again. + Err(BucketTargetError::BucketRemoteTargetsUnreadable { .. }) => { + warn!( + event = EVENT_REPLICATION_CONFIG_LOOKUP_SKIPPED, + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_REPLICATION, + bucket, + reason = "target_config_unreadable", + "Bucket replication targets are unreadable; replication heal queue fails closed" + ); + + return ReplicationQueueAdmission::Missed; + } Err(err) => { debug!( event = EVENT_REPLICATION_CONFIG_LOOKUP_SKIPPED, diff --git a/crates/ecstore/src/bucket/replication/replication_target_boundary.rs b/crates/ecstore/src/bucket/replication/replication_target_boundary.rs index a5cbf6b34..f7531319e 100644 --- a/crates/ecstore/src/bucket/replication/replication_target_boundary.rs +++ b/crates/ecstore/src/bucket/replication/replication_target_boundary.rs @@ -15,7 +15,8 @@ use std::collections::HashMap; use std::sync::Arc; -use crate::bucket::bucket_target_sys::{BucketTargetError, BucketTargetSys}; +pub(crate) use crate::bucket::bucket_target_sys::BucketTargetError; +use crate::bucket::bucket_target_sys::BucketTargetSys; use aws_sdk_s3::operation::head_object::HeadObjectOutput; use aws_sdk_s3::types::{ObjectLockLegalHoldStatus, ObjectLockRetentionMode}; use http::HeaderMap; diff --git a/rustfs/src/admin/handlers/replication.rs b/rustfs/src/admin/handlers/replication.rs index b2b2b80dc..b7bf5ef78 100644 --- a/rustfs/src/admin/handlers/replication.rs +++ b/rustfs/src/admin/handlers/replication.rs @@ -121,6 +121,11 @@ fn map_bucket_target_error(err: BucketTargetError) -> S3Error { | BucketTargetError::BucketRemoteRemoveDisallowed { .. } => { S3Error::with_message(S3ErrorCode::InvalidRequest, err.to_string()) } + // A stored target configuration this node cannot decode is a + // server-side data fault, not a bad request (rustfs/backlog#2282). + BucketTargetError::BucketRemoteTargetsUnreadable { .. } => { + S3Error::with_message(S3ErrorCode::InternalError, err.to_string()) + } BucketTargetError::Io(io_err) => S3Error::with_message(S3ErrorCode::InternalError, io_err.to_string()), } } @@ -753,7 +758,12 @@ impl Operation for ListRemoteTargetHandler { .map_err(ApiError::from)?; let sys = BucketTargetSys::get(); - let targets = sys.list_targets(bucket, "").await; + // An unreadable targets configuration must not be reported as an + // empty target list (rustfs/backlog#2282). + let targets = sys.list_targets(bucket, "").await.map_err(|e| { + error!("list remote targets failed: {}", e); + map_bucket_target_error(e) + })?; let targets: Vec<_> = targets .iter() From 4dbc58887afe4c0e38dc11704f403479a9f7c79e Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 5 Sep 2026 14:00:14 +0800 Subject: [PATCH 2/4] fix(tier): probe legacy transition version state (#7138) --- .../bucket/lifecycle/bucket_lifecycle_ops.rs | 782 +++++++++++++++++- crates/ecstore/src/services/tier/test_util.rs | 5 +- crates/ecstore/src/services/tier/tier.rs | 13 + .../ecstore/src/services/tier/warm_backend.rs | 53 +- .../src/services/tier/warm_backend_s3.rs | 43 +- crates/ecstore/src/store/init.rs | 160 +++- crates/filemeta/src/filemeta/version.rs | 125 ++- 7 files changed, 1133 insertions(+), 48 deletions(-) diff --git a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs index 7649e3ac9..dbb5dd154 100644 --- a/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs +++ b/crates/ecstore/src/bucket/lifecycle/bucket_lifecycle_ops.rs @@ -584,33 +584,173 @@ impl ExpiryOp for FreeVersionTask { } } +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum TransitionDeleteVersionPlan { + Direct { version_id_exact: bool }, + ProbeLegacyUnknown, +} + +fn legacy_transition_version_state_missing(oi: &ObjectInfo) -> Result { + use rustfs_utils::http::metadata_compat::{ + SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, contains_key_str, get_consistent_str, + }; + + if !contains_key_str(&oi.user_defined, SUFFIX_TRANSITIONED_VERSION_STATE) { + let version_key_present = contains_key_str(&oi.user_defined, SUFFIX_TRANSITIONED_VERSION_ID); + if version_key_present { + if oi.transitioned_object.version_id.is_empty() { + let has_non_empty_version = oi.user_defined.iter().any(|(key, value)| { + rustfs_utils::http::metadata_compat::strip_internal_prefix_preserving_case(key) + .is_some_and(|suffix| suffix.eq_ignore_ascii_case(SUFFIX_TRANSITIONED_VERSION_ID)) + && !value.is_empty() + }); + if !has_non_empty_version { + // MinIO writes the transitioned-versionID key with an empty value + // for unversioned tier objects. The backend probe remains the proof. + return Ok(true); + } + } else if get_consistent_str(&oi.user_defined, SUFFIX_TRANSITIONED_VERSION_ID) + == Some(oi.transitioned_object.version_id.as_str()) + { + return Ok(true); + } + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "legacy remote tier version metadata is conflicting or malformed", + )); + } + if !oi.transitioned_object.version_id.is_empty() { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "legacy remote tier version metadata is missing or inconsistent", + )); + } + return Ok(true); + } + let persisted = get_consistent_str(&oi.user_defined, SUFFIX_TRANSITIONED_VERSION_STATE).ok_or_else(|| { + std::io::Error::new( + std::io::ErrorKind::InvalidData, + "remote tier object has conflicting transition version state metadata", + ) + })?; + if persisted != oi.transition_version_state.as_str() { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "remote tier object transition version state metadata changed during decoding", + )); + } + Ok(false) +} + +fn transition_remote_version_delete_plan(oi: &ObjectInfo) -> Result { + match oi.transition_version_state { + rustfs_filemeta::TransitionVersionState::Unknown => { + if legacy_transition_version_state_missing(oi)? { + Ok(TransitionDeleteVersionPlan::ProbeLegacyUnknown) + } else { + validate_transition_remote_version(oi) + .map(|version_id_exact| TransitionDeleteVersionPlan::Direct { version_id_exact }) + } + } + _ => validate_transition_remote_version(oi) + .map(|version_id_exact| TransitionDeleteVersionPlan::Direct { version_id_exact }), + } +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +struct ResolvedTransitionDeleteVersion { + version_id_exact: bool, + remote_already_missing: bool, +} + async fn acquire_free_version_tier_lease( oi: &ObjectInfo, tier_config_mgr: &Arc>, -) -> Result<(TierOperationLease, bool), std::io::Error> { - let version_id_exact = validate_transition_remote_version(oi)?; +) -> Result<(TierOperationLease, TransitionDeleteVersionPlan), std::io::Error> { + let delete_plan = transition_remote_version_delete_plan(oi)?; let identity = tier_destination_id_from_metadata(&oi.user_defined)? .ok_or_else(|| std::io::Error::other("tier free-version has no durable backend identity"))?; let lease = TierConfigMgr::acquire_operation_lease_for_backend_identity(tier_config_mgr, &oi.transitioned_object.tier, identity) .await .map_err(std::io::Error::other)?; - Ok((lease, version_id_exact)) + Ok((lease, delete_plan)) +} + +async fn resolve_transition_delete_version_plan( + oi: &ObjectInfo, + lease: &TierOperationLease, + delete_plan: TransitionDeleteVersionPlan, +) -> Result { + match delete_plan { + TransitionDeleteVersionPlan::Direct { version_id_exact } => Ok(ResolvedTransitionDeleteVersion { + version_id_exact, + remote_already_missing: false, + }), + TransitionDeleteVersionPlan::ProbeLegacyUnknown => { + let expected_version = oi.transitioned_object.version_id.as_str(); + if expected_version.is_empty() { + return Err(std::io::Error::new( + std::io::ErrorKind::WouldBlock, + "remote tier cannot safely delete a legacy object without an exact version ID", + )); + } + let probe = lease + .probe_transition_version(&oi.transitioned_object.name, expected_version) + .await?; + match (expected_version, probe) { + (expected, crate::services::tier::warm_backend::TransitionCandidateProbe::VersionedPresent(actual)) + if expected == actual => + { + lease.validate_remote_version_id(expected)?; + Ok(ResolvedTransitionDeleteVersion { + version_id_exact: true, + remote_already_missing: false, + }) + } + (_, crate::services::tier::warm_backend::TransitionCandidateProbe::Missing) => { + Ok(ResolvedTransitionDeleteVersion { + version_id_exact: false, + remote_already_missing: true, + }) + } + (_, crate::services::tier::warm_backend::TransitionCandidateProbe::Unsupported) => Err(std::io::Error::new( + std::io::ErrorKind::Unsupported, + "remote tier cannot prove legacy transition delete state", + )), + _ => Err(std::io::Error::new( + std::io::ErrorKind::WouldBlock, + "remote tier object version state is unknown", + )), + } + } + } +} + +async fn execute_resolved_transition_delete( + oi: &ObjectInfo, + lease: &TierOperationLease, + resolved: ResolvedTransitionDeleteVersion, +) -> Result<(), std::io::Error> { + if !resolved.remote_already_missing { + delete_object_from_remote_tier_with_lease_idempotent( + &oi.transitioned_object.name, + &oi.transitioned_object.version_id, + lease, + resolved.version_id_exact, + ) + .await?; + } + Ok(()) } async fn delete_free_version_remote_object_with_lease( oi: &ObjectInfo, lease: &TierOperationLease, - version_id_exact: bool, + delete_plan: TransitionDeleteVersionPlan, ) -> Result<(), std::io::Error> { - delete_object_from_remote_tier_with_lease_idempotent( - &oi.transitioned_object.name, - &oi.transitioned_object.version_id, - lease, - version_id_exact, - ) - .await?; - Ok(()) + let resolved = resolve_transition_delete_version_plan(oi, lease, delete_plan).await?; + execute_resolved_transition_delete(oi, lease, resolved).await } fn free_version_physical_topology_generation(api: &ECStore) -> String { @@ -641,6 +781,16 @@ fn free_version_remote_tuple_matches(candidate: &ObjectInfo, expected: &ObjectIn if candidate.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown || expected.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown { + let candidate_legacy_missing = legacy_transition_version_state_missing(candidate)?; + let expected_legacy_missing = legacy_transition_version_state_missing(expected)?; + if candidate.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown + && expected.transition_version_state == rustfs_filemeta::TransitionVersionState::Unknown + && candidate_legacy_missing + && expected_legacy_missing + && candidate.transitioned_object.version_id == expected.transitioned_object.version_id + { + return Ok(true); + } return Err(std::io::Error::new( std::io::ErrorKind::WouldBlock, "tier free-version remote version state is unknown", @@ -716,7 +866,7 @@ async fn cleanup_free_version_exact(api: Arc, oi: &ObjectInfo, cancel: .acquire_bucket_lifecycle_read_lock(&oi.bucket) .await .map_err(std::io::Error::other)?; - let (lease, version_id_exact) = acquire_free_version_tier_lease(oi, &api.tier_config_mgr()).await?; + let (lease, delete_plan) = acquire_free_version_tier_lease(oi, &api.tier_config_mgr()).await?; let local_object = encode_dir_object(&oi.name); let object_guards = api .acquire_all_physical_object_write_locks("tier_free_version_cleanup", &oi.bucket, &local_object) @@ -734,16 +884,30 @@ async fn cleanup_free_version_exact(api: Arc, oi: &ObjectInfo, cancel: "tier free-version cleanup fence is invalid before remote delete", )); } + let resolved = tokio::select! { + _ = cancel.cancelled() => { + return Err(std::io::Error::new(std::io::ErrorKind::Interrupted, "tier free-version cleanup was cancelled")); + } + result = tokio::time::timeout_at(deadline, resolve_transition_delete_version_plan(oi, &lease, delete_plan)) => { + result.map_err(|_| { + std::io::Error::new(std::io::ErrorKind::TimedOut, "tier free-version remote probe timed out") + })?? + } + }; + if !free_version_cleanup_fences_current(&topology_generation, &api, &bucket_guard, &object_guards, &lease, cancel, deadline) { + return Err(std::io::Error::new( + std::io::ErrorKind::WouldBlock, + "tier free-version cleanup fence changed after remote probe", + )); + } tokio::select! { _ = cancel.cancelled() => { return Err(std::io::Error::new(std::io::ErrorKind::Interrupted, "tier free-version cleanup was cancelled")); } - result = tokio::time::timeout_at( - deadline, - delete_free_version_remote_object_with_lease(oi, &lease, version_id_exact), - ) => { - result - .map_err(|_| std::io::Error::new(std::io::ErrorKind::TimedOut, "tier free-version remote delete timed out"))??; + result = tokio::time::timeout_at(deadline, execute_resolved_transition_delete(oi, &lease, resolved)) => { + result.map_err(|_| { + std::io::Error::new(std::io::ErrorKind::TimedOut, "tier free-version remote delete timed out") + })??; } } if !free_version_cleanup_fences_current(&topology_generation, &api, &bucket_guard, &object_guards, &lease, cancel, deadline) { @@ -791,8 +955,8 @@ async fn delete_free_version_remote_object( oi: &ObjectInfo, tier_config_mgr: &Arc>, ) -> Result<(), std::io::Error> { - let (lease, version_id_exact) = acquire_free_version_tier_lease(oi, tier_config_mgr).await?; - delete_free_version_remote_object_with_lease(oi, &lease, version_id_exact).await + let (lease, delete_plan) = acquire_free_version_tier_lease(oi, tier_config_mgr).await?; + delete_free_version_remote_object_with_lease(oi, &lease, delete_plan).await } #[allow( @@ -808,8 +972,8 @@ where F: FnOnce() -> Fut, Fut: std::future::Future, { - let (lease, version_id_exact) = acquire_free_version_tier_lease(oi, tier_config_mgr).await?; - delete_free_version_remote_object_with_lease(oi, &lease, version_id_exact).await?; + let (lease, delete_plan) = acquire_free_version_tier_lease(oi, tier_config_mgr).await?; + delete_free_version_remote_object_with_lease(oi, &lease, delete_plan).await?; let result = delete_local().await; drop(lease); Ok(result) @@ -4688,6 +4852,39 @@ fn validate_transition_remote_version(oi: &ObjectInfo) -> Result Result { + let version = oi.transitioned_object.version_id.as_str(); + match oi.transition_version_state { + rustfs_filemeta::TransitionVersionState::Unknown => { + if !legacy_transition_version_state_missing(oi)? { + return validate_transition_remote_version(oi).map(|_| TransitionReadVersionPlan::Direct); + } + if version.is_empty() { + Ok(TransitionReadVersionPlan::ProbeLegacyUnversioned) + } else { + Ok(TransitionReadVersionPlan::Direct) + } + } + rustfs_filemeta::TransitionVersionState::KnownDisabled if version.is_empty() => Ok(TransitionReadVersionPlan::Direct), + rustfs_filemeta::TransitionVersionState::SuspendedNull if version == "null" => Ok(TransitionReadVersionPlan::Direct), + rustfs_filemeta::TransitionVersionState::Exact if !version.is_empty() && version != "null" => { + Ok(TransitionReadVersionPlan::Direct) + } + _ => Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "remote tier object version state conflicts with its version ID", + )), + } +} + // The resolver joins the tier manager as the second injected port this read // needs; grouping the request half into a struct would churn every call site of // a bug fix. @@ -4702,7 +4899,12 @@ pub(crate) async fn get_transitioned_object_reader_with_tier_manager( tier_config_mgr: &Arc>, resolver: Option<&dyn ObjectEncryptionResolver>, ) -> Result { - validate_transition_remote_version(oi)?; + let read_plan = transition_remote_version_read_plan(oi)?; + // Reject invalid ranges and encryption requests before a compatibility + // probe can amplify them into remote listing work. + let plan = ReadPlan::build_for_request(rs.clone(), oi, opts, h, resolver) + .await + .map_err(|err| std::io::Error::other(format!("building the read plan for {bucket}/{object} failed: {err}")))?; let expected_identity = tier_destination_id_from_metadata(&oi.user_defined)?; let lease = match expected_identity { Some(identity) => { @@ -4716,7 +4918,36 @@ pub(crate) async fn get_transitioned_object_reader_with_tier_manager( Err(err) => return Err(std::io::Error::other(err)), }; - tgt_client.validate_remote_version_id(&oi.transitioned_object.version_id)?; + match read_plan { + TransitionReadVersionPlan::Direct => { + tgt_client.validate_remote_version_id(&oi.transitioned_object.version_id)?; + } + TransitionReadVersionPlan::ProbeLegacyUnversioned => { + // RUSTFS_COMPAT_TODO(backlog#2203): remove operation-time probing + // after an admin reconcile can persist every proven legacy state. + let probe = tokio::time::timeout( + LEGACY_TRANSITION_READ_PROBE_TIMEOUT, + tgt_client.probe_transition_candidate(&oi.transitioned_object.name), + ) + .await + .map_err(|_| std::io::Error::new(std::io::ErrorKind::TimedOut, "legacy remote tier version probe timed out"))??; + match probe { + crate::services::tier::warm_backend::TransitionCandidateProbe::UnversionedPresent => {} + crate::services::tier::warm_backend::TransitionCandidateProbe::Unsupported => { + return Err(std::io::Error::new( + std::io::ErrorKind::Unsupported, + "remote tier cannot prove legacy unversioned transition state", + )); + } + _ => { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + "remote tier object version state is unknown", + )); + } + } + } + } // The same read plan the local path uses, so the tier fetch is positioned in // the object's *stored* coordinate system and the stream is handed the same @@ -4724,9 +4955,6 @@ pub(crate) async fn get_transitioned_object_reader_with_tier_manager( // through a plaintext-coordinate range and skipping the transform is how a // transitioned SSE object used to come back as silently corrupt bytes of the // right length (rustfs/rustfs#6025). - let plan = ReadPlan::build_for_request(rs.clone(), oi, opts, h, resolver) - .await - .map_err(|err| std::io::Error::other(format!("building the read plan for {bucket}/{object} failed: {err}")))?; let (off, length) = (plan.storage_offset() as i64, plan.storage_length()); let mut gopts = WarmBackendGetOpts::default(); @@ -5599,11 +5827,13 @@ mod tests { use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints}; use crate::object_api::{ObjectInfo, ObjectOptions, PutObjReader}; #[cfg(feature = "test-util")] + use crate::services::tier::test_util::MockWarmOp; + #[cfg(feature = "test-util")] use crate::services::tier::test_util::register_mock_tier; #[cfg(feature = "test-util")] use crate::services::tier::tier::TierConfigMgr; #[cfg(feature = "test-util")] - use crate::services::tier::warm_backend::WarmBackend as _; + use crate::services::tier::warm_backend::{TransitionCandidateProbe, WarmBackend as _}; use crate::set_disk::{MultipartCommitBarrier, MultipartCommitPause}; use crate::set_disk::{RUSTFS_MULTIPART_BUCKET_KEY, RUSTFS_MULTIPART_OBJECT_KEY}; use crate::storage_api_contracts::namespace::NamespaceLocking as _; @@ -6299,7 +6529,75 @@ mod tests { #[cfg(feature = "test-util")] #[tokio::test] - async fn transitioned_get_rejects_unknown_version_state_before_backend_io() { + async fn transitioned_get_allows_legacy_unknown_exact_version_for_non_destructive_read() { + let manager = TierConfigMgr::new(); + let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let backend = register_mock_tier(&manager, &tier).await; + let remote_object = format!("remote/{}", Uuid::new_v4()); + let body = Bytes::from_static(b"legacy transitioned object body"); + let remote_version = backend + .put( + &remote_object, + ReaderImpl::Body(body.clone()), + i64::try_from(body.len()).expect("body length should fit"), + ) + .await + .expect("mock remote object should be stored"); + let mut user_defined = HashMap::new(); + insert_legacy_transition_version_id(&mut user_defined, &remote_version); + let object_info = ObjectInfo { + bucket: "bucket".to_string(), + name: "object".to_string(), + size: i64::try_from(body.len()).expect("body length should fit"), + transitioned_object: TransitionedObject { + name: remote_object, + version_id: remote_version, + status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(), + tier: tier.clone(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), + ..Default::default() + }; + + let range = Some(crate::storage_api_contracts::range::HTTPRangeSpec { + is_suffix_length: false, + start: 7, + end: 18, + }); + let mut reader = get_transitioned_object_reader_with_tier_manager( + &object_info.bucket, + &object_info.name, + &range, + &HeaderMap::new(), + &object_info, + &ObjectOptions::default(), + &manager, + None, + ) + .await + .expect("legacy unknown state should still allow a non-destructive read"); + let mut got = Vec::new(); + reader + .stream + .read_to_end(&mut got) + .await + .expect("transitioned reader should drain"); + + assert_eq!(got, &body.as_ref()[7..=18]); + assert_eq!(backend.get_count().await, 1); + assert_eq!(backend.remove_count().await, 0); + assert_eq!( + TierConfigMgr::active_operation_lease_count(&manager, &tier).await, + 0, + "tier generation lease should release after EOF" + ); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn transitioned_get_rejects_explicit_unknown_version_state_before_backend_io() { let manager = TierConfigMgr::new(); let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); let backend = register_mock_tier(&manager, &tier).await; @@ -6315,6 +6613,181 @@ mod tests { ..Default::default() }, transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined_with_transition_version_state(rustfs_filemeta::TransitionVersionState::Unknown).into(), + ..Default::default() + }; + + let err = match get_transitioned_object_reader_with_tier_manager( + &object_info.bucket, + &object_info.name, + &None, + &HeaderMap::new(), + &object_info, + &ObjectOptions::default(), + &manager, + None, + ) + .await + { + Ok(_) => panic!("explicit unknown remote version state must fail before backend IO"), + Err(err) => err, + }; + + assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + assert_eq!(backend.op_log().await, Vec::::new()); + assert_eq!(backend.get_count().await, 0); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn transitioned_get_rejects_present_but_invalid_legacy_version_metadata() { + let manager = TierConfigMgr::new(); + let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let backend = register_mock_tier(&manager, &tier).await; + + for persisted_version in [ + Uuid::nil().to_string(), + "\u{fffd}".to_string(), + "bad\u{0001}version".to_string(), + ] { + let mut user_defined = HashMap::new(); + insert_legacy_transition_version_id(&mut user_defined, &persisted_version); + let object_info = ObjectInfo { + bucket: "bucket".to_string(), + name: "object".to_string(), + size: 1, + transitioned_object: TransitionedObject { + name: "remote/object".to_string(), + version_id: String::new(), + status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(), + tier: tier.clone(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), + ..Default::default() + }; + + let err = match get_transitioned_object_reader_with_tier_manager( + &object_info.bucket, + &object_info.name, + &None, + &HeaderMap::new(), + &object_info, + &ObjectOptions::default(), + &manager, + None, + ) + .await + { + Ok(_) => panic!("present but invalid legacy version metadata must fail before backend IO"), + Err(err) => err, + }; + + assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + } + + assert_eq!(backend.op_log().await, Vec::::new()); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn transitioned_get_probes_legacy_empty_unknown_state_before_unversioned_read() { + let manager = TierConfigMgr::new(); + let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let backend = register_mock_tier(&manager, &tier).await; + backend.set_put_remote_version(Some(String::new())).await; + let remote_object = format!("remote/{}", Uuid::new_v4()); + let body = Bytes::from_static(b"legacy unversioned transitioned object body"); + let remote_version = backend + .put( + &remote_object, + ReaderImpl::Body(body.clone()), + i64::try_from(body.len()).expect("body length should fit"), + ) + .await + .expect("mock remote object should be stored"); + assert!(remote_version.is_empty()); + let object_info = ObjectInfo { + bucket: "bucket".to_string(), + name: "object".to_string(), + size: i64::try_from(body.len()).expect("body length should fit"), + transitioned_object: TransitionedObject { + name: remote_object.clone(), + version_id: String::new(), + status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(), + tier: tier.clone(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: HashMap::from([("x-minio-internal-transitioned-versionID".to_string(), String::new())]).into(), + ..Default::default() + }; + + let mut reader = get_transitioned_object_reader_with_tier_manager( + &object_info.bucket, + &object_info.name, + &None, + &HeaderMap::new(), + &object_info, + &ObjectOptions::default(), + &manager, + None, + ) + .await + .expect("probe-proven legacy unversioned state should allow a non-destructive read"); + let mut got = Vec::new(); + reader + .stream + .read_to_end(&mut got) + .await + .expect("transitioned reader should drain"); + + assert_eq!(got, body.as_ref()); + assert_eq!(backend.remove_count().await, 0); + assert_eq!( + backend.op_log().await, + vec![ + MockWarmOp::Put { + object: remote_object.clone() + }, + MockWarmOp::Probe { + object: remote_object.clone() + }, + MockWarmOp::Get { object: remote_object }, + ] + ); + assert_eq!( + TierConfigMgr::active_operation_lease_count(&manager, &tier).await, + 0, + "tier generation lease should release after EOF" + ); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn transitioned_get_rejects_ambiguous_empty_unknown_state_without_backend_get() { + let manager = TierConfigMgr::new(); + let tier = format!("COLDTIER{}", &Uuid::new_v4().simple().to_string()[..8]).to_uppercase(); + let backend = register_mock_tier(&manager, &tier).await; + let remote_object = format!("remote/{}", Uuid::new_v4()); + backend + .set_transition_candidate_probe_override(Some(TransitionCandidateProbe::VersionedPresent( + "versioned-candidate".to_string(), + ))) + .await; + let object_info = ObjectInfo { + bucket: "bucket".to_string(), + name: "object".to_string(), + size: 1, + transitioned_object: TransitionedObject { + name: remote_object.clone(), + version_id: String::new(), + status: crate::bucket::lifecycle::lifecycle::TRANSITION_COMPLETE.to_string(), + tier, + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, ..Default::default() }; @@ -6330,19 +6803,28 @@ mod tests { ) .await { - Ok(_) => panic!("unknown remote version state must fail before backend IO"), + Ok(_) => panic!("versioned legacy unknown state without stored version must fail before backend GET"), Err(err) => err, }; assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + assert_eq!(backend.op_log().await, vec![MockWarmOp::Probe { object: remote_object }]); assert_eq!(backend.get_count().await, 0); + assert_eq!(backend.remove_count().await, 0); } #[cfg(feature = "test-util")] #[tokio::test] - async fn free_version_delete_rejects_unknown_version_state_before_backend_io() { + async fn free_version_delete_rejects_explicit_unknown_before_backend_io() { let manager = TierConfigMgr::new(); let backend = register_mock_tier(&manager, "WARM").await; + let identity = test_tier_destination_identity(&manager, "WARM").await; + let mut user_defined = user_defined_with_tier_destination_identity(identity); + rustfs_utils::http::metadata_compat::insert_str( + &mut user_defined, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE, + rustfs_filemeta::TransitionVersionState::Unknown.as_str().to_string(), + ); let object_info = ObjectInfo { transitioned_object: TransitionedObject { name: "remote/object".to_string(), @@ -6351,17 +6833,251 @@ mod tests { ..Default::default() }, transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), ..Default::default() }; let err = super::delete_free_version_remote_object(&object_info, &manager) .await - .expect_err("unknown remote version state must fail before backend IO"); + .expect_err("explicit unknown cleanup must fail before backend IO"); assert_eq!(err.kind(), std::io::ErrorKind::InvalidData); + assert!(err.to_string().contains("version state is unknown")); + assert_eq!(backend.op_log().await, Vec::::new()); assert_eq!(backend.remove_count().await, 0); } + #[cfg(feature = "test-util")] + async fn test_tier_destination_identity( + manager: &Arc>, + tier: &str, + ) -> crate::services::tier::tier::TierDestinationId { + TierConfigMgr::acquire_operation_lease(manager, tier) + .await + .expect("test tier lease should be available") + .backend_identity() + } + + #[cfg(feature = "test-util")] + fn user_defined_with_tier_destination_identity( + identity: crate::services::tier::tier::TierDestinationId, + ) -> HashMap { + let mut user_defined = HashMap::new(); + rustfs_utils::http::metadata_compat::insert_str( + &mut user_defined, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID, + rustfs_utils::crypto::hex(identity), + ); + user_defined + } + + #[cfg(feature = "test-util")] + fn user_defined_with_transition_version_state(state: rustfs_filemeta::TransitionVersionState) -> HashMap { + let mut user_defined = HashMap::new(); + rustfs_utils::http::metadata_compat::insert_str( + &mut user_defined, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE, + state.as_str().to_string(), + ); + user_defined + } + + #[cfg(feature = "test-util")] + fn insert_legacy_transition_version_id(user_defined: &mut HashMap, version_id: &str) { + rustfs_utils::http::metadata_compat::insert_str( + user_defined, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_ID, + version_id.to_string(), + ); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn free_version_tuple_rejects_mixed_legacy_missing_and_explicit_unknown() { + let manager = TierConfigMgr::new(); + register_mock_tier(&manager, "WARM").await; + let identity = test_tier_destination_identity(&manager, "WARM").await; + let mut legacy_metadata = user_defined_with_tier_destination_identity(identity); + insert_legacy_transition_version_id(&mut legacy_metadata, "legacy-version"); + let mut explicit_metadata = legacy_metadata.clone(); + rustfs_utils::http::metadata_compat::insert_str( + &mut explicit_metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE, + rustfs_filemeta::TransitionVersionState::Unknown.as_str().to_string(), + ); + let make_info = |user_defined: HashMap| ObjectInfo { + transitioned_object: TransitionedObject { + name: "remote/object".to_string(), + version_id: "legacy-version".to_string(), + tier: "WARM".to_string(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), + ..Default::default() + }; + + let err = super::free_version_remote_tuple_matches(&make_info(legacy_metadata), &make_info(explicit_metadata)) + .expect_err("mixed legacy-missing and explicit unknown provenance must fail closed"); + + assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn free_version_delete_probes_exact_version_hidden_by_current_delete_marker() { + let manager = TierConfigMgr::new(); + let tier = "WARM"; + let backend = register_mock_tier(&manager, tier).await; + let identity = test_tier_destination_identity(&manager, tier).await; + let remote_object = format!("remote/{}", Uuid::new_v4()); + let body = Bytes::from_static(b"legacy exact cleanup body"); + let remote_version = backend + .put( + &remote_object, + ReaderImpl::Body(body), + i64::try_from(b"legacy exact cleanup body".len()).expect("body length should fit"), + ) + .await + .expect("mock remote object should be stored"); + let mut user_defined = user_defined_with_tier_destination_identity(identity); + insert_legacy_transition_version_id(&mut user_defined, &remote_version); + backend + .set_transition_candidate_probe_override(Some(TransitionCandidateProbe::Missing)) + .await; + assert_eq!( + backend + .probe_transition_candidate_state(&remote_object) + .await + .expect("current remote view should be readable"), + TransitionCandidateProbe::Missing, + "a current delete marker must hide the historical data version from an unversioned probe" + ); + backend.clear_op_log().await; + let object_info = ObjectInfo { + transitioned_object: TransitionedObject { + name: remote_object.clone(), + version_id: remote_version, + tier: tier.to_string(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), + ..Default::default() + }; + + super::delete_free_version_remote_object(&object_info, &manager) + .await + .expect("probe-proven legacy exact cleanup should delete the remote version"); + super::delete_free_version_remote_object(&object_info, &manager) + .await + .expect("a retry after the exact remote version is already missing should be idempotent"); + + assert_eq!( + backend.op_log().await, + vec![ + MockWarmOp::Get { + object: remote_object.clone() + }, + MockWarmOp::Remove { + object: remote_object.clone() + }, + MockWarmOp::Get { + object: remote_object.clone() + }, + ] + ); + assert_eq!( + backend.remove_versions().await, + vec![(remote_object, object_info.transitioned_object.version_id)] + ); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn free_version_delete_retains_legacy_unknown_unversioned_object() { + let manager = TierConfigMgr::new(); + let tier = "WARM"; + let backend = register_mock_tier(&manager, tier).await; + backend.set_put_remote_version(Some(String::new())).await; + let identity = test_tier_destination_identity(&manager, tier).await; + let remote_object = format!("remote/{}", Uuid::new_v4()); + let body = Bytes::from_static(b"legacy unversioned cleanup body"); + let remote_version = backend + .put( + &remote_object, + ReaderImpl::Body(body), + i64::try_from(b"legacy unversioned cleanup body".len()).expect("body length should fit"), + ) + .await + .expect("mock remote object should be stored"); + assert!(remote_version.is_empty()); + backend.clear_op_log().await; + let mut user_defined = user_defined_with_tier_destination_identity(identity); + user_defined.insert("x-minio-internal-transitioned-versionID".to_string(), String::new()); + let object_info = ObjectInfo { + transitioned_object: TransitionedObject { + name: remote_object.clone(), + version_id: String::new(), + tier: tier.to_string(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), + ..Default::default() + }; + + let err = super::delete_free_version_remote_object(&object_info, &manager) + .await + .expect_err("legacy unversioned cleanup cannot exclude a versioning-state race"); + + assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock); + assert!(backend.op_log().await.is_empty()); + assert_eq!(backend.remove_count().await, 0); + assert!(backend.remove_versions().await.is_empty()); + } + + #[cfg(feature = "test-util")] + #[tokio::test] + async fn free_version_delete_does_not_remove_a_different_remote_version() { + let manager = TierConfigMgr::new(); + let tier = "WARM"; + let backend = register_mock_tier(&manager, tier).await; + let identity = test_tier_destination_identity(&manager, tier).await; + let remote_object = format!("remote/{}", Uuid::new_v4()); + backend.set_put_remote_version(Some("different-version".to_string())).await; + backend + .put( + &remote_object, + ReaderImpl::Body(Bytes::from_static(b"different remote version")), + i64::try_from(b"different remote version".len()).expect("body length should fit"), + ) + .await + .expect("different remote version should be stored"); + backend.clear_op_log().await; + let mut user_defined = user_defined_with_tier_destination_identity(identity); + insert_legacy_transition_version_id(&mut user_defined, "legacy-version"); + let object_info = ObjectInfo { + transitioned_object: TransitionedObject { + name: remote_object.clone(), + version_id: "legacy-version".to_string(), + tier: tier.to_string(), + ..Default::default() + }, + transition_version_state: rustfs_filemeta::TransitionVersionState::Unknown, + user_defined: user_defined.into(), + ..Default::default() + }; + + super::delete_free_version_remote_object(&object_info, &manager) + .await + .expect("a missing exact legacy version should be an idempotent cleanup success"); + + assert_eq!(backend.op_log().await, vec![MockWarmOp::Get { object: remote_object }]); + assert_eq!(backend.remove_count().await, 0); + assert!(backend.remove_versions().await.is_empty()); + } + #[cfg(feature = "test-util")] #[tokio::test] async fn free_version_remote_delete_requires_persisted_destination_identity() { diff --git a/crates/ecstore/src/services/tier/test_util.rs b/crates/ecstore/src/services/tier/test_util.rs index f2c4ddeb6..59b4d6f4e 100644 --- a/crates/ecstore/src/services/tier/test_util.rs +++ b/crates/ecstore/src/services/tier/test_util.rs @@ -701,7 +701,7 @@ impl WarmBackend for MockWarmBackend { Ok(version) } - async fn get(&self, object: &str, _rv: &str, opts: WarmBackendGetOpts) -> Result { + async fn get(&self, object: &str, rv: &str, opts: WarmBackendGetOpts) -> Result { self.precondition().await?; let barrier = self.inner.get_barrier.lock().await.take(); if let Some(barrier) = barrier { @@ -719,6 +719,9 @@ impl WarmBackend for MockWarmBackend { let Some(stored) = objects.get(object) else { return Err(std::io::Error::new(std::io::ErrorKind::NotFound, "mock object not found")); }; + if !rv.is_empty() && stored.remote_version_id != rv { + return Err(std::io::Error::new(std::io::ErrorKind::NotFound, "NoSuchVersion")); + } let bytes = &stored.bytes; let start = opts.start_offset.max(0) as usize; diff --git a/crates/ecstore/src/services/tier/tier.rs b/crates/ecstore/src/services/tier/tier.rs index 887015a1e..af24423fa 100644 --- a/crates/ecstore/src/services/tier/tier.rs +++ b/crates/ecstore/src/services/tier/tier.rs @@ -2346,6 +2346,10 @@ impl WarmBackend for SharedWarmBackendProxy { self.0.probe_transition_candidate(object).await } + async fn probe_transition_version(&self, object: &str, remote_version_id: &str) -> io::Result { + self.0.probe_transition_version(object, remote_version_id).await + } + async fn in_use(&self) -> io::Result { self.0.in_use().await } @@ -2458,6 +2462,15 @@ impl TierOperationLease { Ok(()) } + pub(crate) async fn probe_transition_version( + &self, + object: &str, + remote_version_id: &str, + ) -> io::Result { + self.validate_remote_version_id(remote_version_id)?; + self.inner.driver.probe_transition_version(object, remote_version_id).await + } + pub(crate) fn is_current_generation(&self) -> bool { lock_unpoisoned(&self.runtime) .generations diff --git a/crates/ecstore/src/services/tier/warm_backend.rs b/crates/ecstore/src/services/tier/warm_backend.rs index ee48c116c..b5cf4ab38 100644 --- a/crates/ecstore/src/services/tier/warm_backend.rs +++ b/crates/ecstore/src/services/tier/warm_backend.rs @@ -40,6 +40,7 @@ use rustfs_s3_client::credentials::{Credentials, SignatureType, Static, Value}; use rustfs_s3_client::transition_api::{BucketLookupType, Options, TransitionClient, TransitionCore}; use rustfs_s3_client::{ admin_handler_utils::AdminError, + api_error_response::to_error_response, api_put_object::{AdvancedPutOptions, PutObjectOptions}, transition_api::{ReadCloser, ReaderImpl}, }; @@ -48,11 +49,14 @@ use rustfs_utils::egress::validate_outbound_url; use rustfs_utils::http::headers::{ CACHE_CONTROL, CONTENT_DISPOSITION, CONTENT_ENCODING, CONTENT_LANGUAGE, CONTENT_TYPE, EXPIRES, HeaderExt as _, }; -use s3s::dto::{ObjectLockLegalHoldStatus, ObjectLockRetentionMode, ReplicationStatus}; use s3s::header::{ X_AMZ_OBJECT_LOCK_LEGAL_HOLD, X_AMZ_OBJECT_LOCK_MODE, X_AMZ_OBJECT_LOCK_RETAIN_UNTIL_DATE, X_AMZ_REPLICATION_STATUS, X_AMZ_STORAGE_CLASS, }; +use s3s::{ + S3ErrorCode, + dto::{ObjectLockLegalHoldStatus, ObjectLockRetentionMode, ReplicationStatus}, +}; use std::collections::HashMap; use std::sync::Arc; use std::time::Duration; @@ -141,6 +145,42 @@ pub trait WarmBackend { async fn probe_transition_candidate(&self, _object: &str) -> Result { Ok(TransitionCandidateProbe::Unsupported) } + async fn probe_transition_version( + &self, + object: &str, + remote_version_id: &str, + ) -> Result { + if remote_version_id.is_empty() { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "an exact tier probe requires a remote version ID", + )); + } + self.validate_remote_version_id(remote_version_id)?; + match self + .get( + object, + remote_version_id, + WarmBackendGetOpts { + start_offset: 0, + length: 1, + }, + ) + .await + { + Ok(_) => Ok(TransitionCandidateProbe::VersionedPresent(remote_version_id.to_string())), + Err(err) if matches!(to_error_response(&err).code, S3ErrorCode::InvalidRange) => { + Ok(TransitionCandidateProbe::VersionedPresent(remote_version_id.to_string())) + } + Err(err) + if err.kind() == std::io::ErrorKind::NotFound + || matches!(to_error_response(&err).code, S3ErrorCode::NoSuchKey | S3ErrorCode::NoSuchVersion) => + { + Ok(TransitionCandidateProbe::Missing) + } + Err(err) => Err(err), + } + } async fn in_use(&self) -> Result; } @@ -437,6 +477,17 @@ impl WarmBackend for MeteredWarmBackend { Self::record(TierRequestOperation::Probe, result) } + async fn probe_transition_version( + &self, + object: &str, + remote_version_id: &str, + ) -> Result { + Self::record( + TierRequestOperation::Probe, + self.inner.probe_transition_version(object, remote_version_id).await, + ) + } + async fn in_use(&self) -> Result { Self::record(TierRequestOperation::InUse, self.inner.in_use().await) } diff --git a/crates/ecstore/src/services/tier/warm_backend_s3.rs b/crates/ecstore/src/services/tier/warm_backend_s3.rs index 5462fc52c..b830ea7f2 100644 --- a/crates/ecstore/src/services/tier/warm_backend_s3.rs +++ b/crates/ecstore/src/services/tier/warm_backend_s3.rs @@ -529,6 +529,10 @@ mod tests { "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\nNoSuchKeymissing", "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 66\r\nConnection: close\r\n\r\nNoSuchObjectmissing", "HTTP/1.1 403 Forbidden\r\nContent-Type: application/xml\r\nContent-Length: 65\r\nConnection: close\r\n\r\nAccessDenieddenied", + "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\nNoSuchKeymissing", + "HTTP/1.1 416 Range Not Satisfiable\r\nContent-Type: application/xml\r\nContent-Length: 72\r\nConnection: close\r\n\r\nInvalidRangeempty version", + "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 67\r\nConnection: close\r\n\r\nNoSuchVersionmissing", + "HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\nNoSuchKeymissing", ]; let mut requests = Vec::new(); for response in responses { @@ -622,15 +626,52 @@ mod tests { .await .expect_err("an authorization failure must not be mistaken for a missing key"); assert_eq!(to_error_response(&err).code, S3ErrorCode::AccessDenied); + assert_eq!( + backend + .probe_transition_candidate("delete-marker-hidden") + .await + .expect("a current delete marker should hide the data version"), + TransitionCandidateProbe::Missing + ); + assert_eq!( + backend + .probe_transition_version("delete-marker-hidden", "historical-version") + .await + .expect("the stored historical version should be probed exactly"), + TransitionCandidateProbe::VersionedPresent("historical-version".to_string()) + ); + assert_eq!( + backend + .probe_transition_version("delete-marker-hidden", "missing-version") + .await + .expect("a missing exact version should be classified"), + TransitionCandidateProbe::Missing + ); + assert_eq!( + backend + .probe_transition_version("missing-object", "historical-version") + .await + .expect("a missing key for an exact version probe should be classified"), + TransitionCandidateProbe::Missing + ); let requests = fixture.await.expect("candidate fixture should join"); - for request in requests { + for request in &requests[..6] { let request = request.to_ascii_lowercase(); assert!(request.starts_with("get /bucket/"), "candidate discovery must use object GET"); assert!(request.contains("\r\nrange: bytes=0-0\r\n")); assert!(!request.contains("?versioning")); assert!(!request.contains("?versions")); } + for request in &requests[6..] { + let request = request.to_ascii_lowercase(); + assert!(request.starts_with("get /bucket/"), "exact discovery must use object GET"); + assert!(request.contains("\r\nrange: bytes=0-0\r\n")); + } + assert!(!requests[5].to_ascii_lowercase().contains("versionid=")); + assert!(requests[6].to_ascii_lowercase().contains("?versionid=historical-version")); + assert!(requests[7].to_ascii_lowercase().contains("?versionid=missing-version")); + assert!(requests[8].to_ascii_lowercase().contains("?versionid=historical-version")); } fn list_versions(versions: &[(&str, &str)], delete_markers: &[(&str, &str)], is_truncated: bool) -> ListVersionsResult { diff --git a/crates/ecstore/src/store/init.rs b/crates/ecstore/src/store/init.rs index 22d7e37c3..463f50a82 100644 --- a/crates/ecstore/src/store/init.rs +++ b/crates/ecstore/src/store/init.rs @@ -11575,6 +11575,7 @@ mod tests { pool_index: usize, bucket: &str, object: &str, + minio_unversioned: bool, ) { for disk_index in 0..4 { let metadata_path = @@ -11608,6 +11609,11 @@ mod tests { ] { rustfs_utils::http::metadata_compat::remove_bytes(&mut object_meta.meta_sys, suffix); } + if minio_unversioned { + object_meta + .meta_sys + .insert("x-minio-internal-transitioned-versionID".to_string(), Vec::new()); + } *shallow = rustfs_filemeta::FileMetaShallowVersion::try_from(version) .expect("legacy transitioned version should re-encode"); } @@ -11618,6 +11624,152 @@ mod tests { } } + #[cfg(feature = "test-util")] + async fn read_store_body( + store: &Arc, + bucket: &str, + object: &str, + range: Option, + opts: &ObjectOptions, + ) -> Vec { + let mut reader = store + .get_object_reader(bucket, object, range, HeaderMap::new(), opts) + .await + .expect("object reader should open"); + let mut body = Vec::new(); + reader.stream.read_to_end(&mut body).await.expect("object body should drain"); + body + } + + #[cfg(feature = "test-util")] + #[tokio::test] + #[serial_test::serial(storage_class_env)] + async fn legacy_unknown_unversioned_transition_supports_head_get_and_range_without_backfill() { + let temp_dir = tempfile::tempdir().expect("create legacy unknown unversioned store dir"); + let (ctx, store, _shutdown) = + without_storage_class_env(build_isolated_test_store(temp_dir.path(), "legacy-unknown-unversioned-read", &[4])).await; + crate::bucket::metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await; + let tier_name = "LEGACY-UNKNOWN-UNVERSIONED-READ"; + let backend = register_mock_tier(&ctx.tier_config_mgr(), tier_name).await; + backend.set_put_remote_version(Some(String::new())).await; + let bucket = "legacy-unknown-unversioned-read-bucket"; + let object = "object.bin"; + let payload = b"legacy unversioned remote tier object remains readable".repeat(1024); + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("legacy source bucket should be created"); + let mut reader = PutObjReader::from_vec(payload.clone()); + let source = store + .put_object(bucket, object, &mut reader, &ObjectOptions::default()) + .await + .expect("legacy source should be written"); + store + .transition_object( + bucket, + object, + &ObjectOptions { + transition: TransitionOptions { + status: TRANSITION_PENDING.to_string(), + tier: tier_name.to_string(), + etag: source.etag.clone().expect("legacy source should have an etag"), + ..Default::default() + }, + mod_time: source.mod_time, + ..Default::default() + }, + ) + .await + .expect("legacy source should transition"); + rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 0, bucket, object, true).await; + backend.clear_op_log().await; + + let opts = ObjectOptions { + metadata_cache_safe: false, + ..Default::default() + }; + let head = store + .get_object_info(bucket, object, &opts) + .await + .expect("legacy transitioned HEAD should use local metadata"); + assert_eq!(head.transition_version_state, rustfs_filemeta::TransitionVersionState::Unknown); + assert!(head.transitioned_object.version_id.is_empty()); + assert_eq!( + head.user_defined + .get("x-minio-internal-transitioned-versionID") + .map(String::as_str), + Some(""), + "the MinIO empty version-key provenance must survive xl.meta decoding" + ); + assert!( + !rustfs_utils::http::metadata_compat::contains_key_str( + &head.user_defined, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE, + ), + "the compatibility read must not synthesize version-state metadata" + ); + + let full_body = read_store_body(&store, bucket, object, None, &opts).await; + assert_eq!(full_body, payload); + + let range = HTTPRangeSpec { + is_suffix_length: false, + start: 7, + end: 38, + }; + let ranged_body = read_store_body(&store, bucket, object, Some(range), &opts).await; + assert_eq!(ranged_body, &payload[7..=38]); + + let after_read = store.pools[0] + .get_disks_by_key(object) + .load_file_info_versions_exact(bucket, object) + .await + .expect("legacy metadata should remain readable after GET") + .expect("legacy object metadata should remain on disk") + .versions + .into_iter() + .find(|version| version.transition_status == rustfs_filemeta::TRANSITION_COMPLETE) + .expect("legacy transitioned source should remain visible after GET"); + assert_eq!(after_read.transition_version_state, rustfs_filemeta::TransitionVersionState::Unknown); + assert!(after_read.transition_version.is_none()); + assert!(after_read.transition_version_id.is_none()); + assert_eq!( + after_read + .metadata + .get("x-minio-internal-transitioned-versionID") + .map(String::as_str), + Some(""), + "the MinIO empty version-key provenance must remain after GET and Range GET" + ); + assert!( + !rustfs_utils::http::metadata_compat::contains_key_str( + &after_read.metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE, + ), + "the compatibility read must remain side-effect free" + ); + + assert_eq!( + backend.op_log().await, + vec![ + MockWarmOp::Probe { + object: after_read.transitioned_objname.clone(), + }, + MockWarmOp::Get { + object: after_read.transitioned_objname.clone(), + }, + MockWarmOp::Probe { + object: after_read.transitioned_objname.clone(), + }, + MockWarmOp::Get { + object: after_read.transitioned_objname, + }, + ], + "legacy reads should probe before each unversioned GET and never mutate local metadata" + ); + assert_eq!(backend.remove_count().await, 0); + } + #[cfg(feature = "test-util")] #[tokio::test] #[serial_test::serial(storage_class_env)] @@ -11658,7 +11810,7 @@ mod tests { ) .await .expect("legacy source should transition"); - rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 0, bucket, object).await; + rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 0, bucket, object, false).await; let legacy = store.pools[0] .get_disks_by_key(object) .load_file_info_versions_exact(bucket, object) @@ -12799,7 +12951,7 @@ mod tests { .expect("merge-loser source should transition"); copy_test_xlmeta_between_pools(temp_dir.path(), 0, 1, bucket, object).await; } - rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 1, bucket, "legacy/item.bin").await; + rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 1, bucket, "legacy/item.bin", false).await; backend.set_remove_failure(true); store.pools[1] .delete_object(bucket, "hidden/item.bin", ObjectOptions::default()) @@ -16866,6 +17018,10 @@ mod tests { .find(|version| version.version_id == history.version_id) .expect("transitioned history should exist"); transitioned.transition_version_state = rustfs_filemeta::TransitionVersionState::Unknown; + rustfs_utils::http::metadata_compat::remove_str( + &mut transitioned.metadata, + rustfs_utils::http::metadata_compat::SUFFIX_TRANSITIONED_VERSION_STATE, + ); metadata .add_version(transitioned) .expect("unknown state should replace the transitioned version"); diff --git a/crates/filemeta/src/filemeta/version.rs b/crates/filemeta/src/filemeta/version.rs index f0969c3cf..42aa5a74d 100644 --- a/crates/filemeta/src/filemeta/version.rs +++ b/crates/filemeta/src/filemeta/version.rs @@ -297,6 +297,20 @@ fn transitioned_version_from_bytes(value: Option<&[u8]>, state: TransitionVersio } } +fn transition_version_metadata_value(raw: &[u8], decoded: Option<&str>) -> String { + decoded.map(str::to_owned).unwrap_or_else(|| { + if raw.is_empty() { + String::new() + } else { + String::from_utf8_lossy(raw).into_owned() + } + }) +} + +fn is_transition_version_metadata_key(key: &str) -> bool { + strip_internal_prefix_preserving_case(key).is_some_and(|suffix| suffix.eq_ignore_ascii_case(SUFFIX_TRANSITIONED_VERSION_ID)) +} + fn validate_transition_version_state(state: TransitionVersionState, version: Option<&str>) -> Result<()> { let valid = match state { TransitionVersionState::Unknown | TransitionVersionState::KnownDisabled => version.is_none(), @@ -366,14 +380,26 @@ impl<'a> DerivedInternalMetadata<'a> { } *slot = Some(value.as_slice()); } + fn merge_consistent<'a>(canonical: Option<&'a [u8]>, legacy: Option<&'a [u8]>) -> Result> { + if let (Some(canonical), Some(legacy)) = (canonical, legacy) + && canonical != legacy + { + return Err(Error::FileCorrupt); + } + Ok(canonical.or(legacy)) + } + Ok(Self { checksum: canonical.checksum.or(legacy.checksum), part_checksums: canonical.part_checksums.or(legacy.part_checksums), - transition_status: canonical.transition_status.or(legacy.transition_status), - transitioned_object: canonical.transitioned_object.or(legacy.transitioned_object), - transitioned_version: canonical.transitioned_version.or(legacy.transitioned_version), - transitioned_version_state: canonical.transitioned_version_state.or(legacy.transitioned_version_state), - transition_tier: canonical.transition_tier.or(legacy.transition_tier), + transition_status: merge_consistent(canonical.transition_status, legacy.transition_status)?, + transitioned_object: merge_consistent(canonical.transitioned_object, legacy.transitioned_object)?, + transitioned_version: merge_consistent(canonical.transitioned_version, legacy.transitioned_version)?, + transitioned_version_state: merge_consistent( + canonical.transitioned_version_state, + legacy.transitioned_version_state, + )?, + transition_tier: merge_consistent(canonical.transition_tier, legacy.transition_tier)?, }) } } @@ -438,8 +464,14 @@ impl FileInfo { } } -fn set_transition_version_state(meta_sys: &mut HashMap>, state: TransitionVersionState) { - if state == TransitionVersionState::Unknown { +fn set_transition_version_state( + meta_sys: &mut HashMap>, + state: TransitionVersionState, + source_metadata: &HashMap, +) { + if state == TransitionVersionState::Unknown + && !rustfs_utils::http::metadata_compat::contains_key_str(source_metadata, SUFFIX_TRANSITIONED_VERSION_STATE) + { remove_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE); } else { insert_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE, state.as_str().as_bytes().to_vec()); @@ -2643,6 +2675,11 @@ impl MetaObject { if derived_metadata.transitioned_version_state.is_some() { validate_transition_version_state(transition_version_state, transition_version.as_deref())?; } + for (key, value) in &self.meta_sys { + if is_transition_version_metadata_key(key) { + metadata.insert(key.to_owned(), transition_version_metadata_value(value, transition_version.as_deref())); + } + } let transition_version_id = transition_version.as_deref().and_then(|value| Uuid::parse_str(value).ok()); let transition_tier = derived_metadata .transition_tier @@ -2689,7 +2726,7 @@ impl MetaObject { } else { remove_bytes(&mut self.meta_sys, SUFFIX_TRANSITIONED_VERSION_ID); } - set_transition_version_state(&mut self.meta_sys, fi.transition_version_state); + set_transition_version_state(&mut self.meta_sys, fi.transition_version_state, &fi.metadata); insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER, fi.transition_tier.as_bytes().to_vec()); if let Some(destination_id) = get_str(&fi.metadata, SUFFIX_TRANSITION_TIER_DESTINATION_ID) { insert_bytes(&mut self.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID, destination_id.into_bytes()); @@ -2830,7 +2867,7 @@ impl From for MetaObject { insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version); } if !value.transition_status.is_empty() { - set_transition_version_state(&mut meta_sys, value.transition_version_state); + set_transition_version_state(&mut meta_sys, value.transition_version_state, &value.metadata); } if !value.transition_tier.is_empty() { @@ -2985,6 +3022,12 @@ impl MetaDeleteMarker { fi.transition_version_state = transition_version_state_from_bytes(derived_metadata.transitioned_version_state)?; fi.transition_version = transitioned_version_from_bytes(derived_metadata.transitioned_version, fi.transition_version_state); + for (key, value) in &self.meta_sys { + if is_transition_version_metadata_key(key) { + fi.metadata + .insert(key.to_owned(), transition_version_metadata_value(value, fi.transition_version.as_deref())); + } + } fi.transition_version_id = fi.transition_version.as_deref().and_then(|value| Uuid::parse_str(value).ok()); if derived_metadata.transitioned_version_state.is_some() { validate_transition_version_state(fi.transition_version_state, fi.transition_version.as_deref())?; @@ -3152,7 +3195,7 @@ impl From for MetaDeleteMarker { insert_bytes(&mut meta_sys, SUFFIX_TRANSITIONED_VERSION_ID, transition_version); } if !value.transition_status.is_empty() || value.tier_free_version() { - set_transition_version_state(&mut meta_sys, value.transition_version_state); + set_transition_version_state(&mut meta_sys, value.transition_version_state, &value.metadata); } if !value.transition_tier.is_empty() { insert_bytes(&mut meta_sys, SUFFIX_TRANSITION_TIER, value.transition_tier.as_bytes().to_vec()); @@ -4574,6 +4617,7 @@ mod tests { .into_fileinfo("b", "k", false) .expect("into_fileinfo"); assert_eq!(fi.transition_version_id, None); + assert_eq!(get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID), Some(String::new())); } #[test] @@ -4585,6 +4629,10 @@ mod tests { .into_fileinfo("b", "k", false) .expect("into_fileinfo"); assert_eq!(fi.transition_version_id, None); + assert!( + get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID).is_some_and(|value| !value.is_empty()), + "nil UUID bytes must remain distinguishable from an empty MinIO version" + ); } #[test] @@ -4598,6 +4646,7 @@ mod tests { assert_eq!(fi.transition_version_id, Some(id)); assert_eq!(fi.transition_version, Some(id.to_string())); assert_eq!(fi.transition_version_state, TransitionVersionState::Unknown); + assert_eq!(get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID), Some(id.to_string())); } #[test] @@ -4637,6 +4686,36 @@ mod tests { assert_eq!(fi.transition_version_state, TransitionVersionState::Unknown); } + #[test] + fn meta_object_transition_version_state_explicit_unknown_is_not_legacy_missing() { + let mut metadata = HashMap::new(); + rustfs_utils::http::metadata_compat::insert_str( + &mut metadata, + SUFFIX_TRANSITIONED_VERSION_STATE, + TransitionVersionState::Unknown.as_str().to_string(), + ); + let fi = FileInfo { + transition_status: "complete".to_string(), + transition_version_state: TransitionVersionState::Unknown, + metadata, + ..Default::default() + }; + + let object = MetaObject::from(fi); + assert_eq!( + get_consistent_bytes(&object.meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE), + Some(b"unknown".as_slice()) + ); + let decoded = object + .into_fileinfo("b", "k", false) + .expect("explicit unknown state should decode"); + assert_eq!(decoded.transition_version_state, TransitionVersionState::Unknown); + assert_eq!( + rustfs_utils::http::metadata_compat::get_consistent_str(&decoded.metadata, SUFFIX_TRANSITIONED_VERSION_STATE,), + Some("unknown") + ); + } + #[test] fn meta_object_transition_version_state_exact_round_trips_dual_keys() { let id = sample_version_id(); @@ -4753,6 +4832,10 @@ mod tests { .expect("invalid transition version bytes must not fail the object read"); assert_eq!(fi.transition_version_id, None); assert_eq!(fi.transition_version, None); + assert!( + get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID).is_some_and(|value| !value.is_empty()), + "invalid raw bytes must remain distinguishable from an empty MinIO version" + ); } #[test] @@ -4795,6 +4878,10 @@ mod tests { .into_fileinfo("b", "k", false) .expect("nil tier version should remain an absent remote version"); assert_eq!(fi.transition_version_id, None); + assert!( + get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID).is_some_and(|value| !value.is_empty()), + "nil UUID bytes must remain distinguishable from an empty MinIO version" + ); } #[test] @@ -4812,6 +4899,7 @@ mod tests { .expect("legacy binary UUID tier version should decode"); assert_eq!(fi.transition_version_id, Some(id)); assert_eq!(fi.transition_version, Some(id.to_string())); + assert_eq!(get_str(&fi.metadata, SUFFIX_TRANSITIONED_VERSION_ID), Some(id.to_string())); } #[test] @@ -4910,6 +4998,23 @@ mod tests { assert_eq!(err, Error::FileCorrupt); } + #[test] + fn meta_object_transition_version_state_mixed_case_alias_conflict_fails_closed() { + let sys = HashMap::from([ + ( + format!("{RUSTFS_INTERNAL_PREFIX}{SUFFIX_TRANSITIONED_VERSION_STATE}"), + b"unknown".to_vec(), + ), + ("X-Minio-Internal-transitioned-version-state".to_string(), b"exact".to_vec()), + ]); + + let err = make_meta_object_with_sys(sys) + .into_fileinfo("b", "k", false) + .expect_err("mixed-case transition state aliases must agree"); + + assert_eq!(err, Error::FileCorrupt); + } + #[test] fn version_header_sorts_before_prefers_object_over_delete_marker_on_equal_mod_time() { let object = FileMetaVersionHeader { From a3b8183be9cc75e9ea3775e8e8f4709a1de90182 Mon Sep 17 00:00:00 2001 From: cxymds Date: Sat, 5 Sep 2026 14:12:45 +0800 Subject: [PATCH 3/4] test(ecstore): narrow barrier re-export cfgs (#7152) Co-authored-by: Zhengchao An --- crates/ecstore/src/set_disk/mod.rs | 2 +- crates/ecstore/src/store/mod.rs | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 1bc87619d..0690d70cb 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -876,7 +876,7 @@ pub use ops::multipart::{MultipartCommitBarrier, MultipartCommitPause}; pub(crate) use ops::object::DeleteObjectCommitBarrier; #[cfg(any(test, feature = "test-util"))] pub(crate) use ops::object::TransitionCleanupStoreBarrier as SetDiskTransitionCleanupStoreBarrier; -#[cfg(test)] +#[cfg(all(test, feature = "test-util"))] pub(crate) use ops::object::TransitionUploadedCommitBarrier as SetDiskTransitionUploadedCommitBarrier; pub(crate) use ops::object::body_cache_plaintext_len; #[cfg(all(test, feature = "test-util"))] diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index 77aaf98c8..8a4579e1b 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -425,7 +425,7 @@ pub(crate) mod init_format; pub(crate) mod list_objects; mod multipart; mod object; -#[cfg(any(test, feature = "test-util"))] +#[cfg(feature = "test-util")] pub use object::DeleteAfterObjectLockSnapshotBarrier; pub(crate) use object::{ DecommissionFixedReadAnchor, ObjectLockDiagGuard, RemoteTuplePublicationCommitGuard, RemoteTuplePublicationFence, From d8c3b1bb26250a6ec9008b7100df436ac7720643 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Sat, 5 Sep 2026 14:13:59 +0800 Subject: [PATCH 4/4] fix(app): fail closed on an unreadable bucket encryption config (#7183) The object write path read the bucket default encryption configuration with `.ok()`, which made "this bucket has no default encryption" and "the encryption configuration cannot be read" the same value. A bucket whose encryption blob is damaged therefore stored plaintext objects the operator had mandated be encrypted, with nothing returned to the client and nothing in the object to tell those writes apart afterwards. PUT, COPY and the snowball extract path now share one resolver: an absent configuration still writes plaintext exactly as before, and every other outcome refuses the write, carrying the accessor's typed error so a damaged blob surfaces as a deterministic InternalError while a transient metadata read failure surfaces as the retryable ServiceUnavailable. A missing bucket and a cold metadata cache both still resolve to "no configuration", so neither becomes a refusal. This matches `prepare_sse_configuration` in `storage::sse`, the resolver the multipart writer has always used, which fails closed on this lookup. --- rustfs/src/app/object/copy.rs | 46 +++++++++- rustfs/src/app/object/extract.rs | 2 +- rustfs/src/app/object/put.rs | 120 ++++++++++++++++++++++++- rustfs/src/app/object/shared.rs | 123 ++++++++++++++++++++++++++ rustfs/src/app/object/test_support.rs | 33 +++++++ 5 files changed, 321 insertions(+), 3 deletions(-) diff --git a/rustfs/src/app/object/copy.rs b/rustfs/src/app/object/copy.rs index fb9496c13..13393f797 100644 --- a/rustfs/src/app/object/copy.rs +++ b/rustfs/src/app/object/copy.rs @@ -394,7 +394,7 @@ impl DefaultObjectUsecase { // Bucket metadata uses the bucket name as its namespace-lock key. Load // every copy-time bucket snapshot before a same-object key can collide // with that key (for example, copying `bucket/bucket` onto itself). - let bucket_sse_config = metadata_sys::get_sse_config(&bucket).await.ok(); + let bucket_sse_config = load_bucket_default_sse_config(&bucket).await?; let object_lock_config_state = load_bucket_object_lock_config_state(&bucket).await?; if cp_src_dst_same && key == bucket { dst_opts.object_lock_config_snapshot = @@ -1388,4 +1388,48 @@ mod tests { .unwrap_err(); assert_eq!(err.code(), &S3ErrorCode::InvalidRequest); } + + #[tokio::test] + #[serial_test::serial] + async fn execute_copy_object_refuses_a_bucket_whose_encryption_config_is_unreadable() { + use crate::app::storage_api::test::contract::bucket::{BucketOperations as _, MakeBucketOptions}; + + let (store, context) = real_store_test_context().await; + let bucket = format!("copy-sse-unreadable-{}", Uuid::new_v4()); + let source = "source.bin"; + let destination = "destination.bin"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("unreadable-encryption copy bucket must be created"); + let mut reader = PutObjReader::from_vec(b"copied while the bucket still had a readable configuration".to_vec()); + store + .put_object(&bucket, source, &mut reader, &ObjectOptions::default()) + .await + .expect("copy source object must be written"); + install_unreadable_bucket_sse_config(&bucket).await; + + let input = CopyObjectInput::builder() + .copy_source(CopySource::Bucket { + bucket: bucket.clone().into(), + key: source.into(), + version_id: None, + }) + .bucket(bucket.clone()) + .key(destination.to_string()) + .build() + .expect("copy input must build"); + let usecase = DefaultObjectUsecase::with_context(Some(Arc::clone(&context))); + + let err = Box::pin(usecase.execute_copy_object(build_request(input, Method::PUT))) + .await + .expect_err("an unreadable bucket encryption configuration must refuse the copy"); + + assert_eq!(err.code(), &S3ErrorCode::InternalError); + let lookup_err = store + .get_object_info(&bucket, destination, &ObjectOptions::default()) + .await + .expect_err("a refused copy must not leave a destination object behind"); + assert!(is_err_object_not_found(&lookup_err), "{lookup_err}"); + } } diff --git a/rustfs/src/app/object/extract.rs b/rustfs/src/app/object/extract.rs index e5b7fb4da..2db042fa3 100644 --- a/rustfs/src/app/object/extract.rs +++ b/rustfs/src/app/object/extract.rs @@ -2037,7 +2037,7 @@ impl DefaultObjectUsecase { let sse_customer_key_md5 = sse_customer_key_md5.or(h_md5); let original_sse = server_side_encryption.or(extract_server_side_encryption_from_headers(&req.headers)?); - let bucket_sse_config = metadata_sys::get_sse_config(&bucket).await.ok(); + let bucket_sse_config = load_bucket_default_sse_config(&bucket).await?; let (mut effective_sse, mut effective_kms_key_id) = resolve_bucket_default_sse( bucket_sse_config.as_ref().map(|(config, _timestamp)| config), original_sse, diff --git a/rustfs/src/app/object/put.rs b/rustfs/src/app/object/put.rs index d108b1885..86d6e255a 100644 --- a/rustfs/src/app/object/put.rs +++ b/rustfs/src/app/object/put.rs @@ -1485,8 +1485,9 @@ impl DefaultObjectUsecase { }; let sse_config_stage_start = put_stage_metrics_enabled.then(Instant::now); - let bucket_sse_config = metadata_sys::get_sse_config(&bucket).await.ok(); + let bucket_sse_config = load_bucket_default_sse_config(&bucket).await; rustfs_io_metrics::record_put_object_stage_duration_from("app_sse_config_lookup", sse_config_stage_start); + let bucket_sse_config = bucket_sse_config?; debug!( target: "rustfs::app::object_usecase", component = "app", @@ -3912,4 +3913,121 @@ mod tests { .expect_err("writes after the zero-byte quota update must be denied"); assert!(matches!(err, StorageError::QuotaExceeded { current: 4096, limit: 0 })); } + + #[tokio::test] + #[serial_test::serial] + async fn execute_put_object_refuses_a_bucket_whose_encryption_config_is_unreadable() { + use crate::app::storage_api::test::contract::bucket::{BucketOperations as _, MakeBucketOptions}; + + let (store, context) = real_store_test_context().await; + let bucket = format!("put-sse-unreadable-{}", Uuid::new_v4()); + let object = "object.bin"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("unreadable-encryption PUT bucket must be created"); + install_unreadable_bucket_sse_config(&bucket).await; + + let payload = Bytes::from_static(b"an operator mandated encryption for this bucket"); + let input = PutObjectInput::builder() + .bucket(bucket.clone()) + .key(object.to_string()) + .body(Some(StreamingBlob::from(s3s::Body::from(payload.clone())))) + .content_length(Some(i64::try_from(payload.len()).expect("test payload length must fit i64"))) + .build() + .expect("PUT input must build"); + let usecase = DefaultObjectUsecase::with_context(Some(Arc::clone(&context))); + + let err = Box::pin(usecase.execute_put_object(&FS::new(), build_request(input, Method::PUT))) + .await + .expect_err("an unreadable bucket encryption configuration must refuse the write"); + + assert_eq!(err.code(), &S3ErrorCode::InternalError); + let lookup_err = store + .get_object_info(&bucket, object, &ObjectOptions::default()) + .await + .expect_err("a refused PUT must not leave an object behind"); + assert!(is_err_object_not_found(&lookup_err), "{lookup_err}"); + } + + #[tokio::test] + #[serial_test::serial] + async fn execute_put_object_still_writes_plaintext_without_bucket_encryption() { + use crate::app::storage_api::test::contract::bucket::{BucketOperations as _, MakeBucketOptions}; + + let (store, context) = real_store_test_context().await; + let bucket = format!("put-sse-absent-{}", Uuid::new_v4()); + let object = "object.bin"; + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("plaintext PUT bucket must be created"); + + let payload = Bytes::from_static(b"no default encryption is configured for this bucket"); + let input = PutObjectInput::builder() + .bucket(bucket.clone()) + .key(object.to_string()) + .body(Some(StreamingBlob::from(s3s::Body::from(payload.clone())))) + .content_length(Some(i64::try_from(payload.len()).expect("test payload length must fit i64"))) + .build() + .expect("PUT input must build"); + let usecase = DefaultObjectUsecase::with_context(Some(Arc::clone(&context))); + + Box::pin(usecase.execute_put_object(&FS::new(), build_request(input, Method::PUT))) + .await + .expect("a bucket without default encryption must still accept a plaintext write"); + + let stored = store + .get_object_info(&bucket, object, &ObjectOptions::default()) + .await + .expect("the plaintext object must be readable"); + assert_eq!(stored.size, i64::try_from(payload.len()).expect("test payload length must fit i64")); + assert!( + !stored + .user_defined + .keys() + .any(|key| key.eq_ignore_ascii_case(AMZ_SERVER_SIDE_ENCRYPTION) + || key.starts_with("x-rustfs-encryption-") + || key.starts_with("x-minio-encryption-")), + "the object must carry no encryption metadata: {:?}", + stored.user_defined + ); + } + + #[tokio::test] + #[serial_test::serial] + async fn execute_put_object_extract_refuses_a_bucket_whose_encryption_config_is_unreadable() { + use crate::app::storage_api::test::contract::bucket::{BucketOperations as _, MakeBucketOptions}; + + let (store, context) = real_store_test_context().await; + let bucket = format!("extract-sse-unreadable-{}", Uuid::new_v4()); + store + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("unreadable-encryption extract bucket must be created"); + install_unreadable_bucket_sse_config(&bucket).await; + + let payload = Bytes::from_static(b"archive bytes that must never be unpacked in plaintext"); + let input = PutObjectInput::builder() + .bucket(bucket.clone()) + .key("archive.tar".to_string()) + .body(Some(StreamingBlob::from(s3s::Body::from(payload.clone())))) + .content_length(Some(i64::try_from(payload.len()).expect("test payload length must fit i64"))) + .build() + .expect("extract PUT input must build"); + let mut req = build_request(input, Method::PUT); + req.headers.insert(AMZ_SNOWBALL_EXTRACT, HeaderValue::from_static("true")); + let usecase = DefaultObjectUsecase::with_context(Some(Arc::clone(&context))); + + let err = Box::pin(usecase.execute_put_object(&FS::new(), req)) + .await + .expect_err("an unreadable bucket encryption configuration must refuse the extract upload"); + + assert_eq!(err.code(), &S3ErrorCode::InternalError); + let lookup_err = store + .get_object_info(&bucket, "archive.tar", &ObjectOptions::default()) + .await + .expect_err("a refused extract upload must not leave an object behind"); + assert!(is_err_object_not_found(&lookup_err), "{lookup_err}"); + } } diff --git a/rustfs/src/app/object/shared.rs b/rustfs/src/app/object/shared.rs index bb382ec11..c35619edb 100644 --- a/rustfs/src/app/object/shared.rs +++ b/rustfs/src/app/object/shared.rs @@ -269,6 +269,129 @@ pub(super) fn resolve_bucket_default_sse( (effective_sse, effective_kms_key_id) } +/// The bucket's default encryption configuration for a write path. +/// +/// `Ok(None)` carries one meaning only — this bucket has no default encryption +/// — and the write proceeds in plaintext exactly as before. Every other +/// outcome refuses the write rather than collapsing onto that same value: an +/// encryption blob that exists but cannot be read fails closed in +/// `get_sse_config` since rustfs/rustfs#7172, and swallowing that error here +/// stores plaintext into a bucket whose operator mandated encryption, with +/// nothing returned to the client and nothing in the object to tell it apart +/// afterwards (rustfs/backlog#2287). +/// +/// The states the lookup can report, and what each one does: +/// +/// * configured and readable — apply the bucket default; +/// * no encryption blob at all, including a bucket that does not exist and a +/// bucket whose metadata document is absent — `ConfigNotFound`, so a cold +/// cache and a missing bucket are never turned into a refusal, and the write +/// still fails later with its own `NoSuchBucket`; +/// * blob present but unparseable — deterministic, so retrying cannot help; +/// surfaces as `InternalError` until an operator repairs or removes it; +/// * the metadata read itself failed (namespace lock, quorum, disk, an +/// uninitialized metadata system) — transient, and the typed error maps to +/// the retryable `ServiceUnavailable`. +/// +/// The last two are distinguished by the typed error the accessor returns, not +/// re-derived here: [`ApiError`] already separates them. This mirrors +/// `prepare_sse_configuration` in `storage::sse`, the resolver the multipart +/// writer uses, which has always failed closed on the same lookup. +pub(super) async fn load_bucket_default_sse_config( + bucket: &str, +) -> S3Result> { + classify_bucket_default_sse_lookup(bucket, metadata_sys::get_sse_config(bucket).await) +} + +fn classify_bucket_default_sse_lookup( + bucket: &str, + lookup: Result<(ServerSideEncryptionConfiguration, OffsetDateTime), StorageError>, +) -> S3Result> { + match lookup { + Ok(config) => Ok(Some(config)), + Err(err) if err == StorageError::ConfigNotFound => Ok(None), + Err(err) => { + let api_error = ApiError::from(err); + error!( + event = "bucket_sse_config_lookup_failed", + component = LOG_COMPONENT_APP, + subsystem = LOG_SUBSYSTEM_OBJECT, + result = "write_refused", + bucket = %bucket, + code = %api_error.code.as_str(), + error = %api_error, + "Bucket default encryption is unreadable; refusing the write instead of storing plaintext" + ); + Err(api_error.into()) + } + } +} + +#[cfg(test)] +mod bucket_default_sse_lookup_tests { + use super::*; + use s3s::dto::{ServerSideEncryptionByDefault, ServerSideEncryptionRule}; + use time::OffsetDateTime; + + fn sse_config() -> ServerSideEncryptionConfiguration { + ServerSideEncryptionConfiguration { + rules: vec![ServerSideEncryptionRule { + apply_server_side_encryption_by_default: Some(ServerSideEncryptionByDefault { + sse_algorithm: ServerSideEncryption::from_static(ServerSideEncryption::AES256), + kms_master_key_id: None, + }), + blocked_encryption_types: None, + bucket_key_enabled: None, + }], + } + } + + #[test] + fn an_absent_configuration_still_writes_plaintext() { + let resolved = classify_bucket_default_sse_lookup("bucket", Err(StorageError::ConfigNotFound)) + .expect("a bucket without default encryption must keep writing plaintext"); + + assert!(resolved.is_none()); + assert_eq!(resolve_bucket_default_sse(None, None, None, false), (None, None)); + } + + #[test] + fn a_readable_configuration_is_returned() { + let resolved = classify_bucket_default_sse_lookup("bucket", Ok((sse_config(), OffsetDateTime::UNIX_EPOCH))) + .expect("a readable configuration must not refuse the write") + .expect("a readable configuration must be applied"); + + assert_eq!(resolved.0.rules.len(), 1); + } + + #[test] + fn an_unreadable_configuration_refuses_the_write() { + let err = classify_bucket_default_sse_lookup( + "bucket", + Err(StorageError::other("persisted bucket encryption configuration is invalid")), + ) + .expect_err("a corrupt encryption blob must never degrade to plaintext"); + + assert_eq!(err.code(), &S3ErrorCode::InternalError); + } + + #[test] + fn an_unavailable_metadata_read_refuses_the_write_as_retryable() { + let err = classify_bucket_default_sse_lookup("bucket", Err(StorageError::ErasureReadQuorum)) + .expect_err("an unreadable metadata subsystem must never degrade to plaintext"); + + assert_eq!(err.code(), &S3ErrorCode::ServiceUnavailable); + } + + #[test] + fn a_missing_bucket_keeps_its_own_error() { + let err = classify_bucket_default_sse_lookup("bucket", Err(StorageError::BucketNotFound("bucket".to_string()))) + .expect_err("a bucket-not-found lookup must not be reported as an encryption failure"); + + assert_eq!(err.code(), &S3ErrorCode::NoSuchBucket); + } +} + #[cfg(test)] mod deadlock_request_guard_tests { use super::DeadlockRequestGuard; diff --git a/rustfs/src/app/object/test_support.rs b/rustfs/src/app/object/test_support.rs index 61b368cf6..55ebb28b2 100644 --- a/rustfs/src/app/object/test_support.rs +++ b/rustfs/src/app/object/test_support.rs @@ -96,3 +96,36 @@ pub(super) fn real_cold_fill_plan( }; plan } + +/// A store with an ambient `AppContext`, for tests that drive a handler end to +/// end without the object-data-cache overrides of +/// [`real_cold_fill_test_context`]. +pub(super) async fn real_store_test_context() -> (Arc, Arc) { + let store = crate::app::gating_test_env::shared_gating_ecstore().await; + if current_app_context().is_none() { + crate::app::runtime_sources::install_test_app_context(Arc::clone(&store)).await; + } + let ambient = current_app_context().expect("real-store tests require an ambient AppContext"); + let context = Arc::new(AppContext::new(Arc::clone(&store), ambient.iam(), ambient.kms())); + (store, context) +} + +/// Leave the bucket in the state a damaged encryption blob produces: the raw +/// document is retained and the typed configuration stays `None`, which is the +/// durable "exists but cannot be read" signal `get_sse_config` fails closed on +/// (rustfs/rustfs#7172). +pub(super) async fn install_unreadable_bucket_sse_config(bucket: &str) { + use crate::app::storage_api::test::{get_global_bucket_metadata_sys, set_bucket_metadata}; + + let sys = get_global_bucket_metadata_sys().expect("bucket metadata system must be initialized"); + let metadata = { + let sys = sys.read().await; + sys.get(bucket).await.expect("bucket metadata must be cached") + }; + let mut metadata = (*metadata).clone(); + metadata.encryption_config_xml = b"truncated".to_vec(); + metadata.sse_config = None; + set_bucket_metadata(bucket.to_string(), metadata) + .await + .expect("unreadable bucket encryption configuration must be installed"); +}