mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-10 14:16:01 +00:00
Compare commits
9 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 7b6e45d372 | |||
| 09c85a5f29 | |||
| 00aeb12914 | |||
| ffb18979f8 | |||
| 97c7b451d2 | |||
| 9ba95cac37 | |||
| cec796328c | |||
| 1aa4fa145f | |||
| d286f3d06c |
@@ -1,2 +1,2 @@
|
||||
sha256-darwin=874c881d7b45f12378a5817c7f42c95c4981960a2ec9ce12dcf4af239ae1f9d5
|
||||
sha256-linux=9351e25b45bf7dfce18b951a5e3740225f457cacc53b8bf9f500f6947763ec0e
|
||||
sha256-darwin=15cb0cf9909bfbfc5a835fb08bd3675516c641db1b92ecc1e170a1cfa0fa2fd5
|
||||
sha256-linux=3163fdd29df5def86cf511ca7db05880d5c0caaa608a7031a71394327f7c217a
|
||||
|
||||
@@ -1 +1 @@
|
||||
sha256=6d18f9cce820c51d5589de944e8cc185f73eeca0ea9a9916651943e3759169d0
|
||||
sha256=5fbb230b89212b7c3d7229d6cef3e7e2d16f0ecfec62237ebc770785706f67d9
|
||||
|
||||
@@ -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
@@ -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
@@ -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());
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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";
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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)]
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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(¤t.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.
|
||||
|
||||
@@ -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(),
|
||||
|
||||
@@ -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"),
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user