Compare commits

..

40 Commits

Author SHA1 Message Date
overtrue 2ea71a8725 fix(ecstore): integrate verified rename preflight evidence 2026-09-05 12:30:42 +08:00
overtrue cbb494beb8 refactor(ecstore): preserve rename observations in commit module 2026-09-05 12:26:32 +08:00
overtrue c332ba3d41 test(ecstore): cover observed rename outer failures 2026-09-05 12:21:08 +08:00
overtrue c084a26a38 Merge remote-tracking branch 'origin/overtrue/fix/ecstore-write-completion' into overtrue/fix/ecstore-write-completion 2026-09-05 12:13:46 +08:00
overtrue 0658b228e1 fix(ecstore): preserve known preflight rename rejections 2026-09-05 12:08:28 +08:00
overtrue 0b8ccae25d Merge remote-tracking branch 'origin/main' into overtrue/fix/ecstore-core-regressions 2026-09-05 12:01:40 +08:00
overtrue 03726fc322 test(ecstore): count decommission faults across retry restarts 2026-09-05 12:01:40 +08:00
Zhengchao An c9202db5e1 Merge branch 'main' into overtrue/fix/ecstore-write-completion 2026-09-05 11:57:29 +08:00
Zhengchao An 918d6af726 Merge branch 'main' into overtrue/refactor/local-commit-boundary 2026-09-05 11:57:24 +08:00
overtrue 7b74dfddb0 fix(odm): resolve pagination Clippy failures 2026-09-05 11:55:42 +08:00
overtrue e1f8cb9b89 Merge remote-tracking branch 'origin/main' into overtrue/fix/ecstore-core-regressions 2026-09-05 10:10:36 +08:00
overtrue 21e955ad29 test(odm): use app facade for listing wire types 2026-09-05 10:10:36 +08:00
overtrue 8fa93227e9 refactor(ecstore): remove moved quota fence import 2026-09-05 10:10:36 +08:00
overtrue 8ddc3cd436 Merge remote-tracking branch 'origin/main' into overtrue/refactor/local-commit-boundary 2026-09-05 10:01:42 +08:00
overtrue ccf8e2362c test(ecstore): match sealed context fixture map type 2026-09-05 10:01:42 +08:00
overtrue 2cb0425380 Merge remote-tracking branch 'origin/main' into overtrue/fix/ecstore-write-completion 2026-09-05 09:57:43 +08:00
overtrue daba4f7b32 test(ecstore): match sealed context fixture map type 2026-09-05 09:57:43 +08:00
overtrue 5513dadb78 test(ecstore): mark rollback fixtures as inline data 2026-09-05 09:57:43 +08:00
overtrue 19690eba7b fix(ecstore): retain indeterminate rename recovery evidence 2026-09-05 09:57:43 +08:00
overtrue dbe0073362 fix(ecstore): retain per-disk rename rollback outcomes 2026-09-05 09:57:43 +08:00
overtrue 2a6b2f31c7 refactor(ecstore): remove moved quota fence import 2026-09-05 09:56:41 +08:00
overtrue fe079cd2e8 refactor(ecstore): isolate local object rename commit 2026-09-05 09:54:53 +08:00
overtrue a732698586 refactor(ecstore): isolate metadata quorum decisions 2026-09-05 09:54:53 +08:00
overtrue 87c90c874d test(odm): match SDK bucket-root listing requests 2026-09-05 09:50:47 +08:00
overtrue bbcc51c124 test(ecstore): mark rollback fixtures as inline data 2026-09-05 09:50:47 +08:00
overtrue dacb617ff1 refactor(ecstore): isolate local object rename commit 2026-09-05 09:37:50 +08:00
overtrue 664f8c92ed fix(ecstore): retain indeterminate rename recovery evidence 2026-09-05 09:32:22 +08:00
overtrue 1c604f9b83 test(ci): require a fresh core JUnit report 2026-09-05 09:30:43 +08:00
overtrue 544bd1d7cf docs(odm): clarify folded source probe pagination 2026-09-05 09:30:26 +08:00
overtrue 286a5bed18 fix(ecstore): drain backfill checkpoint before confirmation 2026-09-05 09:30:26 +08:00
overtrue e83ace9533 fix(ecstore): drain backfill checkpoint before confirmation 2026-09-05 09:29:01 +08:00
overtrue 9e71b92f91 test(ecstore): match sealed context fixture map type 2026-09-05 09:28:46 +08:00
overtrue dd3de6d1c4 fix(odm): reject non-progressing listing cursors 2026-09-05 09:26:58 +08:00
overtrue b3b3eb0d36 fix(ecstore): retain PUT staging after incomplete rollback 2026-09-05 09:26:34 +08:00
overtrue 70be76c54e fix(ecstore): retain PUT staging after incomplete rollback 2026-09-05 09:26:12 +08:00
overtrue 6019391da7 docs(ecstore): define generation authority and recovery boundary 2026-09-05 09:25:58 +08:00
overtrue d8d9c02dcf fix(ecstore): retain per-disk rename rollback outcomes 2026-09-05 09:25:35 +08:00
overtrue 413b880b23 fix(ecstore): drain durable control-plane write tails 2026-09-05 09:24:54 +08:00
overtrue 6967c35074 test(ecstore): require core invariant tests in existing CI lane 2026-09-05 09:24:53 +08:00
overtrue de6ef09635 fix(ecstore): drain durable control-plane write tails 2026-09-05 09:24:29 +08:00
69 changed files with 816 additions and 6249 deletions
+1 -1
View File
@@ -1,2 +1,2 @@
sha256-darwin=a881fd7d3f5cb94654221ca85b8b30cce1b95e608824a55a15339cbc294e6d34
sha256-linux=a2933d83dfe74ffa03410a0959333a1c48288b8469ca9f17273d449d7510c24b
sha256-linux=e9a8d64e73f627c4d26c236dbbba690c9ee03a9e26d42a4244515b4439365535
+2 -26
View File
@@ -18,7 +18,7 @@ on:
workflow_dispatch:
inputs:
from_version:
description: 'OLD RustFS release tag, e.g. 1.0.0-rc.3 (its release must ship a .deb asset). Leave empty for the default.'
description: 'OLD RustFS release tag (must ship a .deb asset, e.g. 1.0.0-rc.3)'
required: false
default: '1.0.0-rc.3'
from_url:
@@ -26,7 +26,7 @@ on:
required: false
type: string
to_version:
description: 'NEW RustFS release tag, e.g. 1.0.0-rc.5 (any version with a .deb asset). Leave empty for latest nightly.'
description: 'NEW RustFS release tag (leave empty for latest nightly)'
required: false
to_url:
description: 'NEW .deb URL. Overrides to_version / nightly default.'
@@ -145,7 +145,6 @@ jobs:
continue-on-error: true
env:
LOG_FILE: /tmp/rustfs-upgrade.log
GH_TOKEN: ${{ secrets.PF_TESTING_GH_TOKEN }}
run: |
set -euo pipefail
chmod +x auto-testing/rustfs-upgrade-test.sh
@@ -176,29 +175,6 @@ jobs:
else
ARGS+=(--to-url "${RUSTFS_NIGHTLY_PACKAGE_URL}")
fi
# Fail fast with a clear message when a requested release tag has
# no .deb asset (e.g. 1.0.0-rc.4 ships only zips), instead of
# letting the suite die mid-run on a 404.
check_release_asset() {
local version="$1" tag asset url
[ -n "${version}" ] && [ "${version}" != "null" ] || return 0
tag="${version#v}"
asset="rustfs_${tag//-/.}_amd64.deb"
url="https://github.com/rustfs/rustfs/releases/download/${tag}/${asset}"
if ! gh api "repos/rustfs/rustfs/releases/tags/${tag}" --jq '.assets[].name' 2>/dev/null | grep -qxF "${asset}"; then
echo "ERROR: release ${tag} has no downloadable asset ${asset}:" >&2
echo " ${url}" >&2
echo "Pick a tag whose release ships a .deb (check its release assets)." >&2
exit 1
fi
echo "resolved ${tag} -> ${url}"
}
if [ -z "${FROM_URL}" ]; then
check_release_asset "${FROM_VERSION}"
fi
if [ -z "${TO_URL}" ]; then
check_release_asset "${TO_VERSION}"
fi
./auto-testing/rustfs-upgrade-test.sh "${ARGS[@]}"
- name: Generate report
-1
View File
@@ -33,7 +33,6 @@ profile.json
*.zst
.secrets
*.go
!crates/zip/tests/fixtures/snowball/**/generate/*.go
*.pb
*.svg
deploy/logs/*.log.*
Generated
+112 -324
View File
File diff suppressed because it is too large Load Diff
+6 -10
View File
@@ -168,7 +168,7 @@ reqwest = "0.13.4"
rustfs-kafka-async = { version = "1.3.1" }
socket2 = { version = "0.6.5" }
tokio = { version = "1.53.1" }
tokio-rustls = { default-features = false, version = "0.26.5" }
tokio-rustls = { default-features = false, version = "0.26.4" }
tokio-stream = { version = "0.1.19" }
tokio-test = "0.4.5"
tokio-util = { version = "0.7.19" }
@@ -234,19 +234,15 @@ tokio-postgres-rustls = "0.14.0"
# Utilities and Tools
anyhow = "1.0.104"
arc-swap = "1.9.2"
# RUSTFS_COMPAT_TODO(tokio-tar-extension-limits): keep the fork pin while Snowball and Swift still depend on it. Remove after Snowball uses a released tar-codec/tar-framing API that exposes precedence-resolved MinIO vendor records, RustFS preserves cancellation-safe ownership of large streamed members, footerless minio-go input is accepted only at an authenticated complete request boundary, the existing resource-limit, cancellation, and error-fuse regressions pass, and Swift no longer needs this fork.
# RUSTFS_COMPAT_TODO(tokio-tar-extension-limits): keep the fork pin until every parser hardening used by Snowball is released upstream. Remove after astral-sh/tokio-tar#118 is merged and a published release includes extension, physical-entry, and sparse limits, cancellation-safe sparse parsing, and error-fused entry streams.
astral-tokio-tar = { git = "https://github.com/cxymds/tokio-tar.git", rev = "603756478b7668436e464519c77ccac22a99ba96" }
# Candidate Snowball parser versions exercised by rustfs-zip compatibility fixtures.
tar-codec = "0.0.14"
tar-framing = "0.0.14"
atoi = "3.1.0"
atomic_enum = "0.3.0"
aws-config = { version = "1.12.0" }
aws-config = { version = "1.11.0" }
aws-credential-types = { version = "1.3.0" }
aws-sdk-kms = { default-features = false, version = "1.118.0" }
aws-sdk-s3 = { default-features = false, version = "1.145.0" }
aws-sdk-sts = { default-features = false, version = "1.114.0" }
aws-smithy-async = { version = "1.3.0" }
aws-sdk-kms = { default-features = false, version = "1.117.0" }
aws-sdk-s3 = { default-features = false, version = "1.144.0" }
aws-sdk-sts = { default-features = false, version = "1.113.0" }
aws-smithy-http-client = { default-features = false, version = "1.4.0" }
aws-smithy-runtime-api = { version = "1.16.0" }
aws-smithy-types = { version = "1.6.3" }
@@ -6743,99 +6743,6 @@ async fn test_site_replication_replicates_object_with_bucket_versioning_real_dua
Ok(())
}
#[tokio::test]
async fn test_site_replication_replays_bucket_created_during_peer_outage_real_dual_node() -> TestResult {
init_logging();
// Keep compilation outside the scenario timeout. Recovery itself waits
// for the production 30-second lightweight retry tick.
let _rustfs_binary = rustfs_binary_path();
match timeout(Duration::from_secs(150), async {
let mut site_env = replication_fast_env();
site_env.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV);
let mut site_a_env = RustFSTestEnvironment::new().await?;
site_a_env.start_rustfs_server_with_env(vec![], &site_env).await?;
let mut site_b_env = RustFSTestEnvironment::new().await?;
site_b_env.start_rustfs_server_without_cleanup_with_env(&site_env).await?;
let site_a_client = site_a_env.create_s3_client();
let site_b_client = site_b_env.create_s3_client();
let bucket = "site-repl-peer-outage";
let key = "after-recovery.txt";
let payload = b"site replication recovered the missed bucket".to_vec();
let add_status = site_replication_add(
&site_a_env,
&[
PeerSite {
name: "outage-site-a".to_string(),
endpoint: site_a_env.url.clone(),
access_key: site_a_env.access_key.clone(),
secret_key: site_a_env.secret_key.clone(),
..Default::default()
},
PeerSite {
name: "outage-site-b".to_string(),
endpoint: site_b_env.url.clone(),
access_key: site_b_env.access_key.clone(),
secret_key: site_b_env.secret_key.clone(),
..Default::default()
},
],
)
.await?;
assert!(add_status.success, "unexpected site add result: {add_status:?}");
wait_for_site_replication_enabled(&site_a_env, 2).await?;
wait_for_site_replication_enabled(&site_b_env, 2).await?;
site_b_env.stop_server();
site_a_client.create_bucket().bucket(bucket).send().await?;
site_a_client.head_bucket().bucket(bucket).send().await?;
let queued = site_replication_info(&site_a_env)
.await?
.retry_stats
.ok_or("peer outage did not persist a site replication retry event")?;
assert!(queued.pending + queued.failed > 0, "peer outage retry queue was unexpectedly empty");
site_b_env.restart_server_preserving_data(vec![], &site_env).await?;
let recovery_deadline = tokio::time::Instant::now() + Duration::from_secs(75);
loop {
let bucket_recovered = site_b_client.head_bucket().bucket(bucket).send().await.is_ok();
let queue_empty = site_replication_info(&site_a_env).await?.retry_stats.is_none();
if bucket_recovered && queue_empty {
break;
}
if tokio::time::Instant::now() >= recovery_deadline {
return Err(format!(
"site replication retry did not settle after peer recovery; bucket_recovered={bucket_recovered}, queue_empty={queue_empty}"
)
.into());
}
sleep(Duration::from_millis(250)).await;
}
site_a_client
.put_object()
.bucket(bucket)
.key(key)
.body(ByteStream::from(payload.clone()))
.send()
.await?;
assert_eq!(wait_for_object_on_target(&site_b_client, bucket, key).await?, payload);
Ok(())
})
.await
{
Ok(result) => result,
Err(_) => Err("site replication peer-outage recovery timed out after 150 seconds".into()),
}
}
/// Re-applying a site's own replication config must not disable the peer's reverse direction.
///
/// `PutBucketReplication` broadcasts the config to every peer — the console's replication
-1
View File
@@ -244,7 +244,6 @@ windows-sys = { workspace = true, features = [
windows-sys = { workspace = true, features = ["Win32_System_Ioctl"] }
[dev-dependencies]
aws-smithy-async.workspace = true
tokio = { workspace = true, features = ["rt-multi-thread", "macros", "test-util", "fs"] }
criterion = { workspace = true, features = ["html_reports"] }
temp-env = { workspace = true, features = ["async_closure"] }
+25 -84
View File
@@ -59,7 +59,7 @@ use rustfs_utils::http::{
insert_header,
};
use serde::{Deserialize, Serialize};
use std::collections::{HashMap, HashSet};
use std::collections::HashMap;
use std::error::Error;
use std::fmt;
use std::str::FromStr as _;
@@ -376,11 +376,6 @@ 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>,
@@ -424,7 +419,6 @@ 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()),
@@ -634,40 +628,30 @@ impl BucketTargetSys {
health_map.clone()
}
/// 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> {
pub async fn list_targets(&self, bucket: &str, arn_type: &str) -> Vec<BucketTarget> {
let health_stats = self.target_health_stats().await;
let mut targets = Vec::new();
if !bucket.is_empty() {
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);
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;
}
targets.push(target);
}
}
Err(BucketTargetError::BucketRemoteTargetNotFound { .. }) => {}
Err(err) => return Err(err),
}
return Ok(targets);
return targets;
}
let targets_map = self.targets_map.read().await;
@@ -690,16 +674,10 @@ impl BucketTargetSys {
}
}
Ok(targets)
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 {
@@ -712,30 +690,13 @@ 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: 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);
// 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).
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;
@@ -1132,11 +1093,6 @@ 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 {
@@ -1144,9 +1100,9 @@ impl BucketTargetSys {
}
}
// 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).
// 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).
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;
@@ -1205,11 +1161,6 @@ 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;
};
@@ -2325,13 +2276,6 @@ 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,
},
@@ -2365,9 +2309,6 @@ 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}")
}
@@ -3315,7 +3256,7 @@ mod tests {
}],
);
let targets = sys.list_targets("", "").await.expect("listing every bucket's targets");
let targets = sys.list_targets("", "").await;
assert_eq!(targets.len(), 1);
assert!(!targets[0].online);
@@ -584,173 +584,33 @@ 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, TransitionDeleteVersionPlan), std::io::Error> {
let delete_plan = transition_remote_version_delete_plan(oi)?;
) -> Result<(TierOperationLease, bool), std::io::Error> {
let version_id_exact = validate_transition_remote_version(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, 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(())
Ok((lease, version_id_exact))
}
async fn delete_free_version_remote_object_with_lease(
oi: &ObjectInfo,
lease: &TierOperationLease,
delete_plan: TransitionDeleteVersionPlan,
version_id_exact: bool,
) -> Result<(), std::io::Error> {
let resolved = resolve_transition_delete_version_plan(oi, lease, delete_plan).await?;
execute_resolved_transition_delete(oi, lease, resolved).await
delete_object_from_remote_tier_with_lease_idempotent(
&oi.transitioned_object.name,
&oi.transitioned_object.version_id,
lease,
version_id_exact,
)
.await?;
Ok(())
}
fn free_version_physical_topology_generation(api: &ECStore) -> String {
@@ -781,16 +641,6 @@ 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",
@@ -866,7 +716,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, delete_plan) = acquire_free_version_tier_lease(oi, &api.tier_config_mgr()).await?;
let (lease, version_id_exact) = 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)
@@ -884,30 +734,16 @@ 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, 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")
})??;
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"))??;
}
}
if !free_version_cleanup_fences_current(&topology_generation, &api, &bucket_guard, &object_guards, &lease, cancel, deadline) {
@@ -955,8 +791,8 @@ async fn delete_free_version_remote_object(
oi: &ObjectInfo,
tier_config_mgr: &Arc<RwLock<TierConfigMgr>>,
) -> Result<(), std::io::Error> {
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 (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
}
#[allow(
@@ -972,8 +808,8 @@ where
F: FnOnce() -> Fut,
Fut: std::future::Future<Output = T>,
{
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 (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 result = delete_local().await;
drop(lease);
Ok(result)
@@ -4852,39 +4688,6 @@ 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.
@@ -4899,12 +4702,7 @@ 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> {
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}")))?;
validate_transition_remote_version(oi)?;
let expected_identity = tier_destination_id_from_metadata(&oi.user_defined)?;
let lease = match expected_identity {
Some(identity) => {
@@ -4918,36 +4716,7 @@ pub(crate) async fn get_transitioned_object_reader_with_tier_manager(
Err(err) => return Err(std::io::Error::other(err)),
};
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",
));
}
}
}
}
tgt_client.validate_remote_version_id(&oi.transitioned_object.version_id)?;
// 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
@@ -4955,6 +4724,9 @@ 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();
@@ -5827,13 +5599,11 @@ 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::{TransitionCandidateProbe, WarmBackend as _};
use crate::services::tier::warm_backend::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 _;
@@ -6529,75 +6299,7 @@ mod tests {
#[cfg(feature = "test-util")]
#[tokio::test]
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() {
async fn transitioned_get_rejects_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;
@@ -6613,7 +6315,6 @@ 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()
};
@@ -6629,202 +6330,19 @@ mod tests {
)
.await
{
Ok(_) => panic!("explicit unknown remote version state must fail before backend IO"),
Ok(_) => panic!("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()
};
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!("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_explicit_unknown_before_backend_io() {
async fn free_version_delete_rejects_unknown_version_state_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(),
@@ -6833,251 +6351,17 @@ 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("explicit unknown cleanup must fail before backend IO");
.expect_err("unknown remote version state 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() {
@@ -15,6 +15,8 @@
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use super::runtime_boundary as runtime_sources;
use crate::bucket::lifecycle::bucket_lifecycle_ops::ExpiryOp;
@@ -70,11 +72,9 @@ static REMOTE_DELETE_BREAKER: LazyLock<Mutex<RemoteDeleteBreaker>> = LazyLock::n
});
#[cfg(test)]
type RemoteTierDeleteTestHook = Box<dyn Fn(&str, &str, &str) -> std::io::Result<()> + Send + Sync>;
#[cfg(test)]
static REMOTE_TIER_DELETE_TEST_HOOK: std::sync::LazyLock<std::sync::Mutex<Option<RemoteTierDeleteTestHook>>> =
std::sync::LazyLock::new(|| std::sync::Mutex::new(None));
static REMOTE_TIER_DELETE_TEST_HOOK: std::sync::LazyLock<
std::sync::Mutex<Option<Box<dyn Fn(&str, &str, &str) -> std::io::Result<()> + Send + Sync>>>,
> = std::sync::LazyLock::new(|| std::sync::Mutex::new(None));
#[derive(Debug)]
struct RemoteDeleteBreaker {
@@ -107,7 +107,7 @@ impl RemoteDeleteBreaker {
fn prune(&mut self, now: Instant) {
while let Some(ts) = self.failures.front().copied() {
if now.duration_since(ts) > self.window {
let _ = self.failures.pop_front();
self.failures.pop_front();
} else {
break;
}
@@ -137,10 +137,10 @@ fn is_signer_header_error(err: &std::io::Error) -> bool {
return false;
}
if let Some(source) = err.get_ref()
&& error_chain_contains_signer_header_marker(source)
{
return true;
if let Some(source) = err.get_ref() {
if error_chain_contains_signer_header_marker(source) {
return true;
}
}
let message = err.to_string().to_ascii_lowercase();
@@ -205,7 +205,7 @@ impl ObjSweeper {
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
pub fn with_version(&mut self, vid: Option<Uuid>) -> &Self {
self.version_id = vid;
self.version_id = vid.clone();
self
}
@@ -219,7 +219,7 @@ impl ObjSweeper {
#[allow(dead_code, reason = "MinIO-parity surface with no caller in this port (backlog#1823)")]
pub fn get_opts(&self) -> lifecycle::ObjectOpts {
let mut opts = ObjectOpts {
version_id: self.version_id,
version_id: self.version_id.clone(),
versioned: self.versioned,
version_suspended: self.suspended,
..Default::default()
@@ -388,8 +388,8 @@ impl Jentry {
impl ExpiryOp for Jentry {
fn op_hash(&self) -> u64 {
let mut hasher = Sha256::new();
hasher.update(self.tier_name.as_bytes());
hasher.update(self.obj_name.as_bytes());
hasher.update(format!("{}", self.tier_name).as_bytes());
hasher.update(format!("{}", self.obj_name).as_bytes());
xxh64::xxh64(hasher.finalize().as_slice(), XXHASH_SEED)
}
@@ -436,7 +436,7 @@ async fn delete_object_from_remote_tier_raw_with_manager(
tier_name: &str,
tier_config_mgr: &Arc<tokio::sync::RwLock<TierConfigMgr>>,
) -> Result<(), std::io::Error> {
let lease = TierConfigMgr::acquire_operation_lease(tier_config_mgr, tier_name)
let lease = TierConfigMgr::acquire_operation_lease(&tier_config_mgr, tier_name)
.await
.map_err(std::io::Error::other)?;
delete_object_from_remote_tier_raw_with_lease(obj_name, rv_id, &lease, false, true).await
+8 -162
View File
@@ -477,18 +477,6 @@ 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
@@ -976,32 +964,7 @@ impl BucketMetadata {
Ok(())
}
/// 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<()> {
fn parse_all_configs(&mut self) -> Result<()> {
if let Err(e) = self.parse_policy_config() {
tracing::warn!(
event = "bucket_metadata_parse_failed",
@@ -1125,26 +1088,20 @@ 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() {
match serde_json::from_slice::<BucketTargets>(&self.bucket_targets_config_json) {
Ok(targets) => self.bucket_target_config = Some(targets),
Err(e) => tracing::error!(
if let Err(e) = serde_json::from_slice::<BucketTargets>(&self.bucket_targets_config_json)
.map(|t| self.bucket_target_config = Some(t))
{
tracing::warn!(
event = "bucket_metadata_parse_failed",
component = "ecstore",
subsystem = "bucket_metadata",
bucket = %self.name,
config = "bucket_targets",
error = %e,
"Bucket replication targets are unreadable; replication for this bucket fails closed"
),
"Failed to parse bucket metadata config"
);
self.bucket_target_config = Some(BucketTargets::default());
}
} else {
self.bucket_target_config = Some(BucketTargets::default());
@@ -1578,117 +1535,6 @@ 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");
+4 -161
View File
@@ -360,16 +360,6 @@ 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;
@@ -2128,9 +2118,7 @@ impl BucketMetadataSys {
pub async fn get_public_access_block_config(&self, bucket: &str) -> Result<(PublicAccessBlockConfiguration, OffsetDateTime)> {
let (bm, _) = self.get_config(bucket).await?;
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 {
if let Some(config) = &bm.public_access_block_config {
Ok((config.clone(), bm.public_access_block_config_updated_at))
} else {
Err(Error::ConfigNotFound)
@@ -2441,9 +2429,7 @@ impl BucketMetadataSys {
pub async fn get_sse_config(&self, bucket: &str) -> Result<(ServerSideEncryptionConfiguration, OffsetDateTime)> {
let (bm, _) = self.get_config(bucket).await?;
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 {
if let Some(config) = &bm.sse_config {
Ok((config.clone(), bm.encryption_config_updated_at))
} else {
Err(Error::ConfigNotFound)
@@ -2514,9 +2500,7 @@ impl BucketMetadataSys {
pub async fn get_quota_config(&self, bucket: &str) -> Result<(BucketQuota, OffsetDateTime)> {
let (bm, _) = self.get_config(bucket).await?;
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 {
if let Some(config) = &bm.quota_config {
Ok((config.clone(), bm.quota_config_updated_at))
} else {
Err(Error::ConfigNotFound)
@@ -2538,9 +2522,7 @@ impl BucketMetadataSys {
pub async fn get_bucket_targets_config(&self, bucket: &str) -> Result<BucketTargets> {
let (bm, _) = self.get_config(bucket).await?;
if bm.bucket_targets_unreadable() {
Err(Error::other("persisted bucket replication target configuration is invalid"))
} else if let Some(config) = &bm.bucket_target_config {
if let Some(config) = &bm.bucket_target_config {
Ok(config.clone())
} else {
Err(Error::ConfigNotFound)
@@ -2611,7 +2593,6 @@ 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,
@@ -2807,36 +2788,6 @@ 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;
@@ -4115,114 +4066,6 @@ 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() {
+1 -170
View File
@@ -652,10 +652,9 @@ async fn build_aws_s3_http_client_from_tls_path() -> Option<SharedHttpClient> {
#[cfg(test)]
mod tests {
use super::*;
use aws_smithy_async::time::TimeSource;
use aws_smithy_runtime_api::http::StatusCode as SmithyStatusCode;
use std::sync::Mutex;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::atomic::{AtomicUsize, Ordering};
fn spec(endpoint: &str, secure: bool) -> RemoteS3EndpointSpec {
RemoteS3EndpointSpec {
@@ -825,174 +824,6 @@ mod tests {
);
}
#[derive(Clone, Debug)]
struct ClockSkewTimeSource(Arc<AtomicU64>);
impl TimeSource for ClockSkewTimeSource {
fn now(&self) -> SystemTime {
SystemTime::UNIX_EPOCH + Duration::from_secs(self.0.load(Ordering::SeqCst))
}
}
#[derive(Clone, Debug)]
struct ClockSkewConnector {
request_headers: RecordedHeaders,
error_code: &'static str,
skew_seconds: i64,
clock: ClockSkewTimeSource,
}
fn recorded_header<'a>(headers: &'a [(String, String)], name: &str) -> &'a str {
headers
.iter()
.find(|(key, _)| key.eq_ignore_ascii_case(name))
.map(|(_, value)| value.as_str())
.unwrap_or_else(|| panic!("signed request must contain {name}"))
}
fn signing_time(headers: &[(String, String)]) -> chrono::NaiveDateTime {
chrono::NaiveDateTime::parse_from_str(recorded_header(headers, "x-amz-date"), "%Y%m%dT%H%M%SZ")
.expect("SDK signing timestamp must use the SigV4 format")
}
impl SmithyHttpConnector for ClockSkewConnector {
fn call(&self, request: HttpRequest) -> HttpConnectorFuture {
let mut headers = self.request_headers.lock().expect("clock skew request capture lock");
assert!(headers.len() < 3, "clock skew fixture must not exceed two GET attempts and one HEAD");
headers.push(
request
.headers()
.iter()
.map(|(key, value)| (key.to_string(), value.to_string()))
.collect(),
);
let server_time = chrono::DateTime::<chrono::Utc>::from(self.clock.now()).naive_utc()
+ chrono::Duration::seconds(self.skew_seconds);
let (status, body) = if headers.len() == 1 {
(
403,
format!("<Error><Code>{}</Code><Message>Clock skew fixture</Message></Error>", self.error_code),
)
} else {
(200, String::new())
};
let response = http::Response::builder()
.status(status)
.header("date", server_time.format("%a, %d %b %Y %H:%M:%S GMT").to_string())
.header("content-type", "application/xml")
.header("content-length", body.len())
.body(SdkBody::from(body))
.expect("clock skew fixture response");
HttpConnectorFuture::ready(Ok(HttpResponse::try_from(response).expect("Smithy fixture response")))
}
}
async fn clock_skew_client(
error_code: &'static str,
skew_seconds: i64,
retry: RemoteS3RetryPolicy,
) -> (S3Client, RecordedHeaders, ClockSkewTimeSource) {
let headers: RecordedHeaders = Arc::new(Mutex::new(Vec::new()));
let clock = ClockSkewTimeSource(Arc::new(AtomicU64::new(1_700_000_000)));
let connector = SharedHttpConnector::new(ClockSkewConnector {
request_headers: Arc::clone(&headers),
error_code,
skew_seconds,
clock: clock.clone(),
});
let mut spec = spec("s3.example.com", true);
spec.retry = retry;
let config = build_remote_s3_config(&spec)
.await
.expect("clock skew fixture uses the production outbound configuration")
.http_client(http_client_fn(move |_settings, _components| connector.clone()))
.time_source(clock.clone())
.build();
(S3Client::from_conf(config), headers, clock)
}
#[tokio::test(start_paused = true)]
async fn remote_s3_clock_skew_retries_resign_and_seed_next_operation() {
for error_code in ["RequestTimeTooSkewed", "SignatureDoesNotMatch"] {
for skew_seconds in [-600, 600] {
let (client, headers, clock) = clock_skew_client(error_code, skew_seconds, REPLICATION_TARGET_RETRY_POLICY).await;
let initial = chrono::DateTime::<chrono::Utc>::from(clock.now()).naive_utc();
client
.get_object()
.bucket("bucket")
.key("object")
.send()
.await
.expect("clock skew GET must retry successfully");
assert_eq!(
headers.lock().expect("captured requests").len(),
2,
"{error_code}: GET needs exactly one retry"
);
clock.0.fetch_add(17, Ordering::SeqCst);
// SDK signing time is independent of Tokio's retry/scheduler clock.
tokio::time::advance(Duration::from_secs(61)).await;
client
.head_bucket()
.bucket("bucket")
.send()
.await
.expect("subsequent HEAD must use the client's cached skew");
let headers = headers.lock().expect("captured signed requests");
assert_eq!(headers.len(), 3, "subsequent operation must succeed on its first attempt");
assert_eq!(signing_time(&headers[0]), initial, "the first attempt must use the injected clock");
assert_eq!(
signing_time(&headers[1]),
initial + chrono::Duration::seconds(skew_seconds),
"{error_code}: retry must apply the measured offset exactly"
);
assert_eq!(
signing_time(&headers[2]),
initial + chrono::Duration::seconds(skew_seconds + 17),
"{error_code}: the next operation must apply cached skew to the advanced signing clock"
);
let signature = |index: usize| {
recorded_header(&headers[index], "authorization")
.rsplit_once("Signature=")
.expect("SigV4 authorization contains a signature")
.1
};
assert_ne!(
signature(0),
signature(1),
"{error_code}: retry must be signed again after adjusting its date"
);
}
}
}
#[tokio::test(start_paused = true)]
async fn remote_s3_clock_skew_respects_one_attempt_policy() {
use aws_smithy_types::error::metadata::ProvideErrorMetadata;
for error_code in ["RequestTimeTooSkewed", "SignatureDoesNotMatch"] {
for retry in [
RemoteS3RetryPolicy::Disabled,
RemoteS3RetryPolicy::Standard { max_attempts: 1 },
] {
let (client, headers, _clock) = clock_skew_client(error_code, 600, retry).await;
let error = client
.get_object()
.bucket("bucket")
.key("object")
.send()
.await
.expect_err("clock skew must not override the caller's one-attempt budget");
assert_eq!(error.as_service_error().and_then(ProvideErrorMetadata::code), Some(error_code));
assert_eq!(
headers.lock().expect("captured requests").len(),
1,
"{error_code}: {retry:?} must send exactly one request"
);
}
}
}
#[test]
fn path_style_auto_and_path_force_path_style() {
assert!(PathStyle::Auto.force_path_style());
@@ -46,7 +46,7 @@ use super::replication_storage_boundary::{
HTTPPreconditions, ObjectInfo, ObjectOptions, ObjectToDelete, ReplicationDeletedObject, ReplicationObjectIO,
ReplicationStorage,
};
use super::replication_target_boundary::{BucketTargetError, ReplicationTargetStore, replication_object_is_ssec_encrypted};
use super::replication_target_boundary::{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,23 +3084,6 @@ 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,8 +15,7 @@
use std::collections::HashMap;
use std::sync::Arc;
pub(crate) use crate::bucket::bucket_target_sys::BucketTargetError;
use crate::bucket::bucket_target_sys::BucketTargetSys;
use crate::bucket::bucket_target_sys::{BucketTargetError, 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,9 +719,6 @@ 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,10 +2346,6 @@ 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
}
@@ -2462,15 +2458,6 @@ 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
@@ -15,6 +15,8 @@
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use serde::{Deserialize, Deserializer, Serialize, Serializer, de};
@@ -143,7 +145,7 @@ mod tests {
assert_eq!(creds.access_key, "access");
assert_eq!(creds.secret_key, "secret");
assert_eq!(creds.creds_json.as_slice(), service_account);
assert_eq!(creds.creds_json.as_slice(), &service_account[..]);
let wire = serde_json::to_value(&creds).expect("madmin tier credentials should encode");
assert_eq!(wire["access"], "access");
@@ -160,7 +162,7 @@ mod tests {
.expect("the former RustFS field names and byte-array encoding should remain readable");
assert_eq!(legacy.access_key, "legacy-access");
assert_eq!(legacy.secret_key, "legacy-secret");
assert_eq!(legacy.creds_json.as_slice(), service_account);
assert_eq!(legacy.creds_json.as_slice(), &service_account[..]);
}
#[test]
@@ -40,7 +40,6 @@ 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},
};
@@ -49,14 +48,11 @@ 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;
@@ -145,42 +141,6 @@ 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>;
}
@@ -477,17 +437,6 @@ 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)
}
@@ -15,6 +15,8 @@
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use std::collections::HashMap;
@@ -15,6 +15,8 @@
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use std::collections::HashMap;
@@ -15,6 +15,8 @@
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use std::collections::{HashMap, HashSet};
use std::future::Future;
@@ -144,11 +146,11 @@ pub struct WarmBackendGCS {
impl WarmBackendGCS {
pub async fn new(conf: &TierGCS, tier: &str) -> Result<Self, std::io::Error> {
if conf.creds.is_empty() {
if conf.creds == "" {
return Err(std::io::Error::other("both access and secret keys are required"));
}
if conf.bucket.is_empty() {
if conf.bucket == "" {
return Err(std::io::Error::other("no bucket name was provided"));
}
@@ -193,11 +195,11 @@ impl WarmBackendGCS {
}
pub fn get_dest(&self, object: &str) -> String {
if self.prefix.is_empty() {
object.to_string()
} else {
format!("{}/{}", self.prefix, object)
let mut dest_obj = object.to_string();
if self.prefix != "" {
dest_obj = format!("{}/{}", &self.prefix, object);
}
return dest_obj;
}
}
@@ -221,7 +223,7 @@ impl WarmBackend for WarmBackendGCS {
let bucket = gcs_bucket_resource_name(&self.bucket);
let Ok(res) = Box::pin(
self.client
.write_object(&bucket, self.get_dest(object), Bytes::from(d))
.write_object(&bucket, &self.get_dest(object), Bytes::from(d))
.send_buffered(),
)
.await
@@ -238,7 +240,7 @@ impl WarmBackend for WarmBackendGCS {
async fn get(&self, object: &str, rv: &str, opts: WarmBackendGetOpts) -> Result<ReadCloser, std::io::Error> {
let bucket = gcs_bucket_resource_name(&self.bucket);
let mut req = self.client.read_object(&bucket, self.get_dest(object));
let mut req = self.client.read_object(&bucket, &self.get_dest(object));
let mut max_response_bytes = None;
if let Some(generation) = parse_generation(rv)? {
req = req.set_generation(generation);
@@ -15,6 +15,8 @@
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use std::collections::HashMap;
@@ -15,6 +15,8 @@
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use std::collections::HashMap;
@@ -15,6 +15,8 @@
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use std::collections::HashMap;
@@ -15,6 +15,8 @@
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use std::collections::HashMap;
@@ -529,10 +529,6 @@ 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 {
@@ -626,52 +622,15 @@ 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[..6] {
for request in requests {
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 {
@@ -15,6 +15,8 @@
#![allow(unused_variables)]
#![allow(unused_mut)]
#![allow(unused_assignments)]
#![allow(unused_must_use)]
#![allow(clippy::all)]
use std::collections::HashMap;
+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(all(test, feature = "test-util"))]
#[cfg(test)]
pub(crate) use ops::object::TransitionUploadedCommitBarrier as SetDiskTransitionUploadedCommitBarrier;
pub(crate) use ops::object::body_cache_plaintext_len;
#[cfg(all(test, feature = "test-util"))]
+2 -158
View File
@@ -11611,7 +11611,6 @@ mod tests {
pool_index: usize,
bucket: &str,
object: &str,
minio_unversioned: bool,
) {
for disk_index in 0..4 {
let metadata_path =
@@ -11645,11 +11644,6 @@ 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");
}
@@ -11660,152 +11654,6 @@ 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)]
@@ -11846,7 +11694,7 @@ mod tests {
)
.await
.expect("legacy source should transition");
rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 0, bucket, object, false).await;
rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 0, bucket, object).await;
let legacy = store.pools[0]
.get_disks_by_key(object)
.load_file_info_versions_exact(bucket, object)
@@ -12987,7 +12835,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", false).await;
rewrite_transitioned_xlmeta_as_legacy_unknown(temp_dir.path(), 1, bucket, "legacy/item.bin").await;
backend.set_remove_failure(true);
store.pools[1]
.delete_object(bucket, "hidden/item.bin", ObjectOptions::default())
@@ -17054,10 +16902,6 @@ 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(feature = "test-util")]
#[cfg(any(test, feature = "test-util"))]
pub use object::DeleteAfterObjectLockSnapshotBarrier;
pub(crate) use object::{
DecommissionFixedReadAnchor, ObjectLockDiagGuard, RemoteTuplePublicationCommitGuard, RemoteTuplePublicationFence,
+10 -115
View File
@@ -297,20 +297,6 @@ 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(),
@@ -380,26 +366,14 @@ 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: 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)?,
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),
})
}
}
@@ -464,14 +438,8 @@ impl FileInfo {
}
}
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)
{
fn set_transition_version_state(meta_sys: &mut HashMap<String, Vec<u8>>, state: TransitionVersionState) {
if state == TransitionVersionState::Unknown {
remove_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE);
} else {
insert_bytes(meta_sys, SUFFIX_TRANSITIONED_VERSION_STATE, state.as_str().as_bytes().to_vec());
@@ -2675,11 +2643,6 @@ 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
@@ -2726,7 +2689,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, &fi.metadata);
set_transition_version_state(&mut self.meta_sys, fi.transition_version_state);
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());
@@ -2867,7 +2830,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, &value.metadata);
set_transition_version_state(&mut meta_sys, value.transition_version_state);
}
if !value.transition_tier.is_empty() {
@@ -3022,12 +2985,6 @@ 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())?;
@@ -3195,7 +3152,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, &value.metadata);
set_transition_version_state(&mut meta_sys, value.transition_version_state);
}
if !value.transition_tier.is_empty() {
insert_bytes(&mut meta_sys, SUFFIX_TRANSITION_TIER, value.transition_tier.as_bytes().to_vec());
@@ -4617,7 +4574,6 @@ 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]
@@ -4629,10 +4585,6 @@ 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]
@@ -4646,7 +4598,6 @@ 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]
@@ -4686,36 +4637,6 @@ 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();
@@ -4832,10 +4753,6 @@ 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]
@@ -4878,10 +4795,6 @@ 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]
@@ -4899,7 +4812,6 @@ 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]
@@ -4998,23 +4910,6 @@ 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 {
-3
View File
@@ -52,9 +52,6 @@ static REMOTE_SCANNER_CYCLE_REFRESH: LazyLock<AsyncMutex<()>> = LazyLock::new(||
mod stream;
#[cfg(test)]
pub(crate) use stream::checkpoint_fixture_partial_return;
pub use stream::{RemoteScannerAdmission, RemoteScannerRequest, serve_remote_scanner_request};
pub(crate) use stream::{RemoteScannerOutcome, RemoteScannerScanSpec, scan_remote_bucket};
use stream::{RemoteScannerReplayCache, RemoteScannerRequestWire, RemoteScannerValidatedCycle};
@@ -1017,48 +1017,6 @@ fn finish_remote_scanner_stream(
#[cfg(test)]
const TEST_NEXT_CYCLE: u64 = 11;
#[cfg(test)]
pub(crate) async fn checkpoint_fixture_partial_return(progress: (u64, u64), entries_visited: u64) {
let request_id = Uuid::new_v4();
let writer_auth = FrameAuthenticator::for_test(request_id);
let reader_auth = FrameAuthenticator::for_test(request_id);
let mut bytes = Vec::new();
write_frame(
&mut bytes,
&writer_auth,
&mut 0,
&RemoteScannerFrame::terminal(
RemoteScannerProgress {
objects_scanned: progress.0,
directories_started: progress.1,
entries_visited,
},
RemoteScannerFrameResult::Partial,
),
)
.await
.expect("checkpoint partial frame must encode");
let frame = read_frame(&mut std::io::Cursor::new(bytes.as_slice()), &reader_auth, &mut 0)
.await
.expect("checkpoint progress frame must authenticate");
assert_eq!(frame.progress.entries_visited, entries_visited);
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new_with_progress_tracking(&parent, Default::default());
let result = consume_remote_scanner_stream(
std::io::Cursor::new(bytes),
parent,
budget.clone(),
"bucket",
DataUsageCacheSource::new(0, 0),
DataUsageScanPlanDigest([17; 32]),
reader_auth,
)
.await
.expect("checkpoint partial frame must decode");
assert!(matches!(result, RemoteScannerOutcome::Partial));
assert_eq!(budget.progress(), progress);
}
#[cfg(test)]
async fn consume_remote_scanner_stream<R>(
reader: R,
@@ -24,8 +24,6 @@ use std::io::Write;
use std::os::unix::fs::{PermissionsExt, symlink};
use std::sync::Mutex;
mod checkpoint_fixture;
/// Reset the process-global alert cooldown map; test-only.
fn reset_alert_cooldowns() {
*SCANNER_ALERT_EMISSION_COOLDOWN
@@ -1,410 +0,0 @@
// Copyright 2026 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use super::*;
use crate::scanner_budget::ScannerCycleBudgetConfig;
use crate::scanner_io::{ScannerDiskScanOutcome, ScannerIODisk};
use crate::storage_api::scanner_io::ObjectIO;
use crate::{DataUsageCacheSource, DataUsageScanPlanDigest};
use std::io::Cursor;
use tokio::io::AsyncReadExt;
const CACHE_NAME: &str = "bucket/checkpoint-fixture.bin";
const STATIC_OBJECTS: u64 = 24;
const MAX_CACHE_BYTES: u64 = 1024 * 1024;
const SOURCE: DataUsageCacheSource = DataUsageCacheSource::new(0, 0);
const PLAN: DataUsageScanPlanDigest = DataUsageScanPlanDigest([17; 32]);
/// Real cache persistence codec and CAS calls, backed by two bounded local files.
#[derive(Debug)]
struct FixtureStore {
root: tempfile::TempDir,
reject_save: AtomicBool,
}
impl FixtureStore {
fn new() -> Arc<Self> {
Arc::new(Self {
root: tempfile::tempdir().expect("checkpoint fixture storage directory"),
reject_save: AtomicBool::new(false),
})
}
fn path(&self, object: &str) -> std::path::PathBuf {
assert!(object.ends_with(CACHE_NAME) || object.ends_with(&format!("{CACHE_NAME}.bkp")));
self.root
.path()
.join(if object.ends_with(".bkp") { "backup" } else { "main" })
}
async fn strict_load(&self) -> DataUsageCache {
let bytes = tokio::fs::read(self.root.path().join("main"))
.await
.expect("saved checkpoint fixture must exist");
decode_fixture(&bytes).expect("saved checkpoint fixture must contain a valid bucket root")
}
}
#[async_trait::async_trait]
impl ObjectIO for FixtureStore {
type Error = crate::EcstoreError;
type RangeSpec = crate::storage_api::scanner_io::HTTPRangeSpec;
type HeaderMap = http::HeaderMap;
type ObjectOptions = crate::ScannerObjectOptions;
type ObjectInfo = crate::ScannerObjectInfo;
type GetObjectReader = crate::ScannerGetObjectReader;
type PutObjectReader = crate::ScannerPutObjReader;
async fn get_object_reader(
&self,
_bucket: &str,
object: &str,
_range: Option<Self::RangeSpec>,
_headers: Self::HeaderMap,
_options: &Self::ObjectOptions,
) -> crate::EcstoreResult<Self::GetObjectReader> {
let bytes = tokio::fs::read(self.path(object)).await.map_err(|error| {
if error.kind() == std::io::ErrorKind::NotFound {
crate::EcstoreError::FileNotFound
} else {
crate::EcstoreError::from(error)
}
})?;
assert!(u64::try_from(bytes.len()).expect("cache length") <= MAX_CACHE_BYTES);
Ok(crate::ScannerGetObjectReader {
stream: Box::new(Cursor::new(bytes)),
object_info: crate::ScannerObjectInfo {
etag: Some("fixture".into()),
..Default::default()
},
buffered_body: None,
body_source: Default::default(),
})
}
async fn put_object(
&self,
_bucket: &str,
object: &str,
data: &mut Self::PutObjectReader,
options: &Self::ObjectOptions,
) -> crate::EcstoreResult<Self::ObjectInfo> {
if self.reject_save.load(Ordering::SeqCst) {
return Err(crate::EcstoreError::PreconditionFailed);
}
let path = self.path(object);
let exists = tokio::fs::try_exists(&path).await?;
let preconditions = options.http_preconditions.as_ref().expect("checkpoint writes must use CAS");
if (exists && preconditions.if_none_match_value() == Some("*"))
|| (!exists && preconditions.if_match_value().is_some())
|| (exists && preconditions.if_match_value() != Some("fixture"))
{
return Err(crate::EcstoreError::PreconditionFailed);
}
let mut bytes = Vec::new();
(&mut data.stream).take(MAX_CACHE_BYTES + 1).read_to_end(&mut bytes).await?;
assert!(u64::try_from(bytes.len()).expect("cache length") <= MAX_CACHE_BYTES);
tokio::fs::write(path, bytes).await?;
Ok(crate::ScannerObjectInfo {
etag: Some("fixture".into()),
..Default::default()
})
}
}
#[async_trait::async_trait]
impl crate::ScannerConfigObjectDelete for FixtureStore {
async fn delete_config_object(
&self,
_bucket: &str,
_object: &str,
_options: crate::ScannerObjectOptions,
) -> crate::EcstoreResult<crate::ScannerObjectInfo> {
Err(crate::EcstoreError::NotImplemented)
}
async fn scanner_data_usage_publication_admission(&self) -> Option<crate::ScannerDataUsagePublicationAdmission> {
Some(crate::ScannerDataUsagePublicationAdmission::unfenced())
}
}
fn decode_fixture(bytes: &[u8]) -> Result<DataUsageCache, &'static str> {
if bytes.is_empty() || bytes.len() > usize::try_from(MAX_CACHE_BYTES).expect("fixture bound") {
return Err("missing or oversized checkpoint fixture");
}
let cache = DataUsageCache::unmarshal(bytes).map_err(|_| "corrupt checkpoint fixture")?;
if cache.info.name != "bucket" || cache.checked_flatten("bucket").is_none() {
return Err("checkpoint fixture has no valid bucket root");
}
Ok(cache)
}
fn retained(cache: &DataUsageCache) -> u64 {
assert!(
!cache.root().is_some_and(|root| root.compacted),
"a compacted bucket root cannot prove static-prefix coverage"
);
cache
.checked_flatten("bucket/static")
.map_or(0, |entry| u64::try_from(entry.objects).expect("fixture object count fits u64"))
}
#[derive(Debug, PartialEq, Eq)]
enum CoverageDiagnosis {
Progress,
NoNewWork,
LostAtPrepare,
LostAtReload,
WalkWithoutRetention,
}
fn diagnose(previous: u64, prepared: u64, walked: u64, scanned: u64, reloaded: u64) -> CoverageDiagnosis {
if reloaded < scanned {
CoverageDiagnosis::LostAtReload
} else if prepared < previous {
CoverageDiagnosis::LostAtPrepare
} else if walked > 0 && reloaded <= previous {
CoverageDiagnosis::WalkWithoutRetention
} else if reloaded > previous {
CoverageDiagnosis::Progress
} else {
CoverageDiagnosis::NoNewWork
}
}
#[test]
fn checkpoint_fixture_diagnosis_rejects_walk_without_retention() {
assert_eq!(diagnose(4, 4, 9, 8, 8), CoverageDiagnosis::Progress);
assert_eq!(diagnose(4, 4, 9, 4, 4), CoverageDiagnosis::WalkWithoutRetention);
assert_eq!(diagnose(4, 0, 9, 4, 4), CoverageDiagnosis::LostAtPrepare);
assert_eq!(diagnose(4, 4, 9, 8, 4), CoverageDiagnosis::LostAtReload);
assert_eq!(diagnose(4, 4, 0, 4, 4), CoverageDiagnosis::NoNewWork);
}
#[test]
fn checkpoint_fixture_missing_and_corrupt_inputs_fail() {
for bytes in [
vec![],
vec![0xc1],
DataUsageCache::default().marshal_msg().expect("empty cache encoding"),
vec![0; usize::try_from(MAX_CACHE_BYTES + 1).expect("oversized fixture")],
] {
assert!(decode_fixture(&bytes).is_err(), "invalid fixture must not become an empty complete root");
}
}
#[test]
fn checkpoint_fixture_compaction_preserves_aggregate_not_child_enumeration() {
let mut cache = DataUsageCache::default();
cache.info.name = "bucket".to_string();
cache.replace("bucket", "", DataUsageEntry::default());
cache.replace("bucket/static", "bucket", DataUsageEntry::default());
for index in 0..4 {
cache.replace(
&format!("bucket/static/{index}"),
"bucket/static",
DataUsageEntry {
objects: 1,
..Default::default()
},
);
}
cache.reduce_children_of(&hash_path("bucket/static"), 1, true);
let decoded = decode_fixture(&cache.marshal_msg().expect("encode compacted cache")).expect("decode compacted fixture");
let entry = decoded
.find("bucket/static")
.expect("compaction must retain the static subtree root");
assert!(entry.compacted);
assert!(entry.children.is_empty());
assert_eq!(
retained(&decoded),
4,
"compaction retains aggregate coverage even when leaf keys are absent"
);
}
#[tokio::test]
#[serial]
async fn checkpoint_fixture_save_reload_resume() {
run_checkpoint_fixture(false).await;
}
#[tokio::test]
#[serial]
async fn checkpoint_fixture_hot_digest_diagnostic() {
run_checkpoint_fixture(true).await;
}
async fn run_checkpoint_fixture(change_digest: bool) {
let (scanner, root) = build_test_scanner().await;
let _guard = TestGuard {
temp_dir: Some(root.clone()),
};
for index in 0..STATIC_OBJECTS {
write_test_object_metadata(&root, "bucket", &format!("static/{index:04}")).await;
}
let store = FixtureStore::new();
let mut previous = 0;
let mut visited = 0;
for round in 0..3_u8 {
write_test_object_metadata(&root, "bucket", "hot/current").await;
let mut cache = DataUsageCache::default();
let revisions = cache
.load_with_revisions(store.clone(), CACHE_NAME)
.await
.expect("load checkpoint revisions");
if round > 0 {
assert_eq!(retained(&store.strict_load().await), previous);
}
let plan = crate::scanner_io::checkpoint_fixture_bucket_digest(PLAN, change_digest.then_some(u64::from(round)));
crate::scanner_io::current_cache_root_or_prepare_with_generation(
&mut cache,
"bucket",
SOURCE,
11,
7,
plan,
crate::scanner_io::DataUsageCacheReuseOptions {
require_source: true,
tier_registry_generation: None,
},
);
let prepared = retained(&cache);
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new_with_progress_tracking(
&parent,
ScannerCycleBudgetConfig {
max_objects: Some(4),
..Default::default()
},
);
let outcome = scanner
.local_disk
.clone()
.nsscanner_disk(
budget.token(),
budget.clone(),
vec![scanner.local_disk.clone()],
cache,
None,
HealScanMode::Normal,
)
.await
.expect("budgeted local disk scan returns partial cache");
let ScannerDiskScanOutcome::Partial(cache) = outcome else {
panic!("budgeted fixture must remain partial")
};
assert!(!cache.info.snapshot_complete, "partial must never publish a complete root");
assert_eq!(budget.reason(), Some(crate::scanner_budget::ScannerCycleBudgetReason::Objects));
let scanned = retained(&cache);
cache
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
.await
.expect("persist partial checkpoint");
let mut loaded = DataUsageCache::default();
loaded
.load(store.clone(), CACHE_NAME)
.await
.expect("reload persisted partial checkpoint");
let reloaded = retained(&loaded);
assert_eq!(reloaded, retained(&store.strict_load().await));
assert_eq!(scanned, reloaded, "save/load must retain static subtree coverage");
assert!(!loaded.info.snapshot_complete);
visited += budget.entries_visited();
let diagnosis = diagnose(previous, prepared, budget.entries_visited(), scanned, reloaded);
eprintln!(
"checkpoint_fixture round={round} hot_digest={change_digest} visited_total={visited} before={previous} prepared={prepared} scanned={scanned} reloaded={reloaded} diagnosis={diagnosis:?}"
);
if !change_digest || std::env::var_os("RUSTFS_CHECKPOINT_REQUIRE_PROGRESS").is_some() {
assert_eq!(
diagnosis,
CoverageDiagnosis::Progress,
"visited growth must produce durable static coverage"
);
}
crate::remote_scanner::checkpoint_fixture_partial_return(budget.progress(), budget.entries_visited()).await;
previous = reloaded;
}
assert!(visited > 0, "fixture must exercise the directory walk");
assert!(previous > 0, "fixture must retain and enumerate static subtree entries");
let mut loaded = DataUsageCache::default();
let revisions = loaded
.load_with_revisions(store.clone(), CACHE_NAME)
.await
.expect("load final checkpoint");
let before = tokio::fs::read(store.root.path().join("main"))
.await
.expect("read durable checkpoint bytes");
let epoch_error = loaded
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 1)
.await
.expect_err("stale publication epoch must reject persistence");
assert!(epoch_error.to_string().contains(crate::SCANNER_PUBLICATION_EPOCH_CHANGED));
store.reject_save.store(true, Ordering::SeqCst);
loaded.info.next_cycle += 1;
loaded
.save_with_revisions_for_epoch(store.clone(), CACHE_NAME, &revisions, 0)
.await
.expect_err("injected save failure must not report durable progress");
assert_eq!(
tokio::fs::read(store.root.path().join("main"))
.await
.expect("read unchanged checkpoint bytes"),
before
);
let parent = CancellationToken::new();
parent.cancel();
let budget = ScannerCycleBudget::new(&parent, Default::default());
let result = scanner
.local_disk
.clone()
.nsscanner_disk(
budget.token(),
budget.clone(),
vec![scanner.local_disk.clone()],
loaded.clone(),
None,
HealScanMode::Normal,
)
.await;
assert!(result.is_err(), "pre-scan cancellation must not produce a complete root");
assert_eq!(budget.reason(), None, "parent cancellation is not object budget exhaustion");
let parent = CancellationToken::new();
let budget = ScannerCycleBudget::new(&parent, Default::default());
let result = scanner
.local_disk
.clone()
.nsscanner_disk(
budget.token(),
budget,
vec![scanner.local_disk.clone()],
loaded,
None,
HealScanMode::Normal,
)
.await
.expect("unbounded scan must complete after durable partial progress");
let ScannerDiskScanOutcome::Complete(cache) = result else {
panic!("unbounded fixture must produce a complete disk cache");
};
assert!(cache.info.snapshot_complete);
assert!(cache.info.scan_checkpoint.is_none());
assert_eq!(
cache.checked_flatten("bucket").expect("complete bucket root").objects,
usize::try_from(STATIC_OBJECTS + 1).expect("fixture object count fits usize")
);
}
-8
View File
@@ -209,14 +209,6 @@ fn scanner_bucket_cache_digest(
DataUsageScanPlanDigest(hasher.finalize().into())
}
#[cfg(test)]
pub(crate) fn checkpoint_fixture_bucket_digest(
scan_plan_digest: DataUsageScanPlanDigest,
dirty_generation: Option<u64>,
) -> DataUsageScanPlanDigest {
scanner_bucket_cache_digest(scan_plan_digest, dirty_generation)
}
fn finalize_nsscanner_result(results: &[DataUsageCache], first_err: Option<Error>) -> Result<()> {
if results.iter().any(|result| result.info.last_update.is_some()) {
return Ok(());
-21
View File
@@ -1048,27 +1048,6 @@ fn scanner_cycle_status_requires_a_clean_complete_snapshot() {
}
}
#[test]
fn checkpoint_fixture_superseded_is_distinct_from_partial_and_cancel() {
for (budget, cancelled, bucket, expected) in [
(false, false, ScannerBucketScanStatus::Complete, ScannerCycleStatus::Superseded),
(true, false, ScannerBucketScanStatus::Partial, ScannerCycleStatus::Incomplete),
(false, true, ScannerBucketScanStatus::Partial, ScannerCycleStatus::Incomplete),
] {
assert_eq!(
classify_nsscanner_cycle(
true,
budget,
cancelled,
bucket,
DirtyUsageSnapshotStatus::Changed,
ScannerCycleActivityStatus::Unchanged
),
expected,
);
}
}
#[test]
fn unverified_activity_defers_partial_and_floor_cycles() {
let expected = ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable);
-9
View File
@@ -49,14 +49,5 @@ rustfs-rio.workspace = true
tokio = { workspace = true, features = ["io-util", "macros", "rt"] }
thiserror = { workspace = true }
[dev-dependencies]
astral-tokio-tar = { workspace = true }
futures = { workspace = true }
serde = { workspace = true, features = ["derive"] }
serde_json = { workspace = true }
sha2 = { workspace = true }
tar-codec = { workspace = true }
tar-framing = { workspace = true }
[lints]
workspace = true
@@ -1,24 +0,0 @@
# minio-go Snowball fixtures
These request bodies are generated by
`github.com/minio/minio-go/v7.Client.PutObjectsSnowball` at the version pinned
in `generate/go.mod`. They cover the raw TAR and S2-compressed forms accepted by
RustFS Snowball extraction.
The decoded TAR intentionally ends immediately after the final padded member
body because minio-go flushes, rather than closes, its TAR writer. The
compatibility test permits that shape only when the authenticated request body
is complete at the exact member boundary; it does not make incomplete TAR
terminators generally valid.
Regenerate them from this directory with Go 1.25:
```console
cd generate
go mod download
go run . -out ..
```
`manifest.json` records the input objects and SHA-256 digest of each captured
request body. Review changes to the manifest and binary fixtures together when
updating minio-go.
@@ -1,26 +0,0 @@
module rustfs.local/snowball-fixture
go 1.25.0
require github.com/minio/minio-go/v7 v7.3.0
require (
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/dustin/go-humanize v1.0.1 // indirect
github.com/google/uuid v1.6.0 // indirect
github.com/klauspost/compress v1.19.2 // indirect
github.com/klauspost/cpuid/v2 v2.4.0 // indirect
github.com/klauspost/crc32 v1.3.0 // indirect
github.com/minio/crc64nvme v1.1.1 // indirect
github.com/minio/md5-simd v1.1.2 // indirect
github.com/philhofer/fwd v1.2.0 // indirect
github.com/rs/xid v1.6.0 // indirect
github.com/tinylib/msgp v1.6.4 // indirect
github.com/zeebo/xxh3 v1.1.0 // indirect
go.yaml.in/yaml/v3 v3.0.5 // indirect
golang.org/x/crypto v0.55.0 // indirect
golang.org/x/net v0.58.0 // indirect
golang.org/x/sys v0.47.0 // indirect
golang.org/x/text v0.41.0 // indirect
gopkg.in/ini.v1 v1.67.3 // indirect
)
@@ -1,59 +0,0 @@
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY=
github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto=
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/klauspost/compress v1.19.2 h1:hMRETovs/pu/dVWN7zIT1PGG8t509MwT6bO7XSi26R8=
github.com/klauspost/compress v1.19.2/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ=
github.com/klauspost/cpuid/v2 v2.0.1/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg=
github.com/klauspost/cpuid/v2 v2.4.0 h1:S6Hrbc7+ywsr0r+RLapfGBHfyefhCTwEh3A0tV913Dw=
github.com/klauspost/cpuid/v2 v2.4.0/go.mod h1:19jmZ9mjzoF//ddRSUsv0zfBTJWh3QJh9FNxZTMrGxU=
github.com/klauspost/crc32 v1.3.0 h1:sSmTt3gUt81RP655XGZPElI0PelVTZ6YwCRnPSupoFM=
github.com/klauspost/crc32 v1.3.0/go.mod h1:D7kQaZhnkX/Y0tstFGf8VUzv2UofNGqCjnC3zdHB0Hw=
github.com/minio/crc64nvme v1.1.1 h1:8dwx/Pz49suywbO+auHCBpCtlW1OfpcLN7wYgVR6wAI=
github.com/minio/crc64nvme v1.1.1/go.mod h1:eVfm2fAzLlxMdUGc0EEBGSMmPwmXD5XiNRpnu9J3bvg=
github.com/minio/md5-simd v1.1.2 h1:Gdi1DZK69+ZVMoNHRXJyNcxrMA4dSxoYHZSQbirFg34=
github.com/minio/md5-simd v1.1.2/go.mod h1:MzdKDxYpY2BT9XQFocsiZf/NKVtR7nkE4RoEpN+20RM=
github.com/minio/minio-go/v7 v7.3.0 h1:HM4pFCSQq/TK+j0/zmorSh5ddh81iDgRgU0BG0Vz/YU=
github.com/minio/minio-go/v7 v7.3.0/go.mod h1:KUPWdecEO1LWyUz+sTGXAuf2jZHrPh5fCsRH86QbPfk=
github.com/philhofer/fwd v1.2.0 h1:e6DnBTl7vGY+Gz322/ASL4Gyp1FspeMvx1RNDoToZuM=
github.com/philhofer/fwd v1.2.0/go.mod h1:RqIHx9QI14HlwKwm98g9Re5prTQ6LdeRQn+gXJFxsJM=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/rs/xid v1.6.0 h1:fV591PaemRlL6JfRxGDEPl69wICngIQ3shQtzfy2gxU=
github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw=
github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo=
github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA=
github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU=
github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo=
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/tinylib/msgp v1.6.4 h1:mOwYbyYDLPj35mkA2BjjYejgJk9BuHxDdvRnb6v2ZcQ=
github.com/tinylib/msgp v1.6.4/go.mod h1:RSp0LW9oSxFut3KzESt5Voq4GVWyS+PSulT77roAqEA=
github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ=
github.com/zeebo/assert v1.3.0/go.mod h1:Pq9JiuJQpG8JLJdtkwrJESF0Foym2/D9XMU5ciN/wJ0=
github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs=
github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s=
go.yaml.in/yaml/v3 v3.0.5 h1:N6y/pJk8buWs9NY5ERU2HSMfm+IuD/OtfdAnq6kESPw=
go.yaml.in/yaml/v3 v3.0.5/go.mod h1:HVTZu1O7/Vkt2N+BFy8Zza+lnLsABggaTM2ZpNIGuKg=
golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M=
golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis=
golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To=
golang.org/x/net v0.58.0/go.mod h1:YwCddHnFlT7eLQqVprV19OnhLGtc5xOKgE0RyqgfWAU=
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8=
golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/ini.v1 v1.67.3 h1:iM9Lhz5MRSGhHVGGwCuzG9KO8PoirCXj/m/qTmOJJQw=
gopkg.in/ini.v1 v1.67.3/go.mod h1:x/cyOwCgZqOkJoDIJ3c1KNHMo10+nLGAhh+kn3Zizss=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
@@ -1,193 +0,0 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package main
import (
"bytes"
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"flag"
"fmt"
"io"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strings"
"time"
"github.com/minio/minio-go/v7"
"github.com/minio/minio-go/v7/pkg/credentials"
)
const minioGoVersion = "v7.3.0"
type fixtureManifest struct {
Generator string `json:"generator"`
MinioGo string `json:"minio_go"`
GeneratedAt string `json:"generated_at"`
Objects []fixtureObject `json:"objects"`
Archives []fixtureArchive `json:"archives"`
}
type fixtureObject struct {
Key string `json:"key"`
Body string `json:"body"`
ModTime string `json:"mod_time"`
VersionID string `json:"version_id,omitempty"`
Headers map[string][]string `json:"headers,omitempty"`
}
type fixtureArchive struct {
File string `json:"file"`
Compressed bool `json:"compressed"`
Length int `json:"length"`
SHA256 string `json:"sha256"`
}
func objects() []fixtureObject {
return []fixtureObject{
{
Key: "alpha.txt",
Body: "alpha-body",
ModTime: "2024-01-02T03:04:05Z",
VersionID: "018cc251-f400-7c22-9e8d-8b1800000001",
Headers: map[string][]string{
"Content-Type": {"text/plain"},
"X-Amz-Meta-Owner": {"snowball-fixture"},
"X-Amz-Tagging": {"project=rustfs&source=minio-go"},
},
},
{
Key: "nested/世界.txt",
Body: "bravo-body",
ModTime: "2024-01-02T03:05:05Z",
Headers: map[string][]string{
"Content-Language": {"zh-CN"},
"X-Amz-Meta-Note": {"unicode-path"},
},
},
}
}
func captureSnowball(compressed bool, specs []fixtureObject) ([]byte, error) {
body := make(chan []byte, 1)
server := httptest.NewServer(http.HandlerFunc(func(writer http.ResponseWriter, request *http.Request) {
payload, err := io.ReadAll(request.Body)
if err != nil {
http.Error(writer, err.Error(), http.StatusInternalServerError)
return
}
body <- payload
writer.Header().Set("ETag", `"snowball-fixture"`)
writer.WriteHeader(http.StatusOK)
}))
defer server.Close()
client, err := minio.New(strings.TrimPrefix(server.URL, "http://"), &minio.Options{
// The S3 authentication layer removes AWS streaming-signature framing
// before Snowball extraction sees the request body. Anonymous signing
// captures those decoded archive bytes directly.
Creds: credentials.NewStatic("", "", "", credentials.SignatureAnonymous),
Secure: false,
Region: "us-east-1",
})
if err != nil {
return nil, fmt.Errorf("construct minio client: %w", err)
}
input := make(chan minio.SnowballObject, len(specs))
for _, spec := range specs {
modTime, err := time.Parse(time.RFC3339, spec.ModTime)
if err != nil {
return nil, fmt.Errorf("parse mod time for %q: %w", spec.Key, err)
}
headers := make(http.Header, len(spec.Headers))
for name, values := range spec.Headers {
headers[name] = append([]string(nil), values...)
}
input <- minio.SnowballObject{
Key: spec.Key,
Size: int64(len(spec.Body)),
ModTime: modTime,
Content: bytes.NewReader([]byte(spec.Body)),
VersionID: spec.VersionID,
Headers: headers,
}
}
close(input)
err = client.PutObjectsSnowball(context.Background(), "fixture-bucket", minio.SnowballOptions{
Opts: minio.PutObjectOptions{
ContentType: "application/octet-stream",
},
InMemory: true,
Compress: compressed,
}, input)
if err != nil {
return nil, fmt.Errorf("generate snowball request: %w", err)
}
return <-body, nil
}
func main() {
outDir := flag.String("out", "..", "fixture output directory")
flag.Parse()
specs := objects()
archives := make([]fixtureArchive, 0, 2)
for _, fixture := range []struct {
name string
compressed bool
}{
{name: "snowball.tar"},
{name: "snowball.tar.s2", compressed: true},
} {
payload, err := captureSnowball(fixture.compressed, specs)
if err != nil {
panic(err)
}
path := filepath.Join(*outDir, fixture.name)
if err := os.WriteFile(path, payload, 0o644); err != nil {
panic(fmt.Errorf("write %s: %w", path, err))
}
digest := sha256.Sum256(payload)
archives = append(archives, fixtureArchive{
File: fixture.name,
Compressed: fixture.compressed,
Length: len(payload),
SHA256: hex.EncodeToString(digest[:]),
})
}
manifest := fixtureManifest{
Generator: "github.com/minio/minio-go/v7.Client.PutObjectsSnowball",
MinioGo: minioGoVersion,
GeneratedAt: "2026-09-05T00:00:00Z",
Objects: specs,
Archives: archives,
}
payload, err := json.MarshalIndent(manifest, "", " ")
if err != nil {
panic(err)
}
payload = append(payload, '\n')
path := filepath.Join(*outDir, "manifest.json")
if err := os.WriteFile(path, payload, 0o644); err != nil {
panic(fmt.Errorf("write %s: %w", path, err))
}
}
@@ -1,51 +0,0 @@
{
"generator": "github.com/minio/minio-go/v7.Client.PutObjectsSnowball",
"minio_go": "v7.3.0",
"generated_at": "2026-09-05T00:00:00Z",
"objects": [
{
"key": "alpha.txt",
"body": "alpha-body",
"mod_time": "2024-01-02T03:04:05Z",
"version_id": "018cc251-f400-7c22-9e8d-8b1800000001",
"headers": {
"Content-Type": [
"text/plain"
],
"X-Amz-Meta-Owner": [
"snowball-fixture"
],
"X-Amz-Tagging": [
"project=rustfs\u0026source=minio-go"
]
}
},
{
"key": "nested/世界.txt",
"body": "bravo-body",
"mod_time": "2024-01-02T03:05:05Z",
"headers": {
"Content-Language": [
"zh-CN"
],
"X-Amz-Meta-Note": [
"unicode-path"
]
}
}
],
"archives": [
{
"file": "snowball.tar",
"compressed": false,
"length": 4096,
"sha256": "f00f2789dcb65b567f722f49cfdac9705e7bdac6c0badae75194327c32193d2e"
},
{
"file": "snowball.tar.s2",
"compressed": true,
"length": 528,
"sha256": "f8a9d9aa9b9ccdfae24ded1bff3741aacb935f1457a252efc9266674ff13c992"
}
]
}
Binary file not shown.
@@ -1,548 +0,0 @@
// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use std::collections::BTreeMap;
use std::fmt::Write as _;
use std::io::Cursor;
use futures::StreamExt;
use rustfs_zip::CompressionFormat;
use serde::Deserialize;
use sha2::{Digest, Sha256};
use tar_codec::{Archive as _, DecodePolicy, Member, MemberPayload as _, PaxDecodePolicy, PaxVendorExtensionPolicy, TarArchive};
use tar_framing::{
FrameError, FrameErrorInner, PaxKeyword, PaxRecord, PaxValue, StreamPolicy, UstarKind,
logical::{MemberExtensions, PaxState, TarReader},
};
use tokio::io::AsyncReadExt;
const FIXTURE_ROOT: &str = "fixtures/snowball/minio-go-v7.3.0";
const RAW_FIXTURE: &[u8] = include_bytes!("fixtures/snowball/minio-go-v7.3.0/snowball.tar");
const S2_FIXTURE: &[u8] = include_bytes!("fixtures/snowball/minio-go-v7.3.0/snowball.tar.s2");
const MANIFEST: &[u8] = include_bytes!("fixtures/snowball/minio-go-v7.3.0/manifest.json");
#[derive(Debug, Deserialize)]
struct FixtureManifest {
generator: String,
minio_go: String,
generated_at: String,
objects: Vec<FixtureObject>,
archives: Vec<FixtureArchive>,
}
#[derive(Debug, Deserialize)]
struct FixtureObject {
key: String,
body: String,
mod_time: String,
#[serde(default)]
version_id: String,
#[serde(default)]
headers: BTreeMap<String, Vec<String>>,
}
#[derive(Debug, Deserialize)]
struct FixtureArchive {
file: String,
compressed: bool,
length: usize,
sha256: String,
}
#[derive(Debug, Eq, PartialEq)]
struct ParsedMember {
path: String,
size: u64,
mtime: Option<u64>,
body: Vec<u8>,
minio_pax: BTreeMap<String, Option<Vec<u8>>>,
}
fn sha256_hex(bytes: &[u8]) -> String {
let mut encoded = String::with_capacity(64);
for byte in Sha256::digest(bytes) {
write!(&mut encoded, "{byte:02x}").expect("writing to a String should not fail");
}
encoded
}
async fn decode_s2(bytes: &[u8]) -> Vec<u8> {
let mut decoder = CompressionFormat::S2
.get_decoder(Cursor::new(bytes.to_vec()))
.expect("S2 fixture decoder should be available");
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).await.expect("S2 fixture should decode");
decoded
}
async fn parse_with_tokio_tar(bytes: &[u8]) -> Vec<ParsedMember> {
let mut archive = tokio_tar::Archive::new(Cursor::new(bytes.to_vec()));
let mut entries = archive.entries().expect("tokio-tar should create an entry stream");
let mut parsed = Vec::new();
while let Some(entry) = entries.next().await {
let mut entry = entry.expect("tokio-tar should parse the fixture member");
let kind = entry.header().entry_type();
if kind == tokio_tar::EntryType::XGlobalHeader {
continue;
}
let path_bytes = entry.path_bytes().expect("tokio-tar should resolve the fixture path");
let path = std::str::from_utf8(path_bytes.as_ref())
.expect("fixture paths should be UTF-8")
.to_owned();
let size = entry.effective_size();
let mtime = entry.header().mtime().ok();
let mut minio_pax = BTreeMap::new();
if let Some(extensions) = entry
.pax_extensions()
.await
.expect("tokio-tar should parse local PAX records")
{
for extension in extensions {
let extension = extension.expect("fixture PAX record should be valid");
let key = extension.key().expect("fixture PAX keys should be UTF-8");
if key.starts_with("minio.") {
minio_pax.insert(key.to_owned(), Some(extension.value_bytes().to_vec()));
}
}
}
let mut body = Vec::new();
entry
.read_to_end(&mut body)
.await
.expect("tokio-tar should read the fixture body");
parsed.push(ParsedMember {
path,
size,
mtime,
body,
minio_pax,
});
}
parsed
}
fn effective_minio_pax(state: &PaxState<'_>, known_keywords: &mut Vec<PaxKeyword>) -> BTreeMap<String, Option<Vec<u8>>> {
for extension in state.extensions() {
for record in extension.records() {
let keyword = record.keyword();
if matches!(&keyword, PaxKeyword::Vendor { vendor, .. } if vendor.as_ref() == "minio")
&& !known_keywords.contains(&keyword)
{
known_keywords.push(keyword);
}
}
}
known_keywords
.iter()
.filter_map(|keyword| {
let record = state.effective_record(keyword)?;
let PaxRecord::Vendor { vendor, name, value } = record else {
return None;
};
let key = format!("{vendor}.{name}");
let value = match value {
PaxValue::Value(value) => Some(value.to_vec()),
PaxValue::Deleted => None,
};
Some((key, value))
})
.collect()
}
fn effective_mtime(header_mtime: Option<u64>, extensions: &MemberExtensions<'_>) -> Option<u64> {
let MemberExtensions::Pax(state) = extensions else {
return header_mtime;
};
match state.effective_record(&PaxKeyword::Mtime) {
Some(PaxRecord::Mtime(PaxValue::Value(value))) => Some(*value),
Some(PaxRecord::Mtime(PaxValue::Deleted)) => None,
_ => header_mtime,
}
}
fn padded_member_end(position: u64, size: u64) -> u64 {
let padded_size = size.checked_add(511).expect("fixture member size should not overflow") / 512 * 512;
position
.checked_add(512)
.and_then(|position| position.checked_add(padded_size))
.expect("fixture member end should not overflow")
}
fn is_authenticated_footerless_end(error: &FrameError, last_member_end: Option<u64>, request_body_complete: bool) -> bool {
// The production gate must source `request_body_complete` from RustFS's
// length, checksum, and trailing-header validation state.
request_body_complete && matches!(&error.inner, FrameErrorInner::MissingEndMarker) && last_member_end == Some(error.position)
}
fn candidate_snowball_decode_policy() -> DecodePolicy {
DecodePolicy::default()
.allow_gnu(true)
.allow_all_nul_numeric_fields(true)
.max_gnu_extension_size(1_048_576)
.pax_policy(
PaxDecodePolicy::default()
.max_extension_size(1_048_576)
.max_global_extensions_size(67_108_864)
.allow_global_pax_extensions(false)
.allow_non_utf8_pax_vendor_values(false)
.allow_duplicate_pax_records(false)
.allow_global_pax_member_metadata(false)
.vendor_extension_policy(PaxVendorExtensionPolicy::ignore(["minio"])),
)
}
async fn parse_with_tar_framing(bytes: &[u8]) -> (Vec<ParsedMember>, Option<FrameError>, Option<u64>) {
let policy = StreamPolicy::default()
.max_pax_extension_size(1024 * 1024)
.max_global_pax_extensions_size(4 * 1024 * 1024)
.max_gnu_extension_size(128 * 1024);
let mut reader = TarReader::new(Cursor::new(bytes.to_vec())).with_policy(policy);
let mut parsed = Vec::new();
let mut known_minio_keywords = Vec::new();
let mut last_member_end = None;
loop {
let mut frame = match reader.next_frame().await {
Ok(Some(frame)) => frame,
Ok(None) => return (parsed, None, last_member_end),
Err(error) => return (parsed, Some(error), last_member_end),
};
assert_eq!(frame.header.kind, UstarKind::Regular);
let path = String::from_utf8(
frame
.effective_path()
.expect("tar-framing should resolve the fixture path")
.into_owned(),
)
.expect("fixture paths should be UTF-8");
let size = frame.header.effective_size;
let mtime = effective_mtime(frame.header.mtime, &frame.extensions);
let minio_pax = match &frame.extensions {
MemberExtensions::Pax(state) => effective_minio_pax(state, &mut known_minio_keywords),
MemberExtensions::Gnu { .. } => BTreeMap::new(),
};
let mut body = Vec::new();
let mut chunk = Vec::new();
while frame
.payload
.next_chunk(&mut chunk, 64 * 1024)
.await
.expect("tar-framing should read the fixture body")
{
body.extend_from_slice(&chunk);
}
last_member_end = Some(padded_member_end(frame.header.position, size));
parsed.push(ParsedMember {
path,
size,
mtime,
body,
minio_pax,
});
}
}
#[test]
fn checked_in_fixtures_match_the_minio_go_manifest() {
let manifest: FixtureManifest = serde_json::from_slice(MANIFEST).expect("fixture manifest should be valid JSON");
assert_eq!(manifest.generator, "github.com/minio/minio-go/v7.Client.PutObjectsSnowball");
assert_eq!(manifest.minio_go, "v7.3.0");
assert_eq!(manifest.generated_at, "2026-09-05T00:00:00Z");
assert_eq!(manifest.objects.len(), 2);
assert_eq!(manifest.objects[0].key, "alpha.txt");
assert_eq!(manifest.objects[0].body, "alpha-body");
assert_eq!(manifest.objects[0].mod_time, "2024-01-02T03:04:05Z");
assert_eq!(manifest.objects[0].version_id, "018cc251-f400-7c22-9e8d-8b1800000001");
assert_eq!(
manifest.objects[0].headers.get("X-Amz-Meta-Owner"),
Some(&vec!["snowball-fixture".to_owned()])
);
for archive in &manifest.archives {
let bytes = match archive.file.as_str() {
"snowball.tar" => RAW_FIXTURE,
"snowball.tar.s2" => S2_FIXTURE,
file => panic!("unexpected archive in {FIXTURE_ROOT}/manifest.json: {file}"),
};
assert_eq!(bytes.len(), archive.length);
assert_eq!(sha256_hex(bytes), archive.sha256);
assert_eq!(archive.compressed, archive.file.ends_with(".s2"));
}
}
#[tokio::test]
async fn minio_go_raw_and_s2_fixtures_have_identical_footerless_tar_data() {
assert_eq!(decode_s2(S2_FIXTURE).await, RAW_FIXTURE);
assert_eq!(RAW_FIXTURE.len() % 512, 0);
assert!(RAW_FIXTURE.len() >= 1024);
assert!(
!RAW_FIXTURE[RAW_FIXTURE.len() - 1024..].iter().all(|byte| *byte == 0),
"minio-go Flush output should not contain the standard two-block terminator"
);
}
#[tokio::test]
async fn tar_framing_matches_tokio_tar_before_rejecting_the_missing_terminator() {
let expected = parse_with_tokio_tar(RAW_FIXTURE).await;
let (actual, error, last_member_end) = parse_with_tar_framing(RAW_FIXTURE).await;
let error = error.expect("footerless minio-go fixture should fail strict termination");
assert_eq!(actual, expected);
assert_eq!(
actual,
[
ParsedMember {
path: "alpha.txt".to_owned(),
size: 10,
mtime: Some(1_704_164_645),
body: b"alpha-body".to_vec(),
minio_pax: BTreeMap::from([
("minio.metadata.Content-Type".to_owned(), Some(b"text/plain".to_vec()),),
("minio.metadata.X-Amz-Meta-Owner".to_owned(), Some(b"snowball-fixture".to_vec()),),
(
"minio.metadata.X-Amz-Tagging".to_owned(),
Some(b"project=rustfs&source=minio-go".to_vec()),
),
("minio.versionId".to_owned(), Some(b"018cc251-f400-7c22-9e8d-8b1800000001".to_vec()),),
]),
},
ParsedMember {
path: "nested/世界.txt".to_owned(),
size: 10,
mtime: Some(1_704_164_705),
body: b"bravo-body".to_vec(),
minio_pax: BTreeMap::from([
("minio.metadata.Content-Language".to_owned(), Some(b"zh-CN".to_vec()),),
("minio.metadata.X-Amz-Meta-Note".to_owned(), Some(b"unicode-path".to_vec()),),
]),
},
]
);
assert!(matches!(&error.inner, FrameErrorInner::MissingEndMarker));
assert_eq!(
error.position,
u64::try_from(RAW_FIXTURE.len()).expect("fixture length should fit in u64")
);
assert_eq!(last_member_end, Some(error.position));
}
#[tokio::test]
async fn footerless_compatibility_requires_authenticated_eof_at_the_member_boundary() {
let (_, error, last_member_end) = parse_with_tar_framing(RAW_FIXTURE).await;
let error = error.expect("the real fixture should be footerless");
assert!(is_authenticated_footerless_end(&error, last_member_end, true));
assert!(!is_authenticated_footerless_end(&error, last_member_end, false));
let mut one_zero_block = RAW_FIXTURE.to_vec();
one_zero_block.extend([0; 512]);
let (_, error, last_member_end) = parse_with_tar_framing(&one_zero_block).await;
let error = error.expect("one zero block is not a valid TAR terminator");
assert!(matches!(&error.inner, FrameErrorInner::MissingEndMarker));
assert_eq!(
last_member_end,
Some(u64::try_from(RAW_FIXTURE.len()).expect("fixture length should fit in u64"))
);
assert_eq!(
error.position,
u64::try_from(one_zero_block.len()).expect("fixture length should fit in u64")
);
assert!(!is_authenticated_footerless_end(&error, last_member_end, true));
}
#[tokio::test]
async fn tar_codec_policy_accepts_only_the_explicit_minio_vendor_namespace() {
let default_error = match TarArchive::new(Cursor::new(RAW_FIXTURE.to_vec())).members().next().await {
Err(error) => error,
Ok(_) => panic!("the default policy should reject minio vendor records"),
};
assert!(default_error.to_string().contains("pax vendor extension minio."));
let mut members = TarArchive::new(Cursor::new(RAW_FIXTURE.to_vec()))
.with_policy(candidate_snowball_decode_policy())
.members();
let mut bodies = Vec::new();
loop {
let member = match members.next().await {
Ok(Some(member)) => member,
Ok(None) => panic!("footerless minio-go fixture should not report a valid archive end"),
Err(error) => {
assert!(error.to_string().contains("missing two-block end-of-archive marker"));
break;
}
};
let Member::File { mut payload, .. } = member else {
panic!("fixture should contain only regular files");
};
let mut body = Vec::new();
let mut chunk = Vec::new();
while payload
.next_chunk(&mut chunk, 64 * 1024)
.await
.expect("tar-codec should read the fixture body")
{
body.extend_from_slice(&chunk);
}
bodies.push(body);
}
assert_eq!(bodies, [b"alpha-body".to_vec(), b"bravo-body".to_vec()]);
assert!(
members
.next()
.await
.expect("the member cursor should be fused after an error")
.is_none()
);
}
fn pax_record(key: &str, value: &str) -> Vec<u8> {
let payload = format!("{key}={value}\n");
let mut len = payload.len() + 3;
loop {
let record = format!("{len} {payload}");
if record.len() == len {
return record.into_bytes();
}
len = record.len();
}
}
async fn append_pax_header(
builder: &mut tokio_tar::Builder<Cursor<Vec<u8>>>,
entry_type: tokio_tar::EntryType,
records: &[(&str, &str)],
) {
let mut payload = Vec::new();
for (key, value) in records {
payload.extend(pax_record(key, value));
}
let mut header = tokio_tar::Header::new_ustar();
header.set_entry_type(entry_type);
header.set_size(u64::try_from(payload.len()).expect("PAX test payload should fit in u64"));
header.set_mode(0o644);
header.set_cksum();
builder
.append_data(&mut header, "PaxHeaders.X/snowball", Cursor::new(payload))
.await
.expect("PAX test header should be written");
}
async fn append_regular(builder: &mut tokio_tar::Builder<Cursor<Vec<u8>>>, path: &str) {
let body = path.as_bytes();
let mut header = tokio_tar::Header::new_ustar();
header.set_entry_type(tokio_tar::EntryType::Regular);
header.set_size(u64::try_from(body.len()).expect("test member body should fit in u64"));
header.set_mode(0o644);
header.set_mtime(1_704_164_645);
header.set_cksum();
builder
.append_data(&mut header, path, Cursor::new(body))
.await
.expect("ordinary test member should be written");
}
async fn archive_with_local_pax(records: &[(&str, &str)]) -> Vec<u8> {
let mut builder = tokio_tar::Builder::new(Cursor::new(Vec::new()));
append_pax_header(&mut builder, tokio_tar::EntryType::XHeader, records).await;
append_regular(&mut builder, "member.txt").await;
builder.into_inner().await.expect("policy archive should finish").into_inner()
}
#[tokio::test]
async fn candidate_policy_rejects_unknown_vendor_and_duplicate_pax_records() {
let unknown_vendor = archive_with_local_pax(&[("acme.metadata.owner", "mallory")]).await;
let error = match TarArchive::new(Cursor::new(unknown_vendor))
.with_policy(candidate_snowball_decode_policy())
.members()
.next()
.await
{
Err(error) => error,
Ok(_) => panic!("the candidate Snowball policy should reject unknown vendors"),
};
assert!(
error
.to_string()
.contains("pax vendor extension acme.metadata.owner is not allowed")
);
let duplicate = archive_with_local_pax(&[
("minio.metadata.x-amz-meta-owner", "first"),
("minio.metadata.x-amz-meta-owner", "second"),
])
.await;
let error = match TarArchive::new(Cursor::new(duplicate))
.with_policy(candidate_snowball_decode_policy())
.members()
.next()
.await
{
Err(error) => error,
Ok(_) => panic!("the candidate Snowball policy should reject duplicate PAX records"),
};
assert!(
error
.to_string()
.contains("pax extended header contains duplicate record minio.metadata.x-amz-meta-owner")
);
}
#[tokio::test]
async fn global_minio_pax_inheritance_is_an_explicit_migration_difference() {
let mut builder = tokio_tar::Builder::new(Cursor::new(Vec::new()));
append_pax_header(
&mut builder,
tokio_tar::EntryType::XGlobalHeader,
&[("minio.metadata.x-amz-meta-owner", "global")],
)
.await;
append_pax_header(
&mut builder,
tokio_tar::EntryType::XHeader,
&[("minio.metadata.x-amz-meta-owner", "local")],
)
.await;
append_regular(&mut builder, "local.txt").await;
append_regular(&mut builder, "inherited.txt").await;
let archive = builder
.into_inner()
.await
.expect("precedence archive should finish")
.into_inner();
let legacy = parse_with_tokio_tar(&archive).await;
let (framing, error, _) = parse_with_tar_framing(&archive).await;
assert!(error.is_none());
assert_eq!(legacy.len(), 2);
assert_eq!(framing.len(), 2);
let owner_key = "minio.metadata.x-amz-meta-owner";
assert_eq!(legacy[0].minio_pax.get(owner_key), Some(&Some(b"local".to_vec())));
assert!(!legacy[1].minio_pax.contains_key(owner_key));
assert_eq!(framing[0].minio_pax.get(owner_key), Some(&Some(b"local".to_vec())));
assert_eq!(framing[1].minio_pax.get(owner_key), Some(&Some(b"global".to_vec())));
let error = match TarArchive::new(Cursor::new(archive))
.with_policy(candidate_snowball_decode_policy())
.members()
.next()
.await
{
Err(error) => error,
Ok(_) => panic!("the candidate Snowball policy should reject global PAX state"),
};
assert!(error.to_string().contains("global pax extended headers are not allowed"));
}
+2 -2
View File
@@ -37,8 +37,8 @@ unknown-git = "deny"
allow-registry = ["https://github.com/rust-lang/crates.io-index"]
allow-git = [
# Temporary tokio-tar fork pinned to the reviewed parser limits,
# cancellation safety, and error-fusing change while Snowball is
# prototyped against tar-codec and Swift retains its current reader.
# cancellation safety, and error-fusing change while
# astral-sh/tokio-tar#118 awaits an upstream release.
# owner: cxymds review: 2026-10
"https://github.com/cxymds/tokio-tar.git",
# Official s3s repository. Temporarily pinned to the merged generic REST
+1 -1
View File
@@ -13,7 +13,7 @@
- `backlog-1337` legacy restore orphan recovery: releases that predate the restore worker-lock marker can leave a valid operation-id and `ongoing-request="true"` after cancellation or process failure, with no durable liveness proof. New servers allow an exact, non-nil legacy generation to be superseded only when its consistently parsed request date is at least 24 hours old. Remove the clock-based legacy fallback after the minimum supported direct-upgrade release writes the v1 worker-lock marker on every restore and operators have resolved every retained pre-v1 ongoing generation.
- `backlog-2133-tier-delete-chunk-parent` bounded tier-delete dispatch compatibility: prefixes at or below the legacy manifest limit keep the byte-compatible v1 single-manifest protocol, while larger prefixes place a chunk-parent sentinel at the original deterministic root path and use operation-scoped child manifests. Older binaries reject the sentinel and child paths, preserving the v6 sole-owner downgrade fence instead of starting a competing local delete. Remove the v1 reader and fail-closed mixed-version sentinel only after every supported rollback release validates the parent/child protocol and migration tooling confirms that no retained v1 dispatch manifest remains.
- `tokio-tar-extension-limits` bounded archive parser hardening: Snowball extraction depends on precedence-resolved MinIO PAX metadata; per-entry and cumulative extension limits; a physical-entry limit; cancellation-safe parsing and ownership of large streamed members; fused streams after errors; and compatibility with minio-go streams that omit the two-block terminator. Swift bulk extraction also uses the same fork. Keep the reviewed pin while the Snowball path is prototyped against tar-codec/tar-framing. Remove it only after a released API exposes the effective allowed vendor records, RustFS provides a cancellation-safe handoff for borrowed member payloads, footerless input is accepted solely when authenticated request framing proves EOF immediately after a complete member, the existing resource-limit, cancellation, error-fuse, and real minio-go fixtures pass against the replacement, and Swift no longer depends on the fork.
- `tokio-tar-extension-limits` bounded archive parser hardening: Snowball extraction depends on per-entry and cumulative GNU long-name, GNU long-link, and PAX extension limits; physical-entry, GNU sparse-map, and sparse-continuation limits; cancellation-safe sparse parsing; and fused entry streams after parser errors. The released tokio-tar API does not provide this complete boundary. Keep the reviewed fork pin until astral-sh/tokio-tar#118 is merged and one published tokio-tar release contains every listed capability with the Snowball regression fixtures passing against that release.
- `backlog-2102` rc.2/rc.3 empty scanner usage floor recovery: old DeleteBucket cleanup could synthesize an empty incomplete v2 usage primary/backup before leadership added an epoch, while newer scanners require a durable authoritative baseline identity. New scanners recognize only that exact serialized empty-fence shape, preserve its epoch through a CAS-protected recovery marker, and rebuild namespace coverage without treating zero usage as authoritative. Remove this recovery path and marker after rc.2 and rc.3 are no longer supported direct-upgrade sources.
- `backlog-2122` rc.1-rc.3 non-empty scanner usage floor recovery: leadership fencing in those releases can stamp scanner_epoch onto a real bucket-usage snapshot before any scanner cycle completed, leaving a non-empty floor with no scanner_cycle and no authoritative baseline identity. New scanners recognize only this consistent incomplete fenced shape, preserve the epoch through the CAS-protected recovery marker, and rebuild namespace coverage without treating the old usage data as authoritative. Remove this recovery path after rc.1, rc.2, and rc.3 are no longer supported direct-upgrade sources.
- `s3gate-metadata-xml` persisted bucket XML migration: mixed-version site-replication peers, retained `.metadata.bin` objects, and backup archives can all carry XML written by the s3s codec, so the gateway migration must keep the legacy codec available until every stored form has crossed a verified rewrite boundary. Remove the legacy s3s parser and serializer only after the minimum supported direct-upgrade release reads and writes every persisted XML configuration family through the gateway codec, every supported mixed-version site-replication topology has completed its writer upgrade, and migration tooling has verified or rewritten every retained bucket metadata object and restorable backup archive.
-2
View File
@@ -21,8 +21,6 @@ Pick the lowest layer that can prove the change; add a higher-layer test only wh
Every script named above is indexed with status and wiring in [`scripts/README.md`](../../scripts/README.md). Fixed GHSA advisories map to named regression tests in [security-regressions.md](security-regressions.md).
The [scanner checkpoint fixture](scanner-checkpoint-fixture.md) diagnoses retained subtree coverage across budget interruption, persistence, reload, and plan invalidation.
## Naming conventions
### Reserved test-name substrings (migration gate)
@@ -1,22 +0,0 @@
# Scanner Checkpoint Fixture
The `checkpoint_fixture` tests exercise a bounded namespace of 24 static objects and one repeatedly updated hot object. Each of three rounds runs the production local disk scanner with an object budget, saves the returned partial cache through the production persistence codec and revision checks to a two-file test backend, and reloads it before preparing the next round. The fixture prints static-subtree coverage at each boundary and cumulative visited entries. This is a diagnostic of retained coverage, not a throughput benchmark.
Run the fixture and confirm the test filter selects a nonzero number of tests:
```sh
cargo test -p rustfs-scanner --lib checkpoint_fixture -- --list
RUST_MIN_STACK=4194304 cargo test -p rustfs-scanner --lib checkpoint_fixture -- --nocapture
```
The unchanged-plan case requires durable static coverage to increase each round. The hot-plan diagnostic changes the bucket plan digest between rounds and reports where coverage is lost without asserting that a particular defect must remain present. To require progress in this diagnostic as well:
```sh
RUST_MIN_STACK=4194304 RUSTFS_CHECKPOINT_REQUIRE_PROGRESS=1 cargo test -p rustfs-scanner --lib checkpoint_fixture_hot_digest_diagnostic -- --nocapture
```
A nonzero exit from the strict command means that walked work did not become additional retained static coverage. `LostAtPrepare` identifies invalidation before traversal; `LostAtReload` identifies loss between the returned cache and persisted data; `WalkWithoutRetention` identifies visited growth without durable coverage growth. Missing, corrupt, empty-root, and oversized checkpoint inputs are rejected by the strict fixture reader. Save failure and publication-epoch rejection must preserve the preceding file bytes. Parent cancellation is checked separately from object-budget exhaustion. Superseded classification is tested separately from either incomplete outcome.
For every saved partial cache, the fixture also passes its progress through the production authenticated remote terminal-frame writer and stream consumer. A remote partial result must remain partial even when its progress reports visited objects. This covers the return-frame contract; it does not execute the remote RPC server, distributed locks, EC quorum persistence, mixed-version peers, process crashes, or fsync durability. The file backend models revision preconditions and persistence errors, not a concurrent object store.
The synthetic namespace contains no customer data. Temporary files are removed with their owning fixture. Production scan semantics and persistent formats are unchanged, so rollback consists of removing these tests and this guide. A passing fixture alone does not establish that the field report in [issue #7108](https://github.com/rustfs/rustfs/issues/7108) has been independently reproduced or fixed. A field diagnosis must separately identify the source capture, cycle and leader identity, and decoded bucket/set caches.
+1 -11
View File
@@ -121,11 +121,6 @@ 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()),
}
}
@@ -758,12 +753,7 @@ impl Operation for ListRemoteTargetHandler {
.map_err(ApiError::from)?;
let sys = BucketTargetSys::get();
// 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 = sys.list_targets(bucket, "").await;
let targets: Vec<_> = targets
.iter()
+246 -420
View File
@@ -190,8 +190,7 @@ fn site_replicator_service_account_policy() -> S3Result<Policy> {
.map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("parse site replicator policy failed: {e}")))
}
// Lock order: lifecycle -> bucket-mutation admission -> per-bucket mutation
// -> bucket operation -> repair admission -> state -> per-bucket metadata.
// Lock order: lifecycle -> bucket operation -> repair admission -> state -> per-bucket metadata.
// "state" is the distributed state-object lock in
// crate::site_replication::state_lock, entered through
// update_site_replication_state (P1-15). There is no process-local state
@@ -435,7 +434,6 @@ pub fn register_site_replication_route(r: &mut S3Router<AdminOperation>) -> std:
// into this module: startup sits below this layer and must not depend upwards. The admin
// router is built before startup reconciles, so the hook is always installed in time.
crate::site_replication_reconcile::register_site_replication_reconciler(reconcile_site_replication_wiring);
crate::site_replication_reconcile::register_site_replication_retry_drainer(reconcile_site_replication_retry_drain);
for (method, path, operation) in [
(Method::PUT, "/v3/site-replication/add", AdminOperation(&SiteReplicationAddHandler {})),
@@ -1805,61 +1803,28 @@ async fn reconcile_site_replication_buckets() -> S3Result<()> {
/// (`SiteReplicationEditHandler`), so a tick landing between them would rewrite the targets
/// from the stale endpoint. The pending marker in the persisted state closes that window.
/// Skipping costs nothing — the timer comes back.
async fn site_replication_reconcile_prerequisites_ready() -> bool {
if current_iam_handle().is_none() || current_object_store_handle().is_none() {
return false;
}
if let Err(err) = migrate_collapsed_retry_queue_paths().await {
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
result = "retry_queue_migration_failed",
error = ?err,
"admin site replication state"
);
return false;
}
true
}
fn reconcile_site_replication_retry_drain() -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>> {
Box::pin(async {
let Some(lifecycle) = SiteReplicationLifecycleGuard::try_acquire() else {
return;
};
if !site_replication_reconcile_prerequisites_ready().await {
return;
}
match load_site_replication_state().await {
Ok(state) => {
if state.pending_endpoint_refresh.is_some() || state.pending_rotation.is_some() || state.pending_remove.is_some()
{
return;
}
}
Err(_) => return,
}
// Admission above observes a lifecycle-stable state. The lightweight
// drain itself handles only idempotent bucket setup, reloads state
// under the distributed repair lock, and shares that lock with bucket
// deletion. Do not hold this process-local guard across peer I/O: an
// outage recovery must not make admin add/edit/remove time out.
drop(lifecycle);
drain_site_replication_retry_queue_lightweight().await;
})
}
fn reconcile_site_replication_wiring() -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>> {
Box::pin(async {
// The scheduler starts before IAM and the object store are guaranteed ready (IAM
// bootstrap may still be recovering), so an early tick returns quietly instead of
// logging a failure for every reconciler.
let Some(lifecycle) = SiteReplicationLifecycleGuard::try_acquire() else {
if current_iam_handle().is_none() || current_object_store_handle().is_none() {
return;
}
let Some(_lifecycle) = SiteReplicationLifecycleGuard::try_acquire() else {
return;
};
if !site_replication_reconcile_prerequisites_ready().await {
if let Err(err) = migrate_collapsed_retry_queue_paths().await {
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
result = "retry_queue_migration_failed",
error = ?err,
"admin site replication state"
);
return;
}
@@ -1913,9 +1878,8 @@ fn reconcile_site_replication_wiring() -> std::pin::Pin<Box<dyn std::future::Fut
"admin site replication state"
);
}
// The retry path re-checks membership from distributed state before
// each request; release lifecycle before a potentially large replay.
drop(lifecycle);
// Failed peer deliveries recorded in the retry queue; runs behind the
// same lifecycle guard and pending_* gates as the reconcilers above.
drain_site_replication_retry_queue().await;
})
}
@@ -3082,7 +3046,6 @@ fn set_pending_endpoint_refresh(state: &mut SiteReplicationState, pending: Pendi
last_error: "endpoint target refresh pending".to_string(),
updated_at: Some(OffsetDateTime::now_utc()),
edit_generation: None,
peer_unreachable: false,
deletions_recorded: false,
});
state.pending_endpoint_refresh = Some(pending);
@@ -3629,15 +3592,16 @@ const PEER_EDIT_FENCE_STALENESS_WINDOW_NANOS: u64 = 24 * 60 * 60 * 1_000_000_000
/// must be a site this state currently replicates with — the same membership
/// rule the load-time mark pruning applies, so every mark recorded behind
/// this check is one a reload would keep — and not this site itself, which
/// never delivers edits to itself. The caller acknowledges an inadmissible
/// fenced request without applying it: after a remove commits, an older
/// in-flight retry from the departed origin must not recreate topology. Old
/// peers remain compatible because their unstamped edits still follow the
/// pre-fence path. The generation itself is NOT bounded here: a genuine
/// origin whose hybrid clock persisted a wall-clock excursion allocates
/// arbitrarily far in the future, and refusing to record its marks would
/// strip the ordering fence from exactly the deliveries that still race —
/// the staleness window on the read side is what defuses forged marks instead.
/// never delivers edits to itself. The caller IGNORES an inadmissible fence
/// rather than failing the request: the delivery applies exactly as an
/// unstamped (pre-fence) delivery would, no high-water mark is read or
/// written, and the worst a forged fence achieves is forfeiting an ordering
/// guarantee its sender was never owed. The generation itself is NOT
/// bounded here: a genuine origin whose hybrid clock persisted a wall-clock
/// excursion allocates arbitrarily far in the future, and refusing to
/// record its marks would strip the ordering fence from exactly the
/// deliveries that still race — the staleness window on the read side is
/// what defuses forged marks instead.
fn peer_edit_fence_is_admissible(state: &SiteReplicationState, local_deployment_id: &str, fence: &(String, u64)) -> bool {
let (origin, generation) = fence;
if origin != local_deployment_id && state.peers.contains_key(origin) {
@@ -4827,135 +4791,105 @@ async fn backfill_existing_buckets_after_add(
let resync_id = Uuid::new_v4().to_string();
for bucket in &buckets {
let operation_name = bucket.name.clone();
let lock_bucket = operation_name.clone();
let operation_state = state.clone();
let operation_local_peer = local_peer.clone();
let operation_resync_id = resync_id.clone();
let operation_bootstrap_token = bootstrap_token.map(str::to_owned);
let bucket_errors = with_site_replication_bucket_mutation_lock(store.clone(), &lock_bucket, move || async move {
let mut errors = SiteReplicationErrorSummary::default();
let name = &operation_name;
let name = &bucket.name;
if let Err(err) = ensure_site_replication_bucket_versioning(name).await {
if let Err(err) = ensure_site_replication_bucket_versioning(name).await {
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
bucket = %name,
result = "backfill_versioning_setup_failed",
error = ?err,
"admin site replication state"
);
errors.push(format!("{name}: versioning setup failed: {err}"));
continue;
}
match ensure_site_replication_bucket_setup(name).await {
Ok(true) => {}
Ok(false) => {
// Runtime targets unavailable: the setup silently no-ops, which would make the
// downstream make-bucket broadcast and resync fail. Record it and skip so the
// operator sees this bucket was not propagated instead of an unqualified success.
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
bucket = %name,
result = "backfill_versioning_setup_failed",
error = ?err,
result = "backfill_bucket_setup_skipped",
"admin site replication state"
);
errors.push(format!("{name}: versioning setup failed: {err}"));
return errors;
errors.push(format!("{name}: replication setup skipped (site replication runtime unavailable)"));
continue;
}
match ensure_site_replication_bucket_setup(name).await {
Ok(true) => {}
Ok(false) => {
// Runtime targets unavailable: the setup silently no-ops, which would make the
// downstream make-bucket broadcast and resync fail. Record it and skip so the
// operator sees this bucket was not propagated instead of an unqualified success.
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
bucket = %name,
result = "backfill_bucket_setup_skipped",
"admin site replication state"
);
errors.push(format!("{name}: replication setup skipped (site replication runtime unavailable)"));
return errors;
}
Err(err) => {
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
bucket = %name,
result = "backfill_bucket_setup_failed",
error = ?err,
"admin site replication state"
);
errors.push(format!("{name}: bucket setup failed: {err}"));
}
}
// Broadcast the bucket to peers so they create it too (idempotent on the peer side).
// Read the real lock_enabled flag so peers recreate the bucket with the same object-lock
// setting — object lock cannot be added after bucket creation.
let lock_enabled = match metadata_sys::get(name).await {
Ok(bm) => bm.lock_enabled,
Err(err) => {
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
bucket = %name,
result = "backfill_bucket_metadata_read_failed",
fallback = "lock_enabled=false",
error = ?err,
"admin site replication state"
);
false
}
};
if let Err(err) =
broadcast_site_replication_make_bucket(name, lock_enabled, None, operation_bootstrap_token.as_deref()).await
{
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
bucket = %name,
result = "backfill_make_bucket_broadcast_failed",
error = ?err,
"admin site replication state"
);
errors.push(format!("{name}: make-bucket broadcast failed: {err}"));
}
// Kick a resync toward every remote peer so existing objects travel across.
for peer in operation_state.peers.values() {
if peer.deployment_id == operation_local_peer.deployment_id
|| same_identity_endpoint(&peer.endpoint, &operation_local_peer.endpoint)
{
continue;
}
let manifest = site_bucket_resync_manifest_entry(name, peer, OffsetDateTime::now_utc()).await;
let result = if manifest.target_arn.is_empty() {
manifest
} else {
start_site_bucket_resync(name, &manifest.target_arn, &operation_resync_id).await
};
if result.status == "failed" {
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
bucket = %name,
peer = %peer.endpoint,
result = "backfill_resync_kick_failed",
detail = %result.err_detail,
"admin site replication state"
);
errors.push(format!("{name} -> {}: resync kick failed: {}", peer.endpoint, result.err_detail));
}
}
errors
})
.await;
match bucket_errors {
Ok(bucket_errors) => errors.extend(bucket_errors),
Err(err) => {
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
bucket = %lock_bucket,
result = "backfill_bucket_mutation_lock_failed",
bucket = %name,
result = "backfill_bucket_setup_failed",
error = ?err,
"admin site replication state"
);
errors.push(format!("{lock_bucket}: bucket mutation lock failed: {err}"));
errors.push(format!("{name}: bucket setup failed: {err}"));
}
}
// Broadcast the bucket to peers so they create it too (idempotent on the peer side).
// Read the real lock_enabled flag so peers recreate the bucket with the same object-lock
// setting — object lock cannot be added after bucket creation.
let lock_enabled = match metadata_sys::get(name).await {
Ok(bm) => bm.lock_enabled,
Err(err) => {
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
bucket = %name,
result = "backfill_bucket_metadata_read_failed",
fallback = "lock_enabled=false",
error = ?err,
"admin site replication state"
);
false
}
};
if let Err(err) = broadcast_site_replication_make_bucket(name, lock_enabled, None, bootstrap_token).await {
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
bucket = %name,
result = "backfill_make_bucket_broadcast_failed",
error = ?err,
"admin site replication state"
);
errors.push(format!("{name}: make-bucket broadcast failed: {err}"));
}
// Kick a resync toward every remote peer so existing objects travel across.
for peer in state.peers.values() {
if peer.deployment_id == local_peer.deployment_id || same_identity_endpoint(&peer.endpoint, &local_peer.endpoint) {
continue;
}
let manifest = site_bucket_resync_manifest_entry(name, peer, OffsetDateTime::now_utc()).await;
let result = if manifest.target_arn.is_empty() {
manifest
} else {
start_site_bucket_resync(name, &manifest.target_arn, &resync_id).await
};
if result.status == "failed" {
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
bucket = %name,
peer = %peer.endpoint,
result = "backfill_resync_kick_failed",
detail = %result.err_detail,
"admin site replication state"
);
errors.push(format!("{name} -> {}: resync kick failed: {}", peer.endpoint, result.err_detail));
}
}
}
@@ -6138,204 +6072,146 @@ fn parse_peer_join_response(body: &[u8], fallback_peer: PeerInfo) -> Result<SRPe
serde_json::from_slice(body)
}
fn ensure_add_bucket_set_matches_preflight(expected: &HashSet<String>, present: &HashSet<String>) -> S3Result<()> {
let mut missing = expected.difference(present).cloned().collect::<Vec<_>>();
if !missing.is_empty() {
missing.sort_unstable();
return Err(S3Error::with_message(
S3ErrorCode::InvalidRequest,
format!(
"bucket `{}` disappeared while site replication was being added; peers may already be joined — re-run replicate add",
missing[0]
),
));
}
let mut unexpected = present.difference(expected).cloned().collect::<Vec<_>>();
if !unexpected.is_empty() {
unexpected.sort_unstable();
return Err(S3Error::with_message(
S3ErrorCode::InvalidRequest,
format!(
"bucket `{}` appeared while site replication was being added; peers may already be joined — re-run replicate add",
unexpected[0]
),
));
}
Ok(())
}
#[async_trait::async_trait]
impl Operation for SiteReplicationAddHandler {
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
let cred = validate_site_replication_admin_request(&req, AdminAction::SiteReplicationAddAction).await?;
reject_site_replicator_on_public_admin(&cred)?;
let replicate_ilm_expiry = sr_add_replicate_ilm_expiry(&req.uri);
let local_endpoint = site_replication_local_endpoint(&req.uri, &req.headers);
let lifecycle_guard = SiteReplicationLifecycleGuard::acquire().await?;
// Everything up to the commit below is preflight: peer probes, IAM
// work and the join fan-out all talk to the network, so none of it may
// run inside the state transaction. The snapshot read here is what the
// `updated_at` CAS in the commit validates.
let current_state = load_site_replication_state().await?;
if pending_endpoint_refresh(&current_state).is_some() {
return Err(s3_error!(InvalidRequest, "endpoint target refresh is pending"));
}
let local_peer = current_local_peer(&req, &current_state);
let mut sites: Vec<PeerSite> = read_site_replication_json(req, &cred.secret_key, true).await?;
let admin_access_key = cred.access_key.clone();
let admission_store = current_object_store_handle()
.ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))?;
let list_store = admission_store.clone();
let (state, edit_generation, local_peer, service_account_secret_key, mut initial_sync_errors, _add_guard) =
with_site_replication_bucket_mutation_admission_lock(admission_store, move || async move {
// The writer starts before the local bucket snapshot and stays
// held through every peer join and the topology commit. A
// delete followed by a same-name create therefore cannot hide
// behind an unchanged final name set. Peer bootstrap callbacks
// use their internal path and do not acquire this public-
// mutation admission lock.
let current_state = load_site_replication_state().await?;
if pending_endpoint_refresh(&current_state).is_some() {
return Err(s3_error!(InvalidRequest, "endpoint target refresh is pending"));
}
let local_peer = local_peer_at_endpoint(local_endpoint, &current_state);
// The web console's "Set Up Site Replication" omits the local deployment from the payload;
// inject it so the add preflight (which requires the local deployment) succeeds. No-op for `mc`.
ensure_local_site_present(&mut sites, &local_peer);
validate_add_sites(&sites, &local_peer)?;
let preflight_infos = add_preflight_infos(&sites, &current_state, &local_peer).await?;
validate_add_preflight_topology(&preflight_infos, &local_peer)?;
let expected_updated_at = current_state.updated_at;
require_add_peer_tls_capability(&sites, &local_peer).await?;
// Early exit on a state that moved under the preflight probes, BEFORE
// the IAM write and the join fan-out change anything remote. Advisory
// only — the binding check is the CAS inside the commit — but it fences
// the common race off the side-effect path and refreshes the merge
// base so the CAS window is only the join round trips.
let latest_state = load_site_replication_state().await?;
ensure_edit_precondition(&latest_state, expected_updated_at, None, "add preflight")?;
let current_state = latest_state;
let (service_account_access_key, service_account_secret_key) =
ensure_site_replicator_service_account(&admin_access_key, false).await?;
let expected_buckets: HashSet<String> =
preflight_infos.iter().flat_map(|info| info.buckets.keys().cloned()).collect();
let bootstrap_buckets: HashSet<String> = preflight_infos
.iter()
.filter(|info| !same_identity_endpoint(&info.endpoint, &local_peer.endpoint))
.flat_map(|info| info.buckets.keys().cloned())
.collect();
let add_in_progress_guard =
SiteReplicationAddInProgressGuard::start(lifecycle_guard, bootstrap_buckets.clone())?;
let mut state = merge_add_sites(
current_state,
local_peer.clone(),
sites.clone(),
service_account_access_key.clone(),
admin_access_key,
replicate_ilm_expiry,
);
state.sync_state_initialized = true;
let join_req = SRPeerJoinEnvelope {
request: SRPeerJoinReq {
svc_acct_access_key: service_account_access_key,
svc_acct_secret_key: service_account_secret_key.clone(),
svc_acct_parent: String::new(),
peers: state.peers.clone(),
updated_at: state.updated_at,
},
defer_sync_state_enable: true,
};
let peer_join_path = with_site_replication_bootstrap_token(
SITE_REPLICATION_PEER_JOIN_PATH,
&add_in_progress_guard.token.to_string(),
);
// The web console's "Set Up Site Replication" omits the local deployment from the payload;
// inject it so the add preflight (which requires the local deployment) succeeds. No-op for `mc`.
ensure_local_site_present(&mut sites, &local_peer);
validate_add_sites(&sites, &local_peer)?;
let preflight_infos = add_preflight_infos(&sites, &current_state, &local_peer).await?;
validate_add_preflight_topology(&preflight_infos, &local_peer)?;
let expected_updated_at = current_state.updated_at;
require_add_peer_tls_capability(&sites, &local_peer).await?;
// Early exit on a state that moved under the preflight probes, BEFORE
// the IAM write and the join fan-out change anything remote. Advisory
// only — the binding check is the CAS inside the commit — but it fences
// the common race off the side-effect path and refreshes the merge
// base so the CAS window is only the join round trips.
let latest_state = load_site_replication_state().await?;
ensure_edit_precondition(&latest_state, expected_updated_at, None, "add preflight")?;
let current_state = latest_state;
let (service_account_access_key, service_account_secret_key) =
ensure_site_replicator_service_account(&cred.access_key, false).await?;
let bootstrap_buckets = preflight_infos
.iter()
.filter(|info| !same_identity_endpoint(&info.endpoint, &local_peer.endpoint))
.flat_map(|info| info.buckets.keys().cloned())
.collect();
let add_in_progress_guard = SiteReplicationAddInProgressGuard::start(lifecycle_guard, bootstrap_buckets)?;
let mut state = merge_add_sites(
current_state,
local_peer.clone(),
sites.clone(),
service_account_access_key.clone(),
cred.access_key.clone(),
replicate_ilm_expiry,
);
state.sync_state_initialized = true;
let join_req = SRPeerJoinEnvelope {
request: SRPeerJoinReq {
svc_acct_access_key: service_account_access_key,
svc_acct_secret_key: service_account_secret_key.clone(),
svc_acct_parent: String::new(),
peers: state.peers.clone(),
updated_at: state.updated_at,
},
defer_sync_state_enable: true,
};
let peer_join_path =
with_site_replication_bootstrap_token(SITE_REPLICATION_PEER_JOIN_PATH, &add_in_progress_guard.token.to_string());
let mut joined_endpoints = HashSet::new();
let mut initial_sync_errors = SiteReplicationErrorSummary::default();
for (site, preflight) in sites.iter().zip(preflight_infos.iter()) {
if same_identity_endpoint(&site.endpoint, &local_peer.endpoint)
|| !joined_endpoints.insert(site_identity_key(&site.endpoint))
{
continue;
}
let mut joined_endpoints = HashSet::new();
let mut initial_sync_errors = SiteReplicationErrorSummary::default();
for (site, preflight) in sites.iter().zip(preflight_infos.iter()) {
if same_identity_endpoint(&site.endpoint, &local_peer.endpoint)
|| !joined_endpoints.insert(site_identity_key(&site.endpoint))
{
continue;
}
let mut peer_join_req = join_req.clone();
peer_join_req.request.svc_acct_parent = site.access_key.clone();
let connection = PeerConnection::try_from(site)?;
let body = PeerAdminRequest::put(&connection, &peer_join_path, &site.access_key)
.send(&site.secret_key, &peer_join_req)
.await?;
let mut fallback_peer = existing_peer_for_endpoint(&state, &site.endpoint)
.unwrap_or_else(|| normalize_peer_site(site.clone(), replicate_ilm_expiry));
fallback_peer.deployment_id = preflight.deployment_id.clone();
let join_response = parse_peer_join_response(&body, fallback_peer).map_err(|e| {
S3Error::with_message(
S3ErrorCode::InternalError,
format!("parse peer join response from {} failed: {e}", site.endpoint),
)
})?;
if !join_response.initial_sync_error_message.is_empty() {
initial_sync_errors.push(format!("{}: {}", site.endpoint, join_response.initial_sync_error_message));
}
// An explicit no-op join. The peer answered 200 but wrote nothing —
// its persisted state is already newer than the snapshot it was
// sent — so the add is only PARTIALLY configured and saying
// "configured successfully" would be a lie (rustfs/rustfs#5963).
// `None` (a MinIO peer, or one older than the field) is not a
// no-op signal and is deliberately not reported.
if join_response.applied == Some(false) {
initial_sync_errors.push(format!(
"{}: peer did not apply the join (its site replication state is newer than the snapshot it was sent); \
the site is not configured against this peer",
site.endpoint
));
}
state = reconcile_peer_with_actual_identity(state, join_response.peer);
let reconciled_peer = existing_peer_for_endpoint(&state, &site.endpoint).ok_or_else(|| {
S3Error::with_message(
S3ErrorCode::InternalError,
format!("peer join response from {} did not identify the requested site", site.endpoint),
)
})?;
validate_proposed_peer(&reconciled_peer).map_err(|err| {
S3Error::with_message(
S3ErrorCode::InvalidRequest,
format!("invalid peer join response from {}: {err}", site.endpoint),
)
})?;
}
mark_unknown_peer_sync_enabled(&mut state.peers);
// Commit. The state transaction's CAS still fences topology
// writers that do not use bucket admission. By this point
// remote sites may already have accepted their joins, so a
// mismatch asks the operator to re-run add and reconverge.
let next_state = state;
let present = list_store
.list_bucket(&BucketOptions::default())
.await
.map_err(ApiError::from)?
.into_iter()
.map(|bucket| bucket.name)
.collect::<HashSet<_>>();
ensure_add_bucket_set_matches_preflight(&expected_buckets, &present)?;
let (state, edit_generation) = update_site_replication_state(move |state| {
if state.updated_at != expected_updated_at || pending_endpoint_refresh(state).is_some() {
return Err(s3_error!(
InvalidRequest,
"site replication state changed during peer join; the peers may already be joined — re-run replicate add"
));
}
adopt_add_commit_state(state, next_state);
let edit_generation = next_peer_edit_generation(state);
Ok((state.clone(), edit_generation))
})
let mut peer_join_req = join_req.clone();
peer_join_req.request.svc_acct_parent = site.access_key.clone();
let connection = PeerConnection::try_from(site)?;
let body = PeerAdminRequest::put(&connection, &peer_join_path, &site.access_key)
.send(&site.secret_key, &peer_join_req)
.await?;
Ok((
state,
edit_generation,
local_peer,
service_account_secret_key,
initial_sync_errors,
add_in_progress_guard,
))
})
.await?;
let mut fallback_peer = existing_peer_for_endpoint(&state, &site.endpoint)
.unwrap_or_else(|| normalize_peer_site(site.clone(), replicate_ilm_expiry));
fallback_peer.deployment_id = preflight.deployment_id.clone();
let join_response = parse_peer_join_response(&body, fallback_peer).map_err(|e| {
S3Error::with_message(
S3ErrorCode::InternalError,
format!("parse peer join response from {} failed: {e}", site.endpoint),
)
})?;
if !join_response.initial_sync_error_message.is_empty() {
initial_sync_errors.push(format!("{}: {}", site.endpoint, join_response.initial_sync_error_message));
}
// An explicit no-op join. The peer answered 200 but wrote nothing —
// its persisted state is already newer than the snapshot it was
// sent — so the add is only PARTIALLY configured and saying
// "configured successfully" would be a lie (rustfs/rustfs#5963).
// `None` (a MinIO peer, or one older than the field) is not a
// no-op signal and is deliberately not reported.
if join_response.applied == Some(false) {
initial_sync_errors.push(format!(
"{}: peer did not apply the join (its site replication state is newer than the snapshot it was sent); \
the site is not configured against this peer",
site.endpoint
));
}
state = reconcile_peer_with_actual_identity(state, join_response.peer);
let reconciled_peer = existing_peer_for_endpoint(&state, &site.endpoint).ok_or_else(|| {
S3Error::with_message(
S3ErrorCode::InternalError,
format!("peer join response from {} did not identify the requested site", site.endpoint),
)
})?;
validate_proposed_peer(&reconciled_peer).map_err(|err| {
S3Error::with_message(
S3ErrorCode::InvalidRequest,
format!("invalid peer join response from {}: {err}", site.endpoint),
)
})?;
}
mark_unknown_peer_sync_enabled(&mut state.peers);
// Commit. The CAS runs inside the transaction, against the state the
// transaction itself loaded — the peer round trips above took however
// long they took, and only this check can tell whether the topology
// this add was planned against is still the current one. The error
// says so: by this point the remote sites already accepted their
// joins, and re-running the add is what reconverges the local side.
let next_state = state;
let (state, edit_generation) = update_site_replication_state(move |state| {
if state.updated_at != expected_updated_at || pending_endpoint_refresh(state).is_some() {
return Err(s3_error!(
InvalidRequest,
"site replication state changed during peer join; the peers may already be joined — re-run replicate add"
));
}
adopt_add_commit_state(state, next_state);
let edit_generation = next_peer_edit_generation(state);
Ok((state.clone(), edit_generation))
})
.await?;
// The finalize fan-out delivers peer-edit payloads, so it carries the
// generation allocated in the commit above: the receiving site orders
@@ -7309,14 +7185,8 @@ impl Operation for SRPeerEditHandler {
// The fence is self-reported — the shared service account means
// the sender cannot be identified — so it is honoured only after
// the admissibility check, against the same state it will gate.
let commit_fence = match commit_fence {
Some(fence) if peer_edit_fence_is_admissible(state, &local_peer.deployment_id, &fence) => Some(fence),
// A fenced edit can only come from a current remote peer. If
// that origin left while the retry was in flight, applying
// its body here would resurrect the removed topology.
Some(_) => return Ok(StateCommit::Unchanged(PeerEditOutcome::Acked)),
None => None,
};
let commit_fence =
commit_fence.filter(|fence| peer_edit_fence_is_admissible(state, &local_peer.deployment_id, fence));
// Ordering fence: the sending site allocates the generation under
// its state-object lock, so a delivery that lost the race carries
// a generation this site has already passed. Applying it would
@@ -9016,41 +8886,6 @@ mod tests {
);
}
#[test]
fn add_admission_starts_before_preflight_and_rejects_bucket_set_changes() {
let expected = HashSet::from(["remote-owned".to_string(), "shared".to_string()]);
let present = HashSet::from(["shared".to_string()]);
let err = ensure_add_bucket_set_matches_preflight(&expected, &present)
.expect_err("a missing bootstrap bucket must reject the topology commit");
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest);
let present = HashSet::from([
"remote-owned".to_string(),
"shared".to_string(),
"created-during-add".to_string(),
]);
let err = ensure_add_bucket_set_matches_preflight(&expected, &present)
.expect_err("a bucket created during add must reject the topology commit");
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest);
let src = include_str!("site_replication.rs");
let add = src
.split("impl Operation for SiteReplicationAddHandler")
.nth(1)
.and_then(|rest| rest.split("pub struct SiteReplicationRemoveHandler").next())
.expect("add handler block");
let admission = add
.find("with_site_replication_bucket_mutation_admission_lock")
.expect("distributed mutation admission");
let preflight = add.find("add_preflight_infos").expect("bucket preflight");
let validation = add
.find("ensure_add_bucket_set_matches_preflight")
.expect("bucket-set validation");
let commit = add.find("adopt_add_commit_state").expect("topology commit");
assert!(admission < preflight && preflight < validation && validation < commit);
}
#[test]
fn test_tls_capability_gates_run_before_add_or_edit_state_side_effects() {
let src = include_str!("site_replication.rs");
@@ -9347,19 +9182,13 @@ mod tests {
);
// Fence hardening: origin and generation are self-reported by a
// caller the shared service account cannot identify, so the handler
// must admit the fence against the same state it gates. An origin
// removed while a retry was in flight is acknowledged without
// applying the stale body; otherwise it could recreate topology.
// must pass the fence through the admissibility check — against the
// same state the fence gates, i.e. inside the transaction — before
// reading or raising any high-water mark.
assert!(
handler_block.contains(
"Some(fence) if peer_edit_fence_is_admissible(state, &local_peer.deployment_id, &fence) => Some(fence)"
),
handler_block.contains(".filter(|fence| peer_edit_fence_is_admissible(state, &local_peer.deployment_id, fence))"),
"SRPeerEditHandler must admit a fence only through peer_edit_fence_is_admissible inside the state transaction"
);
assert!(
handler_block.contains("Some(_) => return Ok(StateCommit::Unchanged(PeerEditOutcome::Acked))"),
"SRPeerEditHandler must not apply a fenced edit after its origin leaves the current topology"
);
// P1-15 PR2: both halves of the fence and the edit they fence share
// ONE transaction. Checking the fence against a state read outside the
// lock would let the check pass on one snapshot and the write land on
@@ -10523,9 +10352,8 @@ mod tests {
/// A fence is self-reported: every site authenticates peer traffic with
/// the same site-replicator credential, so a compromised peer can stamp
/// ANY origin with ANY generation. An origin the receiver does not
/// replicate with — or the receiver itself — is inadmissible and plants
/// no mark; the handler acknowledges such a request without applying its
/// body. A mark a compromised peer plants for a CURRENT origin cannot
/// replicate with — or the receiver itself — is ignored and plants no
/// mark; a mark a compromised peer plants for a CURRENT origin cannot
/// silence that origin, because the staleness window refuses to fence on
/// a mark implausibly far above the genuine deliveries.
#[test]
@@ -12963,7 +12791,6 @@ mod tests {
last_error: "site replication is not enabled".to_string(),
updated_at: Some(OffsetDateTime::now_utc()),
edit_generation: None,
peer_unreachable: false,
deletions_recorded: false,
}],
..Default::default()
@@ -13162,7 +12989,6 @@ mod tests {
last_error: "peer offline".to_string(),
updated_at: Some(OffsetDateTime::now_utc()),
edit_generation: None,
peer_unreachable: false,
deletions_recorded: false,
}],
..Default::default()
+37 -63
View File
@@ -75,8 +75,7 @@ use crate::auth::get_condition_values_with_client_info;
use crate::error::ApiError;
use crate::shared_types::RemoteAddr;
use crate::site_replication::{
cancel_site_replication_delete_bucket, commit_site_replication_delete_bucket, prepare_site_replication_delete_bucket,
site_replication_bucket_meta_hook, site_replication_make_bucket_hook, with_site_replication_bucket_mutation_lock,
site_replication_bucket_meta_hook, site_replication_delete_bucket_hook, site_replication_make_bucket_hook,
};
use crate::storage::storage_api::lock_bucket_targets_metadata;
use http::StatusCode;
@@ -1332,34 +1331,23 @@ impl DefaultBucketUsecase {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
};
// Keep the local namespace mutation and its peer hook ordered across
// every node in this site. Otherwise a delete waiting for repair
// coordination can arrive after this create on remote sites.
let operation_bucket = bucket.clone();
let operation_store = store.clone();
let make_result = with_site_replication_bucket_mutation_lock(store, &bucket, move || async move {
let make_result = operation_store
.make_bucket(
&operation_bucket,
&MakeBucketOptions {
force_create: false,
lock_enabled,
..Default::default()
},
)
.await;
if make_result.is_ok() {
crate::storage::invalidate_bucket_validation_cache(&operation_bucket);
if let Err(err) = site_replication_make_bucket_hook(&operation_bucket, lock_enabled).await {
warn!(bucket = %operation_bucket, error = ?err, "site replication make bucket hook failed");
}
}
make_result
})
.await?;
let make_result = store
.make_bucket(
&bucket,
&MakeBucketOptions {
force_create: false,
lock_enabled,
..Default::default()
},
)
.await;
match make_result {
Ok(()) => {}
Ok(()) => {
// Invalidate the bucket validation cache so subsequent GETs
// see the newly created bucket immediately.
crate::storage::invalidate_bucket_validation_cache(&bucket);
}
Err(StorageError::BucketExists(_)) => {
// Per S3 spec: bucket namespace is global. Owner recreating returns 200 OK;
// non-owner gets 409 BucketAlreadyExists.
@@ -1370,6 +1358,10 @@ impl DefaultBucketUsecase {
Err(e) => return Err(ApiError::from(e).into()),
}
if let Err(err) = site_replication_make_bucket_hook(&bucket, lock_enabled).await {
warn!(bucket = %bucket, error = ?err, "site replication make bucket hook failed");
}
let output = CreateBucketOutput::default();
counter!("rustfs_create_bucket_total").increment(1);
let result = Ok(S3Response::new(output));
@@ -1405,41 +1397,16 @@ impl DefaultBucketUsecase {
authorize_request(&mut req, Action::S3Action(S3Action::ForceDeleteBucketAction)).await?;
}
// Keep the local namespace mutation and its peer hook ordered across
// every node in this site so an older delete cannot overtake a new
// same-name make while it waits for repair coordination.
let operation_bucket = input.bucket.clone();
let operation_store = store.clone();
with_site_replication_bucket_mutation_lock(store, &input.bucket, move || async move {
let intent = prepare_site_replication_delete_bucket(&operation_bucket, force).await?;
let delete_result = operation_store
.delete_bucket(
&operation_bucket,
&DeleteBucketOptions {
force,
..Default::default()
},
)
.await;
match delete_result {
Ok(()) => {
crate::storage::invalidate_bucket_validation_cache(&operation_bucket);
if let Some(intent) = intent
&& let Err(err) = commit_site_replication_delete_bucket(&intent).await
{
warn!(bucket = %operation_bucket, error = ?err, "site replication delete bucket hook failed");
}
Ok::<(), S3Error>(())
}
Err(err) => {
if let Some(intent) = intent {
cancel_site_replication_delete_bucket(intent).await;
}
Err(S3Error::from(ApiError::from(err)))
}
}
})
.await??;
store
.delete_bucket(
&input.bucket,
&DeleteBucketOptions {
force,
..Default::default()
},
)
.await
.map_err(ApiError::from)?;
// Drop every cached object body for the now-deleted bucket so dead
// bytes do not sit resident until TTL. Covers both the normal and the
@@ -1448,9 +1415,16 @@ impl DefaultBucketUsecase {
let cache_adapter = current_object_data_cache_for_context(self.context.as_deref());
let _ = invalidate_object_data_cache_bucket_after_delete(&cache_adapter, &input.bucket).await;
// Invalidate bucket validation cache
crate::storage::invalidate_bucket_validation_cache(&input.bucket);
// Re-evaluate lifecycle and replication after bucket removal.
rustfs_scanner::record_scanner_maintenance_change(&input.bucket);
if let Err(err) = site_replication_delete_bucket_hook(&input.bucket, force).await {
warn!(bucket = %input.bucket, error = ?err, "site replication delete bucket hook failed");
}
// Notify peers to drop their cached metadata for the now-deleted bucket.
let request_context = req.extensions.get::<request_context::RequestContext>().cloned();
notify_bucket_metadata_delete(input.bucket.clone(), request_context);
+1 -45
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 = load_bucket_default_sse_config(&bucket).await?;
let bucket_sse_config = metadata_sys::get_sse_config(&bucket).await.ok();
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,48 +1388,4 @@ 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 = load_bucket_default_sse_config(&bucket).await?;
let bucket_sse_config = metadata_sys::get_sse_config(&bucket).await.ok();
let (mut effective_sse, mut effective_kms_key_id) = resolve_bucket_default_sse(
bucket_sse_config.as_ref().map(|(config, _timestamp)| config),
original_sse,
+1 -119
View File
@@ -1485,9 +1485,8 @@ impl DefaultObjectUsecase {
};
let sse_config_stage_start = put_stage_metrics_enabled.then(Instant::now);
let bucket_sse_config = load_bucket_default_sse_config(&bucket).await;
let bucket_sse_config = metadata_sys::get_sse_config(&bucket).await.ok();
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",
@@ -3913,121 +3912,4 @@ 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,129 +269,6 @@ 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,36 +96,3 @@ 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");
}
+35 -427
View File
@@ -22,57 +22,6 @@ pub(crate) const SITE_REPLICATION_BUCKET_OP_CONFIGURE_REPLICATION: &str = "confi
pub(crate) static SITE_REPLICATION_BUCKET_OP_LOCK: LazyLock<RwLock<()>> = LazyLock::new(|| RwLock::new(()));
const SITE_REPLICATION_BUCKET_MUTATION_LOCK_PREFIX: &str = "config/site-replication/bucket-mutation";
pub(crate) const SITE_REPLICATION_BUCKET_MUTATION_ADMISSION_LOCK_PATH: &str =
"config/site-replication/bucket-mutation-admission.lock";
pub(crate) fn site_replication_bucket_mutation_lock_path(bucket: &str) -> String {
format!("{SITE_REPLICATION_BUCKET_MUTATION_LOCK_PREFIX}/{bucket}.lock")
}
pub(crate) async fn with_site_replication_bucket_mutation_lock<F, Fut, T>(
store: Arc<ECStore>,
bucket: &str,
operation: F,
) -> S3Result<T>
where
F: FnOnce() -> Fut + Send + 'static,
Fut: std::future::Future<Output = T> + Send + 'static,
T: Send + 'static,
{
let mutation_store = store.clone();
let mutation_path = site_replication_bucket_mutation_lock_path(bucket);
with_config_object_read_lock(
store,
SITE_REPLICATION_BUCKET_MUTATION_ADMISSION_LOCK_PATH.to_string(),
move || async move {
with_config_object_write_lock(mutation_store, mutation_path, operation)
.await
.map_err(|err| S3Error::from(ApiError::from(err)))
},
)
.await
.map_err(|err| S3Error::from(ApiError::from(err)))?
}
/// Exclude every local bucket namespace mutation from an add's local preflight
/// snapshot until its topology commit. Peer bootstrap callbacks do not enter
/// this public-mutation admission path, so they can finish while the writer is
/// held; post-commit fan-out and backfill must run after it is released.
pub(crate) async fn with_site_replication_bucket_mutation_admission_lock<F, Fut, T>(
store: Arc<ECStore>,
operation: F,
) -> S3Result<T>
where
F: FnOnce() -> Fut + Send + 'static,
Fut: std::future::Future<Output = S3Result<T>> + Send + 'static,
T: Send + 'static,
{
with_config_object_write_lock(store, SITE_REPLICATION_BUCKET_MUTATION_ADMISSION_LOCK_PATH.to_string(), operation)
.await
.map_err(|err| S3Error::from(ApiError::from(err)))?
}
#[derive(Debug, Default)]
pub(crate) struct SiteReplicationBootstrapPlan {
pub(crate) iam_items: Vec<SRIAMItem>,
@@ -380,91 +329,6 @@ pub(crate) fn site_replication_bootstrap_plan(info: &SRInfo) -> S3Result<SiteRep
Ok(plan)
}
/// Build only the two bucket operations needed by the lightweight retry
/// drain. The full bootstrap plan scans every bucket and IAM record; doing
/// that on a 30-second recovery cadence would make lifecycle admission scale
/// with the whole site instead of the one queued bucket.
pub(crate) fn site_replication_bucket_retry_plan_for(
bucket: &SRBucketInfo,
replicate_ilm_expiry: bool,
) -> S3Result<SiteReplicationBootstrapPlan> {
let mut plan = SiteReplicationBootstrapPlan {
bucket_make_ops: vec![bootstrap_bucket_make_op_path(bucket)],
bucket_configure_ops: vec![bootstrap_bucket_op_path(
&bucket.bucket,
SITE_REPLICATION_BUCKET_OP_CONFIGURE_REPLICATION,
)],
..Default::default()
};
append_bootstrap_bucket_items(&mut plan, bucket, replicate_ilm_expiry)?;
Ok(plan)
}
pub(crate) fn site_replication_bucket_retry_plan_from_info(
bucket: &SRBucketInfo,
replicate_ilm_expiry: bool,
) -> S3Result<SiteReplicationBootstrapPlan> {
let mut plan = site_replication_bucket_retry_plan_for(bucket, replicate_ilm_expiry)?;
// Omit only metadata the make/configure operations can reproduce exactly.
// Non-default versioning fields and operator-authored replication rules
// remain in the plan; their extra request cost intentionally defers the
// event to the complete drain when the lightweight budget is too small.
plan.bucket_items.retain(|item| !retry_bucket_metadata_is_redundant(item));
Ok(plan)
}
fn retry_bucket_metadata_is_redundant(item: &SRBucketMeta) -> bool {
match item.r#type.as_str() {
"version-config" => item.versioning.as_deref().is_some_and(|raw| {
deserialize::<VersioningConfiguration>(&decode_bucket_meta_wire_value(raw)).is_ok_and(|config| {
config
== VersioningConfiguration {
status: Some(BucketVersioningStatus::from_static(BucketVersioningStatus::ENABLED)),
..Default::default()
}
})
}),
"replication-config" => item.replication_config.as_deref().is_some_and(|raw| {
deserialize::<ReplicationConfiguration>(&decode_bucket_meta_wire_value(raw))
.is_ok_and(|config| config.role.trim().is_empty() && config.rules.iter().all(is_derived_site_replication_rule))
}),
// `Some("")` is the in-memory sentinel used when the bucket is lock
// enabled but has no object-lock configuration body. The make query
// carries lockEnabled=true; sending an empty metadata body is neither
// useful nor parseable.
"object-lock-config" => item.object_lock_config.as_deref() == Some(""),
_ => false,
}
}
pub(crate) async fn site_replication_bucket_retry_plan(
bucket: &str,
replicate_ilm_expiry: bool,
) -> S3Result<SiteReplicationBootstrapPlan> {
let Some(store) = current_object_store_handle() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
};
let bucket_info = match store.get_bucket_info(bucket, &BucketOptions::default()).await {
Ok(bucket_info) => bucket_info,
Err(err) if is_err_bucket_not_found(&err) => return Ok(SiteReplicationBootstrapPlan::default()),
Err(err) => return Err(ApiError::from(err).into()),
};
let lock_enabled = bucket_info.object_locking;
let metadata = metadata_sys::get(bucket).await.map_err(ApiError::from)?;
let mut bucket_info = SRBucketInfo {
bucket: bucket.to_string(),
created_at: bucket_info.created,
location: current_region().map(|region| region.to_string()).unwrap_or_default(),
api_version: Some(SITE_REPL_API_VERSION.to_string()),
..Default::default()
};
populate_sr_bucket_info_from_metadata(&mut bucket_info, &metadata).await;
if lock_enabled && bucket_info.object_lock_config.is_none() {
bucket_info.object_lock_config = Some(String::new());
}
site_replication_bucket_retry_plan_from_info(&bucket_info, replicate_ilm_expiry)
}
pub async fn site_replication_make_bucket_hook(bucket: &str, lock_enabled: bool) -> S3Result<()> {
let _bucket_op_guard = SITE_REPLICATION_BUCKET_OP_LOCK.read().await;
let runtime = {
@@ -529,273 +393,20 @@ pub(crate) async fn broadcast_site_replication_make_bucket(
broadcast_site_replication_json_using_runtime(runtime, &configure_path, &serde_json::json!({})).await
}
const SITE_REPLICATION_DELETE_INTENT_PENDING: &str =
"bucket deletion reserved; local completion and peer delivery are not yet known";
#[derive(Clone)]
struct SiteReplicationDeleteBucketReservation {
peer: PeerInfo,
previous: Option<SiteReplicationRetryEvent>,
observed: SiteReplicationRetryEvent,
}
pub(crate) struct SiteReplicationDeleteBucketIntent {
path: String,
reservations: Vec<SiteReplicationDeleteBucketReservation>,
displaced: Vec<SiteReplicationRetryEvent>,
}
fn site_replication_delete_bucket_path(bucket: &str, force_delete: bool) -> String {
pub async fn site_replication_delete_bucket_hook(bucket: &str, force_delete: bool) -> S3Result<()> {
let operation = if force_delete {
"force-delete-bucket"
} else {
"delete-bucket"
};
format!(
let path = format!(
"/rustfs/admin/v3/site-replication/peer/bucket-ops?{}",
form_urlencoded::Serializer::new(String::new())
.append_pair("bucket", bucket)
.append_pair("operation", operation)
.finish()
)
}
/// Reserve every destructive peer delivery before the local namespace is
/// changed. The state transaction either persists the complete set or writes
/// nothing, so a full/unreadable queue fails the S3 delete closed.
pub(crate) async fn prepare_site_replication_delete_bucket(
bucket: &str,
force_delete: bool,
) -> S3Result<Option<SiteReplicationDeleteBucketIntent>> {
let path = site_replication_delete_bucket_path(bucket, force_delete);
let reservation_path = path.clone();
update_site_replication_state_when_changed(move |state| {
if !state.enabled() {
return Ok(StateCommit::Unchanged(None));
}
let local_peer = current_local_runtime_peer(state);
let peers = state
.peers
.values()
.filter(|peer| {
peer.deployment_id != local_peer.deployment_id && !same_identity_endpoint(&peer.endpoint, &local_peer.endpoint)
})
.cloned()
.collect::<Vec<_>>();
if peers.is_empty() {
return Ok(StateCommit::Unchanged(None));
}
let mut reservations = Vec::with_capacity(peers.len());
let mut displaced = Vec::new();
for peer in peers {
let previous = state
.retry_queue
.iter()
.find(|event| retry_event_matches(event, &peer, &reservation_path))
.cloned();
displaced.extend(upsert_site_replication_retry_event(
&mut state.retry_queue,
&peer,
&reservation_path,
SITE_REPLICATION_DELETE_INTENT_PENDING,
None,
)?);
let observed = state
.retry_queue
.iter()
.find(|event| retry_event_matches(event, &peer, &reservation_path))
.cloned()
.ok_or_else(|| {
S3Error::with_message(
S3ErrorCode::InternalError,
"site replication delete reservation disappeared before commit".to_string(),
)
})?;
reservations.push(SiteReplicationDeleteBucketReservation {
peer,
previous,
observed,
});
}
Ok(StateCommit::Changed(Some(SiteReplicationDeleteBucketIntent {
path: reservation_path,
reservations,
displaced,
})))
})
.await
}
/// Roll back a reservation when the local storage delete definitively failed.
/// A concurrently revised reservation is preserved; it belongs to a newer
/// observation and this operation has no authority to settle it.
pub(crate) async fn cancel_site_replication_delete_bucket(intent: SiteReplicationDeleteBucketIntent) {
let path = intent.path.clone();
let result = update_site_replication_state_when_changed(move |state| {
let mut changed = false;
for reservation in intent.reservations {
let Some(index) = state.retry_queue.iter().position(|event| {
retry_event_matches(event, &reservation.peer, &reservation.observed.path)
&& event.id == reservation.observed.id
&& event.updated_at == reservation.observed.updated_at
}) else {
continue;
};
if let Some(previous) = reservation.previous {
state.retry_queue[index] = previous;
} else {
state.retry_queue.remove(index);
}
changed = true;
}
let mut restored_all = true;
for displaced in intent.displaced {
let duplicate = state.retry_queue.iter().any(|event| {
event.id == displaced.id
|| (event.peer_deployment_id == displaced.peer_deployment_id && event.path == displaced.path)
});
if duplicate {
continue;
}
if state.retry_queue.len() >= SITE_REPLICATION_RETRY_QUEUE_LIMIT {
restored_all = false;
continue;
}
state.retry_queue.push(displaced);
changed = true;
}
Ok(if changed {
StateCommit::Changed(restored_all)
} else {
StateCommit::Unchanged(restored_all)
})
})
.await;
match result {
Ok(true) => {}
Ok(false) => warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
path,
result = "delete_intent_cancel_incomplete",
"admin site replication state"
),
Err(err) => warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
path,
result = "delete_intent_cancel_failed",
error = ?err,
"admin site replication state"
),
}
}
async fn broadcast_site_replication_delete_bucket(intent: &SiteReplicationDeleteBucketIntent) -> S3Result<()> {
let sends = intent.reservations.iter().cloned().map(|reservation| {
let request_path = intent.path.clone();
async move {
let fallback_peer = reservation.peer.clone();
let observed = reservation.observed.clone();
let delivery_path = request_path.clone();
let delivery = with_site_replication_state_read_lock(move |state| async move {
let Some(current_peer) = state.peers.get(&fallback_peer.deployment_id).cloned() else {
return Ok(None);
};
let service_account_secret_key =
match site_replicator_service_account_secret(&state.service_account_access_key).await {
Ok(secret) => secret,
Err(err) => {
let Some(secret) = legacy_site_replicator_state_secret(&state) else {
return Err(err);
};
warn!(
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
component = LOG_COMPONENT_ADMIN,
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
result = "legacy_state_service_account_secret_fallback",
error = ?err,
"admin site replication state"
);
secret
}
};
let result = async {
let transport = PeerTransport::for_runtime_peer(&current_peer).await?;
PeerAdminRequest::put(&transport.connection, &delivery_path, &state.service_account_access_key)
.with_client(&transport.client)
.send(&service_account_secret_key, &serde_json::json!({}))
.await
}
.await;
Ok(Some((current_peer, result)))
})
.await;
match delivery {
Ok(Some((current_peer, Ok(_)))) => {
dequeue_observed_site_replication_retry_event(&current_peer, &observed).await;
None
}
Ok(Some((current_peer, Err(err)))) => {
// Keep the failed deletion operator-visible, but never
// replay it automatically: without a bucket-incarnation
// fence, a delayed delete could erase a recreated bucket.
enqueue_site_replication_retry_event(&current_peer, &request_path, &err).await;
Some(err)
}
Ok(None) => {
dequeue_observed_site_replication_retry_event(&reservation.peer, &observed).await;
None
}
Err(err) => {
enqueue_site_replication_retry_event(&reservation.peer, &request_path, &err).await;
Some(err)
}
}
}
});
futures::future::join_all(sends)
.await
.into_iter()
.flatten()
.next()
.map_or(Ok(()), Err)
}
pub(crate) async fn commit_site_replication_delete_bucket(intent: &SiteReplicationDeleteBucketIntent) -> S3Result<()> {
let _bucket_op_guard = SITE_REPLICATION_BUCKET_OP_LOCK.read().await;
let store =
current_object_store_handle().ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()))?;
let retry_peers = intent
.reservations
.iter()
.map(|reservation| reservation.peer.clone())
.collect::<Vec<_>>();
let retry_path = intent.path.clone();
let delivery_intent = SiteReplicationDeleteBucketIntent {
path: intent.path.clone(),
reservations: intent.reservations.clone(),
displaced: Vec::new(),
};
match with_config_object_write_lock(store, SITE_REPLICATION_REPAIR_EXECUTION_LOCK_PATH.to_string(), move || async move {
broadcast_site_replication_delete_bucket(&delivery_intent).await
})
.await
{
Ok(result) => result,
Err(err) => {
let err: S3Error = ApiError::from(err).into();
for peer in &retry_peers {
enqueue_site_replication_retry_event(peer, &retry_path, &err).await;
}
Err(err)
}
}
);
broadcast_site_replication_json(&path, &serde_json::json!({})).await
}
pub async fn site_replication_bucket_meta_hook(mut item: SRBucketMeta) -> S3Result<()> {
@@ -904,39 +515,6 @@ pub(crate) fn maybe_time(value: OffsetDateTime) -> Option<OffsetDateTime> {
(value != OffsetDateTime::UNIX_EPOCH).then_some(value)
}
async fn populate_sr_bucket_info_from_metadata(entry: &mut SRBucketInfo, metadata: &BucketMetadata) {
entry.policy = raw_config_to_string(&metadata.policy_config_json).and_then(|raw| serde_json::from_str(&raw).ok());
entry.versioning = raw_config_to_base64(&metadata.versioning_config_xml);
entry.tags = raw_config_to_base64(&metadata.tagging_config_xml);
entry.object_lock_config = raw_config_to_base64(&metadata.object_lock_config_xml);
entry.sse_config = raw_config_to_base64(&metadata.encryption_config_xml);
entry.replication_config = raw_config_to_base64(&metadata.replication_config_xml);
entry.quota_config = raw_config_to_base64(&metadata.quota_config_json);
// Expiry subset only: this entry feeds both the bootstrap/repair plan
// (peers must not receive transition rules) and cross-site consistency
// views (transition rules are site-local and would read as false
// mismatches). A deleted expiry state is a `None` value with the
// deletion's axis so repair can converge peers that missed the live
// delete.
let expiry_statement = lifecycle_expiry_statement(metadata);
entry.expiry_lc_config = expiry_statement.as_ref().and_then(|(subset, _)| subset.clone());
entry.cors_config = raw_config_to_base64(&metadata.cors_config_xml);
entry.policy_updated_at = maybe_time(metadata.policy_config_updated_at);
entry.tag_config_updated_at = maybe_time(metadata.tagging_config_updated_at);
entry.object_lock_config_updated_at = maybe_time(metadata.object_lock_config_updated_at);
entry.sse_config_updated_at = maybe_time(metadata.encryption_config_updated_at);
entry.versioning_config_updated_at = maybe_time(metadata.versioning_config_updated_at);
entry.replication_config_updated_at = maybe_time(metadata.replication_config_updated_at);
entry.quota_config_updated_at = maybe_time(metadata.quota_config_updated_at);
// The expiry axis, not the whole-config write time: local transition-only
// edits inflate the latter, and a repair item stamped with it could
// out-rank a newer real expiry edit on a third site.
entry.expiry_lc_config_updated_at = expiry_statement.map(|(_, axis)| axis);
entry.cors_config_updated_at = maybe_time(metadata.cors_config_updated_at);
entry.replication_targets_online =
Some(site_replication_targets_online(&entry.bucket, &metadata.replication_config_xml).await);
}
pub(crate) async fn build_sr_info(state: &SiteReplicationState, local_peer: &PeerInfo) -> S3Result<SRInfo> {
let Some(store) = current_object_store_handle() else {
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
@@ -968,7 +546,37 @@ pub(crate) async fn build_sr_info(state: &SiteReplicationState, local_peer: &Pee
};
if let Some(metadata) = metadata {
populate_sr_bucket_info_from_metadata(&mut entry, &metadata).await;
entry.policy = raw_config_to_string(&metadata.policy_config_json).and_then(|raw| serde_json::from_str(&raw).ok());
entry.versioning = raw_config_to_base64(&metadata.versioning_config_xml);
entry.tags = raw_config_to_base64(&metadata.tagging_config_xml);
entry.object_lock_config = raw_config_to_base64(&metadata.object_lock_config_xml);
entry.sse_config = raw_config_to_base64(&metadata.encryption_config_xml);
entry.replication_config = raw_config_to_base64(&metadata.replication_config_xml);
entry.quota_config = raw_config_to_base64(&metadata.quota_config_json);
// Expiry subset only: this entry feeds both the bootstrap/repair
// plan (peers must not receive transition rules) and cross-site
// consistency views (transition rules are site-local and would
// read as false mismatches). A deleted expiry state is a `None`
// value with the deletion's axis so repair can converge peers
// that missed the live delete.
let expiry_statement = lifecycle_expiry_statement(&metadata);
entry.expiry_lc_config = expiry_statement.as_ref().and_then(|(subset, _)| subset.clone());
entry.cors_config = raw_config_to_base64(&metadata.cors_config_xml);
entry.policy_updated_at = maybe_time(metadata.policy_config_updated_at);
entry.tag_config_updated_at = maybe_time(metadata.tagging_config_updated_at);
entry.object_lock_config_updated_at = maybe_time(metadata.object_lock_config_updated_at);
entry.sse_config_updated_at = maybe_time(metadata.encryption_config_updated_at);
entry.versioning_config_updated_at = maybe_time(metadata.versioning_config_updated_at);
entry.replication_config_updated_at = maybe_time(metadata.replication_config_updated_at);
entry.quota_config_updated_at = maybe_time(metadata.quota_config_updated_at);
// The expiry axis, not the whole-config write time: local
// transition-only edits inflate the latter, and a repair item
// stamped with it could out-rank a newer real expiry edit on a
// third site.
entry.expiry_lc_config_updated_at = expiry_statement.map(|(_, axis)| axis);
entry.cors_config_updated_at = maybe_time(metadata.cors_config_updated_at);
entry.replication_targets_online =
Some(site_replication_targets_online(&bucket.name, &metadata.replication_config_xml).await);
}
info.buckets.insert(bucket.name, entry);
+6 -7
View File
@@ -47,7 +47,6 @@ use self::identity::{
canonical_endpoint, deployment_id_for_endpoint, mark_unknown_peer_sync_enabled, normalize_peer_map_by_identity_with,
same_identity_endpoint,
};
pub(crate) use self::state_lock::with_site_replication_state_read_lock;
use self::state_lock::{SITE_REPLICATION_STATE_PATH, with_site_replication_state_lock};
use crate::auth::constant_time_eq;
use crate::config::get_config_snapshot;
@@ -65,12 +64,12 @@ use crate::storage_api::site_replication::s3::{
#[cfg(test)]
use crate::storage_api::site_replication::save_config as save_admin_config;
use crate::storage_api::site_replication::{
ARN, BUCKET_REPLICATION_CONFIG, BUCKET_TARGETS_FILE, BUCKET_VERSIONING_CONFIG, BucketMetadata, BucketOperations,
BucketOptions, BucketTarget, BucketTargetSys, BucketTargetType, BucketTargets, Credentials, ECStore, OperatorRuleContract,
StorageError, VersioningApi as _, assign_site_replication_rule_priorities, delete_config_no_lock, deserialize,
is_err_bucket_not_found, is_site_replication_role, lock_bucket_targets_metadata, metadata_sys,
read_config as read_admin_config, read_config_no_lock, replication_target_arn_deployment_id, save_config_no_lock, serialize,
site_replication_rule_deployment_id, with_config_object_read_lock, with_config_object_write_lock,
ARN, BUCKET_REPLICATION_CONFIG, BUCKET_TARGETS_FILE, BUCKET_VERSIONING_CONFIG, BucketOperations, BucketOptions, BucketTarget,
BucketTargetSys, BucketTargetType, BucketTargets, Credentials, ECStore, OperatorRuleContract, StorageError,
VersioningApi as _, assign_site_replication_rule_priorities, delete_config_no_lock, deserialize, is_site_replication_role,
lock_bucket_targets_metadata, metadata_sys, read_config as read_admin_config, read_config_no_lock,
replication_target_arn_deployment_id, save_config_no_lock, serialize, site_replication_rule_deployment_id,
with_config_object_read_lock, with_config_object_write_lock,
};
use base64_simd::STANDARD as BASE64_STANDARD;
use base64_simd::URL_SAFE_NO_PAD;
+1 -3
View File
@@ -649,9 +649,7 @@ pub(crate) async fn persist_site_replication_repair_task(
let path = path.to_string();
update_site_replication_state(move |state| {
match failure.as_deref() {
Some(error) => {
upsert_site_replication_retry_event(&mut state.retry_queue, &peer, &path, error, None)?;
}
Some(error) => upsert_site_replication_retry_event(&mut state.retry_queue, &peer, &path, error, None),
None => {
dequeue_site_replication_retry_events_including_escalated(&mut state.retry_queue, &peer, &path);
// A repair is the operator's accountability transfer for the
File diff suppressed because it is too large Load Diff
+6 -31
View File
@@ -28,18 +28,14 @@
//! process-local lock must never be reintroduced in front of it as if it
//! added protection. All IO inside the closure must use the `*_no_lock`
//! config helpers — the locked variants would self-deadlock on the same
//! object lock. Write-lock closures must not perform peer network calls or
//! take other config locks. A read-lock closure may carry bounded peer
//! delivery only when the receiver cannot write this state. Peer-edit
//! delivery must run after the read lock is released.
//! object lock. Do not perform peer network calls or take other config locks
//! inside the closure.
//!
//! Lock order: lifecycle -> bucket-mutation admission -> per-bucket mutation
//! -> bucket operation -> repair admission -> state object lock ->
//! per-bucket metadata. A path may skip levels, but must not acquire an
//! earlier level while holding a later one.
//! Lock order: lifecycle -> bucket operation -> repair admission
//! -> state object lock -> per-bucket metadata.
use super::{S3Error, S3ErrorCode, S3Result, SiteReplicationState, load_site_replication_state_no_lock};
use crate::storage_api::site_replication::{ECStore, with_config_object_read_lock, with_config_object_write_lock};
use super::{S3Error, S3ErrorCode, S3Result};
use crate::storage_api::site_replication::{ECStore, with_config_object_write_lock};
use std::sync::Arc;
use crate::runtime_sources::current_object_store_handle;
@@ -61,27 +57,6 @@ where
with_site_replication_state_lock_on(store, operation).await
}
/// Hold the distributed state-object read lock while `operation` validates a
/// topology snapshot. The closure may carry a bounded peer delivery only when
/// its receiver cannot write site replication state; peer-edit delivery must
/// run after this lock is released. Topology writers use the matching write
/// lock through [`with_site_replication_state_lock`].
pub(crate) async fn with_site_replication_state_read_lock<T, F, Fut>(operation: F) -> S3Result<T>
where
T: Send + 'static,
F: FnOnce(SiteReplicationState) -> Fut + Send + 'static,
Fut: std::future::Future<Output = S3Result<T>> + Send + 'static,
{
let store = current_object_store_handle().ok_or_else(|| S3Error::with_message(S3ErrorCode::InternalError, "Not init"))?;
let read_store = store.clone();
with_config_object_read_lock(store, SITE_REPLICATION_STATE_PATH.to_string(), move || async move {
let state = load_site_replication_state_no_lock(read_store).await?;
operation(state).await
})
.await
.map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("lock site replication state failed: {e}")))?
}
/// Context-store variant for callers that resolve their store from an
/// explicit [`AppContext`] (the service-side reload driven over node RPC).
///
+41 -568
View File
@@ -33,18 +33,6 @@ use temp_env::with_var;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpListener;
#[test]
fn test_bucket_mutation_lock_path_is_bucket_scoped() {
assert_eq!(
site_replication_bucket_mutation_lock_path("photos"),
"config/site-replication/bucket-mutation/photos.lock"
);
assert_ne!(
site_replication_bucket_mutation_lock_path("photos"),
site_replication_bucket_mutation_lock_path("videos")
);
}
fn valid_test_ca_pem(name: &str) -> String {
rcgen::generate_simple_self_signed(vec![name.to_string()])
.expect("generate test CA")
@@ -393,30 +381,6 @@ async fn peer_clients_do_not_follow_redirects() {
assert!(tls_server.await.expect("custom redirect TLS server task"));
}
#[tokio::test]
async fn peer_http_error_body_cannot_spoof_an_unreachable_peer() {
let (endpoint, ca_pem, server) = spawn_test_tls_server_with_response(
b"HTTP/1.1 500 Internal Server Error\r\ncontent-length: 27\r\nconnection: close\r\n\r\ndownstream failed (connect)",
)
.await;
let connection = validate_peer_connection_inner(&endpoint, false, &ca_pem, true).expect("custom CA peer connection");
let client =
build_custom_site_replication_peer_client(&empty_outbound_tls_state(), &connection).expect("custom CA peer client");
let err = PeerAdminRequest::post(&connection, SITE_REPLICATION_PEER_DEVNULL_PATH, "access-key")
.with_client(&client)
.send("secret-key", &serde_json::json!({}))
.await
.expect_err("HTTP 500 must fail");
let detail = err.to_string();
assert!(detail.contains("downstream failed (connect)"));
assert!(
!retry_error_indicates_peer_unreachable(&detail),
"an untrusted response body must not enable the fast reachability probe"
);
assert!(server.await.expect("HTTP error TLS server task"));
}
fn peer(name: &str, endpoint: &str) -> PeerInfo {
PeerInfo {
name: name.to_string(),
@@ -455,7 +419,6 @@ fn drain_event(peer: &str, path: &str, retry_count: u32, updated_at: Option<Offs
last_error: "remote-operation-failed".to_string(),
updated_at,
edit_generation: None,
peer_unreachable: false,
deletions_recorded: false,
}
}
@@ -569,25 +532,25 @@ fn test_record_failed_iam_delivery_records_deletions_and_flags_entry() {
// Non-deletion failure: entry flagged, no record.
let mut user_update = user_delete_item("alice");
user_update.iam_user.as_mut().expect("iam user").is_delete_req = false;
record_failed_iam_delivery(&mut state, &target, &user_update, "peer offline").expect("record failure");
record_failed_iam_delivery(&mut state, &target, &user_update, "peer offline");
assert_eq!(state.retry_queue.len(), 1);
assert_eq!(state.retry_queue[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH);
assert!(state.retry_queue[0].deletions_recorded);
assert!(state.iam_deletion_replays.is_empty());
// Deletion failure: recorded for replay, entry stays flagged.
record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline").expect("record failure");
record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline");
assert_eq!(state.iam_deletion_replays.len(), 1);
assert_eq!(state.iam_deletion_replays[0].entity, "iam-user:alice");
assert!(state.retry_queue[0].deletions_recorded);
assert_eq!(state.retry_queue.len(), 1, "IAM failures stay collapsed per peer");
// Same entity again: newest body replaces the record.
record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline").expect("record failure");
record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline");
assert_eq!(state.iam_deletion_replays.len(), 1);
// Different entity: second record.
record_failed_iam_delivery(&mut state, &target, &policy_delete_item("readonly"), "peer offline").expect("record failure");
record_failed_iam_delivery(&mut state, &target, &policy_delete_item("readonly"), "peer offline");
assert_eq!(state.iam_deletion_replays.len(), 2);
// A legacy entry (created without recording) is never stamped.
@@ -602,9 +565,8 @@ fn test_record_failed_iam_delivery_records_deletions_and_flags_entry() {
SITE_REPLICATION_PEER_IAM_ITEM_WIRE_PATH,
"peer offline",
None,
)
.expect("upsert retry event");
record_failed_iam_delivery(&mut state, &legacy, &user_delete_item("bob"), "peer offline").expect("record failure");
);
record_failed_iam_delivery(&mut state, &legacy, &user_delete_item("bob"), "peer offline");
let legacy_event = state
.retry_queue
.iter()
@@ -627,13 +589,12 @@ fn test_record_failed_iam_delivery_overflow_degrades_to_escalation() {
};
let mut state = deletion_replay_state(&target);
for index in 0..SITE_REPLICATION_IAM_DELETION_REPLAY_LIMIT_PER_PEER {
record_failed_iam_delivery(&mut state, &target, &policy_delete_item(&format!("p{index}")), "peer offline")
.expect("record failure");
record_failed_iam_delivery(&mut state, &target, &policy_delete_item(&format!("p{index}")), "peer offline");
}
assert!(state.retry_queue[0].deletions_recorded);
assert_eq!(state.iam_deletion_replays.len(), SITE_REPLICATION_IAM_DELETION_REPLAY_LIMIT_PER_PEER);
record_failed_iam_delivery(&mut state, &target, &policy_delete_item("one-too-many"), "peer offline").expect("record failure");
record_failed_iam_delivery(&mut state, &target, &policy_delete_item("one-too-many"), "peer offline");
assert_eq!(
state.iam_deletion_replays.len(),
SITE_REPLICATION_IAM_DELETION_REPLAY_LIMIT_PER_PEER,
@@ -659,23 +620,33 @@ fn test_settle_replayed_iam_retry_events_settles_or_escalates() {
// Fully recorded: settles.
let mut state = deletion_replay_state(&target);
record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline").expect("record failure");
record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline");
state.retry_queue[0].updated_at = Some(snapshot_at);
let observed = state.retry_queue[0].clone();
let replayed: Vec<String> = state.iam_deletion_replays.iter().map(|record| record.id.clone()).collect();
assert!(settle_replayed_iam_retry_events(&mut state, &target, &observed, &replayed));
assert!(settle_replayed_iam_retry_events(
&mut state,
&target,
SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH,
Some(snapshot_at),
&replayed,
));
assert!(state.retry_queue.is_empty());
assert!(state.iam_deletion_replays.is_empty());
// Not fully recorded: replayed records are still removed, but the entry
// escalates instead of settling.
let mut state = deletion_replay_state(&target);
record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline").expect("record failure");
record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline");
state.retry_queue[0].updated_at = Some(snapshot_at);
state.retry_queue[0].deletions_recorded = false;
let observed = state.retry_queue[0].clone();
let replayed: Vec<String> = state.iam_deletion_replays.iter().map(|record| record.id.clone()).collect();
assert!(!settle_replayed_iam_retry_events(&mut state, &target, &observed, &replayed));
assert!(!settle_replayed_iam_retry_events(
&mut state,
&target,
SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH,
Some(snapshot_at),
&replayed,
));
assert!(state.iam_deletion_replays.is_empty());
assert_eq!(state.retry_queue.len(), 1);
assert_eq!(state.retry_queue[0].last_error, SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER);
@@ -683,13 +654,17 @@ fn test_settle_replayed_iam_retry_events_settles_or_escalates() {
// Newer failure since the snapshot: entry untouched and drain-eligible,
// residual (unreplayed) record kept for the next pass.
let mut state = deletion_replay_state(&target);
record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline").expect("record failure");
state.retry_queue[0].updated_at = Some(snapshot_at);
let observed = state.retry_queue[0].clone();
record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline");
let replayed: Vec<String> = state.iam_deletion_replays.iter().map(|record| record.id.clone()).collect();
record_failed_iam_delivery(&mut state, &target, &user_delete_item("bob"), "peer offline").expect("record failure");
record_failed_iam_delivery(&mut state, &target, &user_delete_item("bob"), "peer offline");
state.retry_queue[0].updated_at = Some(snapshot_at + time::Duration::seconds(5));
assert!(!settle_replayed_iam_retry_events(&mut state, &target, &observed, &replayed));
assert!(!settle_replayed_iam_retry_events(
&mut state,
&target,
SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH,
Some(snapshot_at),
&replayed,
));
assert_eq!(state.retry_queue.len(), 1);
assert_ne!(state.retry_queue[0].last_error, SITE_REPLICATION_RETRY_SNAPSHOT_REPLAYED_MARKER);
assert!(
@@ -698,24 +673,6 @@ fn test_settle_replayed_iam_retry_events_settles_or_escalates() {
);
assert_eq!(state.iam_deletion_replays.len(), 1);
assert_eq!(state.iam_deletion_replays[0].entity, "iam-user:bob");
// A newer deletion of the same entity gets a fresh replay-record id. An
// older settlement therefore removes neither its body nor its queue
// revision, even if the persisted timestamps happen to be equal.
let mut state = deletion_replay_state(&target);
record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "first failure").expect("record failure");
state.retry_queue[0].updated_at = Some(snapshot_at);
let observed = state.retry_queue[0].clone();
let replayed = vec![state.iam_deletion_replays[0].id.clone()];
record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "newer failure").expect("record newer failure");
state.retry_queue[0].updated_at = Some(snapshot_at);
assert_ne!(state.iam_deletion_replays[0].id, replayed[0]);
assert!(!settle_replayed_iam_retry_events(&mut state, &target, &observed, &replayed));
assert_eq!(state.retry_queue.len(), 1);
assert_ne!(state.retry_queue[0].id, observed.id);
assert_eq!(state.iam_deletion_replays.len(), 1);
assert_eq!(state.iam_deletion_replays[0].entity, "iam-user:alice");
}
/// Merging legacy wire-path rows into the collapsed entry must not launder an
@@ -791,281 +748,6 @@ fn test_classify_site_replication_retry_event_actions() {
assert_eq!(classify("/rustfs/admin/v3/site-replication/peer/unknown"), None);
}
#[test]
fn test_bucket_make_retry_replays_matching_configure_before_settlement() {
let make_photos =
"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning".to_string();
let configure_photos =
"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=configure-replication".to_string();
let plan = SiteReplicationBootstrapPlan {
bucket_make_ops: vec![
make_photos.clone(),
"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=videos&operation=make-with-versioning".to_string(),
],
bucket_configure_ops: vec![
"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=videos&operation=configure-replication".to_string(),
configure_photos.clone(),
],
bucket_items: vec![
SRBucketMeta {
bucket: "videos".to_string(),
r#type: "tags".to_string(),
..Default::default()
},
SRBucketMeta {
bucket: "photos".to_string(),
r#type: "policy".to_string(),
..Default::default()
},
],
..Default::default()
};
let tasks = bucket_op_retry_replay_tasks(&plan, SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING, "photos")
.expect("make retry plan should include its configure follow-up");
assert_eq!(
tasks.iter().map(SiteReplicationRepairTask::path).collect::<Vec<_>>(),
vec![
make_photos.as_str(),
"/rustfs/admin/v3/site-replication/peer/bucket-meta",
configure_photos.as_str()
]
);
assert!(matches!(tasks[0], SiteReplicationRepairTask::BucketMake(_)));
assert!(matches!(&tasks[1], SiteReplicationRepairTask::BucketMetadata(item) if item.bucket == "photos"));
assert!(matches!(tasks[2], SiteReplicationRepairTask::Replication(_)));
}
#[test]
fn test_bucket_make_retry_without_matching_configure_fails_closed() {
let plan = SiteReplicationBootstrapPlan {
bucket_make_ops: vec![
"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning".to_string(),
],
..Default::default()
};
let err = match bucket_op_retry_replay_tasks(&plan, SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING, "photos") {
Ok(_) => panic!("make retry must not settle without a matching configure operation"),
Err(err) => err,
};
assert_eq!(err.code(), &S3ErrorCode::InternalError);
}
#[test]
fn test_retry_drain_bounds_each_peer_round_to_one_small_request_chain() {
let plan = SiteReplicationBootstrapPlan {
iam_items: vec![SRIAMItem::default(); 3],
bucket_make_ops: vec![
"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning".to_string(),
],
bucket_configure_ops: vec![
"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=configure-replication".to_string(),
],
bucket_items: vec![SRBucketMeta {
bucket: "photos".to_string(),
r#type: "tags".to_string(),
..Default::default()
}],
..Default::default()
};
let make = RetryDrainAction::BucketOpReplay {
operation: SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING.to_string(),
bucket: "photos".to_string(),
};
assert!(is_lightweight_retry_drain_action(&make));
assert!(!is_lightweight_retry_drain_action(&RetryDrainAction::IamSnapshot));
assert!(!is_lightweight_retry_drain_action(&RetryDrainAction::PeerEdit));
assert_eq!(retry_drain_request_count(&make, Some(&plan)), 3);
assert!(retry_drain_request_count(&make, Some(&plan)) <= SITE_REPLICATION_RETRY_DRAIN_MAX_REQUESTS_PER_PEER);
assert!(
retry_drain_request_count(&RetryDrainAction::IamSnapshot, Some(&plan))
> SITE_REPLICATION_RETRY_DRAIN_MAX_REQUESTS_PER_PEER
);
assert!(
retry_drain_request_count(&RetryDrainAction::PeerEdit, Some(&plan)) > SITE_REPLICATION_RETRY_DRAIN_MAX_REQUESTS_PER_PEER
);
}
#[test]
fn test_lightweight_retry_peer_rotation_covers_all_queued_peers() {
let limit = SITE_REPLICATION_RETRY_DRAIN_PEER_CONCURRENCY;
for peer_count in 1..=(limit * 3 + 1) {
let rounds = peer_count.div_ceil(limit);
let mut seen = HashSet::new();
for round in 7..(7 + rounds as i64) {
let start = lightweight_retry_peer_rotation(peer_count, round);
for offset in 0..limit.min(peer_count) {
seen.insert((start + offset) % peer_count);
}
}
assert_eq!(
seen.len(),
peer_count,
"every peer must enter the bounded lightweight window within {rounds} rounds"
);
}
}
#[test]
fn test_lightweight_bucket_retry_plan_is_targeted_and_preserves_make_options() {
let created_at = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
let bucket = SRBucketInfo {
bucket: "photos".to_string(),
created_at: Some(created_at),
object_lock_config: Some(String::new()),
tags: Some("dGFncy14bWw=".to_string()),
tag_config_updated_at: Some(created_at),
..Default::default()
};
let plan = site_replication_bucket_retry_plan_for(&bucket, false).expect("targeted retry plan");
assert!(plan.iam_items.is_empty());
assert_eq!(plan.bucket_make_ops.len(), 1);
assert!(plan.bucket_make_ops[0].contains("bucket=photos"));
assert!(plan.bucket_make_ops[0].contains("lockEnabled=true"));
assert!(plan.bucket_make_ops[0].contains("createdAt="));
assert_eq!(plan.bucket_items.len(), 2);
assert_eq!(plan.bucket_items[0].r#type, "tags");
assert_eq!(plan.bucket_items[1].r#type, "object-lock-config");
assert_eq!(plan.bucket_configure_ops.len(), 1);
assert!(plan.bucket_configure_ops[0].contains("operation=configure-replication"));
let tasks =
bucket_op_retry_replay_tasks(&plan, SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING, "photos").expect("retry task chain");
assert!(matches!(tasks[0], SiteReplicationRepairTask::BucketMake(_)));
assert!(matches!(tasks[1], SiteReplicationRepairTask::BucketMetadata(_)));
assert!(matches!(tasks[2], SiteReplicationRepairTask::BucketMetadata(_)));
assert!(matches!(tasks[3], SiteReplicationRepairTask::Replication(_)));
}
#[test]
fn test_lightweight_bucket_retry_plan_orders_real_metadata_and_counts_it() {
let versioning = bucket_versioning_xml().expect("canonical versioning config");
let replication = serialize(&site_repl_config("remote-dep")).expect("derived replication config");
let bucket = SRBucketInfo {
bucket: "photos".to_string(),
policy: Some(serde_json::json!({"Version":"2012-10-17","Statement":[]})),
tags: Some(BASE64_STANDARD.encode_to_string("<Tagging/>")),
versioning: Some(BASE64_STANDARD.encode_to_string(&versioning)),
replication_config: Some(BASE64_STANDARD.encode_to_string(&replication)),
..Default::default()
};
let plan = site_replication_bucket_retry_plan_from_info(&bucket, false).expect("targeted retry plan");
let tasks = bucket_op_retry_replay_tasks(&plan, SITE_REPLICATION_BUCKET_OP_MAKE_WITH_VERSIONING, "photos")
.expect("bucket replay tasks");
assert!(matches!(tasks.first(), Some(SiteReplicationRepairTask::BucketMake(_))));
assert!(matches!(tasks.last(), Some(SiteReplicationRepairTask::Replication(_))));
assert!(
tasks[1..tasks.len() - 1]
.iter()
.all(|task| matches!(task, SiteReplicationRepairTask::BucketMetadata(_)))
);
assert_eq!(tasks.len(), 4, "make + policy + tags + configure must all count against the budget");
assert!(
tasks.len() <= SITE_REPLICATION_RETRY_DRAIN_MAX_REQUESTS_PER_PEER,
"the complete metadata chain must fit the bounded lightweight replay"
);
let mut operator_replication = site_repl_config("remote-dep");
operator_replication.rules.push(operator_rule("operator-backup"));
let mut bucket_with_operator_rule = bucket;
bucket_with_operator_rule.replication_config =
Some(BASE64_STANDARD.encode_to_string(&serialize(&operator_replication).expect("operator replication config")));
let plan = site_replication_bucket_retry_plan_from_info(&bucket_with_operator_rule, false).expect("targeted retry plan");
assert!(
plan.bucket_items.iter().any(|item| item.r#type == "replication-config"),
"operator-authored replication rules cannot be replaced by configure-replication"
);
}
#[test]
fn test_delete_bucket_broadcast_fences_target_membership_through_delivery() {
let hooks = include_str!("hooks.rs");
let delete_broadcast = hooks
.split("async fn broadcast_site_replication_delete_bucket")
.nth(1)
.and_then(|rest| rest.split("pub(crate) async fn commit_site_replication_delete_bucket").next())
.expect("delete-bucket broadcast should exist");
assert!(
delete_broadcast.contains("with_site_replication_state_read_lock(move |state| async move {")
&& delete_broadcast.contains("state.peers.get(&fallback_peer.deployment_id)")
&& delete_broadcast.contains("site_replicator_service_account_secret(&state.service_account_access_key)")
&& delete_broadcast
.contains("PeerAdminRequest::put(&transport.connection, &delivery_path, &state.service_account_access_key)"),
"a destructive bucket delivery must resolve current topology and credentials under the distributed state read lock"
);
assert!(
delete_broadcast.contains("enqueue_site_replication_retry_event(&current_peer, &request_path, &err).await"),
"a failed destructive delivery must remain visible for operator repair"
);
let usecase = include_str!("../app/bucket_usecase.rs");
let delete = usecase
.split("async fn execute_delete_bucket_inner")
.nth(1)
.and_then(|rest| rest.split("pub async fn execute_head_bucket").next())
.expect("delete bucket usecase");
assert!(
delete
.find("prepare_site_replication_delete_bucket")
.expect("durable reservation")
< delete.find(".delete_bucket(").expect("local delete"),
"destructive peer liabilities must be persisted before the local bucket is deleted"
);
}
#[test]
fn test_bucket_retry_settlement_preserves_a_newer_same_path_failure() {
let peer = peer("remote", "https://remote.example.com");
let path = "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=configure-replication";
let observed_at = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
let observed = drain_event("remote", path, 1, Some(observed_at));
let mut queue = vec![observed.clone()];
queue[0].id = "evt-remote-new-revision".to_string();
queue[0].retry_count += 1;
assert_eq!(settle_observed_site_replication_retry_event(&mut queue, &peer, &observed), 0);
assert_eq!(queue.len(), 1, "a newer same-timestamp failure must survive stale settlement");
let current = queue[0].clone();
assert_eq!(settle_observed_site_replication_retry_event(&mut queue, &peer, &current), 1);
assert!(queue.is_empty());
}
#[test]
fn test_reachable_probe_promotion_is_fenced_by_the_observed_event() {
let now = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
let path = "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning";
let mut event = drain_event("remote", path, 3, Some(now));
event.peer_unreachable = true;
let recovered = event.clone();
let mut state = SiteReplicationState {
retry_queue: vec![event],
..Default::default()
};
state
.peers
.insert("remote".to_string(), peer("remote", "https://remote.example.com"));
assert_eq!(mark_reachable_deferred_retry_events(&mut state, &[recovered.clone()]), 1);
assert_eq!(state.retry_queue[0].updated_at, None);
assert!(!state.retry_queue[0].peer_unreachable);
assert_eq!(
actionable_site_replication_retry_events(&state, now).len(),
1,
"a successful probe must make the event replayable in the same drain tick"
);
state.retry_queue[0].updated_at = Some(now + time::Duration::seconds(1));
state.retry_queue[0].peer_unreachable = true;
assert_eq!(mark_reachable_deferred_retry_events(&mut state, &[recovered]), 0);
assert_eq!(state.retry_queue[0].updated_at, Some(now + time::Duration::seconds(1)));
assert!(state.retry_queue[0].peer_unreachable);
}
#[test]
fn test_retry_snapshot_fingerprint_detects_concurrent_iam_change() {
let old = SRIAMItem {
@@ -1146,100 +828,6 @@ fn test_site_replication_retry_backoff_schedule() {
assert!(elapsed(30, 86_401));
}
#[test]
fn test_retry_error_marks_peer_unreachable_only_for_connection_failures() {
let mut queue = Vec::new();
let peer = peer("remote", "https://remote.example.com");
let bucket_make = "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning";
upsert_site_replication_retry_event(
&mut queue,
&peer,
bucket_make,
"peer request to https://remote.example.com failed (connect): connection refused",
None,
)
.expect("upsert retry event");
assert!(queue[0].peer_unreachable);
upsert_site_replication_retry_event(
&mut queue,
&peer,
bucket_make,
"peer request to https://remote.example.com failed (timeout): request exceeded 10 seconds",
None,
)
.expect("upsert retry event");
assert!(
!queue[0].peer_unreachable,
"a whole-request timeout does not prove the peer is unreachable"
);
upsert_site_replication_retry_event(
&mut queue,
&peer,
bucket_make,
"peer request to https://remote.example.com failed with 500 Internal Server Error: downstream failed (connect)",
None,
)
.expect("upsert retry event");
assert!(
!queue[0].peer_unreachable,
"application failures and their untrusted bodies must keep the normal replay backoff"
);
upsert_site_replication_retry_event(
&mut queue,
&peer,
bucket_make,
"peer request to https://remote.example.com failed with 500 Internal Server Error: backend failed (connect): spoofed",
None,
)
.expect("upsert retry event");
assert!(!queue[0].peer_unreachable, "peer response bodies must not spoof transport failures");
}
#[test]
fn test_connect_timeout_is_classified_as_a_connection_failure() {
assert_eq!(classify_peer_transport_error(true, true, "tcp connect timed out"), "connect");
assert_eq!(classify_peer_transport_error(false, true, "request timed out"), "timeout");
assert_eq!(
classify_peer_transport_error(false, true, "request timed out for https://tls-gateway.example"),
"timeout"
);
assert_eq!(classify_peer_transport_error(true, false, "tls handshake failed"), "tls handshake");
}
#[test]
fn test_retry_event_peer_unreachable_is_legacy_serde_default() {
let json = r#"{
"id":"evt-legacy",
"peer_deployment_id":"remote",
"peer_endpoint":"https://remote.example.com",
"path":"/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning",
"retry_count":1,
"failed":false,
"last_error":"peer request to https://remote.example.com failed (connect): connection refused"
}"#;
let mut event: SiteReplicationRetryEvent = serde_json::from_str(json).expect("legacy retry event decodes");
assert!(!event.peer_unreachable);
let now = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
event.updated_at = Some(now - time::Duration::seconds(30));
let mut state = SiteReplicationState::default();
state
.peers
.insert("remote".to_string(), peer("remote", "https://remote.example.com"));
state.retry_queue.push(event);
assert_eq!(
deferred_site_replication_retry_events(&state, now).len(),
1,
"rolling-upgrade records must retain fast recovery from their trusted outer error shape"
);
}
/// The actionable subset respects classification, peer membership and
/// backoff; everything else stays untouched in the queue.
#[test]
@@ -1327,51 +915,6 @@ fn test_deferred_site_replication_retry_events_partition() {
assert_eq!(actionable[0].path, "/rustfs/admin/v3/site-replication/peer/bucket-meta");
}
#[test]
fn test_deferred_retry_events_probe_fresh_peer_transport_failures() {
let now = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
let mut state = SiteReplicationState::default();
state
.peers
.insert("remote".to_string(), peer("remote", "https://remote.example.com"));
let bucket_make = "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning";
let mut fresh_transport_failure = drain_event("remote", bucket_make, 1, Some(now - time::Duration::seconds(30)));
fresh_transport_failure.peer_unreachable = true;
state.retry_queue.push(fresh_transport_failure);
let deferred = deferred_site_replication_retry_events(&state, now);
assert_eq!(
deferred.len(),
1,
"fresh transport failures must be eligible for a cheap reachability probe"
);
assert_eq!(deferred[0].path, bucket_make);
let actionable = actionable_site_replication_retry_events(&state, now);
assert!(actionable.is_empty(), "the event is still protected from direct replay by normal backoff");
}
#[test]
fn test_deferred_retry_events_do_not_probe_fresh_application_failures() {
let now = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
let mut state = SiteReplicationState::default();
state
.peers
.insert("remote".to_string(), peer("remote", "https://remote.example.com"));
let bucket_make = "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning";
state
.retry_queue
.push(drain_event("remote", bucket_make, 1, Some(now - time::Duration::seconds(30))));
assert!(
deferred_site_replication_retry_events(&state, now).is_empty(),
"reachable peers that reject an operation must keep the base replay backoff"
);
assert!(actionable_site_replication_retry_events(&state, now).is_empty());
}
/// The drain settles a peer-edit success under a freshly allocated
/// generation; legacy queue entries carry `edit_generation: None` and
/// must be cleared by that generation-scoped settlement (`(Some, None)`
@@ -1453,7 +996,7 @@ fn test_escalate_up_to_marks_snapshot_replayed_and_keeps_newer_failures() {
// successful Bob update on the shared wire path cannot erase it even
// before the drain runs.
let mut queue = Vec::new();
upsert_site_replication_retry_event(&mut queue, &target, path, "alice delete failed", None).expect("upsert retry event");
upsert_site_replication_retry_event(&mut queue, &target, path, "alice delete failed", None);
assert_eq!(queue[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH);
assert_eq!(dequeue_site_replication_retry_events(&mut queue, &target, path), 0);
assert_eq!(queue.len(), 1);
@@ -1462,7 +1005,7 @@ fn test_escalate_up_to_marks_snapshot_replayed_and_keeps_newer_failures() {
// A later hook failure overwrites the marker and re-arms the drain.
let mut queue = vec![drain_event("remote", path, 2, Some(snapshot_at))];
escalate_site_replication_retry_events_up_to(&mut queue, &target, path, Some(snapshot_at));
upsert_site_replication_retry_event(&mut queue, &target, path, "peer offline", None).expect("upsert retry event");
upsert_site_replication_retry_event(&mut queue, &target, path, "peer offline", None);
assert!(classify_site_replication_retry_event(&queue[0]).is_some());
// Legacy entry without a timestamp: escalated.
@@ -2238,85 +1781,17 @@ fn test_retry_event_upsert_marks_repeated_failures() {
};
let mut queue = Vec::new();
upsert_site_replication_retry_event(&mut queue, &peer, "/rustfs/admin/v3/site-replication/peer/iam-item", "first", None)
.expect("upsert retry event");
let first_revision = queue[0].id.clone();
upsert_site_replication_retry_event(&mut queue, &peer, "/rustfs/admin/v3/site-replication/peer/iam-item", "second", None)
.expect("upsert retry event");
let second_revision = queue[0].id.clone();
upsert_site_replication_retry_event(&mut queue, &peer, "/rustfs/admin/v3/site-replication/peer/iam-item", "third", None)
.expect("upsert retry event");
upsert_site_replication_retry_event(&mut queue, &peer, "/rustfs/admin/v3/site-replication/peer/iam-item", "first", None);
upsert_site_replication_retry_event(&mut queue, &peer, "/rustfs/admin/v3/site-replication/peer/iam-item", "second", None);
upsert_site_replication_retry_event(&mut queue, &peer, "/rustfs/admin/v3/site-replication/peer/iam-item", "third", None);
assert_eq!(queue.len(), 1);
assert_ne!(first_revision, second_revision);
assert_ne!(second_revision, queue[0].id, "each failure must advance the settlement revision");
assert_eq!(queue[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH);
assert_eq!(queue[0].retry_count, SITE_REPLICATION_RETRY_FAILED_AFTER);
assert!(queue[0].failed);
assert_eq!(queue[0].last_error, "third");
}
#[test]
fn retry_queue_capacity_never_evicts_destructive_bucket_liabilities() {
let target = PeerInfo {
deployment_id: "remote-dep".to_string(),
..peer("remote", "https://remote.example.com")
};
let destructive = |index: usize| SiteReplicationRetryEvent {
id: format!("delete-{index}"),
peer_deployment_id: target.deployment_id.clone(),
peer_endpoint: target.endpoint.clone(),
path: format!("{SITE_REPLICATION_PEER_BUCKET_OPS_PATH}?bucket=bucket-{index}&operation=delete-bucket"),
..Default::default()
};
let mut queue = (0..SITE_REPLICATION_RETRY_QUEUE_LIMIT).map(destructive).collect::<Vec<_>>();
let original_ids = queue.iter().map(|event| event.id.clone()).collect::<HashSet<_>>();
let new_path = format!("{SITE_REPLICATION_PEER_BUCKET_OPS_PATH}?bucket=overflow&operation=force-delete-bucket");
let err = upsert_site_replication_retry_event(&mut queue, &target, &new_path, "reserve delete", None)
.expect_err("an all-destructive full queue must fail closed");
assert_eq!(err.code(), &S3ErrorCode::ServiceUnavailable);
assert_eq!(queue.len(), SITE_REPLICATION_RETRY_QUEUE_LIMIT);
assert_eq!(queue.iter().map(|event| event.id.clone()).collect::<HashSet<_>>(), original_ids);
queue[0] = SiteReplicationRetryEvent {
id: "iam-snapshot".to_string(),
peer_deployment_id: target.deployment_id.clone(),
peer_endpoint: target.endpoint.clone(),
path: SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH.to_string(),
deletions_recorded: true,
..Default::default()
};
upsert_site_replication_retry_event(&mut queue, &target, &new_path, "reserve delete", None)
.expect_err("a collapsed IAM liability may contain a deletion and must not be evicted");
assert!(queue.iter().any(|event| event.id == "iam-snapshot"));
queue[0] = SiteReplicationRetryEvent {
id: "rebuildable".to_string(),
peer_deployment_id: target.deployment_id.clone(),
peer_endpoint: target.endpoint.clone(),
path: SITE_REPLICATION_PEER_EDIT_PATH.to_string(),
..Default::default()
};
let preserved_delete_ids = queue
.iter()
.filter(|event| is_destructive_bucket_retry_path(&event.path))
.map(|event| event.id.clone())
.collect::<HashSet<_>>();
let evicted = upsert_site_replication_retry_event(&mut queue, &target, &new_path, "reserve delete", None)
.expect("a rebuildable row may make room for a destructive liability");
assert_eq!(evicted.len(), 1);
assert_eq!(evicted[0].id, "rebuildable");
assert_eq!(queue.len(), SITE_REPLICATION_RETRY_QUEUE_LIMIT);
assert!(
preserved_delete_ids
.iter()
.all(|id| queue.iter().any(|event| &event.id == id))
);
assert!(queue.iter().any(|event| event.path == new_path));
}
/// P1-15 review follow-up: a successful peer-edit delivery only proves the
/// peer reached the state THAT delivery carried. Settling it must not
/// erase a retry event a newer edit left behind, or the local site sits on
@@ -2332,8 +1807,7 @@ fn retry_settlement_must_not_erase_a_newer_generation_failure() {
// Edit A (generation 5) delivered successfully and is stalled before
// settling. Edit B (generation 6) commits meanwhile, fails delivery to
// the same peer, and enqueues.
upsert_site_replication_retry_event(&mut queue, &peer, SITE_REPLICATION_PEER_EDIT_PATH, "peer offline", Some(6))
.expect("upsert retry event");
upsert_site_replication_retry_event(&mut queue, &peer, SITE_REPLICATION_PEER_EDIT_PATH, "peer offline", Some(6));
// A resumes: its own settlement must leave B's retry alone.
assert_eq!(
@@ -2344,8 +1818,7 @@ fn retry_settlement_must_not_erase_a_newer_generation_failure() {
assert_eq!(queue[0].edit_generation, Some(6));
// An even older delivery failing afterwards must not lower the fence.
upsert_site_replication_retry_event(&mut queue, &peer, SITE_REPLICATION_PEER_EDIT_PATH, "still offline", Some(4))
.expect("upsert retry event");
upsert_site_replication_retry_event(&mut queue, &peer, SITE_REPLICATION_PEER_EDIT_PATH, "still offline", Some(4));
assert_eq!(queue[0].edit_generation, Some(6));
// B's own delivery succeeding is what clears it.
@@ -2358,7 +1831,7 @@ fn retry_settlement_must_not_erase_a_newer_generation_failure() {
// Collapsed broadcast failures live under an internal snapshot path;
// an unrelated success on their shared wire path cannot settle them.
let iam_path = "/rustfs/admin/v3/site-replication/peer/iam-item";
upsert_site_replication_retry_event(&mut queue, &peer, iam_path, "peer offline", None).expect("upsert retry event");
upsert_site_replication_retry_event(&mut queue, &peer, iam_path, "peer offline", None);
assert_eq!(dequeue_site_replication_retry_events(&mut queue, &peer, iam_path), 0);
assert_eq!(queue[0].path, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH);
}
+13 -17
View File
@@ -331,7 +331,6 @@ pub(crate) fn runtime_peer_connection(peer: &PeerInfo) -> S3Result<PeerConnectio
})
}
#[derive(Clone)]
pub(crate) struct PeerTransport {
pub(crate) connection: PeerConnection,
pub(crate) client: reqwest::Client,
@@ -698,7 +697,19 @@ impl<'a> PeerAdminRequest<'a> {
}
let response = req.send().await.map_err(|e| {
let classify = classify_peer_transport_error(e.is_connect(), e.is_timeout(), &e.to_string());
let classify = if e.is_timeout() {
"timeout"
} else if e.is_connect() && e.to_string().to_ascii_lowercase().contains("dns") {
"dns resolution"
} else if e.to_string().to_ascii_lowercase().contains("certificate")
|| e.to_string().to_ascii_lowercase().contains("tls")
{
"tls handshake"
} else if e.is_connect() {
"connect"
} else {
"request"
};
S3Error::with_message(S3ErrorCode::InternalError, format!("peer request to {url} failed ({classify}): {e}"))
})?;
@@ -814,21 +825,6 @@ impl<'a> PeerAdminRequest<'a> {
}
}
pub(crate) fn classify_peer_transport_error(is_connect: bool, is_timeout: bool, detail: &str) -> &'static str {
let detail = detail.to_ascii_lowercase();
if is_connect && detail.contains("dns") {
"dns resolution"
} else if is_connect && (detail.contains("certificate") || detail.contains("tls")) {
"tls handshake"
} else if is_connect {
"connect"
} else if is_timeout {
"timeout"
} else {
"request"
}
}
pub(crate) fn peer_error_may_be_secret_mismatch(detail: &str) -> bool {
let detail = detail.to_ascii_lowercase();
detail.contains("signaturedoesnotmatch")
+3 -45
View File
@@ -28,19 +28,16 @@ use std::pin::Pin;
use std::sync::OnceLock;
use std::time::Duration;
use tokio::time::Instant;
use tokio_util::sync::CancellationToken;
use tracing::warn;
const RECONCILE_INTERVAL: Duration = Duration::from_secs(600);
pub(crate) const RETRY_DRAIN_INTERVAL: Duration = Duration::from_secs(30);
/// A reconciler reports its own failures; the outcome carries no value because neither
/// caller can act on one — a site that cannot repair its replication wiring still serves S3.
type ReconcileHook = fn() -> Pin<Box<dyn Future<Output = ()> + Send>>;
static RECONCILER: OnceLock<ReconcileHook> = OnceLock::new();
static RETRY_DRAINER: OnceLock<ReconcileHook> = OnceLock::new();
/// Install the admin layer's reconciler. Idempotent: a second call is ignored, which keeps
/// repeated router construction (tests, the embedded server) from panicking.
@@ -48,12 +45,6 @@ pub(crate) fn register_site_replication_reconciler(reconcile: ReconcileHook) {
let _ = RECONCILER.set(reconcile);
}
/// Install the admin layer's lightweight retry drain. Idempotent for the same
/// reason as [`register_site_replication_reconciler`].
pub(crate) fn register_site_replication_retry_drainer(drain: ReconcileHook) {
let _ = RETRY_DRAINER.set(drain);
}
/// Repair drifted site-replication wiring, immediately and then on a timer.
///
/// The first pass runs inside the spawned task rather than on the caller's path: it walks
@@ -71,38 +62,16 @@ pub(crate) fn spawn_site_replication_reconcile_task(ctx: CancellationToken) {
return;
}
spawn_reconcile_loop(ctx.clone(), RECONCILE_INTERVAL, &RECONCILER, true);
if RETRY_DRAINER.get().is_none() {
warn!("site replication retry drainer is not registered; periodic retry drain disabled");
return;
}
spawn_reconcile_loop(ctx, RETRY_DRAIN_INTERVAL, &RETRY_DRAINER, false);
}
fn spawn_reconcile_loop(
ctx: CancellationToken,
interval: Duration,
hook: &'static OnceLock<ReconcileHook>,
run_immediately: bool,
) {
tokio::spawn(async move {
let first_tick = if run_immediately {
Instant::now()
} else {
Instant::now() + interval
};
let mut ticker = tokio::time::interval_at(first_tick, interval);
let mut ticker = tokio::time::interval(RECONCILE_INTERVAL);
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
_ = ctx.cancelled() => break,
// The heavy reconciler owns the startup repair pass. The lightweight
// retry drain starts on its normal cadence so it cannot steal that
// first lifecycle lock and defer bucket/IAM repair for a full interval.
// The first tick fires immediately, which is the startup repair pass.
_ = ticker.tick() => {
if let Some(reconcile) = hook.get() {
if let Some(reconcile) = RECONCILER.get() {
reconcile().await;
}
}
@@ -110,14 +79,3 @@ fn spawn_reconcile_loop(
}
});
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn retry_drain_runs_faster_than_heavy_reconcile() {
assert!(RETRY_DRAIN_INTERVAL < RECONCILE_INTERVAL);
assert!(RETRY_DRAIN_INTERVAL <= Duration::from_secs(60));
}
}
+2 -2
View File
@@ -244,8 +244,8 @@ pub(crate) mod site_replication {
pub(crate) use crate::storage::storage_api::{Endpoint, Endpoints, PoolEndpoints};
pub(crate) use crate::storage::storage_api::{
ECStore, EndpointServerPools, StorageError, delete_config_no_lock, is_err_bucket_not_found, lock_bucket_targets_metadata,
read_config, read_config_no_lock, save_config_no_lock, with_config_object_read_lock, with_config_object_write_lock,
ECStore, EndpointServerPools, StorageError, delete_config_no_lock, lock_bucket_targets_metadata, read_config,
read_config_no_lock, save_config_no_lock, with_config_object_read_lock, with_config_object_write_lock,
};
pub(crate) mod metadata_sys {
+20
View File
@@ -20,6 +20,8 @@
crates/ecstore/src/bucket/lifecycle/tier_last_day_stats.rs|clippy::all
crates/ecstore/src/bucket/lifecycle/tier_last_day_stats.rs|unused_must_use
crates/ecstore/src/bucket/lifecycle/tier_last_day_stats.rs|unused_variables
crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs|clippy::all
crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs|unused_must_use
crates/ecstore/src/bucket/lifecycle/tier_sweeper.rs|unused_variables
crates/s3-client/src/api_error_response.rs|clippy::all
crates/s3-client/src/api_error_response.rs|unused_must_use
@@ -72,18 +74,36 @@ crates/ecstore/src/services/event_notification.rs|unused_variables
crates/ecstore/src/services/tier/tier.rs|clippy::all
crates/ecstore/src/services/tier/tier.rs|unused_must_use
crates/ecstore/src/services/tier/tier.rs|unused_variables
crates/ecstore/src/services/tier/tier_admin.rs|clippy::all
crates/ecstore/src/services/tier/tier_admin.rs|unused_must_use
crates/ecstore/src/services/tier/tier_admin.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend.rs|clippy::all
crates/ecstore/src/services/tier/warm_backend.rs|unused_must_use
crates/ecstore/src/services/tier/warm_backend.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_aliyun.rs|clippy::all
crates/ecstore/src/services/tier/warm_backend_aliyun.rs|unused_must_use
crates/ecstore/src/services/tier/warm_backend_aliyun.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_azure.rs|clippy::all
crates/ecstore/src/services/tier/warm_backend_azure.rs|unused_must_use
crates/ecstore/src/services/tier/warm_backend_azure.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_gcs.rs|clippy::all
crates/ecstore/src/services/tier/warm_backend_gcs.rs|unused_must_use
crates/ecstore/src/services/tier/warm_backend_gcs.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_huaweicloud.rs|clippy::all
crates/ecstore/src/services/tier/warm_backend_huaweicloud.rs|unused_must_use
crates/ecstore/src/services/tier/warm_backend_huaweicloud.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_minio.rs|clippy::all
crates/ecstore/src/services/tier/warm_backend_minio.rs|unused_must_use
crates/ecstore/src/services/tier/warm_backend_minio.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_r2.rs|clippy::all
crates/ecstore/src/services/tier/warm_backend_r2.rs|unused_must_use
crates/ecstore/src/services/tier/warm_backend_r2.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_rustfs.rs|clippy::all
crates/ecstore/src/services/tier/warm_backend_rustfs.rs|unused_must_use
crates/ecstore/src/services/tier/warm_backend_rustfs.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_s3.rs|clippy::all
crates/ecstore/src/services/tier/warm_backend_s3.rs|unused_must_use
crates/ecstore/src/services/tier/warm_backend_s3.rs|unused_variables
crates/ecstore/src/services/tier/warm_backend_tencent.rs|clippy::all
crates/ecstore/src/services/tier/warm_backend_tencent.rs|unused_must_use
crates/ecstore/src/services/tier/warm_backend_tencent.rs|unused_variables