Compare commits

...

9 Commits

Author SHA1 Message Date
cxymds 7b6e45d372 fix(heal): preserve null-version tags and encoded object paths (#7644) 2026-09-10 21:10:32 +08:00
cxymds 09c85a5f29 test(ilm): restore noncurrent compensation tests to serial CI (#7645) 2026-09-10 20:36:32 +08:00
cxymds 00aeb12914 fix(ilm): preserve cleanup ownership on tiered overwrites (#7639)
* fix(ilm): preserve cleanup ownership on tiered overwrites

* fix(ci): refresh E2E selection for tier overwrite regression
2026-09-10 20:14:07 +08:00
houseme ffb18979f8 fix(scanner): find nextest junit fallback for evidence cases (#7643)
Handle cargo-nextest writing JUnit reports under the workspace target directory even when the Rust build uses CARGO_TARGET_DIR. This keeps measured Scanner/Heal evidence cases from passing the real test but failing final receipt packaging.

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 19:58:00 +08:00
houseme 97c7b451d2 test(scanner): bound G14 multi-pool evidence size (#7640)
Keep the multi-pool evidence runner within the registry object budget while retaining deferred outage-write diagnostics for release-gate validation.

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 18:08:02 +08:00
houseme 9ba95cac37 test(e2e): record deferred G14 outage writes (#7637)
Include optional outage-write diagnostics in G14 evidence oracles so deferred multi-pool outage writes bind their down-window refusals and post-rejoin acceptance.

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 17:48:44 +08:00
houseme cec796328c test(e2e): keep G14 multi-set restart graceful (#7633)
test(e2e): keep multi-set restart graceful

Include the EC8+4 multi-set restart scenario in the graceful interruption lane so the unclean-shutdown marker assertion matches the scenario semantics.

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 17:10:46 +08:00
houseme 1aa4fa145f fix(ecstore): use pool lock quorum for set writes (#7630)
Use the pool-wide namespace lock client domain for every set in a pool so degraded EC8+4 multi-set writes are gated by node-level lock quorum instead of the narrower per-set endpoint host slice.

Keep namespace-lock domain deduplication tied to both the pool namespace and shared clients, and add regression coverage for three-locker degraded writes plus cross-pool domain separation.

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 16:42:21 +08:00
houseme d286f3d06c chore(deps): refresh release dependencies (#7632)
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 16:41:31 +08:00
23 changed files with 1499 additions and 114 deletions
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-darwin=874c881d7b45f12378a5817c7f42c95c4981960a2ec9ce12dcf4af239ae1f9d5
sha256-linux=9351e25b45bf7dfce18b951a5e3740225f457cacc53b8bf9f500f6947763ec0e
sha256-darwin=15cb0cf9909bfbfc5a835fb08bd3675516c641db1b92ecc1e170a1cfa0fa2fd5
sha256-linux=3163fdd29df5def86cf511ca7db05880d5c0caaa608a7031a71394327f7c217a
+1 -1
View File
@@ -1 +1 @@
sha256=6d18f9cce820c51d5589de944e8cc185f73eeca0ea9a9916651943e3759169d0
sha256=5fbb230b89212b7c3d7229d6cef3e7e2d16f0ecfec62237ebc770785706f67d9
+2 -6
View File
@@ -335,11 +335,7 @@ jobs:
# re-enabled by backlog#1304 (restore accepts serialize on a short CAS
# guard; the copy-back no longer holds the #4877 whole-copy-back lock,
# so the mid-restore ongoing read and fast 409 rejection it asserts are
# the implemented contract). The remaining exclusions each hit a
# DIFFERENT, independent issue (all tracked under rustfs/backlog#1148;
# they keep #[ignore] with a backlog reference):
# - test_noncurrent_{expiry,transition}_still_works_after_immediate_compensation_transition:
# noncurrent transition/expiry after an immediate compensation transition.
# the implemented contract).
- name: Run ignored ILM integration tests serially
env:
# Match the measured Test and Lint link budget. The default exposed
@@ -352,7 +348,7 @@ jobs:
NEXTEST_HIDE_PROGRESS_BAR=1 timeout --verbose --signal=TERM --kill-after=30s 80m \
cargo nextest run -j1 --run-ignored ignored-only \
-p rustfs-scanner -p rustfs \
-E '(binary(lifecycle_integration_test) or (package(rustfs) and test(lifecycle_transition_api_test))) and not (test(test_noncurrent_expiry_still_works_after_immediate_compensation_transition) or test(test_noncurrent_transition_still_works_after_immediate_compensation_transition))' \
-E 'binary(lifecycle_integration_test) or (package(rustfs) and test(lifecycle_transition_api_test))' \
--status-level all --final-status-level all \
2>&1 | tee artifacts/ilm-integration/nextest.log
status=${PIPESTATUS[0]}
Generated
+46 -46
View File
@@ -271,7 +271,7 @@ version = "1.1.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "40c48f72fd53cd289104fc64099abca73db4166ad86ea0b4341abe65af83dadc"
dependencies = [
"windows-sys 0.61.2",
"windows-sys 0.60.2",
]
[[package]]
@@ -282,7 +282,7 @@ checksum = "291e6a250ff86cd4a820112fb8898808a366d8f9f58ce16d1f538353ad55747d"
dependencies = [
"anstyle",
"once_cell_polyfill",
"windows-sys 0.61.2",
"windows-sys 0.60.2",
]
[[package]]
@@ -1641,9 +1641,9 @@ checksum = "bef38d45163c2f1dde094a7dfd33ccf595c92905c8f8f4fdc18d06fb1037718a"
[[package]]
name = "bitflags"
version = "2.13.1"
version = "2.13.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b588b76d00fde79687d7646a9b5bdf3cc0f655e0bbd080335a95d7e96f3587da"
checksum = "3ded4057c258ba199e2d26386d3af3780957ecaee6c4ef4041c6b4b8b97c0b06"
dependencies = [
"serde_core",
]
@@ -1906,7 +1906,7 @@ dependencies = [
"maybe-owned",
"rustix",
"rustix-linux-procfs",
"windows-sys 0.61.2",
"windows-sys 0.60.2",
"winx",
]
@@ -2270,9 +2270,9 @@ dependencies = [
[[package]]
name = "console"
version = "0.16.4"
version = "0.16.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4fe5f465a4f6fee88fad41b85d990f84c835335e85b5d9e6e63e0d06d28cba7c"
checksum = "e96a4956774c13c126a8b5af4daa79384f4d826534c95a02d76afb39e2ab64e3"
dependencies = [
"encode_unicode",
"libc",
@@ -3942,7 +3942,7 @@ dependencies = [
"libc",
"option-ext",
"redox_users",
"windows-sys 0.61.2",
"windows-sys 0.59.0",
]
[[package]]
@@ -3951,7 +3951,7 @@ version = "0.3.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1e0e367e4e7da84520dedcac1901e4da967309406d1e51017ae1abfb97adbd38"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"objc2",
]
@@ -4287,7 +4287,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb"
dependencies = [
"libc",
"windows-sys 0.61.2",
"windows-sys 0.52.0",
]
[[package]]
@@ -4416,7 +4416,7 @@ version = "25.12.19"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "35f6839d7b3b98adde531effaf34f0c2badc6f4735d26fe74709d8e513a96ef3"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"rustc_version",
]
@@ -5706,7 +5706,7 @@ version = "0.7.15"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ed3bd0ecfbb87805f538bb7b32e5239ca0763890c623e349860ecba69469f2bb"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"cfg-if",
"libc",
]
@@ -6212,7 +6212,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c9f8ff371890db2cf65a0758dba9a79f9cd965de369f6dbdc6581a22780af45e"
dependencies = [
"async-trait",
"bitflags 2.13.1",
"bitflags 2.13.2",
"bytes",
"chrono",
"dashmap",
@@ -6794,7 +6794,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0f27695f286b461da077b8c2f72f47feaa04ce3c3f9c0976257410e90e21208a"
dependencies = [
"base64 0.22.1",
"bitflags 2.13.1",
"bitflags 2.13.2",
"btoi",
"byteorder",
"bytes",
@@ -6826,7 +6826,7 @@ version = "0.7.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "22f9786d56d972959e1408b6a93be6af13b9c1392036c5c1fafa08a1b0c6ee87"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"byteorder",
"derive_builder",
"getset",
@@ -6874,7 +6874,7 @@ version = "0.29.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "71e2746dc3a24dd78b3cfcb7be93368c6de9963d30f43a6a73998a9cf4b17b46"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"cfg-if",
"cfg_aliases",
"libc",
@@ -6887,7 +6887,7 @@ version = "0.30.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "74523f3a35e05aba87a1d978330aef40f67b0304ac79c1c00b294c9830543db6"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"cfg-if",
"cfg_aliases",
"libc",
@@ -6899,7 +6899,7 @@ version = "0.31.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "cf20d2fde8ff38632c426f1165ed7436270b44f199fc55284c38276f9db47c3d"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"cfg-if",
"cfg_aliases",
"libc",
@@ -6963,7 +6963,7 @@ version = "0.50.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5"
dependencies = [
"windows-sys 0.61.2",
"windows-sys 0.59.0",
]
[[package]]
@@ -7100,7 +7100,7 @@ version = "0.13.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d164abbde0b3c03edb9edb9cb8d31a7f5b79015c692b7c771f6e0840e9106b9f"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"libloading",
"nvml-wrapper-sys",
"static_assertions",
@@ -7151,7 +7151,7 @@ version = "0.3.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2a180dd8642fa45cdb7dd721cd4c11b1cadd4929ce112ebd8b9f5803cc79d536"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"dispatch2",
"objc2",
]
@@ -7168,7 +7168,7 @@ version = "0.3.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e3e0adef53c21f888deb4fa59fc59f7eb17404926ee8a6f59f5df0fd7f9f3272"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"objc2",
]
@@ -7721,7 +7721,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e33c6dbf1a8fb7f71742cd70f5e9f0986e60b2d19dc0b28d9ca0d1323259274a"
dependencies = [
"arrayvec",
"bitflags 2.13.1",
"bitflags 2.13.2",
"thiserror 2.0.20",
"zerocopy",
"zerocopy-derive",
@@ -7777,7 +7777,7 @@ version = "0.1.8"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "575828d9d7d205188048eb1508560607a03d21eafdbba47b8cade1736c1c28e1"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"c-enum",
"perf-event-open-sys2",
]
@@ -8294,7 +8294,7 @@ checksum = "4b45fcc2344c680f5025fe57779faef368840d0bd1f42f216291f0dc4ace4744"
dependencies = [
"bit-set",
"bit-vec 0.8.0",
"bitflags 2.13.1",
"bitflags 2.13.2",
"num-traits",
"rand 0.9.5",
"rand_chacha 0.9.0",
@@ -8436,7 +8436,7 @@ version = "0.13.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e9f068eba8e7071c5f9511831b44f32c740d5adf574e990f946ddb53db2f314e"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"memchr",
"unicase",
]
@@ -8712,7 +8712,7 @@ dependencies = [
"once_cell",
"socket2",
"tracing",
"windows-sys 0.61.2",
"windows-sys 0.52.0",
]
[[package]]
@@ -8874,7 +8874,7 @@ version = "11.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "498cd0dc59d73224351ee52a95fee0f1a617a2eae0e7d9d720cc622c73a54186"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
]
[[package]]
@@ -8983,7 +8983,7 @@ version = "0.5.18"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ed2bf2547551a7053d6fdfafda3f938979645c44812fbfcda098faae3f1a362d"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
]
[[package]]
@@ -9310,7 +9310,7 @@ checksum = "036204edbd199552a5b3832f63c60dcdf395dc44c7f06b4af1c0e8139cc11bce"
dependencies = [
"aes 0.9.3",
"aws-lc-rs",
"bitflags 2.13.1",
"bitflags 2.13.2",
"block-padding 0.4.2",
"byteorder",
"bytes",
@@ -9391,7 +9391,7 @@ version = "3.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "093197e526668d92bba562e2bbbe98d1af9831bf080b619c736316ca1fa35101"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"bytes",
"chrono",
"dashmap",
@@ -11061,11 +11061,11 @@ version = "1.1.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b6fe4565b9518b83ef4f91bb47ce29620ca828bd32cb7e408f0062e9930ba190"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"errno",
"libc",
"linux-raw-sys",
"windows-sys 0.61.2",
"windows-sys 0.52.0",
]
[[package]]
@@ -11148,7 +11148,7 @@ dependencies = [
"security-framework",
"security-framework-sys",
"webpki-root-certs",
"windows-sys 0.61.2",
"windows-sys 0.52.0",
]
[[package]]
@@ -11431,7 +11431,7 @@ version = "3.7.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b7f4bc775c73d9a02cde8bf7b2ec4c9d12743edf609006c7facc23998404cd1d"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"core-foundation 0.10.1",
"core-foundation-sys",
"libc",
@@ -11949,7 +11949,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c3d1e2c7f27f8d4cb10542a02c49005dbd6e93095799d6f3be745fae9f8fedd4"
dependencies = [
"libc",
"windows-sys 0.61.2",
"windows-sys 0.60.2",
]
[[package]]
@@ -12100,7 +12100,7 @@ dependencies = [
"cfg-if",
"libc",
"psm",
"windows-sys 0.61.2",
"windows-sys 0.60.2",
]
[[package]]
@@ -12321,7 +12321,7 @@ version = "0.7.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a13f3d0daba03132c0aa9767f98351b3488edc2c100cda2d2ec2b04f3d8d3c8b"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"core-foundation 0.9.4",
"system-configuration-sys",
]
@@ -12405,10 +12405,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd"
dependencies = [
"fastrand",
"getrandom 0.4.3",
"getrandom 0.3.4",
"once_cell",
"rustix",
"windows-sys 0.61.2",
"windows-sys 0.52.0",
]
[[package]]
@@ -12876,7 +12876,7 @@ version = "0.6.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4cfcf7e2740e6fc6d4d688b4ef00650406bb94adf4731e43c096c3a19fe40840"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"bytes",
"futures-util",
"http 1.5.0",
@@ -12895,7 +12895,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "08a05a66a4fdd61cbbe0a1d755ffe0ca6aba159dd4820936a0ff8a8278245b9c"
dependencies = [
"async-compression",
"bitflags 2.13.1",
"bitflags 2.13.2",
"bytes",
"futures-core",
"futures-util",
@@ -13244,9 +13244,9 @@ checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821"
[[package]]
name = "uuid"
version = "1.26.0"
version = "1.26.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b5772d71c9be8a8a6ac2117d949c5b224c1b72241bb611d9a3012edcf8af7812"
checksum = "2ef6dac1e96601b4fb3acccccff2139741fcb757cb9a36089bf5be91cfb285ce"
dependencies = [
"getrandom 0.4.3",
"js-sys",
@@ -13510,7 +13510,7 @@ version = "0.1.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22"
dependencies = [
"windows-sys 0.61.2",
"windows-sys 0.52.0",
]
[[package]]
@@ -13820,7 +13820,7 @@ version = "0.36.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3f3fd376f71958b862e7afb20cfe5a22830e1963462f3a17f49d82a6c1d1f42d"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"windows-sys 0.59.0",
]
+1 -1
View File
@@ -335,7 +335,7 @@ tracing-subscriber = { version = "0.3.23" }
transform-stream = "0.3.1"
url = "2.5.8"
urlencoding = "2.1.3"
uuid = { version = "1.26.0" }
uuid = { version = "1.26.1" }
vaultrs = { version = "0.8.0" }
tar = "0.4.46"
walkdir = "2.5.0"
@@ -1336,7 +1336,7 @@ mod tests {
replacement_format,
expected_pool_metadata,
} = select_replacement_drive(&cluster, 1, background_enabled)?;
let default_online_object_count = if !outage_target_manifest_required { 96 } else { 24 };
let default_online_object_count = if !outage_target_manifest_required { 64 } else { 24 };
let online_object_count = std::env::var("RUSTFS_HEAL_CHAOS_OBJECT_COUNT")
.ok()
.and_then(|value| value.parse::<usize>().ok())
@@ -1823,6 +1823,7 @@ mod tests {
scenario,
InterruptionScenario::BackgroundTargetRestart
| InterruptionScenario::BackgroundTargetRestartEc84
| InterruptionScenario::BackgroundTargetRestartEc84MultiSet
| InterruptionScenario::BackgroundCoordinatorRestart
) {
cluster.stop_node_gracefully(interruption_node).await?;
@@ -2123,7 +2124,21 @@ mod tests {
evidence_context.run.binary.sha256,
"server build changed during restart"
);
let evidence = serde_json::json!({
let outage_write_diagnostic = (!outage_target_manifest_required).then(|| {
serde_json::json!({
"attempted": true,
"required": false,
"accepted": true,
"attempts": if outage_write_deferred_until_rejoin {
max_outage_write_attempts
} else {
service_unavailable_outage_writes + 1
},
"service_unavailable": service_unavailable_outage_writes,
"deferred_until_rejoin": outage_write_deferred_until_rejoin,
})
});
let mut evidence = serde_json::json!({
"schema": 1, "case": evidence_context.case.id, "evidence": evidence_context.case.evidence,
"run_id": evidence_context.run.run_id, "source_revision": evidence_context.run.source_revision,
"test_build": compiled_test_identity(),
@@ -2146,6 +2161,9 @@ mod tests {
"unclean_shutdown_marker": unclean_shutdown_marker_observed.unwrap_or(false),
"objects": evidence_objects, "node_listings": node_listings,
});
if let Some(outage_write) = outage_write_diagnostic {
evidence["outage_write"] = outage_write;
}
let data = serde_json::to_vec(&evidence)?;
if data.len() > 1024 * 1024 {
return Err("scanner/heal oracle exceeds the 1 MiB artifact budget".into());
+77 -2
View File
@@ -52,8 +52,8 @@ use aws_sdk_s3::error::ProvideErrorMetadata;
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{
BucketLifecycleConfiguration, BucketVersioningStatus, CompletedMultipartUpload, CompletedPart, ExpirationStatus,
LifecycleRule, LifecycleRuleFilter, NoncurrentVersionTransition, RestoreRequest, Transition, TransitionStorageClass,
VersioningConfiguration,
LifecycleRule, LifecycleRuleFilter, MetadataDirective, NoncurrentVersionTransition, RestoreRequest, Transition,
TransitionStorageClass, VersioningConfiguration,
};
use http::Method;
use serde::Deserialize;
@@ -963,6 +963,81 @@ async fn test_hermetic_transition_main_path() -> TestResult {
Ok(())
}
/// PUT and materialized self-copy must retain cleanup ownership of a replaced
/// transitioned null version while publishing the new bytes and metadata.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_hermetic_transition_overwrite_and_self_copy() -> TestResult {
let mut cold = RustFSTestEnvironment::new().await?;
cold.access_key = "coldtieradmin".to_string();
cold.secret_key = "coldtiersecret".to_string();
cold.start_rustfs_server_without_cleanup(vec![]).await?;
let cold_client = cold.create_s3_client();
cold_client.create_bucket().bucket(TIER_BUCKET).send().await?;
let mut hot = RustFSTestEnvironment::new().await?;
start_tier_source(&mut hot, crate::common::FAST_DATA_USAGE_SCANNER_ENV).await?;
let hot_client = hot.create_s3_client();
add_rustfs_tier(&hot, &cold).await?;
hot_client.create_bucket().bucket(SOURCE_BUCKET).send().await?;
let data = payload();
for self_copy in [false, true] {
hot_client
.put_bucket_lifecycle_configuration()
.bucket(SOURCE_BUCKET)
.lifecycle_configuration(BucketLifecycleConfiguration::builder().rules(transition_rule()?).build()?)
.send()
.await?;
put_multipart_object(&hot_client, SOURCE_BUCKET, OBJECT_KEY, &data).await?;
wait_for_transition(&hot_client, SOURCE_BUCKET, OBJECT_KEY, StdDuration::from_secs(90)).await?;
assert_eq!(cold_tier_object_count(&cold_client).await?, 1);
// Keep the replacement local so disappearance of the old remote
// object cannot be confused with another automatic transition.
hot_client.delete_bucket_lifecycle().bucket(SOURCE_BUCKET).send().await?;
let expected = if self_copy { data.clone() } else { vec![0x73; 513] };
if self_copy {
hot_client
.copy_object()
.bucket(SOURCE_BUCKET)
.key(OBJECT_KEY)
.copy_source(format!("{SOURCE_BUCKET}/{}", urlencoding::encode(OBJECT_KEY)))
.metadata_directive(MetadataDirective::Replace)
.content_type("text/plain")
.metadata("replacement", "kept")
.send()
.await?;
} else {
hot_client
.put_object()
.bucket(SOURCE_BUCKET)
.key(OBJECT_KEY)
.body(ByteStream::from(expected.clone()))
.content_type("text/plain")
.metadata("replacement", "kept")
.send()
.await?;
}
wait_for_cold_tier_empty(&cold_client, StdDuration::from_secs(90)).await?;
let current = hot_client.get_object().bucket(SOURCE_BUCKET).key(OBJECT_KEY).send().await?;
assert_eq!(current.content_type(), Some("text/plain"));
assert_eq!(current.metadata().and_then(|m| m.get("replacement")).map(String::as_str), Some("kept"));
assert!(
current
.metadata()
.is_none_or(|metadata| !metadata.contains_key(USER_META_KEY))
);
assert_eq!(current.body.collect().await?.into_bytes().as_ref(), expected.as_slice());
hot_client
.delete_object()
.bucket(SOURCE_BUCKET)
.key(OBJECT_KEY)
.send()
.await?;
}
Ok(())
}
/// Restore a transitioned object through a real RustFS remote tier.
///
/// The test covers the externally visible copy-back contract that a mock tier
@@ -27,7 +27,9 @@ mod tests {
use aws_sdk_s3::Client;
use aws_sdk_s3::error::ProvideErrorMetadata;
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{BucketVersioningStatus, CompletedMultipartUpload, CompletedPart, VersioningConfiguration};
use aws_sdk_s3::types::{
BucketVersioningStatus, CompletedMultipartUpload, CompletedPart, Tag, Tagging, VersioningConfiguration,
};
use tracing::info;
fn create_s3_client(env: &RustFSTestEnvironment) -> Client {
@@ -83,6 +85,156 @@ mod tests {
Ok(())
}
async fn assert_version_tags(client: &Client, bucket: &str, key: &str, version: Option<&str>, value: Option<&str>) {
let tags = client
.get_object_tagging()
.bucket(bucket)
.key(key)
.set_version_id(version.map(str::to_owned))
.send()
.await
.expect("GetObjectTagging must accept the exact version selector");
let expected = value
.map(|value| vec![Tag::builder().key("generation").value(value).build().expect("valid tag")])
.unwrap_or_default();
assert_eq!(tags.tag_set(), expected, "version selector: {version:?}");
}
async fn assert_null_tagging_across_versioning_changes(client: &Client, bucket: &str, key: &str) {
assert_version_tags(client, bucket, key, None, None).await;
assert_version_tags(client, bucket, key, Some("null"), None).await;
client
.put_object_tagging()
.bucket(bucket)
.key(key)
.version_id("null")
.tagging(
Tagging::builder()
.tag_set(
Tag::builder()
.key("generation")
.value("original-null")
.build()
.expect("valid tag"),
)
.build()
.expect("valid tagging"),
)
.send()
.await
.expect("tag the original null version");
enable_versioning(client, bucket)
.await
.expect("enable versioning over a null version");
let versioned = client
.put_object()
.bucket(bucket)
.key(key)
.tagging("generation=versioned")
.body(ByteStream::from_static(b"new version"))
.send()
.await
.expect("write a newer UUID version");
let version = versioned.version_id().expect("versioned PUT must return a UUID");
assert_ne!(version, "null");
assert_version_tags(client, bucket, key, None, Some("versioned")).await;
assert_version_tags(client, bucket, key, Some(version), Some("versioned")).await;
assert_version_tags(client, bucket, key, Some("null"), Some("original-null")).await;
client
.put_object_tagging()
.bucket(bucket)
.key(key)
.version_id("null")
.tagging(
Tagging::builder()
.tag_set(
Tag::builder()
.key("generation")
.value("updated-null")
.build()
.expect("valid tag"),
)
.build()
.expect("valid tagging"),
)
.send()
.await
.expect("update tags on the noncurrent null version");
assert_version_tags(client, bucket, key, Some("null"), Some("updated-null")).await;
assert_version_tags(client, bucket, key, None, Some("versioned")).await;
client
.delete_object_tagging()
.bucket(bucket)
.key(key)
.version_id("null")
.send()
.await
.expect("delete only the noncurrent null version tags");
assert_version_tags(client, bucket, key, Some("null"), None).await;
assert_version_tags(client, bucket, key, Some(version), Some("versioned")).await;
suspend_versioning(client, bucket).await.expect("suspend versioning");
client
.put_object()
.bucket(bucket)
.key(key)
.tagging("generation=suspended-null")
.body(ByteStream::from_static(b"replacement null version"))
.send()
.await
.expect("replace the null version while suspended");
assert_version_tags(client, bucket, key, None, Some("suspended-null")).await;
assert_version_tags(client, bucket, key, Some("null"), Some("suspended-null")).await;
assert_version_tags(client, bucket, key, Some(version), Some("versioned")).await;
enable_versioning(client, bucket).await.expect("re-enable versioning");
client
.put_object()
.bucket(bucket)
.key(key)
.tagging("generation=latest")
.body(ByteStream::from_static(b"latest version"))
.send()
.await
.expect("write a new latest version");
assert_version_tags(client, bucket, key, Some("null"), Some("suspended-null")).await;
assert_version_tags(client, bucket, key, None, Some("latest")).await;
client
.delete_object()
.bucket(bucket)
.key(key)
.version_id("null")
.send()
.await
.expect("remove only the null version");
assert_version_tags(client, bucket, key, None, Some("latest")).await;
let absent_version = uuid::Uuid::new_v4().to_string();
for (missing_key, selector, expected_code) in [
(key, Some("null"), "NoSuchVersion"),
(key, Some(absent_version.as_str()), "NoSuchVersion"),
("never-created", None, "NoSuchKey"),
("never-created", Some("null"), "NoSuchVersion"),
("never-created", Some(absent_version.as_str()), "NoSuchVersion"),
] {
let missing = client
.get_object_tagging()
.bucket(bucket)
.key(missing_key)
.set_version_id(selector.map(str::to_owned))
.send()
.await
.expect_err("a missing version must not fall back to latest");
assert_eq!(
missing.as_service_error().and_then(ProvideErrorMetadata::code),
Some(expected_code),
"key: {missing_key}, version selector: {selector:?}"
);
}
assert_version_tags(client, bucket, key, None, Some("latest")).await;
}
/// Test 1: PutObject should return version_id when versioning is enabled
/// This directly addresses the Veeam issue from #1066
#[tokio::test]
@@ -262,7 +414,9 @@ mod tests {
info!("🧪 TEST: PutObject behavior without versioning (no regression)");
let mut env = RustFSTestEnvironment::new().await.expect("Failed to create test environment");
env.start_rustfs_server(vec![]).await.expect("Failed to start RustFS");
env.start_rustfs_server_without_cleanup(vec![])
.await
.expect("Failed to start isolated RustFS");
let client = create_s3_client(&env);
let bucket = "test-no-versioning";
@@ -290,6 +444,9 @@ mod tests {
output.version_id().is_none() || output.version_id() == Some("null"),
"non-versioned PUT must omit version ID or return the S3 null version"
);
// Reuse this unversioned fixture to prove explicit null never becomes
// an implicit latest-version read after enable/suspend transitions.
assert_null_tagging_across_versioning_changes(&client, bucket, key).await;
info!("✅ PASSED: PutObject works correctly without versioning");
}
@@ -822,7 +822,7 @@ async fn scan_exact_free_version_targets(
let mut targets = Vec::new();
for pool in &api.pools {
for set in &pool.disk_set {
let versions = match set.load_file_info_versions_exact(&oi.bucket, &oi.name).await {
let versions = match set.load_file_info_versions_for_tier_cleanup(&oi.bucket, &oi.name).await {
Ok(Some(versions)) => versions,
Ok(None) => continue,
Err(err) if is_err_strict_volume_not_found(&err) => continue,
@@ -8100,6 +8100,75 @@ mod tests {
}
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial]
async fn tier_overwrite_cleanup_retains_a_minority_live_remote_reference() {
let (disk_paths, ecstore) = setup_test_env().await;
let bucket = format!("overwrite-minority-{}", Uuid::new_v4());
let object = "still-referenced";
create_test_bucket(&ecstore, &bucket).await;
let (backend, identity) = register_recovery_mock_tier(&ecstore).await;
seed_recoverable_free_version(&disk_paths, &bucket, object, None, Some(identity.clone())).await;
let page = list_tier_free_versions(Arc::clone(&ecstore), 100, None, None, CancellationToken::new())
.await
.expect("list persisted cleanup owner");
let owner = page.items.into_iter().find(|oi| oi.bucket == bucket).expect("seeded owner");
backend
.set_put_remote_version(Some(owner.transitioned_object.version_id.clone()))
.await;
let lease = TierConfigMgr::acquire_operation_lease(&ecstore.tier_config_mgr(), "WARM")
.await
.expect("remote fixture lease");
lease
.put(
&owner.transitioned_object.name,
rustfs_s3_client::transition_api::ReaderImpl::Body(bytes::Bytes::from_static(b"old")),
3,
)
.await
.expect("seed referenced remote bytes");
drop(lease);
let path = disk_paths[0].join(&bucket).join(object).join(STORAGE_FORMAT_FILE);
let cleanup_metadata = fs::read(&path).await.expect("save completed replica");
let mut live = FileInfo::new(object, 2, 2);
live.volume = bucket.clone();
live.erasure.index = 1;
live.data_dir = Some(Uuid::new_v4());
live.mod_time = Some(OffsetDateTime::now_utc());
live.size = 3;
live.add_object_part(1, "149603e6c03516362a8da23f624db945".to_string(), 3, live.mod_time, 3, None, None);
live.transition_status = TRANSITION_COMPLETE.to_string();
live.transition_tier = "WARM".to_string();
live.transitioned_objname = owner.transitioned_object.name.clone();
live.transition_version = Some(owner.transitioned_object.version_id.clone());
live.transition_version_state = rustfs_filemeta::TransitionVersionState::Exact;
rustfs_utils::http::insert_str(&mut live.metadata, rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID, identity);
let mut old_metadata = FileMeta::new();
old_metadata.add_version(live).expect("prepare minority live source");
fs::write(&path, old_metadata.marshal_msg().expect("encode live source"))
.await
.expect("model one replica retained by an interrupted overwrite");
let err = super::cleanup_free_version_exact(Arc::clone(&ecstore), &owner, &CancellationToken::new())
.await
.expect_err("quorum free versions cannot erase a minority live reference");
assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock);
assert_eq!(backend.remove_count().await, 0);
assert!(backend.contains(&owner.transitioned_object.name).await);
fs::write(&path, cleanup_metadata)
.await
.expect("complete replica convergence");
assert!(
super::cleanup_free_version_exact(Arc::clone(&ecstore), &owner, &CancellationToken::new())
.await
.expect("converged cleanup can delete the exact remote owner")
);
assert_eq!(backend.remove_count().await, 1);
assert!(!backend.contains(&owner.transitioned_object.name).await);
}
#[cfg(feature = "test-util")]
#[tokio::test]
#[serial]
+5 -6
View File
@@ -219,7 +219,10 @@ impl Sets {
let mut disk_set = Vec::with_capacity(set_count);
let lock_registry = runtime_sources::lock_registry();
let pool_lockers = runtime_sources::lock_registry()
.as_ref()
.map(|registry| registry.clients_for_endpoints(endpoints.endpoints.as_ref()))
.unwrap_or_default();
for i in 0..set_count {
let mut set_drive = Vec::with_capacity(set_drive_count);
@@ -270,10 +273,6 @@ impl Sets {
}
}
let lockers = lock_registry
.as_ref()
.map(|registry| registry.clients_for_endpoints(&set_endpoints))
.unwrap_or_default();
let set_disks = SetDisks::new_with_instance_ctx(
runtime_sources::local_node_name().await,
Arc::new(RwLock::new(set_drive)),
@@ -283,7 +282,7 @@ impl Sets {
pool_idx,
set_endpoints,
fm.clone(),
lockers,
pool_lockers.clone(),
instance_ctx.clone(),
)
.await;
@@ -3076,6 +3076,23 @@ impl SetDisks {
&self,
bucket: &str,
object: &str,
) -> Result<Option<rustfs_filemeta::FileInfoVersions>> {
self.load_file_info_versions_for_cleanup(bucket, object, false).await
}
pub(crate) async fn load_file_info_versions_for_tier_cleanup(
&self,
bucket: &str,
object: &str,
) -> Result<Option<rustfs_filemeta::FileInfoVersions>> {
self.load_file_info_versions_for_cleanup(bucket, object, true).await
}
async fn load_file_info_versions_for_cleanup(
&self,
bucket: &str,
object: &str,
retain_unconfirmed_tier_references: bool,
) -> Result<Option<rustfs_filemeta::FileInfoVersions>> {
let disk_object = rustfs_utils::path::encode_dir_object(object);
let disks = self.get_disks_internal().await;
@@ -3152,12 +3169,24 @@ impl SetDisks {
)));
}
let file_info_versions = FileMeta {
let mut file_info_versions = FileMeta {
versions,
..Default::default()
}
.get_all_file_info_versions(bucket, object, true)
.map_err(decode_error)?;
if retain_unconfirmed_tier_references {
// A failed overwrite may leave its live source on a minority
// of disks. Preserve that reference even if quorum merging
// selects only the replacement and its cleanup owner.
file_info_versions.versions.extend(
transition_copies
.into_values()
.flatten()
.map(|(version, _)| version)
.filter(|version| !version.tier_free_version()),
);
}
for file_info in file_info_versions
.versions
@@ -12214,6 +12243,64 @@ mod tests {
assert!(result.is_err(), "missing disks must prevent metadata write quorum");
}
#[tokio::test]
async fn tier_overwrite_cleanup_rejects_unreadable_disk_despite_metadata_quorum() {
let bucket = "tier-unreadable-disk";
let object = "object";
let mut dirs = Vec::new();
let mut disks = Vec::new();
let mut fi = metadata_test_fileinfo(object);
fi.mod_time = Some(OffsetDateTime::now_utc());
for index in 1..=3 {
let (dir, disk) = read_multiple_test_disk(bucket, &[]).await;
fi.erasure.index = index;
disk.write_metadata(bucket, bucket, object, fi.clone())
.await
.expect("seed metadata quorum");
dirs.push(dir);
disks.push(Some(disk));
}
disks.push(None);
let set = io_primitives_test_set(disks, 2).await;
assert!(
set.load_file_info_versions_exact(bucket, object).await.is_err(),
"exact reads must preserve release's unreadable-replica fence"
);
assert!(
set.load_file_info_versions_for_tier_cleanup(bucket, object).await.is_err(),
"unreadable replica may still reference the old remote object"
);
}
#[tokio::test]
async fn tier_overwrite_cleanup_rejects_minority_metadata_in_an_absent_set() {
let bucket = "tier-minority-metadata";
let object = "object";
let mut dirs = Vec::new();
let mut disks = Vec::new();
for index in 1..=4 {
let (dir, disk) = read_multiple_test_disk(bucket, &[]).await;
if index == 1 {
let mut fi = metadata_test_fileinfo(object);
fi.mod_time = Some(OffsetDateTime::now_utc());
disk.write_metadata(bucket, bucket, object, fi)
.await
.expect("seed minority metadata");
}
dirs.push(dir);
disks.push(Some(disk));
}
let set = io_primitives_test_set(disks, 2).await;
assert!(
set.load_file_info_versions_exact(bucket, object).await.is_err(),
"exact reads must preserve release's minority-ownership fence"
);
assert!(
set.load_file_info_versions_for_tier_cleanup(bucket, object).await.is_err(),
"absence on a majority cannot prove this physical set has no remote reference"
);
}
#[tokio::test]
async fn load_file_info_versions_exact_returns_versions_from_read_quorum() {
let bucket = "exact-versions-bucket";
+78 -4
View File
@@ -3852,7 +3852,7 @@ pub struct SetDisks {
pub default_parity_count: usize,
pub set_index: usize,
pub pool_index: usize,
/// Stable namespace shared by every object lock created for this set.
/// Stable namespace shared by every object lock created for this pool.
set_lock_namespace: Arc<str>,
pub format: FormatV3,
#[allow(dead_code, reason = "asserted by this file's tests (backlog#1823)")]
@@ -4491,7 +4491,7 @@ impl SetDisks {
instance_ctx: Arc<InstanceContext>,
) -> Arc<Self> {
let ctx = instance_ctx;
let set_lock_namespace: Arc<str> = format!("set-{pool_index}-{set_index}").into();
let set_lock_namespace: Arc<str> = format!("pool-{pool_index}").into();
let shared_lockers = Arc::from(lockers.to_vec());
Arc::new(SetDisks {
locker_owner,
@@ -4605,7 +4605,9 @@ impl SetDisks {
pub(crate) async fn shares_namespace_lock_domain(&self, other: &Self) -> bool {
match (self.ctx.is_dist_erasure().await, other.ctx.is_dist_erasure().await) {
(false, false) => Arc::ptr_eq(&self.local_lock_manager, &other.local_lock_manager),
(true, true) => same_distributed_lock_domain(&self.lockers, &other.lockers),
(true, true) => {
self.set_lock_namespace == other.set_lock_namespace && same_distributed_lock_domain(&self.lockers, &other.lockers)
}
_ => false,
}
}
@@ -7123,7 +7125,7 @@ mod tests {
ctx.update_erasure_type(SetupType::Erasure).await;
let set = make_test_set_disks_with_ctx(Vec::new(), ctx).await;
assert_eq!(&*set.set_lock_namespace, "set-0-0");
assert_eq!(&*set.set_lock_namespace, "pool-0");
let before = Arc::strong_count(&set.set_lock_namespace);
let lock = set
.new_ns_lock("bucket", "object")
@@ -8348,6 +8350,78 @@ mod tests {
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn test_new_ns_lock_distributed_write_succeeds_with_three_lockers_one_offline() {
let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await;
let manager_a = Arc::new(rustfs_lock::GlobalLockManager::new());
let manager_b = Arc::new(rustfs_lock::GlobalLockManager::new());
let healthy_a: Arc<dyn LockClient> = Arc::new(LocalClient::with_manager(manager_a));
let healthy_b: Arc<dyn LockClient> = Arc::new(LocalClient::with_manager(manager_b));
let failing_client: Arc<dyn LockClient> = Arc::new(FailingClient);
let set_disks = make_test_set_disks(vec![healthy_a, failing_client, healthy_b]).await;
let guard = set_disks
.new_ns_lock("bucket", "object")
.await
.expect("namespace lock should be created")
.get_write_lock(Duration::from_millis(500))
.await
.expect("two healthy lockers should satisfy the three-locker write quorum");
match guard {
NamespaceLockGuard::Standard(_) => {}
NamespaceLockGuard::Fast(_) => panic!("Expected distributed guard for dist-erasure"),
}
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn namespace_lock_domain_includes_pool_namespace() {
let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await;
let first: Arc<dyn LockClient> = Arc::new(LocalClient::with_manager(Arc::new(rustfs_lock::GlobalLockManager::new())));
let second: Arc<dyn LockClient> = Arc::new(LocalClient::with_manager(Arc::new(rustfs_lock::GlobalLockManager::new())));
let lockers = vec![first, second];
let same_pool_first_set = make_test_set_disks_with_ctx(lockers.clone(), bootstrap_ctx()).await;
let same_pool_second_set = SetDisks::new_with_instance_ctx(
"test-owner".to_string(),
Arc::new(RwLock::new(vec![None, None])),
2,
1,
1,
0,
same_pool_first_set.set_endpoints.clone(),
FormatV3::new(2, 2),
lockers.clone(),
bootstrap_ctx(),
)
.await;
let other_pool_set = SetDisks::new_with_instance_ctx(
"test-owner".to_string(),
Arc::new(RwLock::new(vec![None, None])),
2,
1,
0,
1,
same_pool_first_set.set_endpoints.clone(),
FormatV3::new(1, 2),
lockers,
bootstrap_ctx(),
)
.await;
assert!(
same_pool_first_set.shares_namespace_lock_domain(&same_pool_second_set).await,
"sets in the same pool share the object namespace lock domain"
);
assert!(
!same_pool_first_set.shares_namespace_lock_domain(&other_pool_set).await,
"different pool namespaces must not be deduplicated solely by identical clients"
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn streaming_reader_holds_read_lock_until_eof() {
+103
View File
@@ -3821,6 +3821,13 @@ impl SetDisks {
}
fi.metadata = user_defined;
if fi.version_id.is_none_or(|id| id.is_nil()) && !opts.data_movement && expected_restore_operation_id.is_none() {
// Every disk must publish the same cleanup owner alongside a
// replaced null version. This transient key is not persisted
// on the new object; recovery discovers the free-version in
// the committed xl.meta even if this request is cancelled.
fi.set_tier_free_version_id(&Uuid::new_v4().to_string());
}
fi.mod_time = mod_time;
fi.size = w_size as i64;
fi.versioned = opts.versioned || opts.version_suspended;
@@ -18174,6 +18181,102 @@ mod put_object_tmp_cleanup_tests {
drop(temp_dirs);
}
#[tokio::test]
#[serial_test::serial(capacity_dirty_scope)]
async fn tier_overwrite_failed_quorum_and_cancellation_preserve_live_source() {
for cancel_before_rename in [false, true] {
let (dirs, disks, set) = hermetic_set_disks(4).await;
let bucket = "tier-overwrite-failure";
let object = "still-live";
make_completion_test_bucket(&disks, bucket).await;
let old_body = vec![0x31; TEST_OBJECT_SIZE];
let mut metadata = HashMap::from([(
"x-amz-restore".to_string(),
"ongoing-request=\"false\", expiry-date=\"2099-01-01T00:00:00Z\"".to_string(),
)]);
for (suffix, value) in [
(rustfs_utils::http::SUFFIX_TRANSITION_STATUS, "complete".to_string()),
(rustfs_utils::http::SUFFIX_TRANSITION_TIER, "WARM".to_string()),
(rustfs_utils::http::SUFFIX_TRANSITIONED_OBJECTNAME, "remote/still-live".to_string()),
(rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID, "exact-live-version".to_string()),
(rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_STATE, "exact".to_string()),
(rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID, "ab".repeat(32)),
] {
rustfs_utils::http::insert_str(&mut metadata, suffix, value);
}
set.put_object(
bucket,
object,
&mut PutObjReader::from_vec(old_body.clone()),
&ObjectOptions {
user_defined: metadata,
write_completion: WriteCompletion::TailDrained,
..Default::default()
},
)
.await
.expect("seed live transitioned source");
wait_for_tmp_workspace_to_drain(&dirs, "seed write must drain").await;
let before = set
.load_file_info_versions_exact(bucket, object)
.await
.expect("read original metadata")
.expect("original exists");
if cancel_before_rename {
let barrier = PutObjectCommitBarrier::install(bucket, object, PutObjectCommitPause::AfterQuotaReservation);
let writer = Arc::clone(&set);
let put = tokio::spawn(async move {
writer
.put_object(
bucket,
object,
&mut PutObjReader::from_vec(vec![0x32; TEST_OBJECT_SIZE]),
&ObjectOptions::default(),
)
.await
});
barrier.wait_until_paused().await;
put.abort();
assert!(put.await.expect_err("cancel paused replacement").is_cancelled());
wait_for_tmp_workspace_to_drain(&dirs, "cancelled replacement must roll back").await;
drop(barrier);
} else {
let _fault = rename_fault_injection::fail_rename_on(object, &[2, 3]);
let err = set
.put_object(
bucket,
object,
&mut PutObjReader::from_vec(vec![0x32; TEST_OBJECT_SIZE]),
&ObjectOptions {
write_completion: WriteCompletion::TailDrained,
..Default::default()
},
)
.await
.expect_err("two disk commits cannot satisfy write quorum three");
assert!(matches!(err, Error::ErasureWriteQuorum | Error::InsufficientWriteQuorum(_, _)), "{err}");
}
let after = set
.load_file_info_versions_exact(bucket, object)
.await
.expect("read rolled-back metadata")
.expect("live source must survive");
assert_eq!(after.versions, before.versions, "failed replacement must preserve the live version");
assert_eq!(
after.free_versions, before.free_versions,
"failed replacement must not publish a cleanup owner"
);
let mut reader = set
.get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
.await
.expect("live source remains readable");
let mut actual = Vec::new();
reader.stream.read_to_end(&mut actual).await.expect("read original bytes");
assert_eq!(actual, old_body);
}
}
#[tokio::test]
async fn cooperative_cancellation_while_waiting_for_namespace_lock_cleans_tmp_workspace() {
let (temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
+218
View File
@@ -2832,6 +2832,224 @@ mod tests {
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn tier_overwrite_put_and_self_copy_recover_persisted_cleanup_owners() {
use crate::bucket::lifecycle::bucket_lifecycle_ops::ExpiryState;
use crate::bucket::lifecycle::tier_free_version_recovery::recover_tier_free_versions;
use rustfs_filemeta::TransitionVersionState::{Exact, KnownDisabled, SuspendedNull};
use rustfs_s3_client::transition_api::ReaderImpl;
use rustfs_utils::http::{
SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID,
SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, insert_str,
};
let temp_dir = tempfile::tempdir().expect("create tier overwrite store");
let (mut ctx, mut store, mut shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-overwrite", &[4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
let tier = "OVERWRITE-TIER";
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier).await;
let lease = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier)
.await
.expect("tier identity");
let identity = rustfs_utils::crypto::hex(lease.backend_identity());
drop(lease);
for state in [Exact, KnownDisabled, SuspendedNull] {
for suspended in [false, true] {
for self_copy in [false, true] {
let bucket = format!("tier-overwrite-{}", Uuid::new_v4());
let object = "object";
let remote = format!("remote/{bucket}");
let version = match state {
Exact => "opaque-overwrite-version",
SuspendedNull => "null",
_ => "",
};
let payload = vec![0x5b; if suspended { 512 * 1024 } else { 257 }];
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create bucket");
backend.set_put_remote_version(Some(version.to_string())).await;
let lease = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier)
.await
.expect("seed tier lease");
lease
.put(
&remote,
ReaderImpl::Body(bytes::Bytes::from(payload.clone())),
payload.len().try_into().expect("payload size"),
)
.await
.expect("seed remote bytes");
drop(lease);
let mut metadata = HashMap::from([
("content-type".to_string(), "application/octet-stream".to_string()),
(
"x-amz-restore".to_string(),
"ongoing-request=\"false\", expiry-date=\"2099-01-01T00:00:00Z\"".to_string(),
),
]);
for (suffix, value) in [
(SUFFIX_TRANSITION_STATUS, "complete"),
(SUFFIX_TRANSITION_TIER, tier),
(SUFFIX_TRANSITION_TIER_DESTINATION_ID, identity.as_str()),
(SUFFIX_TRANSITIONED_OBJECTNAME, remote.as_str()),
(SUFFIX_TRANSITIONED_VERSION_STATE, state.as_str()),
] {
insert_str(&mut metadata, suffix, value.to_string());
}
if !version.is_empty() {
insert_str(&mut metadata, SUFFIX_TRANSITIONED_VERSION_ID, version.to_string());
}
let options = ObjectOptions {
version_suspended: suspended,
..Default::default()
};
store
.put_object(
&bucket,
object,
&mut PutObjReader::from_vec(payload.clone()),
&ObjectOptions {
user_defined: metadata,
..options.clone()
},
)
.await
.expect("seed transitioned source with locally restored bytes");
let expected = if self_copy {
payload.clone()
} else {
vec![0x73; payload.len()]
};
let new_metadata = HashMap::from([
("content-type".to_string(), "text/plain".to_string()),
("x-amz-meta-replacement".to_string(), "kept".to_string()),
]);
if self_copy {
let mut source = store
.get_object_info(&bucket, object, &options)
.await
.expect("self-copy source");
source.metadata_only = false;
source.user_defined = Arc::new(new_metadata);
source.put_object_reader = Some(PutObjReader::from_vec(expected.clone()));
store
.copy_object(&bucket, object, &bucket, object, &mut source, &options, &options)
.await
.expect("materialized self-copy");
} else {
store
.put_object(
&bucket,
object,
&mut PutObjReader::from_vec(expected.clone()),
&ObjectOptions {
user_defined: new_metadata,
..options.clone()
},
)
.await
.expect("overwrite transitioned null version");
}
let set = store.pools[0].get_disks_by_key(object);
let versions = set
.load_file_info_versions_exact(&bucket, object)
.await
.expect("read committed disk metadata")
.expect("replacement metadata exists");
let free: Vec<_> = versions
.versions
.iter()
.chain(versions.free_versions.iter())
.filter(|fi| fi.tier_free_version())
.collect();
assert_eq!(free.len(), 1, "{state:?}, suspended={suspended}, copy={self_copy}");
assert_eq!(free[0].transitioned_objname, remote);
assert_eq!(free[0].transition_version_state, state);
assert!(backend.contains(&remote).await, "commit must not delete remote bytes before cleanup");
let removed_before = backend.remove_count().await;
// Restart before queue delivery. The new runtime must
// reconstruct ownership solely from the committed xl.meta.
let tier_config = ctx
.tier_config_mgr()
.read()
.await
.tiers
.get(tier)
.expect("tier configuration survives restart")
.clone_with_credentials();
drop(set);
shutdown.cancel();
drop(store);
drop(ctx);
(ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-overwrite-restart", &[4]))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
{
let manager = ctx.tier_config_mgr();
let mut manager = manager.write().await;
manager.tiers.insert(tier.to_string(), tier_config);
manager
.install_test_driver(tier, Box::new(backend.clone()))
.expect("rebind the same remote destination after restart");
}
let set = store.pools[0].get_disks_by_key(object);
ExpiryState::resize_workers(1, Arc::clone(&store)).await;
let recovered = recover_tier_free_versions(Arc::clone(&store), 100, None, None)
.await
.expect("recover persisted cleanup owner");
assert!(recovered.enqueued >= 1);
tokio::time::timeout(Duration::from_secs(30), async {
loop {
let versions = set
.load_file_info_versions_exact(&bucket, object)
.await
.expect("read cleanup progress")
.expect("new object must survive cleanup");
if versions
.versions
.iter()
.chain(versions.free_versions.iter())
.all(|fi| !fi.tier_free_version())
{
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("cleanup must converge");
assert!(!backend.contains(&remote).await);
assert_eq!(backend.remove_count().await, removed_before + 1, "one remote DELETE per owner");
assert_eq!(backend.remove_versions().await.last(), Some(&(remote.clone(), version.to_string())));
let mut reader = store
.get_object_reader(&bucket, object, None, HeaderMap::new(), &options)
.await
.expect("replacement remains readable");
let mut actual = Vec::new();
reader.stream.read_to_end(&mut actual).await.expect("read replacement bytes");
assert_eq!(actual, expected);
let current = store
.get_object_info(&bucket, object, &options)
.await
.expect("replacement metadata");
assert_eq!(current.user_defined.get("content-type").map(String::as_str), Some("text/plain"));
assert_eq!(current.user_defined.get("x-amz-meta-replacement").map(String::as_str), Some("kept"));
assert!(current.transitioned_object.status.is_empty());
}
}
}
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
+38 -18
View File
@@ -50,7 +50,7 @@ use crate::services::notification_sys::{
use crate::services::tier::tier::{TierConfigMgr, TierDestinationId, TierOperationLease, tier_destination_id_from_metadata};
use crate::set_disk::{
SetDisks, get_lock_acquire_timeout, get_object_lock_diag_slow_acquire_threshold, get_object_lock_diag_slow_hold_threshold,
is_lock_optimization_enabled, is_object_lock_diag_enabled, same_distributed_lock_domain,
is_lock_optimization_enabled, is_object_lock_diag_enabled,
};
use crate::storage_api_contracts::{
list::ListOperations as _,
@@ -3440,10 +3440,15 @@ impl ECStore {
for pool in &self.pools {
let hashed_set = pool.get_disks_by_key(object);
let lock_domain_already_held = !distributed
|| locked_sets
.iter()
.any(|locked_set| same_distributed_lock_domain(&locked_set.lockers, &hashed_set.lockers));
let mut lock_domain_already_held = !distributed;
if !lock_domain_already_held {
for locked_set in &locked_sets {
if locked_set.shares_namespace_lock_domain(&hashed_set).await {
lock_domain_already_held = true;
break;
}
}
}
if lock_domain_already_held {
continue;
}
@@ -3500,10 +3505,15 @@ impl ECStore {
let mut locked_sets = vec![fixed_set];
for pool in &self.pools {
for set in &pool.disk_set {
let lock_domain_already_held = !distributed
|| locked_sets
.iter()
.any(|locked_set| same_distributed_lock_domain(&locked_set.lockers, &set.lockers));
let mut lock_domain_already_held = !distributed;
if !lock_domain_already_held {
for locked_set in &locked_sets {
if locked_set.shares_namespace_lock_domain(set).await {
lock_domain_already_held = true;
break;
}
}
}
if lock_domain_already_held {
continue;
}
@@ -3581,10 +3591,15 @@ impl ECStore {
let mut locked_sets = vec![fixed_set];
for pool in &self.pools {
for set in &pool.disk_set {
let lock_domain_already_held = !distributed
|| locked_sets
.iter()
.any(|locked_set| same_distributed_lock_domain(&locked_set.lockers, &set.lockers));
let mut lock_domain_already_held = !distributed;
if !lock_domain_already_held {
for locked_set in &locked_sets {
if locked_set.shares_namespace_lock_domain(set).await {
lock_domain_already_held = true;
break;
}
}
}
if lock_domain_already_held {
continue;
}
@@ -3667,11 +3682,15 @@ impl ECStore {
.get(pool_idx)
.ok_or_else(|| Error::other(format!("invalid data movement publication pool {pool_idx}")))?;
let set = pool.get_disks_by_key(object);
let lock_domain_already_held = !locked_sets.is_empty()
&& (!distributed
|| locked_sets.iter().any(|locked_set: &Arc<crate::set_disk::SetDisks>| {
same_distributed_lock_domain(&locked_set.lockers, &set.lockers)
}));
let mut lock_domain_already_held = !locked_sets.is_empty() && !distributed;
if !lock_domain_already_held {
for locked_set in &locked_sets {
if locked_set.shares_namespace_lock_domain(&set).await {
lock_domain_already_held = true;
break;
}
}
}
if lock_domain_already_held {
continue;
}
@@ -5717,6 +5736,7 @@ mod tests {
GetObjectBodyCacheHook, GetObjectBodyCacheHookLookup, GetObjectBodySource, clear_get_object_body_cache_hook,
lookup_get_object_body_cache_hook, register_get_object_body_cache_hook,
};
use crate::set_disk::same_distributed_lock_domain;
use crate::set_disk::{SetDisks, disk_call_counters};
use crate::storage_api_contracts::bucket::MakeBucketOptions;
use crate::storage_api_contracts::lifecycle::TransitionedObject;
+333 -1
View File
@@ -480,7 +480,110 @@ impl FileMeta {
Ok(())
}
pub fn add_version(&mut self, mut fi: FileInfo) -> Result<()> {
pub fn add_version(&mut self, fi: FileInfo) -> Result<()> {
if let Some(free_version) = self.overwritten_tier_free_version(&fi)? {
// The replacement and its cleanup owner must share one xl.meta
// commit. Keep the original intact if either insertion fails.
let mut next = self.clone();
next.add_version_inner(fi)?;
next.add_version_filemata(free_version)?;
*self = next;
return Ok(());
}
self.add_version_inner(fi)
}
fn overwritten_tier_free_version(&self, fi: &FileInfo) -> Result<Option<FileMetaVersion>> {
use rustfs_utils::http::{
SUFFIX_TIER_FV_ID, SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID,
SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE,
get_consistent_bytes, get_consistent_str, has_internal_suffix, strip_internal_prefix_preserving_case,
};
if fi.version_id.is_some_and(|id| !id.is_nil()) || !contains_key_str(&fi.metadata, SUFFIX_TIER_FV_ID) {
return Ok(None);
}
let Some(existing) = self
.versions
.iter()
.find(|v| v.header.version_id.is_none_or(|id| id.is_nil()))
else {
return Ok(None);
};
let old = existing.parse_version_meta()?;
let Some(mut object) = old.object else {
return Ok(None);
};
let status = get_consistent_bytes(&object.meta_sys, SUFFIX_TRANSITION_STATUS);
if status.is_none()
&& object
.meta_sys
.keys()
.any(|key| has_internal_suffix(key, SUFFIX_TRANSITION_STATUS))
{
// Empty status is a valid local object. The reader distinguishes
// it from conflicting aliases before the ordinary overwrite.
object.into_fileinfo(&fi.volume, &fi.name, false)?;
return Ok(None);
}
if status != Some(TRANSITION_COMPLETE.as_bytes()) {
return Ok(None);
}
// Reuse the reader's alias/state validation. A legacy empty remote
// version is valid and must not be mistaken for conflicting aliases.
object.into_fileinfo(&fi.volume, &fi.name, false)?;
if object
.meta_sys
.keys()
.any(|key| has_internal_suffix(key, SUFFIX_TRANSITION_TIER_DESTINATION_ID))
&& get_consistent_bytes(&object.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID).is_none()
{
return Err(Error::FileCorrupt);
}
let transition_suffixes = [
SUFFIX_TRANSITION_STATUS,
SUFFIX_TRANSITION_TIER,
SUFFIX_TRANSITION_TIER_DESTINATION_ID,
SUFFIX_TRANSITIONED_OBJECTNAME,
SUFFIX_TRANSITIONED_VERSION_ID,
SUFFIX_TRANSITIONED_VERSION_STATE,
];
let replacement = MetaObject::from(fi.clone());
if transition_suffixes
.iter()
.all(|suffix| get_consistent_bytes(&object.meta_sys, suffix) == get_consistent_bytes(&replacement.meta_sys, suffix))
{
return Ok(None);
}
let id = get_consistent_str(&fi.metadata, SUFFIX_TIER_FV_ID).ok_or(Error::FileCorrupt)?;
let id = Uuid::parse_str(id)?;
if id.is_nil() || self.versions.iter().any(|version| version.header.version_id == Some(id)) {
return Err(Error::FileCorrupt);
}
// The reader also accepts legacy key casing. Canonicalize only this
// cleanup source so init_free_version preserves every accepted field,
// including an explicitly empty unversioned remote version.
for suffix in transition_suffixes {
let value = object
.meta_sys
.iter()
.find(|(key, _)| {
strip_internal_prefix_preserving_case(key).is_some_and(|found| found.eq_ignore_ascii_case(suffix))
})
.map(|(_, value)| value.clone());
if let Some(value) = value {
rustfs_utils::http::insert_bytes(&mut object.meta_sys, suffix, value);
}
}
let (free_version, created) = object.init_free_version(fi)?;
if !created {
return Err(Error::FileCorrupt);
}
Ok(Some(free_version))
}
fn add_version_inner(&mut self, mut fi: FileInfo) -> Result<()> {
rustfs_utils::http::remove_str(&mut fi.metadata, rustfs_utils::http::SUFFIX_TIER_FV_ID);
// empty version_id means "null" (versioning disabled/suspended)
if fi.version_id.is_none() {
fi.version_id = Some(Uuid::nil());
@@ -1463,6 +1566,235 @@ mod test {
});
}
fn tier_overwrite_fixture(state: crate::TransitionVersionState) -> (FileMeta, FileInfo) {
let mut source = FileInfo::new("object", 2, 2);
source.erasure.index = 1;
source.mod_time = Some(OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("fixture timestamp"));
source.data_dir = Some(Uuid::new_v4());
source.transition_status = TRANSITION_COMPLETE.to_string();
source.transition_tier = "WARM".to_string();
source.transitioned_objname = "remote/old-object".to_string();
source.transition_version_state = state;
source.transition_version = match state {
crate::TransitionVersionState::Exact => Some("opaque-provider-version".to_string()),
crate::TransitionVersionState::SuspendedNull => Some("null".to_string()),
_ => None,
};
rustfs_utils::http::insert_str(
&mut source.metadata,
rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
"ab".repeat(32),
);
if state == crate::TransitionVersionState::KnownDisabled {
rustfs_utils::http::insert_str(
&mut source.metadata,
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID,
String::new(),
);
}
let mut meta = FileMeta::new();
meta.add_version(source.clone()).expect("seed transitioned null version");
(meta, source)
}
#[test]
fn tier_overwrite_preserves_exact_cleanup_owner_across_reload() {
use crate::TransitionVersionState::{Exact, KnownDisabled, SuspendedNull, Unknown};
use rustfs_utils::http::{MINIO_INTERNAL_PREFIX, RUSTFS_INTERNAL_PREFIX, SUFFIX_TIER_FV_ID};
for state in [Exact, KnownDisabled, SuspendedNull, Unknown] {
for inline in [false, true] {
let (mut meta, _) = tier_overwrite_fixture(state);
let old = meta.versions[0].parse_version_meta().expect("old metadata");
let old = old.object.expect("old object");
let id = Uuid::new_v4();
let mut replacement = FileInfo::new("object", 2, 2);
replacement.version_id = inline.then_some(Uuid::nil());
replacement.mod_time = Some(OffsetDateTime::from_unix_timestamp(1_700_000_001).expect("fixture timestamp"));
replacement.data_dir = Some(Uuid::new_v4());
replacement.size = 3;
if inline {
replacement.data = Some(Bytes::from_static(b"new"));
}
replacement.set_tier_free_version_id(&id.to_string());
meta.add_version(replacement.clone())
.expect("replace transitioned null version");
let bytes = meta.marshal_msg().expect("persist replacement and cleanup owner");
let mut reopened = FileMeta::load(&bytes).expect("reopen committed metadata");
assert_eq!(reopened.versions.len(), 2);
let (_, free) = reopened.find_version(Some(id)).expect("durable cleanup owner");
assert!(free.free_version());
let marker = free.delete_marker.expect("cleanup marker");
for suffix in [
rustfs_utils::http::SUFFIX_TRANSITION_TIER,
rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
rustfs_utils::http::SUFFIX_TRANSITIONED_OBJECTNAME,
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID,
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_STATE,
] {
for prefix in [RUSTFS_INTERNAL_PREFIX, MINIO_INTERNAL_PREFIX] {
let key = format!("{prefix}{suffix}");
assert_eq!(marker.meta_sys.get(&key), old.meta_sys.get(&key), "{state:?}: {key}");
}
}
let (_, current) = reopened.find_version(None).expect("replacement survives restart");
let current = current.object.expect("replacement object");
assert_eq!(current.size, 3);
assert!(!rustfs_utils::http::contains_key_bytes(&current.meta_sys, SUFFIX_TIER_FV_ID));
reopened
.add_version(replacement)
.expect("replaying replacement is idempotent");
assert_eq!(reopened.versions.len(), 2);
}
}
}
#[test]
fn tier_overwrite_rejects_cleanup_failure_without_mutating_source() {
for id in ["not-a-uuid".to_string(), Uuid::nil().to_string()] {
let (mut meta, _) = tier_overwrite_fixture(crate::TransitionVersionState::Exact);
let before = meta.clone();
let mut replacement = FileInfo::new("object", 2, 2);
replacement.mod_time = Some(OffsetDateTime::now_utc());
replacement.set_tier_free_version_id(&id);
assert!(meta.add_version(replacement).is_err());
assert_eq!(meta, before, "invalid cleanup identity must preserve source");
}
with_object_max_versions_for_test(1, || {
let (mut meta, _) = tier_overwrite_fixture(crate::TransitionVersionState::KnownDisabled);
let before = meta.clone();
let mut replacement = FileInfo::new("object", 2, 2);
replacement.mod_time = Some(OffsetDateTime::now_utc());
replacement.data = Some(Bytes::from_static(b"new"));
replacement.set_tier_free_version_id(&Uuid::new_v4().to_string());
assert_eq!(
meta.add_version(replacement)
.expect_err("cleanup owner exceeds version limit"),
Error::MaxVersionsExceeded
);
assert_eq!(meta, before, "failed cleanup insertion must also preserve inline bytes");
});
}
#[test]
fn tier_overwrite_allows_empty_transition_status_on_local_source() {
let mut source = FileInfo::new("object", 2, 2);
source.mod_time = Some(OffsetDateTime::now_utc());
source.data = Some(Bytes::from_static(b"old"));
rustfs_utils::http::insert_str(&mut source.metadata, rustfs_utils::http::SUFFIX_TRANSITION_STATUS, String::new());
let mut meta = FileMeta::new();
meta.add_version(source)
.expect("seed readable local metadata with empty status");
let mut replacement = FileInfo::new("object", 2, 2);
replacement.mod_time = Some(OffsetDateTime::now_utc());
replacement.size = 3;
replacement.data = Some(Bytes::from_static(b"new"));
replacement.set_tier_free_version_id(&Uuid::new_v4().to_string());
meta.add_version(replacement)
.expect("an ordinary overwrite must still succeed");
assert_eq!(meta.versions.len(), 1);
let (_, current) = meta.find_version(None).expect("replacement remains visible");
assert_eq!(current.object.expect("ordinary object").size, 3);
assert!(!meta.versions[0].header.free_version());
}
#[test]
fn tier_overwrite_preserves_legacy_metadata_casing() {
use rustfs_utils::http::{MINIO_INTERNAL_PREFIX, RUSTFS_INTERNAL_PREFIX};
for state in [
crate::TransitionVersionState::Exact,
crate::TransitionVersionState::KnownDisabled,
] {
let (mut meta, _) = tier_overwrite_fixture(state);
let mut source = meta.versions[0].parse_version_meta().expect("seeded source");
let object = source.object.as_mut().expect("transitioned source");
let expected = object.meta_sys.clone();
object.meta_sys = object
.meta_sys
.drain()
.map(|(key, value)| (key.to_ascii_uppercase(), value))
.collect();
meta.versions[0] = FileMetaShallowVersion::try_from(source).expect("legacy key casing");
let id = Uuid::new_v4();
let mut replacement = FileInfo::new("object", 2, 2);
replacement.mod_time = Some(OffsetDateTime::now_utc());
replacement.set_tier_free_version_id(&id.to_string());
meta.add_version(replacement).expect("overwrite readable legacy source");
let reopened = FileMeta::load(&meta.marshal_msg().expect("persist overwrite")).expect("reopen overwrite");
let (_, owner) = reopened
.find_version(Some(id))
.expect("legacy source must retain cleanup ownership");
let marker = owner.delete_marker.expect("cleanup marker");
for suffix in [
rustfs_utils::http::SUFFIX_TRANSITION_TIER,
rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
rustfs_utils::http::SUFFIX_TRANSITIONED_OBJECTNAME,
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID,
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_STATE,
] {
for prefix in [RUSTFS_INTERNAL_PREFIX, MINIO_INTERNAL_PREFIX] {
let key = format!("{prefix}{suffix}");
assert_eq!(marker.meta_sys.get(&key), expected.get(&key), "legacy {state:?}: {key}");
}
}
}
}
#[test]
fn tier_overwrite_rejects_conflicting_remote_metadata_aliases() {
use rustfs_utils::http::{
MINIO_INTERNAL_PREFIX, SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID,
SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE,
};
for suffix in [
SUFFIX_TRANSITION_STATUS,
SUFFIX_TRANSITION_TIER,
SUFFIX_TRANSITION_TIER_DESTINATION_ID,
SUFFIX_TRANSITIONED_OBJECTNAME,
SUFFIX_TRANSITIONED_VERSION_ID,
SUFFIX_TRANSITIONED_VERSION_STATE,
] {
let (mut meta, _) = tier_overwrite_fixture(crate::TransitionVersionState::Exact);
let mut old = meta.versions[0].parse_version_meta().expect("seeded source metadata");
old.object
.as_mut()
.expect("transitioned source")
.meta_sys
.insert(format!("{MINIO_INTERNAL_PREFIX}{suffix}"), b"conflicting-value".to_vec());
meta.versions[0] = FileMetaShallowVersion::try_from(old).expect("encode conflicting aliases");
let before = meta.clone();
let mut replacement = FileInfo::new("object", 2, 2);
replacement.mod_time = Some(OffsetDateTime::now_utc());
replacement.set_tier_free_version_id(&Uuid::new_v4().to_string());
assert_eq!(
meta.add_version(replacement)
.expect_err("ambiguous ownership must fail closed"),
Error::FileCorrupt
);
assert_eq!(meta, before, "conflicting {suffix} must not erase the old remote tuple");
}
}
#[test]
fn tier_overwrite_keeps_retained_remote_and_versioned_copy_ownership() {
let (mut meta, mut source) = tier_overwrite_fixture(crate::TransitionVersionState::Exact);
source.set_tier_free_version_id(&Uuid::new_v4().to_string());
meta.add_version(source.clone())
.expect("restore retains the same remote owner");
assert_eq!(meta.versions.len(), 1);
source.version_id = Some(Uuid::new_v4());
source.transition_status.clear();
source.transition_tier.clear();
source.transitioned_objname.clear();
source.transition_version = None;
source.transition_version_state = crate::TransitionVersionState::Unknown;
meta.add_version(source).expect("versioned write retains historical source");
assert_eq!(meta.versions.len(), 2);
assert!(meta.versions.iter().all(|version| !version.header.free_version()));
}
#[test]
fn add_version_filemata_uses_canonical_equal_time_order() {
let mod_time = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("valid test timestamp");
@@ -1691,7 +1691,7 @@ mod serial_tests {
#[tokio::test(flavor = "multi_thread", worker_threads = 1)]
#[serial]
#[ignore = "FAILING on main: excluded from the serial ILM lane pending a fix, see rustfs/backlog#1148 (ilm-1 partial)"]
#[ignore = "global-state ILM integration test: runs serialized in the CI ILM Integration (serial) lane, see ci.yml test-ilm-integration-serial and rustfs/backlog#1148 (ilm-1)"]
async fn test_noncurrent_expiry_still_works_after_immediate_compensation_transition() {
let (disk_paths, ecstore) = setup_isolated_test_env(true).await;
@@ -1775,7 +1775,7 @@ mod serial_tests {
#[tokio::test(flavor = "multi_thread", worker_threads = 1)]
#[serial]
#[ignore = "FAILING on main: excluded from the serial ILM lane pending a fix, see rustfs/backlog#1148 (ilm-1 partial)"]
#[ignore = "global-state ILM integration test: runs serialized in the CI ILM Integration (serial) lane, see ci.yml test-ilm-integration-serial and rustfs/backlog#1148 (ilm-1)"]
async fn test_noncurrent_transition_still_works_after_immediate_compensation_transition() {
let (disk_paths, ecstore) = setup_isolated_test_env(true).await;
@@ -30,11 +30,23 @@ These are approved-target invariants. A protocol's explicitly labeled current ex
| Remote PUT is in flight or its response is unknown | Transition transaction | Only cleanup of its own canonical candidate, subject to the transaction recovery predicate | Durable transaction identity plus a known remote-version state; the approved target also requires expiry and durable takeover of the creator fence |
| Local transition commit is complete | Exact transitioned version in `xl.meta` | No | Recovery finds the transaction's logical bucket/object/version and requires the complete recorded source identity (version ID, data directory, modification time, size, and ETag), `TRANSITION_COMPLETE`, and the same remote object, tier, and remote version before removing only the transaction record |
| An ordinary delete removes that transitioned version | Hidden `xl.meta` free-version | Yes | Metadata quorum atomically removes the visible version and preserves its exact tier tuple in the free-version |
| PUT or materialized self-copy replaces a transitioned null version | Hidden `xl.meta` free-version | Yes, after replacement commit and complete physical reference checks | The coordinator supplies one cleanup UUID to every disk; replacement metadata and the old remote tuple are written in the same `xl.meta` commit |
| A recursive prefix/delete-all operation cannot preserve per-object markers | v6 journal bound to an immutable single dispatch manifest or a chunk-parent-bound child manifest | Yes, but only after child/manifest completion and all-pool absence proof | `DispatchAuthorized`, exact local destructive mutation, every journal `Committed`, then child/manifest `Completed`; a chunk parent advances only after that child completion |
| Tier configuration mutation, manual job, or decommission receipt | Intent/admission/copy proof only | No | These records gate configuration, scheduling, or migration; they never become remote-object cleanup owners |
An old journal and a free-version can coexist during compatibility recovery. That coexistence is evidence of multiple possible owners, not permission to choose one: the journal path must retain its record until the version-specific recovery rule proves which owner is authoritative.
Null-version replacement uses the existing free-version format and recovery
worker. Failed metadata preparation preserves both the old version and its
inline bytes; rename rollback restores the complete previous metadata. Recovery
must retain cleanup while any physical replica still references the remote tuple,
including a minority version omitted by quorum merging, or any disk cannot be
checked. A successful replacement needs no in-memory queue receipt to survive
restart: the normal free-version sweep discovers its committed owner. Restores
that retain the same remote tuple and ordinary versioned writes retain their
existing ownership. Older binaries can read this format, but all writers and
cleanup workers need the overwrite fix before these guarantees cover the fleet.
## Persisted record inventory
All keys below are objects in the internal metadata bucket. The table gives the canonical target form. Transition-transaction runtime recovery extracts the final 32 hexadecimal characters and UUID while ignoring shard directories and accepting uppercase hex. Manual-job runtime recovery requires exactly two shards matching the filename prefix, but accepts an uppercase UUID when the shards use the same uppercase text; it then loads the lowercase canonical job by UUID. The decommission validator recomputes and rejects a noncanonical manual-job path, but currently inherits the weaker transition parser. Exact runtime canonical-path validation for both protocols is an approved target.
+146 -10
View File
@@ -28,6 +28,7 @@ use futures_util::future::join_all;
use http::{HeaderMap, HeaderValue, Uri};
use hyper::{Method, StatusCode};
use matchit::Params;
use percent_encoding::percent_decode_str;
use rustfs_config::MAX_HEAL_REQUEST_SIZE;
use rustfs_heal::heal::utils::format_set_disk_id;
use rustfs_heal_contracts::heal_channel::{
@@ -35,13 +36,11 @@ use rustfs_heal_contracts::heal_channel::{
};
use rustfs_policy::policy::action::{Action, AdminAction};
use rustfs_scanner::scanner::{BackgroundHealInfo, read_background_heal_info};
use rustfs_utils::path::path_join;
use s3s::header::{CONTENT_LENGTH, CONTENT_TYPE};
use s3s::{Body, S3Request, S3Response, S3Result, s3_error};
use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, BTreeSet, HashSet};
use std::future::Future;
use std::path::PathBuf;
use std::sync::Arc;
use time::{OffsetDateTime, format_description::well_known::Rfc3339};
use tokio::time::{Duration, timeout};
@@ -71,9 +70,17 @@ struct HealInitParams {
}
fn extract_heal_init_params(body: &Bytes, uri: &Uri, params: Params<'_, '_>) -> S3Result<HealInitParams> {
// matchit captures the original URI bytes. Decode once before validation
// so literal %2F keys remain distinct from actual path separators.
let mut hip = HealInitParams {
bucket: params.get("bucket").map(|s| s.to_string()).unwrap_or_default(),
obj_prefix: params.get("prefix").map(|s| s.to_string()).unwrap_or_default(),
bucket: percent_decode_str(params.get("bucket").unwrap_or_default())
.decode_utf8()
.map_err(|_| s3_error!(InvalidRequest, "invalid bucket name encoding"))?
.into_owned(),
obj_prefix: percent_decode_str(params.get("prefix").unwrap_or_default())
.decode_utf8()
.map_err(|_| s3_error!(InvalidRequest, "invalid object name encoding"))?
.into_owned(),
..Default::default()
};
validate_heal_target(&hip.bucket, &hip.obj_prefix)?;
@@ -164,13 +171,13 @@ fn validate_heal_target(bucket: &str, obj_prefix: &str) -> S3Result<()> {
}
fn encode_heal_control_path(bucket: &str, obj_prefix: &str) -> String {
if bucket.is_empty() && obj_prefix.is_empty() {
return String::new();
if obj_prefix.is_empty() {
return bucket.to_owned();
}
path_join(&[PathBuf::from(bucket), PathBuf::from(obj_prefix)])
.to_string_lossy()
.into_owned()
// This identifies an S3 target, not a filesystem path. In particular,
// a leading slash in the object must not alias the sibling without it.
format!("{bucket}/{obj_prefix}")
}
fn heal_control_response_id(heal_path: &str, client_token: &str) -> String {
@@ -200,7 +207,7 @@ pub fn register_heal_route(r: &mut S3Router<AdminOperation>) -> std::io::Result<
r.insert(
Method::POST,
format!("{}{}", ADMIN_PREFIX, "/v3/heal/{bucket}/{prefix}").as_str(),
format!("{}{}", ADMIN_PREFIX, "/v3/heal/{bucket}/{*prefix}").as_str(),
AdminOperation(&HealHandler {}),
)?;
@@ -1616,6 +1623,135 @@ mod tests {
use tokio::sync::mpsc;
use tokio::time::Duration;
fn parse_registered_heal_request(uri: &Uri) -> s3s::S3Result<HealInitParams> {
let mut registered = super::S3Router::new(false);
super::register_heal_route(&mut registered).expect("register production Heal routes");
let mut router = Router::new();
for route in registered.registered_routes() {
router.insert(route.clone(), ()).expect("replay production route");
}
let path = format!("POST|{}", uri.path());
let matched = router.at(&path).expect("request must match a production Heal route");
let body = Bytes::from_static(
br#"{"recursive":false,"dryRun":true,"remove":false,"recreate":false,"scanMode":2,"updateParity":false,"nolock":false,"readRepair":false,"pool":0,"set":0}"#,
);
extract_heal_init_params(&body, uri, matched.params)
}
#[test]
fn test_heal_routes_accept_nested_and_encoded_object_paths() {
let mut router = super::S3Router::new(false);
super::register_heal_route(&mut router).expect("register production Heal routes");
for prefix in ["/rustfs/admin", "/minio/admin"] {
for target in [
"",
"test-bucket",
"test-bucket/object.bin",
"test-bucket/dir/sub/object.bin",
"test-bucket/dir%2Fobject.bin",
] {
let path = format!("{prefix}/v3/heal/{target}");
assert!(router.contains_compatible_route(http::Method::POST, &path), "{path}");
assert!(!router.contains_compatible_route(http::Method::GET, &path), "{path}");
}
}
}
#[test]
fn test_heal_target_decodes_once_and_keeps_start_status_stop_identity() {
for (wire, object) in [
("object.bin", "object.bin"),
("dir/sub/object.bin", "dir/sub/object.bin"),
("dir%2Fsub%2Fobject.bin", "dir/sub/object.bin"),
("dir%2fsub/object.bin", "dir/sub/object.bin"),
("%2Fobject.bin", "/object.bin"),
("dir/", "dir/"),
("dir%2F", "dir/"),
("literal%252Fslash", "literal%2Fslash"),
("space%20key%2Bplus", "space key+plus"),
("literal+plus", "literal+plus"),
("%E4%B8%AD%E6%96%87%2F%E6%96%87%E4%BB%B6", "中文/文件"),
("query%3Fhash%23percent%25", "query?hash#percent%"),
] {
for query in ["", "?clientToken=task", "?clientToken=task&forceStop=true"] {
let uri = format!("/rustfs/admin/v3/heal/test%2Dbucket/{wire}{query}")
.parse()
.expect("valid encoded URI");
let parsed = parse_registered_heal_request(&uri).expect("valid Heal target");
assert_eq!(parsed.bucket, "test-bucket");
assert_eq!(parsed.obj_prefix, object, "wire target: {wire}");
assert_eq!(
encode_heal_control_path(&parsed.bucket, &parsed.obj_prefix),
format!("test-bucket/{object}")
);
assert_eq!(parsed.client_token, if query.is_empty() { "" } else { "task" });
assert_eq!(parsed.force_stop, query.ends_with("forceStop=true"));
if query.is_empty() {
let request = build_heal_channel_request(&parsed);
assert_eq!(request.bucket, "test-bucket");
assert_eq!(request.object_prefix.as_deref(), Some(object));
assert_eq!(request.pool_index, Some(0));
assert_eq!(request.set_index, Some(0));
assert_eq!(request.dry_run, Some(true));
assert_eq!(request.scan_mode, Some(HealScanMode::Deep));
}
}
}
}
#[test]
fn test_heal_target_validates_decoded_paths_before_admission() {
for target in [
"test%2Fbucket/object",
"test%00bucket/object",
"test%FFbucket/object",
"test-bucket/dir%2F..%2Fobject",
"test-bucket/dir/%2e/object",
"test-bucket/dir%5C..%5Cobject",
"test-bucket/dir%2F%2Fobject",
"test-bucket/object%00",
"test-bucket/object%FF",
] {
let uri = format!("/rustfs/admin/v3/heal/{target}").parse().expect("encoded URI");
let err = parse_registered_heal_request(&uri).expect_err("decoded invalid target must fail closed");
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest, "target: {target}");
}
}
#[tokio::test]
async fn test_nested_heal_routes_still_require_authentication() {
use s3s::route::S3Route;
let mut router = super::S3Router::new(false);
super::register_heal_route(&mut router).expect("register production Heal routes");
for prefix in ["/rustfs/admin", "/minio/admin"] {
for object in ["dir/object.bin", "dir%2Fobject.bin", "literal%252Fslash"] {
let mut req = s3s::S3Request {
input: s3s::Body::empty(),
method: http::Method::POST,
uri: format!("{prefix}/v3/heal/test-bucket/{object}").parse().expect("Heal URI"),
headers: http::HeaderMap::new(),
extensions: http::Extensions::new(),
credentials: None,
region: None,
service: None,
trailing_headers: None,
};
let err = router
.check_access(&mut req)
.await
.expect_err("router must require a signature");
assert_eq!(err.code(), &S3ErrorCode::AccessDenied);
let err = router
.call(req)
.await
.expect_err("handler must independently require authentication");
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest);
assert!(err.to_string().contains("authentication required"));
}
}
}
fn replacement_record(task_id: &str) -> rustfs_heal::ReplacementRecoveryRecord {
rustfs_heal::ReplacementRecoveryRecord {
task_id: task_id.to_string(),
+1 -1
View File
@@ -198,7 +198,7 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
admin_route(Method::POST, "/v3/rebalance/stop"),
admin_route(Method::POST, "/v3/heal/"),
admin_route_sample(Method::POST, "/v3/heal/{bucket}", "/v3/heal/test-bucket"),
admin_route_sample(Method::POST, "/v3/heal/{bucket}/{prefix}", "/v3/heal/test-bucket/prefix"),
admin_route_sample(Method::POST, "/v3/heal/{bucket}/{*prefix}", "/v3/heal/test-bucket/prefix"),
admin_route(Method::POST, "/v3/background-heal/status"),
admin_route(Method::GET, "/v4/heal/replacement-recovery"),
admin_route(Method::GET, "/v3/tier"),
+9 -1
View File
@@ -377,6 +377,10 @@ impl FS {
pub(crate) fn parse_object_version_id(version_id: Option<String>) -> S3Result<Option<Uuid>> {
if let Some(vid) = version_id {
if vid == "null" {
// A nil UUID selects the stored null version; None selects latest.
return Ok(Some(Uuid::nil()));
}
let uuid = Uuid::parse_str(&vid).map_err(|e| {
error!("Invalid version ID: {}", e);
s3_error!(InvalidArgument, "Invalid version ID")
@@ -1183,7 +1187,11 @@ impl S3 for FS {
error = %e,
"Object tags not found"
);
return Err(s3_error!(NoSuchKey));
return Err(if opts.version_id.is_some() {
s3_error!(NoSuchVersion)
} else {
s3_error!(NoSuchKey)
});
}
error!(
component = LOG_COMPONENT_STORAGE,
+33 -1
View File
@@ -17,7 +17,9 @@ mod tests {
use crate::config::WorkloadProfile;
use crate::server::cors;
use crate::storage::StorageError;
use crate::storage::ecfs::{FS, propagate_object_lock_peer_reload, validate_object_lock_configuration_input};
use crate::storage::ecfs::{
FS, parse_object_version_id, propagate_object_lock_peer_reload, validate_object_lock_configuration_input,
};
use crate::storage::ecfs_extend::{apply_bucket_default_lock_retention, map_bucket_object_lock_config_state};
use crate::storage::s3_api::common::{rustfs_initiator, rustfs_owner};
use crate::storage::storage_api::ecstore_bucket::metadata_sys::ObjectLockConfigState;
@@ -641,6 +643,36 @@ mod tests {
assert_eq!(metadata.get("content-type"), Some(&"application/octet-stream".to_string()));
}
#[test]
fn test_tagging_version_id_preserves_explicit_null_and_latest_selection() {
assert_eq!(parse_object_version_id(None).expect("latest version selector"), None);
assert_eq!(
parse_object_version_id(Some("null".to_owned())).expect("explicit null version selector"),
Some(uuid::Uuid::nil())
);
for version in [uuid::Uuid::nil(), uuid::Uuid::new_v4()] {
assert_eq!(
parse_object_version_id(Some(version.to_string())).expect("UUID version selector"),
Some(version)
);
}
}
#[test]
fn test_tagging_version_id_rejects_invalid_values_instead_of_selecting_latest() {
for version in [
"",
"NULL",
" null ",
"not-a-version",
"null/other",
"00000000-0000-0000-0000-00000000000g",
] {
let err = parse_object_version_id(Some(version.to_owned())).expect_err("invalid version must fail closed");
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument, "version: {version:?}");
}
}
#[tokio::test]
async fn test_get_object_tagging_returns_internal_error_when_store_uninitialized() {
if !store_uninitialized_premise_holds() {
+55 -6
View File
@@ -112,7 +112,7 @@ apply_runtime_profile() {
export RUSTFS_HEAL_CHAOS_PARTIAL_TIMEOUT_SECS="${RUSTFS_HEAL_CHAOS_PARTIAL_TIMEOUT_SECS:-180}"
;;
background-ec8-4-multi-pool)
export RUSTFS_HEAL_CHAOS_OBJECT_COUNT="${RUSTFS_HEAL_CHAOS_OBJECT_COUNT:-96}"
export RUSTFS_HEAL_CHAOS_OBJECT_COUNT="${RUSTFS_HEAL_CHAOS_OBJECT_COUNT:-64}"
export RUSTFS_HEAL_CHAOS_OBJECT_SIZE_BYTES="${RUSTFS_HEAL_CHAOS_OBJECT_SIZE_BYTES:-4194304}"
export RUSTFS_HEAL_CHAOS_PARTIAL_TIMEOUT_SECS="${RUSTFS_HEAL_CHAOS_PARTIAL_TIMEOUT_SECS:-240}"
;;
@@ -165,6 +165,40 @@ if not status.get("pending_gates"):
PY
}
nextest_junit_candidates() {
local target_dir="$1"
local profile="$2"
local primary="$target_dir/nextest/$profile/junit.xml"
local fallback="$ROOT/target/nextest/$profile/junit.xml"
printf '%s\n' "$primary"
if [[ "$fallback" != "$primary" ]]; then
printf '%s\n' "$fallback"
fi
}
remove_nextest_junit_candidates() {
local target_dir="$1"
local profile="$2"
local candidate
while IFS= read -r candidate; do
rm -f "$candidate"
done < <(nextest_junit_candidates "$target_dir" "$profile")
}
copy_nextest_junit() {
local target_dir="$1"
local profile="$2"
local run_dir="$3"
local candidate
while IFS= read -r candidate; do
if [[ -f "$candidate" ]]; then
cp "$candidate" "$run_dir/junit.xml"
return 0
fi
done < <(nextest_junit_candidates "$target_dir" "$profile")
return 1
}
run_self_test() {
if "$0" --case release --plan-only >/dev/null 2>&1; then
echo "self-test failed: release pseudo-case must not be runnable" >&2
@@ -186,6 +220,24 @@ run_self_test() {
return 1
fi
done < <(case_ids)
local junit_profile="scanner-heal-junit-self-test-$$"
local junit_run_dir="$ROOT/target/scanner-heal-junit-self-test-$$"
local junit_target_dir="$ROOT/target/scanner-heal-junit-custom-target-$$"
local junit_fallback="$ROOT/target/nextest/$junit_profile/junit.xml"
mkdir -p "$(dirname "$junit_fallback")" "$junit_run_dir"
printf '<testsuites />\n' >"$junit_fallback"
if ! copy_nextest_junit "$junit_target_dir" "$junit_profile" "$junit_run_dir"; then
echo "self-test failed: nextest junit fallback was not copied" >&2
rm -rf "$junit_run_dir" "$junit_target_dir" "$(dirname "$junit_fallback")"
return 1
fi
if ! cmp -s "$junit_fallback" "$junit_run_dir/junit.xml"; then
echo "self-test failed: copied nextest junit fallback changed content" >&2
rm -rf "$junit_run_dir" "$junit_target_dir" "$(dirname "$junit_fallback")"
return 1
fi
rm -rf "$junit_run_dir" "$junit_target_dir" "$(dirname "$junit_fallback")"
}
while [[ $# -gt 0 ]]; do
@@ -300,8 +352,7 @@ export RUSTFS_E2E_LOG_DIR="${RUSTFS_E2E_LOG_DIR:-$RUN_DIR/e2e-logs}"
export RUSTFS_HEAL_CHAOS_LOG_DIR="${RUSTFS_HEAL_CHAOS_LOG_DIR:-$RUSTFS_E2E_LOG_DIR}"
mkdir -p "$RUSTFS_E2E_LOG_DIR"
JUNIT_PATH="$TARGET_DIR/nextest/$PROFILE/junit.xml"
rm -f "$JUNIT_PATH"
remove_nextest_junit_candidates "$TARGET_DIR" "$PROFILE"
set +e
NO_PROXY="${NO_PROXY:-127.0.0.1,localhost}" \
HTTP_PROXY= \
@@ -311,9 +362,7 @@ cargo nextest run --profile "$PROFILE" -p e2e_test -E "$TEST_FILTER" --no-tests=
STATUS=$?
set -e
if [[ -f "$JUNIT_PATH" ]]; then
cp "$JUNIT_PATH" "$RUN_DIR/junit.xml"
fi
copy_nextest_junit "$TARGET_DIR" "$PROFILE" "$RUN_DIR" || true
"$PYTHON_BIN" "$ROOT/scripts/check_test_wiring.py" --finish-scanner-heal "$RUN_DIR" "$STATUS"
if [[ "$STATUS" -ne 0 ]]; then