Merge remote-tracking branch 'origin/main' into codex/test-lifecycle-adversarial

This commit is contained in:
loverustfs
2026-09-05 14:23:10 +08:00
20 changed files with 1894 additions and 93 deletions
+84 -25
View File
@@ -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<RwLock<HashMap<String, SsecPassthroughRecord>>>,
pub targets_map: Arc<RwLock<HashMap<String, Vec<BucketTarget>>>>,
/// 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<RwLock<HashSet<String>>>,
pub h_mutex: Arc<RwLock<HashMap<String, EpHealth>>>,
target_h_mutex: Arc<RwLock<HashMap<String, EpHealth>>>,
pub hc_client: Arc<HttpClient>,
@@ -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<BucketTarget> {
/// 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<Vec<BucketTarget>, 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<BucketTargets, BucketTargetError> {
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);
@@ -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<bool, std::io::Error> {
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<TransitionDeleteVersionPlan, std::io::Error> {
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<RwLock<TierConfigMgr>>,
) -> 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<ResolvedTransitionDeleteVersion, std::io::Error> {
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<ECStore>, 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<ECStore>, 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<RwLock<TierConfigMgr>>,
) -> 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<Output = T>,
{
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<bool, std::io::
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum TransitionReadVersionPlan {
Direct,
ProbeLegacyUnversioned,
}
const LEGACY_TRANSITION_READ_PROBE_TIMEOUT: StdDuration = StdDuration::from_secs(30);
fn transition_remote_version_read_plan(oi: &ObjectInfo) -> Result<TransitionReadVersionPlan, std::io::Error> {
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<RwLock<TierConfigMgr>>,
resolver: Option<&dyn ObjectEncryptionResolver>,
) -> Result<GetObjectReader, std::io::Error> {
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::<MockWarmOp>::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::<MockWarmOp>::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::<MockWarmOp>::new());
assert_eq!(backend.remove_count().await, 0);
}
#[cfg(feature = "test-util")]
async fn test_tier_destination_identity(
manager: &Arc<tokio::sync::RwLock<TierConfigMgr>>,
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<String, String> {
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<String, String> {
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<String, String>, 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<String, String>| 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() {
+162 -8
View File
@@ -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::<BucketTargets>(&self.bucket_targets_config_json)
.map(|t| self.bucket_target_config = Some(t))
{
tracing::warn!(
match serde_json::from_slice::<BucketTargets>(&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"<not-a-valid-document".to_vec();
let malformed_json = b"{not-json".to_vec();
let mut bm = BucketMetadata::new("all-configs-malformed");
bm.policy_config_json = malformed_json.clone();
bm.quota_config_json = malformed_json.clone();
bm.bucket_targets_config_json = malformed_json.clone();
bm.notification_config_xml = malformed_xml.clone();
bm.lifecycle_config_xml = malformed_xml.clone();
bm.object_lock_config_xml = malformed_xml.clone();
bm.versioning_config_xml = malformed_xml.clone();
bm.encryption_config_xml = malformed_xml.clone();
bm.tagging_config_xml = malformed_xml.clone();
bm.replication_config_xml = malformed_xml.clone();
bm.cors_config_xml = malformed_xml.clone();
bm.logging_config_xml = malformed_xml.clone();
bm.website_config_xml = malformed_xml.clone();
bm.accelerate_config_xml = malformed_xml.clone();
bm.request_payment_config_xml = malformed_xml.clone();
bm.public_access_block_config_xml = malformed_xml.clone();
// `bucket_acl_config_json` is only checked for UTF-8, so only invalid
// UTF-8 exercises its failure branch.
bm.bucket_acl_config_json = vec![0xff, 0xfe];
bm.parse_all_configs()
.expect("a bucket whose every config is corrupt must still load its metadata");
let cleared: [(&str, bool); 17] = [
("policy", bm.policy_config.is_none()),
("quota", bm.quota_config.is_none()),
("bucket_targets", bm.bucket_target_config.is_none()),
("notification", bm.notification_config.is_none()),
("lifecycle", bm.lifecycle_config.is_none()),
("object_lock", bm.object_lock_config.is_none()),
("versioning", bm.versioning_config.is_none()),
("encryption", bm.sse_config.is_none()),
("tagging", bm.tagging_config.is_none()),
("replication", bm.replication_config.is_none()),
("cors", bm.cors_config.is_none()),
("logging", bm.logging_config.is_none()),
("website", bm.website_config.is_none()),
("accelerate", bm.accelerate_config.is_none()),
("request_payment", bm.request_payment_config.is_none()),
("public_access_block", bm.public_access_block_config.is_none()),
("bucket_acl", bm.bucket_acl_config.is_none()),
];
for (config, is_cleared) in cleared {
assert!(is_cleared, "{config}: a corrupt payload must not be replaced by a default");
}
assert_eq!(bm.bucket_targets_config_json, malformed_json, "raw bytes are retained");
assert_eq!(bm.lifecycle_config_xml, malformed_xml, "raw bytes are retained");
}
#[test]
fn lifecycle_update_config_clears_parsed_config_on_delete() {
let mut bm = BucketMetadata::new("test-bucket");
+161 -4
View File
@@ -360,6 +360,16 @@ async fn refresh_buckets_metadata_once(sys: Arc<RwLock<BucketMetadataSys>>) {
}
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<BucketTargets> {
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"<ServerSideEncryptionConfiguration".to_vec();
metadata.public_access_block_config_xml = b"<PublicAccessBlockConfiguration".to_vec();
metadata.quota_config_json = b"{not-json".to_vec();
metadata.bucket_targets_config_json = b"{not-json".to_vec();
metadata
.parse_all_configs()
.expect("a corrupt sub-config must not fail the load");
sys.set(bucket.to_string(), Arc::new(metadata)).await;
for (config, result) in [
("encryption", sys.get_sse_config(bucket).await.err()),
("public access block", sys.get_public_access_block_config(bucket).await.err()),
("quota", sys.get_quota_config(bucket).await.err()),
("bucket targets", sys.get_bucket_targets_config(bucket).await.err()),
] {
let err = result.unwrap_or_else(|| panic!("malformed {config} metadata must not read as a value"));
assert_ne!(err, Error::ConfigNotFound, "malformed {config} metadata must not be reported as absent");
}
}
#[tokio::test]
async fn config_states_distinguish_authoritative_absence_from_fabricated_metadata() {
use std::sync::atomic::Ordering;
@@ -4066,6 +4115,114 @@ mod tests {
target_sys.delete(bucket).await;
}
/// rustfs/backlog#2282: an unreadable `bucket-targets.json` reaches every
/// targets reader as a typed error; it neither withdraws a snapshot a
/// previous readable load published, nor collapses into the "no targets
/// configured" state that a bucket with an absent configuration reports.
#[tokio::test]
#[serial]
async fn unreadable_bucket_targets_fail_closed_and_stay_distinct_from_absent() {
let (_dirs, ecstore) = isolated_store_over_temp_disks().await;
let sys = BucketMetadataSys::new(ecstore);
let target_sys = BucketTargetSys::get();
let unreadable = "targets-unreadable";
let absent = "targets-absent";
target_sys.delete(unreadable).await;
target_sys.delete(absent).await;
// A readable load publishes this bucket's targets.
let mut readable = BucketMetadata::new(unreadable);
readable.bucket_target_config = Some(BucketTargets {
targets: vec![target(unreadable, "live")],
});
sync_bucket_target_sys(unreadable, &readable).await;
assert_eq!(
target_sys
.list_bucket_targets(unreadable)
.await
.expect("readable targets publish")
.targets
.len(),
1
);
// The same bucket reloaded with a blob that cannot be decoded.
let mut corrupt = BucketMetadata::new(unreadable);
corrupt.bucket_targets_config_json = br#"{"targets":[{"endpoint":"#.to_vec();
corrupt
.parse_all_configs()
.expect("an unreadable targets blob must not fail the metadata load");
sys.set(unreadable.to_string(), Arc::new(corrupt)).await;
assert!(
matches!(
target_sys.list_bucket_targets(unreadable).await,
Err(BucketTargetError::BucketRemoteTargetsUnreadable { .. })
),
"an unreadable configuration must not read as an empty or a missing target set"
);
assert!(
target_sys.list_targets(unreadable, "").await.is_err(),
"the admin listing must surface the fault instead of an empty list"
);
let err = sys
.get_bucket_targets_config(unreadable)
.await
.expect_err("an unreadable targets configuration must not read as a value");
assert_ne!(err, Error::ConfigNotFound, "unreadable must not be reported as absent");
// A bucket that never configured a target keeps its previous behavior.
let mut no_targets = BucketMetadata::new(absent);
no_targets.parse_all_configs().expect("absent targets parse");
sys.set(absent.to_string(), Arc::new(no_targets)).await;
assert!(
matches!(
target_sys.list_bucket_targets(absent).await,
Err(BucketTargetError::BucketRemoteTargetNotFound { .. })
),
"an absent configuration must still report as a missing target set"
);
assert!(
target_sys
.list_targets(absent, "")
.await
.expect("an absent configuration lists no targets")
.is_empty()
);
assert!(
sys.get_bucket_targets_config(absent)
.await
.expect("an absent targets configuration still reads as an empty set")
.is_empty(),
"the absent path must keep returning an empty target set, exactly as before"
);
// One bucket's unreadable configuration does not reach another bucket.
assert!(!matches!(
target_sys.list_bucket_targets(absent).await,
Err(BucketTargetError::BucketRemoteTargetsUnreadable { .. })
));
// A repaired configuration takes effect on the next load, no restart.
let mut repaired = BucketMetadata::new(unreadable);
repaired.bucket_target_config = Some(BucketTargets {
targets: vec![target(unreadable, "repaired")],
});
sync_bucket_target_sys(unreadable, &repaired).await;
assert_eq!(
target_sys
.list_bucket_targets(unreadable)
.await
.expect("a repaired configuration clears the unreadable marker")
.targets
.len(),
1
);
target_sys.delete(unreadable).await;
target_sys.delete(absent).await;
}
#[tokio::test]
#[serial]
async fn metadata_reload_clears_stale_bucket_targets_when_config_is_removed() {
@@ -46,7 +46,7 @@ use super::replication_storage_boundary::{
HTTPPreconditions, ObjectInfo, ObjectOptions, ObjectToDelete, ReplicationDeletedObject, ReplicationObjectIO,
ReplicationStorage,
};
use super::replication_target_boundary::{ReplicationTargetStore, replication_object_is_ssec_encrypted};
use super::replication_target_boundary::{BucketTargetError, ReplicationTargetStore, replication_object_is_ssec_encrypted};
use super::replication_versioning_boundary::ReplicationVersioningStore;
use super::runtime_boundary as runtime_sources;
use futures_util::stream::{self, StreamExt};
@@ -3084,6 +3084,23 @@ pub async fn queue_replication_heal(bucket: &str, oi: ObjectInfo, retry_count: u
let tgts = match ReplicationTargetStore::list_bucket_targets(bucket).await {
Ok(targets) => 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,
@@ -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;
@@ -701,7 +701,7 @@ impl WarmBackend for MockWarmBackend {
Ok(version)
}
async fn get(&self, object: &str, _rv: &str, opts: WarmBackendGetOpts) -> Result<ReadCloser, std::io::Error> {
async fn get(&self, object: &str, rv: &str, opts: WarmBackendGetOpts) -> Result<ReadCloser, std::io::Error> {
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;
+13
View File
@@ -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<TransitionCandidateProbe> {
self.0.probe_transition_version(object, remote_version_id).await
}
async fn in_use(&self) -> io::Result<bool> {
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<TransitionCandidateProbe> {
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
@@ -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<TransitionCandidateProbe, std::io::Error> {
Ok(TransitionCandidateProbe::Unsupported)
}
async fn probe_transition_version(
&self,
object: &str,
remote_version_id: &str,
) -> Result<TransitionCandidateProbe, std::io::Error> {
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<bool, std::io::Error>;
}
@@ -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<TransitionCandidateProbe, std::io::Error> {
Self::record(
TierRequestOperation::Probe,
self.inner.probe_transition_version(object, remote_version_id).await,
)
}
async fn in_use(&self) -> Result<bool, std::io::Error> {
Self::record(TierRequestOperation::InUse, self.inner.in_use().await)
}
@@ -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\n<Error><Code>NoSuchKey</Code><Message>missing</Message></Error>",
"HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 66\r\nConnection: close\r\n\r\n<Error><Code>NoSuchObject</Code><Message>missing</Message></Error>",
"HTTP/1.1 403 Forbidden\r\nContent-Type: application/xml\r\nContent-Length: 65\r\nConnection: close\r\n\r\n<Error><Code>AccessDenied</Code><Message>denied</Message></Error>",
"HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\n<Error><Code>NoSuchKey</Code><Message>missing</Message></Error>",
"HTTP/1.1 416 Range Not Satisfiable\r\nContent-Type: application/xml\r\nContent-Length: 72\r\nConnection: close\r\n\r\n<Error><Code>InvalidRange</Code><Message>empty version</Message></Error>",
"HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 67\r\nConnection: close\r\n\r\n<Error><Code>NoSuchVersion</Code><Message>missing</Message></Error>",
"HTTP/1.1 404 Not Found\r\nContent-Type: application/xml\r\nContent-Length: 63\r\nConnection: close\r\n\r\n<Error><Code>NoSuchKey</Code><Message>missing</Message></Error>",
];
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 {
+1 -1
View File
@@ -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"))]
+158 -2
View File
@@ -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<crate::store::ECStore>,
bucket: &str,
object: &str,
range: Option<HTTPRangeSpec>,
opts: &ObjectOptions,
) -> Vec<u8> {
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");
+1 -1
View File
@@ -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,
+115 -10
View File
@@ -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<Option<&'a [u8]>> {
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<String, Vec<u8>>, state: TransitionVersionState) {
if state == TransitionVersionState::Unknown {
fn set_transition_version_state(
meta_sys: &mut HashMap<String, Vec<u8>>,
state: TransitionVersionState,
source_metadata: &HashMap<String, String>,
) {
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<FileInfo> 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<FileInfo> 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 {
+11 -1
View File
@@ -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()
+45 -1
View File
@@ -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}");
}
}
+1 -1
View File
@@ -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,
+119 -1
View File
@@ -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}");
}
}
+123
View File
@@ -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<Option<(ServerSideEncryptionConfiguration, OffsetDateTime)>> {
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<Option<(ServerSideEncryptionConfiguration, OffsetDateTime)>> {
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;
+33
View File
@@ -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<ECStore>, Arc<AppContext>) {
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"<ServerSideEncryptionConfiguration>truncated".to_vec();
metadata.sse_config = None;
set_bucket_metadata(bucket.to_string(), metadata)
.await
.expect("unreadable bucket encryption configuration must be installed");
}