mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-10 14:16:01 +00:00
Compare commits
172 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 7b6e45d372 | |||
| 09c85a5f29 | |||
| 00aeb12914 | |||
| ffb18979f8 | |||
| 97c7b451d2 | |||
| 9ba95cac37 | |||
| cec796328c | |||
| 1aa4fa145f | |||
| d286f3d06c | |||
| 3391528025 | |||
| 0d1490da3c | |||
| 6421a7ff60 | |||
| 27141f6312 | |||
| 940e221988 | |||
| 8b66e9b62a | |||
| 6a879be1fb | |||
| 0b40d47a8a | |||
| e292637ee0 | |||
| 8b6b1e53a3 | |||
| 358bf6832d | |||
| 70fb8bf504 | |||
| 110f630a5c | |||
| 54bd3c8e24 | |||
| 1643cb29f1 | |||
| bf8d5a32b6 | |||
| 230eeb5fb5 | |||
| 2c6f5f22c0 | |||
| d7b6a8c10d | |||
| e8ffd72575 | |||
| 0b03d535c3 | |||
| 3afb389247 | |||
| 86b6569071 | |||
| 0940fbe1b2 | |||
| 9e7c5dc932 | |||
| 45806bf295 | |||
| 42f00d8532 | |||
| 3e46d61a91 | |||
| 88183daf10 | |||
| 8be1e9b2c1 | |||
| 9b4b366209 | |||
| 43436ad5a7 | |||
| 0c7b6188b2 | |||
| 3ffc3704af | |||
| de3ac27a6a | |||
| 2d0e82e9f7 | |||
| c92fabea3a | |||
| b3b0d892bb | |||
| d9ae1654b6 | |||
| 15949180c9 | |||
| f65305e97f | |||
| 28ac3ab4e7 | |||
| cf41b1d3cf | |||
| ca876930b3 | |||
| c05b3ee01e | |||
| 1076e2fb4f | |||
| 51e6733a9c | |||
| ed2b2cdd19 | |||
| 084477e079 | |||
| bb1b5dea16 | |||
| 1c87413788 | |||
| e1fee0569a | |||
| a6e5bdcbd2 | |||
| 923f51cf12 | |||
| e16274ab14 | |||
| 5eecd657e7 | |||
| 929d906131 | |||
| 6b802a91eb | |||
| 14b3cffe8f | |||
| c62384e58d | |||
| 3f30f6c139 | |||
| 584a52ce3d | |||
| 76e025a50c | |||
| b40eb193df | |||
| 023c987397 | |||
| a4b265bf77 | |||
| c176220a34 | |||
| 82ac2dff05 | |||
| e24eae9eaa | |||
| 8617f2701b | |||
| 8d5592f83e | |||
| b5c9229e77 | |||
| a974e50b1d | |||
| 081910e825 | |||
| 7f7e4fe40b | |||
| 684e6f313a | |||
| b3a02f305e | |||
| 0ea9c19569 | |||
| ae662da256 | |||
| b5d33a1f4e | |||
| d288496752 | |||
| 40baedba4f | |||
| 0c3fd48c22 | |||
| 3e8b0e7c83 | |||
| dd2af9ff88 | |||
| 5db0e14d04 | |||
| 9429225cf4 | |||
| 857584b3c2 | |||
| cfda9b451c | |||
| 354e49de4c | |||
| 8e987ce0a6 | |||
| 929f836e40 | |||
| c7dec044eb | |||
| 14a59f7770 | |||
| 084e9c4b82 | |||
| 921cc7ad94 | |||
| bfa8df00e7 | |||
| 33a8469e27 | |||
| 128080ceb9 | |||
| ae9fe62fb1 | |||
| 929a9f27c4 | |||
| 15376a7fa6 | |||
| 5d789006dd | |||
| f4bd53a63c | |||
| 21fee7f7ff | |||
| 23ed00fc85 | |||
| f49ffa2fec | |||
| ee5f287324 | |||
| 158c613d0e | |||
| a950f914ca | |||
| 1de9b40a30 | |||
| 831468b2e9 | |||
| 600d037b25 | |||
| 3fa2b334be | |||
| 1b549d5907 | |||
| 11c4ce96eb | |||
| 4508a0985d | |||
| ed1b9f25d6 | |||
| 274c2bf402 | |||
| 9aebcefa9c | |||
| 1748814bbf | |||
| ee5f76c180 | |||
| 4d68c32b75 | |||
| 8c6689ff13 | |||
| 595f9f662d | |||
| 4a2b15cb82 | |||
| 2ba7f95547 | |||
| 7b3dad6bae | |||
| c0754f5b1c | |||
| 5dc0e3b402 | |||
| cc5487e7de | |||
| 0b05b6c6ff | |||
| c507da8f75 | |||
| 550dabeffd | |||
| 4d7f0344d3 | |||
| 084338add6 | |||
| 8d339da706 | |||
| 817ad0a682 | |||
| 80321e5bb4 | |||
| e7475cfa4d | |||
| 8c15025a5a | |||
| 6f2ec66263 | |||
| 62542ddc57 | |||
| a8bb53218d | |||
| 5355d9f8f8 | |||
| 50d7a049ee | |||
| 3149c87cf2 | |||
| b7fa6a4615 | |||
| 6fe83f87a4 | |||
| df6981d88e | |||
| 590adad5ae | |||
| be5d14c985 | |||
| 703086d71d | |||
| a159f312f0 | |||
| 2f6f095298 | |||
| efd8ef005f | |||
| 60fa33773a | |||
| ee752b0b03 | |||
| 99c1f4418b | |||
| 7ce0ac72cf | |||
| 944e26d432 | |||
| f0b0a99260 | |||
| 33ddc10ffd |
@@ -1,2 +1,2 @@
|
||||
sha256-linux=775825dcb2b4997c4fa24bd9ba9c0546316c503f4d5e369c39abff0678c95e8e
|
||||
sha256-darwin=775825dcb2b4997c4fa24bd9ba9c0546316c503f4d5e369c39abff0678c95e8e
|
||||
sha256-linux=563bff8f1171d6dbe166ff8440310dbe98430e466aa3ecd8dc39e3c872b320f7
|
||||
sha256-darwin=563bff8f1171d6dbe166ff8440310dbe98430e466aa3ecd8dc39e3c872b320f7
|
||||
|
||||
@@ -1,2 +1,2 @@
|
||||
sha256-darwin=874c881d7b45f12378a5817c7f42c95c4981960a2ec9ce12dcf4af239ae1f9d5
|
||||
sha256-linux=9351e25b45bf7dfce18b951a5e3740225f457cacc53b8bf9f500f6947763ec0e
|
||||
sha256-darwin=15cb0cf9909bfbfc5a835fb08bd3675516c641db1b92ecc1e170a1cfa0fa2fd5
|
||||
sha256-linux=3163fdd29df5def86cf511ca7db05880d5c0caaa608a7031a71394327f7c217a
|
||||
|
||||
@@ -1 +1 @@
|
||||
sha256=6d18f9cce820c51d5589de944e8cc185f73eeca0ea9a9916651943e3759169d0
|
||||
sha256=5fbb230b89212b7c3d7229d6cef3e7e2d16f0ecfec62237ebc770785706f67d9
|
||||
|
||||
@@ -259,6 +259,9 @@ jobs:
|
||||
} > artifacts/test-and-lint/doctest-diagnostics.txt
|
||||
exit "${status}"
|
||||
|
||||
- name: Check offline enrollment E2E root boundary
|
||||
run: ./scripts/check_offline_enrollment_e2e.sh
|
||||
|
||||
- name: Upload test reports and diagnostics
|
||||
if: always()
|
||||
uses: actions/upload-artifact@b7c566a772e6b6bfb58ed0dc250532a479d7789f # v6
|
||||
@@ -293,36 +296,6 @@ jobs:
|
||||
- name: Run rebalance/decommission migration proofs
|
||||
run: ./scripts/check_migration_gate_count.sh
|
||||
|
||||
# The root boundary requires fresh CLI and integration-test builds. Give it
|
||||
# its own time budget instead of sharing the workspace lint/test budget.
|
||||
offline-enrollment-root-boundary:
|
||||
name: Offline Enrollment Root Boundary
|
||||
if: needs.classify-changes.outputs.mode == 'full' && (github.event_name != 'pull_request' || github.event.action != 'closed')
|
||||
needs: [ quick-checks, classify-changes ]
|
||||
runs-on: sm-standard-4
|
||||
timeout-minutes: 90
|
||||
env:
|
||||
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
|
||||
steps:
|
||||
- name: Checkout repository
|
||||
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
|
||||
with:
|
||||
persist-credentials: false
|
||||
|
||||
- name: Setup Rust environment
|
||||
uses: ./.github/actions/setup
|
||||
with:
|
||||
rust-version: stable
|
||||
cache-shared-key: ci-dev
|
||||
cache-save-if: 'false'
|
||||
install-build-packaging-tools: 'false'
|
||||
|
||||
- name: Protect Connect test home
|
||||
run: chmod go-w "$(realpath "$HOME")"
|
||||
|
||||
- name: Check offline enrollment E2E root boundary
|
||||
run: ./scripts/check_offline_enrollment_e2e.sh
|
||||
|
||||
# Dedicated serial lane for the ILM / lifecycle integration tests. These tests
|
||||
# drive the object layer through process-global singletons (the GLOBAL_ENV
|
||||
# ECStore, the global tier-config manager, background-expiry workers) and bind
|
||||
@@ -362,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
|
||||
@@ -379,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]}
|
||||
@@ -997,7 +966,6 @@ jobs:
|
||||
# debug binary; each test spawns its own rustfs server on a random port.
|
||||
- name: Run e2e full suite
|
||||
env:
|
||||
CARGO_BIN_EXE_rustfs: ${{ runner.temp }}/rustfs-startup-cas-input/rustfs
|
||||
RUSTFS_E2E_STARTUP_CAS_BINARY: ${{ runner.temp }}/rustfs-startup-cas-input/rustfs
|
||||
RUSTFS_E2E_STARTUP_CAS_BUILD_MANIFEST: ${{ runner.temp }}/rustfs-startup-cas-input/rustfs.e2e-startup-cas-build.json
|
||||
RUSTFS_E2E_STARTUP_CAS_ARTIFACT_DIR: ${{ runner.temp }}/rustfs-startup-cas-evidence
|
||||
@@ -1217,7 +1185,6 @@ jobs:
|
||||
- typos
|
||||
- quick-checks
|
||||
- test-and-lint
|
||||
- offline-enrollment-root-boundary
|
||||
- test-ilm-integration-serial
|
||||
- test-and-lint-rio-v2
|
||||
- connect-short-credential-boundary
|
||||
@@ -1251,7 +1218,6 @@ jobs:
|
||||
- typos
|
||||
- quick-checks
|
||||
- test-and-lint
|
||||
- offline-enrollment-root-boundary
|
||||
- test-ilm-integration-serial
|
||||
- test-and-lint-rio-v2
|
||||
- test-and-lint-protocols
|
||||
|
||||
@@ -156,7 +156,6 @@ jobs:
|
||||
|
||||
- name: Run cluster fault e2e nightly suite
|
||||
env:
|
||||
CARGO_BIN_EXE_rustfs: ${{ github.workspace }}/target/debug/rustfs
|
||||
RUSTFS_E2E_LOG_DIR: ${{ runner.temp }}/rustfs-e2e-nightly-logs
|
||||
run: cargo nextest run --profile e2e-nightly -p e2e_test
|
||||
|
||||
|
||||
Generated
+51
-51
@@ -3986,7 +3986,7 @@ checksum = "d0881ea181b1df73ff77ffaaf9c7544ecc11e82fba9b5f27b262a3c73a332555"
|
||||
|
||||
[[package]]
|
||||
name = "e2e_test"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"astral-tokio-tar",
|
||||
@@ -9476,7 +9476,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"aes-gcm",
|
||||
"anyhow",
|
||||
@@ -9624,7 +9624,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-audit"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"const-str",
|
||||
"futures",
|
||||
@@ -9646,7 +9646,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-checksums"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"base64-simd",
|
||||
"bytes",
|
||||
@@ -9662,7 +9662,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-common"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"hotpath",
|
||||
"metrics",
|
||||
@@ -9675,7 +9675,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-concurrency"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"hotpath",
|
||||
"insta",
|
||||
@@ -9688,7 +9688,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-config"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"const-str",
|
||||
"hotpath",
|
||||
@@ -9698,7 +9698,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-credentials"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"base64-simd",
|
||||
"hmac 0.13.0",
|
||||
@@ -9712,7 +9712,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-crypto"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"aes-gcm",
|
||||
"argon2",
|
||||
@@ -9733,7 +9733,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-data-usage"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"hotpath",
|
||||
"rmp-serde",
|
||||
@@ -9743,7 +9743,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-ecstore"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"async-channel",
|
||||
@@ -9879,7 +9879,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-extension-schema"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"hotpath",
|
||||
"serde",
|
||||
@@ -9889,7 +9889,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-filemeta"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"byteorder",
|
||||
@@ -9917,7 +9917,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-heal"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"base64-simd",
|
||||
@@ -9954,7 +9954,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-heal-contracts"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"serde",
|
||||
"serde_json",
|
||||
@@ -9964,7 +9964,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-iam"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"async-trait",
|
||||
@@ -10013,7 +10013,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-io-core"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"bytes",
|
||||
"hotpath",
|
||||
@@ -10025,7 +10025,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-io-metrics"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"criterion",
|
||||
"hotpath",
|
||||
@@ -10089,7 +10089,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-keystone"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"bytes",
|
||||
"futures",
|
||||
@@ -10116,7 +10116,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-kms"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"aes-gcm",
|
||||
"anyhow",
|
||||
@@ -10166,14 +10166,14 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-license"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"thiserror 2.0.20",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-lifecycle"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"hotpath",
|
||||
@@ -10195,7 +10195,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-lock"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"compact_str",
|
||||
@@ -10218,7 +10218,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-log-analyzer"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"chrono",
|
||||
"flate2",
|
||||
@@ -10237,7 +10237,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-madmin"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"hotpath",
|
||||
"http 1.5.0",
|
||||
@@ -10275,7 +10275,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-notify"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"async-trait",
|
||||
@@ -10310,7 +10310,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-object-capacity"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"criterion",
|
||||
"futures",
|
||||
@@ -10329,7 +10329,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-object-data-cache"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"bytes",
|
||||
"criterion",
|
||||
@@ -10346,7 +10346,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-obs"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"chrono",
|
||||
"crossbeam-channel",
|
||||
@@ -10404,7 +10404,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-policy"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"base64-simd",
|
||||
@@ -10435,7 +10435,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-protocols"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"astral-tokio-tar",
|
||||
"async-compression",
|
||||
@@ -10497,7 +10497,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-protos"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"flatbuffers",
|
||||
"hotpath",
|
||||
@@ -10522,7 +10522,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-replication"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"byteorder",
|
||||
"bytes",
|
||||
@@ -10540,7 +10540,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-rio"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"aes-gcm",
|
||||
"arc-swap",
|
||||
@@ -10584,7 +10584,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-rio-v2"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"aes-gcm",
|
||||
"bytes",
|
||||
@@ -10607,7 +10607,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-s3-client"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"base64-simd",
|
||||
"bytes",
|
||||
@@ -10651,7 +10651,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-s3-ops"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"hotpath",
|
||||
"rustfs-s3-types",
|
||||
@@ -10659,7 +10659,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-s3-types"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"hotpath",
|
||||
"serde",
|
||||
@@ -10668,7 +10668,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-s3select-api"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"async-compression",
|
||||
@@ -10703,7 +10703,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-s3select-query"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"async-recursion",
|
||||
"async-trait",
|
||||
@@ -10724,7 +10724,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-scanner"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"bytes",
|
||||
@@ -10770,7 +10770,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-scanner-metrics"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"chrono",
|
||||
"jiff",
|
||||
@@ -10785,7 +10785,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-security-governance"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"hotpath",
|
||||
"thiserror 2.0.20",
|
||||
@@ -10793,7 +10793,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-signer"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"base64-simd",
|
||||
"bytes",
|
||||
@@ -10811,7 +10811,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-storage-api"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"hotpath",
|
||||
@@ -10826,7 +10826,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-targets"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"async-nats",
|
||||
@@ -10880,7 +10880,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-test-utils"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"hotpath",
|
||||
"rustfs-data-usage",
|
||||
@@ -10896,7 +10896,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-tls-runtime"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"hotpath",
|
||||
@@ -10917,7 +10917,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-trusted-proxies"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"async-trait",
|
||||
"axum",
|
||||
@@ -10954,7 +10954,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-utils"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"base64-simd",
|
||||
"blake2",
|
||||
@@ -10996,7 +10996,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rustfs-zip"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
dependencies = [
|
||||
"astral-tokio-tar",
|
||||
"async-compression",
|
||||
|
||||
+51
-51
@@ -73,7 +73,7 @@ edition = "2024"
|
||||
license = "Apache-2.0"
|
||||
repository = "https://github.com/rustfs/rustfs"
|
||||
rust-version = "1.98.0"
|
||||
version = "1.0.0-rc.6"
|
||||
version = "1.0.0-rc.5"
|
||||
homepage = "https://rustfs.com"
|
||||
description = "RustFS is a high-performance distributed object storage software built using Rust, one of the most popular languages worldwide. "
|
||||
keywords = ["RustFS", "Minio", "object-storage", "filesystem", "s3"]
|
||||
@@ -90,56 +90,56 @@ redundant_clone = "warn"
|
||||
|
||||
[workspace.dependencies]
|
||||
# RustFS Internal Crates
|
||||
rustfs = { path = "./rustfs", version = "1.0.0-rc.6" }
|
||||
rustfs-heal = { path = "crates/heal", version = "1.0.0-rc.6" }
|
||||
rustfs-heal-contracts = { path = "crates/heal-contracts", version = "1.0.0-rc.6" }
|
||||
rustfs-scanner-metrics = { path = "crates/scanner-metrics", version = "1.0.0-rc.6" }
|
||||
rustfs-audit = { path = "crates/audit", version = "1.0.0-rc.6" }
|
||||
rustfs-checksums = { path = "crates/checksums", version = "1.0.0-rc.6" }
|
||||
rustfs-common = { path = "crates/common", version = "1.0.0-rc.6" }
|
||||
rustfs-data-usage = { path = "crates/data-usage", version = "1.0.0-rc.6" }
|
||||
rustfs-config = { path = "./crates/config", version = "1.0.0-rc.6" }
|
||||
rustfs-concurrency = { path = "./crates/concurrency", version = "1.0.0-rc.6" }
|
||||
rustfs-credentials = { path = "crates/credentials", version = "1.0.0-rc.6" }
|
||||
rustfs-crypto = { path = "crates/crypto", version = "1.0.0-rc.6" }
|
||||
rustfs-ecstore = { path = "crates/ecstore", version = "1.0.0-rc.6" }
|
||||
rustfs-filemeta = { path = "crates/filemeta", version = "1.0.0-rc.6" }
|
||||
rustfs-iam = { path = "crates/iam", version = "1.0.0-rc.6" }
|
||||
rustfs-keystone = { path = "crates/keystone", version = "1.0.0-rc.6" }
|
||||
rustfs-license = { path = "crates/license", version = "1.0.0-rc.6" }
|
||||
rustfs-lifecycle = { path = "crates/lifecycle", version = "1.0.0-rc.6" }
|
||||
rustfs-kms = { path = "crates/kms", version = "1.0.0-rc.6" }
|
||||
rustfs-lock = { path = "crates/lock", version = "1.0.0-rc.6" }
|
||||
rustfs-madmin = { path = "crates/madmin", version = "1.0.0-rc.6" }
|
||||
rustfs-notify = { path = "crates/notify", version = "1.0.0-rc.6" }
|
||||
rustfs-io-metrics = { path = "crates/io-metrics", version = "1.0.0-rc.6" }
|
||||
rustfs-io-core = { path = "crates/io-core", version = "1.0.0-rc.6" }
|
||||
rustfs-object-capacity = { path = "crates/object-capacity", version = "1.0.0-rc.6" }
|
||||
rustfs-object-data-cache = { path = "crates/object-data-cache", version = "1.0.0-rc.6", default-features = false }
|
||||
rustfs-log-analyzer = { path = "crates/log-analyzer", version = "1.0.0-rc.6" }
|
||||
rustfs-obs = { path = "crates/obs", version = "1.0.0-rc.6" }
|
||||
rustfs-policy = { path = "crates/policy", version = "1.0.0-rc.6" }
|
||||
rustfs-protos = { path = "crates/protos", version = "1.0.0-rc.6" }
|
||||
rustfs-protocols = { path = "crates/protocols", version = "1.0.0-rc.6" }
|
||||
rustfs-replication = { path = "crates/replication", version = "1.0.0-rc.6" }
|
||||
rustfs-rio = { path = "crates/rio", version = "1.0.0-rc.6" }
|
||||
rustfs-rio-v2 = { path = "crates/rio-v2", version = "1.0.0-rc.6" }
|
||||
rustfs-s3-client = { path = "crates/s3-client", version = "1.0.0-rc.6" }
|
||||
rustfs-s3-types = { path = "crates/s3-types", version = "1.0.0-rc.6" }
|
||||
rustfs-s3-ops = { path = "crates/s3-ops", version = "1.0.0-rc.6" }
|
||||
rustfs-s3select-api = { path = "crates/s3select-api", version = "1.0.0-rc.6" }
|
||||
rustfs-s3select-query = { path = "crates/s3select-query", version = "1.0.0-rc.6" }
|
||||
rustfs-scanner = { path = "crates/scanner", version = "1.0.0-rc.6" }
|
||||
rustfs-security-governance = { path = "crates/security-governance", version = "1.0.0-rc.6" }
|
||||
rustfs-extension-schema = { path = "crates/extension-schema", version = "1.0.0-rc.6" }
|
||||
rustfs-signer = { path = "crates/signer", version = "1.0.0-rc.6" }
|
||||
rustfs-storage-api = { path = "crates/storage-api", version = "1.0.0-rc.6" }
|
||||
rustfs-trusted-proxies = { path = "crates/trusted-proxies", version = "1.0.0-rc.6" }
|
||||
rustfs-targets = { path = "crates/targets", version = "1.0.0-rc.6" }
|
||||
rustfs-test-utils = { path = "crates/test-utils", version = "1.0.0-rc.6" }
|
||||
rustfs-tls-runtime = { path = "crates/tls-runtime", version = "1.0.0-rc.6" }
|
||||
rustfs-utils = { path = "crates/utils", version = "1.0.0-rc.6" }
|
||||
rustfs-zip = { path = "./crates/zip", version = "1.0.0-rc.6" }
|
||||
rustfs = { path = "./rustfs", version = "1.0.0-rc.5" }
|
||||
rustfs-heal = { path = "crates/heal", version = "1.0.0-rc.5" }
|
||||
rustfs-heal-contracts = { path = "crates/heal-contracts", version = "1.0.0-rc.5" }
|
||||
rustfs-scanner-metrics = { path = "crates/scanner-metrics", version = "1.0.0-rc.5" }
|
||||
rustfs-audit = { path = "crates/audit", version = "1.0.0-rc.5" }
|
||||
rustfs-checksums = { path = "crates/checksums", version = "1.0.0-rc.5" }
|
||||
rustfs-common = { path = "crates/common", version = "1.0.0-rc.5" }
|
||||
rustfs-data-usage = { path = "crates/data-usage", version = "1.0.0-rc.5" }
|
||||
rustfs-config = { path = "./crates/config", version = "1.0.0-rc.5" }
|
||||
rustfs-concurrency = { path = "./crates/concurrency", version = "1.0.0-rc.5" }
|
||||
rustfs-credentials = { path = "crates/credentials", version = "1.0.0-rc.5" }
|
||||
rustfs-crypto = { path = "crates/crypto", version = "1.0.0-rc.5" }
|
||||
rustfs-ecstore = { path = "crates/ecstore", version = "1.0.0-rc.5" }
|
||||
rustfs-filemeta = { path = "crates/filemeta", version = "1.0.0-rc.5" }
|
||||
rustfs-iam = { path = "crates/iam", version = "1.0.0-rc.5" }
|
||||
rustfs-keystone = { path = "crates/keystone", version = "1.0.0-rc.5" }
|
||||
rustfs-license = { path = "crates/license", version = "1.0.0-rc.5" }
|
||||
rustfs-lifecycle = { path = "crates/lifecycle", version = "1.0.0-rc.5" }
|
||||
rustfs-kms = { path = "crates/kms", version = "1.0.0-rc.5" }
|
||||
rustfs-lock = { path = "crates/lock", version = "1.0.0-rc.5" }
|
||||
rustfs-madmin = { path = "crates/madmin", version = "1.0.0-rc.5" }
|
||||
rustfs-notify = { path = "crates/notify", version = "1.0.0-rc.5" }
|
||||
rustfs-io-metrics = { path = "crates/io-metrics", version = "1.0.0-rc.5" }
|
||||
rustfs-io-core = { path = "crates/io-core", version = "1.0.0-rc.5" }
|
||||
rustfs-object-capacity = { path = "crates/object-capacity", version = "1.0.0-rc.5" }
|
||||
rustfs-object-data-cache = { path = "crates/object-data-cache", version = "1.0.0-rc.5", default-features = false }
|
||||
rustfs-log-analyzer = { path = "crates/log-analyzer", version = "1.0.0-rc.5" }
|
||||
rustfs-obs = { path = "crates/obs", version = "1.0.0-rc.5" }
|
||||
rustfs-policy = { path = "crates/policy", version = "1.0.0-rc.5" }
|
||||
rustfs-protos = { path = "crates/protos", version = "1.0.0-rc.5" }
|
||||
rustfs-protocols = { path = "crates/protocols", version = "1.0.0-rc.5" }
|
||||
rustfs-replication = { path = "crates/replication", version = "1.0.0-rc.5" }
|
||||
rustfs-rio = { path = "crates/rio", version = "1.0.0-rc.5" }
|
||||
rustfs-rio-v2 = { path = "crates/rio-v2", version = "1.0.0-rc.5" }
|
||||
rustfs-s3-client = { path = "crates/s3-client", version = "1.0.0-rc.5" }
|
||||
rustfs-s3-types = { path = "crates/s3-types", version = "1.0.0-rc.5" }
|
||||
rustfs-s3-ops = { path = "crates/s3-ops", version = "1.0.0-rc.5" }
|
||||
rustfs-s3select-api = { path = "crates/s3select-api", version = "1.0.0-rc.5" }
|
||||
rustfs-s3select-query = { path = "crates/s3select-query", version = "1.0.0-rc.5" }
|
||||
rustfs-scanner = { path = "crates/scanner", version = "1.0.0-rc.5" }
|
||||
rustfs-security-governance = { path = "crates/security-governance", version = "1.0.0-rc.5" }
|
||||
rustfs-extension-schema = { path = "crates/extension-schema", version = "1.0.0-rc.5" }
|
||||
rustfs-signer = { path = "crates/signer", version = "1.0.0-rc.5" }
|
||||
rustfs-storage-api = { path = "crates/storage-api", version = "1.0.0-rc.5" }
|
||||
rustfs-trusted-proxies = { path = "crates/trusted-proxies", version = "1.0.0-rc.5" }
|
||||
rustfs-targets = { path = "crates/targets", version = "1.0.0-rc.5" }
|
||||
rustfs-test-utils = { path = "crates/test-utils", version = "1.0.0-rc.5" }
|
||||
rustfs-tls-runtime = { path = "crates/tls-runtime", version = "1.0.0-rc.5" }
|
||||
rustfs-utils = { path = "crates/utils", version = "1.0.0-rc.5" }
|
||||
rustfs-zip = { path = "./crates/zip", version = "1.0.0-rc.5" }
|
||||
|
||||
# Async Runtime and Networking
|
||||
async-channel = "2.5.0"
|
||||
|
||||
@@ -141,7 +141,7 @@ chown -R 10001:10001 data logs
|
||||
docker run -d -p 9000:9000 -p 9001:9001 -v $(pwd)/data:/data -v $(pwd)/logs:/logs rustfs/rustfs:latest
|
||||
|
||||
# Using specific version
|
||||
docker run -d -p 9000:9000 -p 9001:9001 -v $(pwd)/data:/data -v $(pwd)/logs:/logs rustfs/rustfs:1.0.0-rc.6
|
||||
docker run -d -p 9000:9000 -p 9001:9001 -v $(pwd)/data:/data -v $(pwd)/logs:/logs rustfs/rustfs:1.0.0-rc.5
|
||||
```
|
||||
|
||||
If you use [podman](https://github.com/containers/podman) instead of docker, you can install the RustFS with the below command
|
||||
|
||||
+1
-1
@@ -138,7 +138,7 @@ chown -R 10001:10001 data logs
|
||||
docker run -d -p 9000:9000 -p 9001:9001 -v $(pwd)/data:/data -v $(pwd)/logs:/logs rustfs/rustfs:latest
|
||||
|
||||
# 使用指定版本运行
|
||||
docker run -d -p 9000:9000 -p 9001:9001 -v $(pwd)/data:/data -v $(pwd)/logs:/logs rustfs/rustfs:1.0.0-rc.6
|
||||
docker run -d -p 9000:9000 -p 9001:9001 -v $(pwd)/data:/data -v $(pwd)/logs:/logs rustfs/rustfs:1.0.0-rc.5
|
||||
```
|
||||
|
||||
如果您通过绑定挂载启用 TLS 证书目录,也请用同样方式准备该目录:
|
||||
|
||||
@@ -62,37 +62,6 @@ Current guidance:
|
||||
- `RUSTFS_CORS_ALLOWED_ORIGINS` defaults to empty, so the S3 endpoint emits no generic CORS headers unless configured. Set `*` for wildcard origins without credentials, or a comma-separated allow-list for credentialed explicit origins.
|
||||
- `RUSTFS_CONSOLE_CORS_ALLOWED_ORIGINS` defaults to `*` for the console service.
|
||||
|
||||
## Console URL prefix
|
||||
|
||||
`RUSTFS_CONSOLE_PREFIX` changes the embedded console URL prefix. The default is
|
||||
`/rustfs/console`. For example, `RUSTFS_CONSOLE_PREFIX=/console` serves the UI at
|
||||
`http://localhost:9001/console/`. Nested prefixes such as `/management/console`
|
||||
are supported; one trailing slash is removed. Restart the server after changing it.
|
||||
|
||||
The prefix must be a non-root absolute path of at most 256 bytes, with nonempty
|
||||
segments containing only ASCII letters, digits, `-`, `_`, `.`, or `~`. Dot
|
||||
segments, encoded characters, and overlaps with reserved admin, RPC, health,
|
||||
profiling, browser entry, and icon routes are rejected at startup. `/` is not supported.
|
||||
Choose a prefix that does not collide with S3 bucket paths.
|
||||
|
||||
The console routes, embedded frontend asset URLs, browser redirects, and OIDC
|
||||
console redirects use this prefix. Admin API paths and the identity provider's
|
||||
`/rustfs/admin/v3/oidc/callback/...` URL remain unchanged. `RUSTFS_CONSOLE_ADDRESS`
|
||||
continues to control only the listening address and port. The server adapts bundled
|
||||
console asset references from their build-time base path to the runtime prefix.
|
||||
|
||||
OEM builds can set `RUSTFS_CONSOLE_BASE_PATH` when compiling RustFS to embed a
|
||||
different default, such as `/nuofans/console`. Build the bundled console with the
|
||||
same `NEXT_PUBLIC_BASE_PATH`. An unset or empty build variable retains
|
||||
`/rustfs/console`. The build path must satisfy the validation rules above and must
|
||||
not have a trailing slash.
|
||||
|
||||
At startup, `RUSTFS_CONSOLE_PREFIX` takes precedence over the compiled default.
|
||||
Changing `RUSTFS_CONSOLE_BASE_PATH` when starting an existing binary has no effect;
|
||||
rebuild both components to change the embedded default. If a runtime prefix is
|
||||
configured, asset adaptation uses the compiled base path as its source, including
|
||||
when restoring `/rustfs/console` for a custom OEM build.
|
||||
|
||||
## Browser redirect environment variables
|
||||
|
||||
- `RUSTFS_BROWSER_REDIRECT_URL` sets the externally reachable browser origin used for OIDC callback, console success redirect, and logout fallback URLs. Configure it to the public scheme and authority without a path, for example `https://console.example.com`. In load-balancer deployments, keep OIDC authorize and callback requests on the same backend node because the in-flight OIDC `state` is local to the RustFS node.
|
||||
|
||||
@@ -213,11 +213,6 @@ pub const ENV_RUSTFS_CONSOLE_ENABLE: &str = "RUSTFS_CONSOLE_ENABLE";
|
||||
/// Environment variable for console server address.
|
||||
pub const ENV_RUSTFS_CONSOLE_ADDRESS: &str = "RUSTFS_CONSOLE_ADDRESS";
|
||||
|
||||
/// URL path prefix for the embedded console, read once at server startup.
|
||||
pub const ENV_RUSTFS_CONSOLE_PREFIX: &str = "RUSTFS_CONSOLE_PREFIX";
|
||||
/// Default embedded console URL path prefix.
|
||||
pub const DEFAULT_CONSOLE_PREFIX: &str = "/rustfs/console";
|
||||
|
||||
/// Public browser entrypoint used to build OIDC callback and console redirects.
|
||||
///
|
||||
/// This should be the externally reachable scheme and authority, without a path.
|
||||
|
||||
@@ -64,14 +64,6 @@ pub const MAX_HEAL_REQUEST_SIZE: usize = 1024 * 1024; // 1 MB
|
||||
/// memory exhaustion from malicious or misconfigured remote services.
|
||||
pub const MAX_S3_CLIENT_RESPONSE_SIZE: usize = 10 * 1024 * 1024; // 10 MB
|
||||
|
||||
/// Maximum body size accepted by a single `PutObject` or `UploadPart` request (5 GiB).
|
||||
/// Used for: the s3s streaming-body limit and the request-header admission check.
|
||||
/// Rationale: matches the AWS S3 single-PUT / single-part ceiling. Larger objects
|
||||
/// must use multipart upload. The header check rejects an oversize
|
||||
/// `Content-Length` before any body byte is read so the client gets
|
||||
/// `EntityTooLarge` immediately instead of streaming 5 GiB into a mid-stream failure.
|
||||
pub const MAX_SINGLE_PUT_OBJECT_SIZE: u64 = 5 * 1024 * 1024 * 1024; // 5 GiB
|
||||
|
||||
/// Maximum size for OIDC provider response bodies (1 MB)
|
||||
/// Used for: discovery documents, JWKS documents and token endpoint responses
|
||||
/// Rationale: a hostile or compromised identity provider must not be able to exhaust
|
||||
|
||||
@@ -38,13 +38,6 @@ All commands assume repo root. `cargo test` triggers an on-demand build of the
|
||||
`rustfs` binary from [`src/common.rs`](src/common.rs) (`rustfs_binary_path`) on
|
||||
first use — the first invocation is slow, later ones reuse the binary.
|
||||
|
||||
Root-heal interruption scenarios use a test-only commit barrier. Prebuild with `e2e-test-hooks` and pin that binary so concurrent cases do not replace it through on-demand builds:
|
||||
|
||||
```bash
|
||||
cargo build -p rustfs --bin rustfs --features e2e-test-hooks
|
||||
CARGO_BIN_EXE_rustfs="$PWD/target/debug/rustfs" cargo nextest run -p e2e_test -E 'test(heal_erasure_disk_rebuild_test)'
|
||||
```
|
||||
|
||||
```bash
|
||||
# Whole crate (default = ignored tests skipped)
|
||||
cargo nextest run -p e2e_test
|
||||
|
||||
@@ -28,7 +28,6 @@ mod harness;
|
||||
mod heal_test;
|
||||
mod object_lock_test;
|
||||
mod observability_test;
|
||||
mod replication_delete_marker_test;
|
||||
mod replication_quota_test;
|
||||
mod s3_basic_test;
|
||||
mod s3_during_data_movement_test;
|
||||
|
||||
@@ -1,137 +0,0 @@
|
||||
// Copyright 2026 RustFS Team
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
//! Functional REP-105 (rustfs/backlog#2195 item 4): a delete marker created
|
||||
//! on a multi-node source cluster must replicate to the bucket-replication
|
||||
//! target. Objects converged in seconds while delete markers did not arrive
|
||||
//! within 180 s on the shared 3-node functional environment; the single-node
|
||||
//! e2e never saw it.
|
||||
|
||||
use super::harness::{
|
||||
DistCluster, DistLayout, TestResult, enable_versioning, put_bucket_replication, put_object, set_remote_target, unique_bucket,
|
||||
wait_for_replicated_bytes, wait_until,
|
||||
};
|
||||
use crate::common::{FAST_DATA_USAGE_SCANNER_ENV, RustFSTestEnvironment, init_logging, replication_fast_env, signed_request};
|
||||
use crate::replication_extension_test::LOOPBACK_REPLICATION_TARGET_ENV;
|
||||
use aws_sdk_s3::Client;
|
||||
use http::{Method, StatusCode};
|
||||
use std::time::Duration;
|
||||
|
||||
async fn target_has_delete_marker(client: &Client, bucket: &str, key: &str) -> TestResult<bool> {
|
||||
let versions = client.list_object_versions().bucket(bucket).prefix(key).send().await?;
|
||||
Ok(versions.delete_markers().iter().any(|marker| marker.key() == Some(key)))
|
||||
}
|
||||
|
||||
async fn delete_marker_replicates(
|
||||
source: &DistCluster,
|
||||
source_bucket: &str,
|
||||
target_client: &Client,
|
||||
target_bucket: &str,
|
||||
) -> TestResult {
|
||||
let key = "delete-marker/object.bin";
|
||||
let body = b"delete marker replication payload".to_vec();
|
||||
// Write through one node, delete through another: behind a load
|
||||
// balancer consecutive requests land on different nodes.
|
||||
put_object(&source.client(1)?, source_bucket, key, body.clone()).await?;
|
||||
wait_for_replicated_bytes(target_client, target_bucket, key, &body, Duration::from_secs(60)).await?;
|
||||
|
||||
let delete = source
|
||||
.client(2)?
|
||||
.delete_object()
|
||||
.bucket(source_bucket)
|
||||
.key(key)
|
||||
.send()
|
||||
.await?;
|
||||
assert_eq!(
|
||||
delete.delete_marker(),
|
||||
Some(true),
|
||||
"a versioned DELETE without versionId must create a marker"
|
||||
);
|
||||
|
||||
wait_until(
|
||||
Duration::from_secs(90),
|
||||
|| async { target_has_delete_marker(target_client, target_bucket, key).await },
|
||||
"delete marker replicated to the target bucket",
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn four_node_bucket_replication_replicates_delete_marker_to_peer_cluster() -> TestResult {
|
||||
init_logging();
|
||||
let (source, target) = DistCluster::start_replication_pair().await?;
|
||||
let source_bucket = unique_bucket("dm-src");
|
||||
let target_bucket = unique_bucket("dm-dst");
|
||||
source.create_bucket(&source_bucket).await?;
|
||||
target.create_bucket(&target_bucket).await?;
|
||||
enable_versioning(&source.client(0)?, &source_bucket).await?;
|
||||
enable_versioning(&target.client(0)?, &target_bucket).await?;
|
||||
|
||||
let arn = set_remote_target(&source.cluster, &source_bucket, &target.cluster, &target_bucket).await?;
|
||||
put_bucket_replication(&source.cluster, &source_bucket, &arn).await?;
|
||||
|
||||
delete_marker_replicates(&source, &source_bucket, &target.client(0)?, &target_bucket).await
|
||||
}
|
||||
|
||||
/// The functional environment replicates from a 3-node site to a single-node
|
||||
/// target; keep that shape as its own case.
|
||||
#[tokio::test]
|
||||
async fn four_node_bucket_replication_replicates_delete_marker_to_single_node_target() -> TestResult {
|
||||
init_logging();
|
||||
let mut extra: Vec<(&str, &str)> = replication_fast_env();
|
||||
extra.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV);
|
||||
extra.extend_from_slice(FAST_DATA_USAGE_SCANNER_ENV);
|
||||
let source = DistCluster::start_with_env(DistLayout::FourNodeFourDisk, &extra).await?;
|
||||
let mut target = RustFSTestEnvironment::new().await?;
|
||||
target.start_rustfs_server_without_cleanup(vec![]).await?;
|
||||
|
||||
let source_bucket = unique_bucket("dm-src");
|
||||
let target_bucket = unique_bucket("dm-dst");
|
||||
source.create_bucket(&source_bucket).await?;
|
||||
let target_client = target.create_s3_client();
|
||||
target_client.create_bucket().bucket(&target_bucket).send().await?;
|
||||
enable_versioning(&source.client(0)?, &source_bucket).await?;
|
||||
enable_versioning(&target_client, &target_bucket).await?;
|
||||
|
||||
let body = serde_json::json!({
|
||||
"endpoint": target.address,
|
||||
"credentials": { "accessKey": target.access_key, "secretKey": target.secret_key },
|
||||
"targetbucket": target_bucket,
|
||||
"secure": false,
|
||||
"type": "replication"
|
||||
});
|
||||
let url = format!(
|
||||
"{}/rustfs/admin/v3/set-remote-target?bucket={}",
|
||||
source.cluster.nodes[0].url,
|
||||
urlencoding::encode(&source_bucket)
|
||||
);
|
||||
let response = signed_request(
|
||||
Method::PUT,
|
||||
&url,
|
||||
&source.cluster.access_key,
|
||||
&source.cluster.secret_key,
|
||||
Some(body.to_string().into_bytes()),
|
||||
Some("application/json"),
|
||||
)
|
||||
.await?;
|
||||
if response.status() != StatusCode::OK {
|
||||
let status = response.status();
|
||||
let body = response.text().await.unwrap_or_default();
|
||||
return Err(format!("set remote target failed: {status} {body}").into());
|
||||
}
|
||||
let arn: String = serde_json::from_slice(&response.bytes().await?)?;
|
||||
put_bucket_replication(&source.cluster, &source_bucket, &arn).await?;
|
||||
|
||||
delete_marker_replicates(&source, &source_bucket, &target_client, &target_bucket).await
|
||||
}
|
||||
@@ -126,165 +126,3 @@ async fn four_node_site_replication_replicates_object_to_peer_site() -> TestResu
|
||||
wait_for_replicated_bytes(&site_a.client(3)?, &bucket, reverse_key, &reverse_body, Duration::from_secs(60)).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn node_admin(
|
||||
cluster: &crate::common::RustFSTestClusterEnvironment,
|
||||
node_idx: usize,
|
||||
method: Method,
|
||||
path_and_query: &str,
|
||||
body: Option<String>,
|
||||
) -> TestResult<(StatusCode, String)> {
|
||||
crate::common::admin_request(
|
||||
&cluster.nodes[node_idx].url,
|
||||
method,
|
||||
path_and_query,
|
||||
body,
|
||||
&cluster.access_key,
|
||||
&cluster.secret_key,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Pair two clusters through site A's first node and wait until both report
|
||||
/// the two-site topology as enabled.
|
||||
async fn pair_sites(site_a: &DistCluster, site_b: &DistCluster) -> TestResult {
|
||||
let sites = vec![
|
||||
PeerSite {
|
||||
name: "site-a".to_string(),
|
||||
endpoint: site_a.cluster.nodes[0].url.clone(),
|
||||
access_key: site_a.cluster.access_key.clone(),
|
||||
secret_key: site_a.cluster.secret_key.clone(),
|
||||
..Default::default()
|
||||
},
|
||||
PeerSite {
|
||||
name: "site-b".to_string(),
|
||||
endpoint: site_b.cluster.nodes[0].url.clone(),
|
||||
access_key: site_b.cluster.access_key.clone(),
|
||||
secret_key: site_b.cluster.secret_key.clone(),
|
||||
..Default::default()
|
||||
},
|
||||
];
|
||||
let add_status = site_replication_add(&site_a.cluster, &sites).await?;
|
||||
assert!(
|
||||
add_status.success && add_status.err_detail.is_empty() && add_status.initial_sync_error_message.is_empty(),
|
||||
"site replication add reported failure: {add_status:?}"
|
||||
);
|
||||
wait_for_site_replication_enabled(&site_a.cluster).await?;
|
||||
wait_for_site_replication_enabled(&site_b.cluster).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn list_users_contains(
|
||||
cluster: &crate::common::RustFSTestClusterEnvironment,
|
||||
node_idx: usize,
|
||||
access_key: &str,
|
||||
) -> TestResult<bool> {
|
||||
let (status, body) = node_admin(cluster, node_idx, Method::GET, "/rustfs/admin/v3/list-users", None).await?;
|
||||
if !status.is_success() {
|
||||
return Err(format!("list-users on node {node_idx} failed: {status} {body}").into());
|
||||
}
|
||||
let users: serde_json::Value = serde_json::from_str(&body)?;
|
||||
Ok(users.get(access_key).is_some())
|
||||
}
|
||||
|
||||
/// backlog#2367 A-7 / functional SITE-102: an IAM change handled by a node
|
||||
/// other than the one that ran `site-replication/add` must still reach the
|
||||
/// peer site. Behind a load balancer every admin call may land on a
|
||||
/// different node, so the coordinator node is not special.
|
||||
#[tokio::test]
|
||||
async fn four_node_site_replication_converges_iam_user_created_on_a_non_coordinator_node() -> TestResult {
|
||||
init_logging();
|
||||
let (site_a, site_b) = DistCluster::start_replication_pair().await?;
|
||||
pair_sites(&site_a, &site_b).await?;
|
||||
|
||||
let user = format!("siteuser-{}", &uuid::Uuid::new_v4().simple().to_string()[..8]);
|
||||
let body = serde_json::json!({ "secretKey": "siteuser-secret-key-1234", "status": "enabled" }).to_string();
|
||||
let (status, response) = node_admin(
|
||||
&site_a.cluster,
|
||||
1,
|
||||
Method::PUT,
|
||||
&format!("/rustfs/admin/v3/add-user?accessKey={user}"),
|
||||
Some(body),
|
||||
)
|
||||
.await?;
|
||||
assert!(status.is_success(), "add-user on site A node 1 failed: {status} {response}");
|
||||
|
||||
let site_b_cluster = &site_b.cluster;
|
||||
let user_ref = user.as_str();
|
||||
wait_until(
|
||||
Duration::from_secs(90),
|
||||
|| async move { list_users_contains(site_b_cluster, 0, user_ref).await },
|
||||
"user created on site A node 1 visible on site B",
|
||||
)
|
||||
.await?;
|
||||
assert!(
|
||||
list_users_contains(&site_a.cluster, 2, &user).await?,
|
||||
"the user must be visible on every site A node"
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// backlog#2367 A-5 / functional SITE-105: a resync started right after
|
||||
/// pairing must not report buckets as failed. The bucket carrying an
|
||||
/// operator-configured bucket-replication target to the peer (the shape the
|
||||
/// functional suite leaves behind) and a plain versioned bucket are both
|
||||
/// wired by the pairing itself.
|
||||
#[tokio::test]
|
||||
async fn four_node_site_replication_resync_start_right_after_pairing_reports_no_failed_bucket() -> TestResult {
|
||||
init_logging();
|
||||
let (site_a, site_b) = DistCluster::start_replication_pair().await?;
|
||||
|
||||
let pre_src = unique_bucket("pre-src");
|
||||
let pre_dst = unique_bucket("pre-dst");
|
||||
let plain = unique_bucket("plain");
|
||||
site_a.create_bucket(&pre_src).await?;
|
||||
site_b.create_bucket(&pre_dst).await?;
|
||||
site_a.create_bucket(&plain).await?;
|
||||
enable_versioning(&site_a.client(0)?, &pre_src).await?;
|
||||
enable_versioning(&site_b.client(0)?, &pre_dst).await?;
|
||||
enable_versioning(&site_a.client(0)?, &plain).await?;
|
||||
let arn = super::harness::set_remote_target(&site_a.cluster, &pre_src, &site_b.cluster, &pre_dst).await?;
|
||||
super::harness::put_bucket_replication(&site_a.cluster, &pre_src, &arn).await?;
|
||||
|
||||
pair_sites(&site_a, &site_b).await?;
|
||||
|
||||
let (status, info) = node_admin(&site_a.cluster, 1, Method::GET, "/rustfs/admin/v3/site-replication/info", None).await?;
|
||||
assert!(status.is_success(), "site-replication/info failed: {status} {info}");
|
||||
let info: serde_json::Value = serde_json::from_str(&info)?;
|
||||
let peer = info["sites"]
|
||||
.as_array()
|
||||
.and_then(|sites| sites.iter().find(|site| site["name"] == "site-b"))
|
||||
.cloned()
|
||||
.ok_or_else(|| format!("site-b peer missing from info: {info}"))?;
|
||||
|
||||
// Through a non-coordinator node, like a load-balanced admin call.
|
||||
let (status, response) = node_admin(
|
||||
&site_a.cluster,
|
||||
1,
|
||||
Method::PUT,
|
||||
"/rustfs/admin/v3/site-replication/resync/op?operation=start",
|
||||
Some(peer.to_string()),
|
||||
)
|
||||
.await?;
|
||||
assert!(status.is_success(), "resync start failed: {status} {response}");
|
||||
let resync: rustfs_madmin::SRResyncOpStatus = serde_json::from_str(&response)?;
|
||||
let failed: Vec<String> = resync
|
||||
.buckets
|
||||
.iter()
|
||||
.filter(|bucket| bucket.status == "failed")
|
||||
.map(|bucket| format!("{}: {}", bucket.bucket, bucket.err_detail))
|
||||
.collect();
|
||||
assert!(
|
||||
failed.is_empty(),
|
||||
"resync right after pairing reported failed buckets: {failed:?} (status={}, detail={})",
|
||||
resync.status,
|
||||
resync.err_detail
|
||||
);
|
||||
assert!(
|
||||
resync.buckets.iter().any(|bucket| bucket.bucket == pre_src)
|
||||
&& resync.buckets.iter().any(|bucket| bucket.bucket == plain),
|
||||
"both buckets must be part of the resync: {:?}",
|
||||
resync.buckets
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -1530,16 +1530,6 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
// Keep the partial-repair checkpoint stable across readiness and admin
|
||||
// requests. Endpoint-blackhole tests must prove their own network stall.
|
||||
let commit_barrier = if scenario != InterruptionScenario::TargetEndpointBlackhole {
|
||||
let barrier = replaced_disk.join(".rustfs.sys/e2e-heal-commit-barrier");
|
||||
std::fs::create_dir_all(barrier.parent().ok_or("commit barrier has no parent")?)?;
|
||||
std::fs::write(&barrier, format!("{bucket}/cluster/online/"))?;
|
||||
Some(barrier)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
cluster.start_node_from_binary(1, &server_binary).await?;
|
||||
|
||||
let status_url = format!("{}/rustfs/admin/v3/background-heal/status", cluster.nodes[0].url);
|
||||
@@ -1673,12 +1663,6 @@ mod tests {
|
||||
sleep(Duration::from_millis(10)).await;
|
||||
};
|
||||
|
||||
if let Some(barrier) = &commit_barrier {
|
||||
assert!(
|
||||
barrier.with_extension("admitted").is_file(),
|
||||
"interruption tests require a server built with e2e-test-hooks"
|
||||
);
|
||||
}
|
||||
let pre_interrupt_status_body = signed_admin_post(&status_url, None, &cluster.access_key, &cluster.secret_key).await?;
|
||||
let pre_interrupt_status: serde_json::Value = serde_json::from_str(&pre_interrupt_status_body)
|
||||
.map_err(|err| format!("pre-interrupt background heal status is not JSON ({err}): {pre_interrupt_status_body}"))?;
|
||||
@@ -1884,9 +1868,6 @@ mod tests {
|
||||
}
|
||||
}
|
||||
}
|
||||
if let Some(barrier) = &commit_barrier {
|
||||
std::fs::remove_file(barrier)?;
|
||||
}
|
||||
cluster.start_node_from_binary(interruption_node, &server_binary).await?;
|
||||
if interruption_node == 0 {
|
||||
let target = cluster.nodes[1]
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -6409,10 +6409,6 @@ impl LocalDisk {
|
||||
// A missing or still-populated directory is benign here; see
|
||||
// is_benign_object_rmdir_error (handles the illumos/Solaris EEXIST
|
||||
// convention, rustfs/rustfs#4978).
|
||||
if is_dir_not_empty_error(&err) {
|
||||
// A populated directory keeps its ancestors populated; no further pruning is needed.
|
||||
return Ok(());
|
||||
}
|
||||
if !is_benign_object_rmdir_error(&err) {
|
||||
warn!(
|
||||
event = EVENT_DISK_LOCAL_DELETE_FAILED,
|
||||
@@ -11226,176 +11222,6 @@ mod test {
|
||||
(disk, dir)
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn delete_pruning_stops_at_live_metadata_below_a_guarded_ancestor() {
|
||||
// Tuple fields drop in order, releasing the disk's root handle before the temporary directory.
|
||||
let fixture = new_disk().await;
|
||||
let (disk, _dir) = &fixture;
|
||||
let base = disk.get_bucket_path(RUSTFS_META_BUCKET).expect("resolve metadata volume");
|
||||
let shared = base.join("buckets");
|
||||
let guard = Arc::new(
|
||||
os::mkdir_all_below_existing_base_std(&shared, &base, &disk.publication_root)
|
||||
.expect("retain the shared publication directory"),
|
||||
);
|
||||
|
||||
for owned in [false, true] {
|
||||
for missing_backup in [false, true] {
|
||||
let transaction = Uuid::new_v4();
|
||||
let object = shared.join(".bloomcycle.bin");
|
||||
let rollback = object.join(transaction.to_string());
|
||||
let metadata = object.join(STORAGE_FORMAT_FILE);
|
||||
let backup = rollback.join(STORAGE_FORMAT_FILE_BACKUP);
|
||||
fs::create_dir_all(&rollback).await.expect("create rollback directory");
|
||||
fs::write(&metadata, b"committed metadata")
|
||||
.await
|
||||
.expect("write live metadata");
|
||||
if !missing_backup {
|
||||
fs::write(&backup, b"old metadata").await.expect("write rollback backup");
|
||||
}
|
||||
let owner: Option<Arc<dyn Send + Sync>> = if owned { Some(guard.clone()) } else { None };
|
||||
let result = disk
|
||||
.delete_with_namespace_owner(
|
||||
RUSTFS_META_BUCKET,
|
||||
&format!("buckets/.bloomcycle.bin/{transaction}/{STORAGE_FORMAT_FILE_BACKUP}"),
|
||||
DeleteOptions::default(),
|
||||
owner,
|
||||
)
|
||||
.await;
|
||||
|
||||
assert!(!backup.exists(), "backup must be absent, owned={owned}, missing={missing_backup}");
|
||||
assert!(!rollback.exists(), "empty rollback directory must be pruned");
|
||||
assert_eq!(fs::read(&metadata).await.expect("read committed metadata"), b"committed metadata");
|
||||
result.expect("a nonempty object must stop pruning before the guarded ancestor");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn delete_pruning_removes_empty_and_missing_ancestors_but_keeps_the_volume() {
|
||||
let fixture = new_disk().await;
|
||||
let (disk, _dir) = &fixture;
|
||||
ensure_test_volume(disk, "pruning").await;
|
||||
let base = disk.get_bucket_path("pruning").expect("resolve test volume");
|
||||
|
||||
for missing in [false, true] {
|
||||
let parent = base.join("parent");
|
||||
let rollback = parent.join("object/transaction");
|
||||
fs::create_dir_all(&rollback).await.expect("create empty ancestor chain");
|
||||
let path = if missing {
|
||||
"parent/object/transaction/missing/xl.meta.bkp"
|
||||
} else {
|
||||
fs::write(rollback.join(STORAGE_FORMAT_FILE_BACKUP), b"backup")
|
||||
.await
|
||||
.expect("create backup");
|
||||
"parent/object/transaction/xl.meta.bkp"
|
||||
};
|
||||
|
||||
disk.delete("pruning", path, DeleteOptions::default())
|
||||
.await
|
||||
.expect("empty and missing ancestors should be pruned");
|
||||
assert!(!parent.exists(), "the whole empty chain should be removed");
|
||||
assert!(base.is_dir(), "pruning must stop at the volume boundary");
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn delete_pruning_does_not_remove_the_base_or_an_outside_path() {
|
||||
let fixture = new_disk().await;
|
||||
let (disk, dir) = &fixture;
|
||||
let base = dir.path().join("base");
|
||||
let outside = dir.path().join("outside");
|
||||
fs::create_dir(&base).await.expect("create base");
|
||||
fs::write(&outside, b"outside data").await.expect("create outside file");
|
||||
|
||||
disk.delete_file(&base, &base, false, false)
|
||||
.await
|
||||
.expect("base path is protected");
|
||||
disk.delete_file(&base, &outside, false, false)
|
||||
.await
|
||||
.expect("outside path is protected");
|
||||
assert!(base.is_dir(), "the base must not be removed even when empty");
|
||||
assert_eq!(fs::read(&outside).await.expect("read outside file"), b"outside data");
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
#[tokio::test]
|
||||
async fn delete_pruning_propagates_a_locked_backup_error() {
|
||||
use std::os::windows::fs::OpenOptionsExt;
|
||||
use windows_sys::Win32::{Foundation::ERROR_SHARING_VIOLATION, Storage::FileSystem::FILE_SHARE_READ};
|
||||
|
||||
let fixture = new_disk().await;
|
||||
let (disk, _dir) = &fixture;
|
||||
ensure_test_volume(disk, "pruning").await;
|
||||
let base = disk.get_bucket_path("pruning").expect("resolve test volume");
|
||||
let backup = base.join(STORAGE_FORMAT_FILE_BACKUP);
|
||||
fs::write(&backup, b"backup").await.expect("write backup");
|
||||
let guard = std::fs::OpenOptions::new()
|
||||
.read(true)
|
||||
.share_mode(FILE_SHARE_READ)
|
||||
.open(&backup)
|
||||
.expect("hold the backup without delete sharing");
|
||||
|
||||
let err = disk
|
||||
.delete("pruning", STORAGE_FORMAT_FILE_BACKUP, DeleteOptions::default())
|
||||
.await
|
||||
.expect_err("a genuine target-file deletion failure must propagate");
|
||||
let DiskError::Io(err) = err else {
|
||||
panic!("expected contextual I/O error, got {err:?}");
|
||||
};
|
||||
let context = err
|
||||
.get_ref()
|
||||
.and_then(|err| err.downcast_ref::<FileAccessDeniedWithContext>())
|
||||
.expect("preserve the failing path and original OS error");
|
||||
assert_eq!(context.path, backup);
|
||||
assert_eq!(
|
||||
context.source.raw_os_error(),
|
||||
Some(i32::try_from(ERROR_SHARING_VIOLATION).expect("OS code fits"))
|
||||
);
|
||||
assert_eq!(fs::read(&backup).await.expect("backup remains readable"), b"backup");
|
||||
drop(guard);
|
||||
}
|
||||
|
||||
#[cfg(windows)]
|
||||
#[tokio::test]
|
||||
async fn delete_pruning_propagates_a_locked_empty_parent_error() {
|
||||
use windows_sys::Win32::Foundation::ERROR_SHARING_VIOLATION;
|
||||
|
||||
let fixture = new_disk().await;
|
||||
let (disk, _dir) = &fixture;
|
||||
ensure_test_volume(disk, "pruning").await;
|
||||
let base = disk.get_bucket_path("pruning").expect("resolve test volume");
|
||||
let parent = base.join("parent");
|
||||
let guard = os::mkdir_all_below_existing_base_std(&parent, &base, &disk.publication_root)
|
||||
.expect("retain an empty parent without delete sharing");
|
||||
let backup = parent.join(STORAGE_FORMAT_FILE_BACKUP);
|
||||
fs::write(&backup, b"backup").await.expect("write backup");
|
||||
|
||||
let err = disk
|
||||
.delete("pruning", "parent/xl.meta.bkp", DeleteOptions::default())
|
||||
.await
|
||||
.expect_err("a real parent failure without a nonempty boundary must still propagate");
|
||||
let DiskError::Io(err) = err else {
|
||||
panic!("expected contextual I/O error, got {err:?}");
|
||||
};
|
||||
let context = err
|
||||
.get_ref()
|
||||
.and_then(|err| err.downcast_ref::<FileAccessDeniedWithContext>())
|
||||
.expect("preserve parent failure context");
|
||||
assert_eq!(context.path, parent);
|
||||
assert_eq!(
|
||||
context.source.raw_os_error(),
|
||||
Some(i32::try_from(ERROR_SHARING_VIOLATION).expect("OS code fits"))
|
||||
);
|
||||
assert!(!backup.exists(), "the target was removed before the parent error");
|
||||
assert!(parent.is_dir(), "the guarded parent remains");
|
||||
drop(guard);
|
||||
disk.delete("pruning", "parent/xl.meta.bkp", DeleteOptions::default())
|
||||
.await
|
||||
.expect("pruning should succeed once the actual guard is released");
|
||||
assert!(!parent.exists());
|
||||
assert!(base.is_dir());
|
||||
}
|
||||
|
||||
// #948: a genuinely missing source is benign and must still return Ok.
|
||||
#[tokio::test]
|
||||
async fn windows_and_unix_move_to_trash_missing_source_is_ok() {
|
||||
|
||||
@@ -47,43 +47,6 @@ use tokio::fs;
|
||||
use tracing::{info, warn};
|
||||
use uuid::Uuid;
|
||||
|
||||
/// Hold later repair publications after admitting one baseline object. The
|
||||
/// fixture arms this on one replacement disk before rejoining the cluster.
|
||||
#[cfg(feature = "e2e-test-hooks")]
|
||||
async fn wait_for_heal_commit_test_barrier(root: &Path, bucket: &str, object: &str) -> Result<()> {
|
||||
use tokio::io::AsyncWriteExt;
|
||||
|
||||
let barrier = root.join(".rustfs.sys/e2e-heal-commit-barrier");
|
||||
let prefix = match fs::read_to_string(&barrier).await {
|
||||
Ok(prefix) => prefix,
|
||||
Err(error) if error.kind() == ErrorKind::NotFound => return Ok(()),
|
||||
Err(error) => return Err(error.into()),
|
||||
};
|
||||
let key = format!("{bucket}/{object}");
|
||||
if prefix.is_empty() || !key.starts_with(&prefix) {
|
||||
return Ok(());
|
||||
}
|
||||
let admitted = barrier.with_extension("admitted");
|
||||
match fs::OpenOptions::new().write(true).create_new(true).open(&admitted).await {
|
||||
Ok(mut file) => {
|
||||
file.write_all(key.as_bytes()).await?;
|
||||
return Ok(());
|
||||
}
|
||||
Err(error) if error.kind() == ErrorKind::AlreadyExists => {}
|
||||
Err(error) => return Err(error.into()),
|
||||
}
|
||||
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(120);
|
||||
loop {
|
||||
if !fs::try_exists(&barrier).await? || fs::read_to_string(&admitted).await? == key {
|
||||
return Ok(());
|
||||
}
|
||||
if tokio::time::Instant::now() >= deadline {
|
||||
return Err(std::io::Error::new(ErrorKind::TimedOut, "heal commit test barrier was not released").into());
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
|
||||
}
|
||||
}
|
||||
|
||||
fn rollback_committed_rename_std(
|
||||
dst_file_path: &Path,
|
||||
new_data_path: Option<&Path>,
|
||||
@@ -290,10 +253,6 @@ impl LocalDisk {
|
||||
state: &mut RenameDataState,
|
||||
) -> Result<RenameDataResp> {
|
||||
crate::hp_guard!("LocalDisk::rename_data");
|
||||
#[cfg(feature = "e2e-test-hooks")]
|
||||
if fi.is_healing() {
|
||||
wait_for_heal_commit_test_barrier(&self.root, dst_volume, dst_path).await?;
|
||||
}
|
||||
let mut fi = fi;
|
||||
// A non-force DeleteBucket must not remove a directory while a local
|
||||
// object commit is publishing into it. The peer's empty scan remains
|
||||
|
||||
@@ -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";
|
||||
|
||||
@@ -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)]
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -548,10 +548,7 @@ fn retry_budget_for_result(task: &HealTask, result: &Result<()>, retryable_batch
|
||||
}
|
||||
|
||||
let error = err.to_string();
|
||||
// Batch aggregation preserves the typed classification in its counters,
|
||||
// while the returned task error retains only the first error's display text.
|
||||
let retryable_batch_result = retryable_batch_failure && matches!(err, Error::TaskExecutionFailed { .. });
|
||||
if !retryable_batch_result && !err.is_recoverable_heal() {
|
||||
if !err.is_recoverable_heal() {
|
||||
return None;
|
||||
}
|
||||
|
||||
|
||||
@@ -2614,62 +2614,27 @@ fn test_retry_request_for_recoverable_error_stops_at_limit() {
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_retry_request_rescans_batch_when_all_exhausted_objects_are_retryable() {
|
||||
for source_error in [
|
||||
Error::Disk(DiskError::FaultyDisk),
|
||||
Error::Disk(DiskError::FaultyRemoteDisk),
|
||||
Error::Storage(EcstoreError::SlowDown),
|
||||
Error::TaskExecutionFailed {
|
||||
message: "Lock acquisition timeout".to_string(),
|
||||
},
|
||||
] {
|
||||
assert!(source_error.is_recoverable_heal());
|
||||
let first_error = source_error.to_string();
|
||||
let storage: Arc<dyn HealStorageAPI> = Arc::new(MockStorage);
|
||||
let task = HealTask::from_request(HealRequest::bucket("bucket".to_string()), storage);
|
||||
let result = Err(task
|
||||
.record_batch_failure(BatchHealFailure {
|
||||
scope: "bucket:bucket".to_string(),
|
||||
failed: 1,
|
||||
retryable: 1,
|
||||
permanent: 0,
|
||||
first_object: "object".to_string(),
|
||||
first_error: first_error.clone(),
|
||||
})
|
||||
.await);
|
||||
let storage: Arc<dyn HealStorageAPI> = Arc::new(MockStorage);
|
||||
let task = HealTask::from_request(HealRequest::bucket("bucket".to_string()), storage);
|
||||
let result = Err(task
|
||||
.record_batch_failure(BatchHealFailure {
|
||||
scope: "bucket:bucket".to_string(),
|
||||
failed: 1,
|
||||
retryable: 1,
|
||||
permanent: 0,
|
||||
first_object: "object".to_string(),
|
||||
first_error: "Lock acquisition timeout".to_string(),
|
||||
})
|
||||
.await);
|
||||
|
||||
let (retry_request, retry_delay, error) = retry_request_for_result_with_budget(&task, &result)
|
||||
.await
|
||||
.expect("all-retryable batch failure should rescan within the manager retry budget");
|
||||
let (retry_request, retry_delay, error) = retry_request_for_result_with_budget(&task, &result)
|
||||
.await
|
||||
.expect("all-retryable batch failure should rescan within the manager retry budget");
|
||||
|
||||
assert_eq!(retry_request.id, task.id);
|
||||
assert_eq!(retry_request.retry_attempts, 1);
|
||||
assert!(retry_delay > Duration::ZERO);
|
||||
assert!(error.contains(&first_error));
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_retry_request_does_not_rescan_cancelled_or_timed_out_retryable_batch() {
|
||||
for terminal_error in [Error::TaskCancelled, Error::TaskTimeout] {
|
||||
let storage: Arc<dyn HealStorageAPI> = Arc::new(MockStorage);
|
||||
let task = HealTask::from_request(HealRequest::bucket("bucket".to_string()), storage);
|
||||
let _ = task
|
||||
.record_batch_failure(BatchHealFailure {
|
||||
scope: "bucket:bucket".to_string(),
|
||||
failed: 1,
|
||||
retryable: 1,
|
||||
permanent: 0,
|
||||
first_object: "object".to_string(),
|
||||
first_error: Error::Disk(DiskError::FaultyDisk).to_string(),
|
||||
})
|
||||
.await;
|
||||
|
||||
assert!(
|
||||
retry_request_for_result_with_budget(&task, &Err(terminal_error))
|
||||
.await
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
assert_eq!(retry_request.id, task.id);
|
||||
assert_eq!(retry_request.retry_attempts, 1);
|
||||
assert!(retry_delay > Duration::ZERO);
|
||||
assert!(error.contains("Lock acquisition timeout"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
|
||||
@@ -730,16 +730,12 @@ impl ReplicationConfigurationExt for ReplicationConfiguration {
|
||||
}
|
||||
}
|
||||
|
||||
// Highest priority first, like MinIO's `FilterActionableRules`. The
|
||||
// tie-breakers make this a total order: a comparator that only
|
||||
// orders same-destination pairs is not transitive, and the standard
|
||||
// library sort panics on such inputs past its insertion-sort
|
||||
// threshold (backlog#2367 C-1).
|
||||
rules.sort_by(|a, b| {
|
||||
b.priority
|
||||
.cmp(&a.priority)
|
||||
.then_with(|| a.destination.bucket.cmp(&b.destination.bucket))
|
||||
.then_with(|| a.id.cmp(&b.id))
|
||||
if a.destination == b.destination {
|
||||
b.priority.cmp(&a.priority)
|
||||
} else {
|
||||
std::cmp::Ordering::Equal
|
||||
}
|
||||
});
|
||||
|
||||
rules
|
||||
@@ -817,19 +813,24 @@ impl ReplicationConfigurationExt for ReplicationConfiguration {
|
||||
return vec![role.to_string()];
|
||||
}
|
||||
|
||||
// Rule order (priority descending) is the ARN order: callers that
|
||||
// iterate targets see the highest-priority destination first.
|
||||
let mut arns: Vec<String> = Vec::new();
|
||||
for rule in self.filter_actionable_rules(obj) {
|
||||
let mut arns = Vec::new();
|
||||
let mut targets_map: HashSet<String> = HashSet::new();
|
||||
let rules = self.filter_actionable_rules(obj);
|
||||
|
||||
for rule in rules {
|
||||
if rule.status == ReplicationRuleStatus::from_static(ReplicationRuleStatus::DISABLED) {
|
||||
continue;
|
||||
}
|
||||
|
||||
let arn = rule.destination.bucket.trim();
|
||||
if !arn.is_empty() && !arns.iter().any(|seen| seen == arn) {
|
||||
arns.push(arn.to_string());
|
||||
if !arn.is_empty() && !targets_map.contains(arn) {
|
||||
targets_map.insert(arn.to_string());
|
||||
}
|
||||
}
|
||||
|
||||
for arn in targets_map {
|
||||
arns.push(arn);
|
||||
}
|
||||
arns
|
||||
}
|
||||
|
||||
@@ -1907,84 +1908,6 @@ mod tests {
|
||||
assert_eq!(decisions, vec![(target_a.to_string(), false), (target_b.to_string(), true)]);
|
||||
}
|
||||
|
||||
// backlog#2367 C-1: the actionable-rule sort must be a total order. A
|
||||
// comparator that answers `Equal` for different destinations but orders
|
||||
// same-destination rules by priority is not transitive, and the standard
|
||||
// library sort panics on such inputs once the slice is past the
|
||||
// insertion-sort threshold (> 20 rules).
|
||||
#[test]
|
||||
fn actionable_rule_sort_is_a_total_order_across_destinations() {
|
||||
let targets = ["arn:target:a", "arn:target:b", "arn:target:c"];
|
||||
let mut seed: u64 = 0x2367;
|
||||
for _ in 0..200 {
|
||||
let rule_count = 21 + (seed % 200) as usize;
|
||||
let rules = (0..rule_count)
|
||||
.map(|index| {
|
||||
seed = seed.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407);
|
||||
let target = targets[(seed >> 33) as usize % targets.len()];
|
||||
delete_marker_rule(&format!("r{index}"), target, "", index as i32, true)
|
||||
})
|
||||
.collect();
|
||||
let config = ReplicationConfiguration {
|
||||
role: String::new(),
|
||||
rules,
|
||||
};
|
||||
let ordered = config.filter_actionable_rules(&ObjectOpts {
|
||||
name: "logs/app.log".to_string(),
|
||||
op_type: ReplicationType::Object,
|
||||
..Default::default()
|
||||
});
|
||||
assert_eq!(ordered.len(), rule_count);
|
||||
assert!(
|
||||
ordered.windows(2).all(|pair| pair[0].priority >= pair[1].priority),
|
||||
"actionable rules must be ordered by descending priority"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// backlog#2367 C-2: a V1 rule carries its prefix at the top level (no
|
||||
// <Filter>). Ignoring it made `<Prefix>logs/</Prefix>` match every object.
|
||||
#[test]
|
||||
fn top_level_rule_prefix_scopes_matching_without_a_filter() {
|
||||
let arn = "arn:target:a";
|
||||
let config = ReplicationConfiguration {
|
||||
role: String::new(),
|
||||
rules: vec![delete_marker_rule("v1-prefix", arn, "logs/", 1, true)],
|
||||
};
|
||||
assert_eq!(config.rules[0].prefix(), "logs/");
|
||||
|
||||
let matching = config.filter_actionable_rules(&ObjectOpts {
|
||||
name: "logs/app.log".to_string(),
|
||||
op_type: ReplicationType::Object,
|
||||
..Default::default()
|
||||
});
|
||||
assert_eq!(matching.len(), 1);
|
||||
|
||||
let outside = config.filter_actionable_rules(&ObjectOpts {
|
||||
name: "data/app.log".to_string(),
|
||||
op_type: ReplicationType::Object,
|
||||
..Default::default()
|
||||
});
|
||||
assert!(outside.is_empty(), "an object outside the V1 prefix must not match: {outside:?}");
|
||||
assert!(
|
||||
config
|
||||
.filter_target_arns(&ObjectOpts {
|
||||
name: "data/app.log".to_string(),
|
||||
op_type: ReplicationType::Object,
|
||||
..Default::default()
|
||||
})
|
||||
.is_empty()
|
||||
);
|
||||
|
||||
// A <Filter> still wins over the deprecated top-level element.
|
||||
let mut filtered = delete_marker_rule("filtered", arn, "logs/", 1, true);
|
||||
filtered.filter = Some(s3s::dto::ReplicationRuleFilter {
|
||||
prefix: Some("photos/".to_string()),
|
||||
..Default::default()
|
||||
});
|
||||
assert_eq!(filtered.prefix(), "photos/");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn force_delete_targets_use_overlapping_rules_and_highest_priority_switch() {
|
||||
let target_a = "arn:target:a";
|
||||
|
||||
@@ -22,10 +22,6 @@ pub trait ReplicationRuleExt {
|
||||
}
|
||||
|
||||
impl ReplicationRuleExt for ReplicationRule {
|
||||
/// The rule's key prefix: `Filter.Prefix`, else `Filter.And.Prefix`, else
|
||||
/// the deprecated top-level `Prefix` of a V1 rule written without a
|
||||
/// `<Filter>` (backlog#2367 C-2). A rule that carries both keeps AWS's
|
||||
/// precedence: the `<Filter>` is authoritative.
|
||||
fn prefix(&self) -> &str {
|
||||
if let Some(filter) = &self.filter {
|
||||
if let Some(prefix) = &filter.prefix {
|
||||
@@ -36,7 +32,7 @@ impl ReplicationRuleExt for ReplicationRule {
|
||||
""
|
||||
}
|
||||
} else {
|
||||
self.prefix.as_deref().unwrap_or("")
|
||||
""
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -77,7 +77,7 @@
|
||||
|
||||
rustfs = rustPlatform.buildRustPackage {
|
||||
pname = "rustfs";
|
||||
version = "1.0.0-rc.6";
|
||||
version = "1.0.0-rc.5";
|
||||
|
||||
src = ./.;
|
||||
|
||||
|
||||
@@ -2,8 +2,8 @@ apiVersion: v2
|
||||
name: rustfs
|
||||
description: RustFS helm chart to deploy RustFS on kubernetes cluster.
|
||||
type: application
|
||||
version: "1.0.0-rc.6"
|
||||
appVersion: "1.0.0-rc.6"
|
||||
version: "1.0.0-rc.5"
|
||||
appVersion: "1.0.0-rc.5"
|
||||
home: https://rustfs.com
|
||||
icon: https://media.sys.truenas.net/apps/rustfs/icons/icon.svg
|
||||
maintainers:
|
||||
|
||||
+2
-5
@@ -1,9 +1,9 @@
|
||||
%global _enable_debug_packages 0
|
||||
%global _empty_manifest_terminate_build 0
|
||||
%global prerelease rc.6
|
||||
%global prerelease rc.5
|
||||
Name: rustfs
|
||||
Version: 1.0.0
|
||||
Release: rc.6
|
||||
Release: rc.5
|
||||
Summary: High-performance distributed object storage for MinIO alternative
|
||||
|
||||
License: Apache-2.0
|
||||
@@ -58,9 +58,6 @@ install %_builddir/%{name}-%{version}-%{prerelease}/target/%_arch/%_arch-unknown
|
||||
%_bindir/rustfs
|
||||
|
||||
%changelog
|
||||
* Thu Sep 10 2026 overtrue <anzhengchao@gmail.com>
|
||||
- Update RPM package to RustFS 1.0.0-rc.6
|
||||
|
||||
* Mon Aug 31 2026 overtrue <anzhengchao@gmail.com>
|
||||
- Update RPM package to RustFS 1.0.0-rc.5
|
||||
|
||||
|
||||
+21
-168
@@ -22,9 +22,9 @@ use crate::server::rate_limit::{
|
||||
apply_throttle_headers, client_ip,
|
||||
};
|
||||
use crate::server::{
|
||||
APPLE_TOUCH_ICON_PATH, APPLE_TOUCH_ICON_PRECOMPOSED_PATH, FAVICON_PATH, HEALTH_PREFIX, HEALTH_READY_PATH, HeaderMapCarrier,
|
||||
HealthProbe, LICENSE, RUSTFS_ADMIN_PREFIX, RequestContextLayer, VERSION, build_health_response_parts,
|
||||
collect_probe_readiness, console_prefix,
|
||||
APPLE_TOUCH_ICON_PATH, APPLE_TOUCH_ICON_PRECOMPOSED_PATH, CONSOLE_PREFIX, FAVICON_PATH, HEALTH_PREFIX, HEALTH_READY_PATH,
|
||||
HeaderMapCarrier, HealthProbe, LICENSE, RUSTFS_ADMIN_PREFIX, RequestContextLayer, VERSION, build_health_response_parts,
|
||||
collect_probe_readiness,
|
||||
};
|
||||
use crate::version::{self, build};
|
||||
use axum::{
|
||||
@@ -83,7 +83,7 @@ async fn static_handler(uri: Uri) -> impl IntoResponse {
|
||||
return Response::builder()
|
||||
.status(StatusCode::OK)
|
||||
.header("Content-Type", mime_type.to_string())
|
||||
.body(Body::from(rewrite_console_asset(path, file.data, crate::server::console_prefix())))
|
||||
.body(Body::from(file.data))
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
@@ -95,7 +95,7 @@ async fn static_handler(uri: Uri) -> impl IntoResponse {
|
||||
return Response::builder()
|
||||
.status(StatusCode::OK)
|
||||
.header("Content-Type", mime_type.to_string())
|
||||
.body(Body::from(rewrite_console_asset(&index_path, file.data, crate::server::console_prefix())))
|
||||
.body(Body::from(file.data))
|
||||
.unwrap();
|
||||
}
|
||||
}
|
||||
@@ -106,11 +106,7 @@ async fn static_handler(uri: Uri) -> impl IntoResponse {
|
||||
Response::builder()
|
||||
.status(StatusCode::OK)
|
||||
.header("Content-Type", mime_type.to_string())
|
||||
.body(Body::from(rewrite_console_asset(
|
||||
"index.html",
|
||||
file.data,
|
||||
crate::server::console_prefix(),
|
||||
)))
|
||||
.body(Body::from(file.data))
|
||||
.unwrap()
|
||||
} else {
|
||||
Response::builder()
|
||||
@@ -120,57 +116,6 @@ async fn static_handler(uri: Uri) -> impl IntoResponse {
|
||||
}
|
||||
}
|
||||
|
||||
// Next exports bake the base path into HTML, chunk loaders, and RSC text payloads.
|
||||
// Rewrite only path references, preserving external URLs such as the source repository.
|
||||
fn rewrite_console_asset<'a>(path: &str, data: std::borrow::Cow<'a, [u8]>, prefix: &str) -> std::borrow::Cow<'a, [u8]> {
|
||||
use std::borrow::Cow;
|
||||
|
||||
if prefix == crate::server::CONSOLE_PREFIX
|
||||
|| !matches!(
|
||||
path.rsplit('.').next(),
|
||||
Some("html" | "js" | "css" | "json" | "txt" | "webmanifest" | "svg")
|
||||
)
|
||||
{
|
||||
return data;
|
||||
}
|
||||
let Ok(text) = std::str::from_utf8(&data) else {
|
||||
return data;
|
||||
};
|
||||
let mut rewritten = Cow::Borrowed(text);
|
||||
let escaped_default = crate::server::CONSOLE_PREFIX.replace('/', "\\/");
|
||||
let escaped_prefix = prefix.replace('/', "\\/");
|
||||
for (source, target) in [
|
||||
(crate::server::CONSOLE_PREFIX, prefix),
|
||||
(escaped_default.as_str(), escaped_prefix.as_str()),
|
||||
] {
|
||||
let mut output = String::new();
|
||||
let mut copied = 0;
|
||||
for (offset, _) in rewritten.match_indices(source) {
|
||||
let end = offset + source.len();
|
||||
let before = rewritten.as_bytes().get(offset.wrapping_sub(1)).copied();
|
||||
let after = rewritten.as_bytes().get(end).copied();
|
||||
let starts_path = before
|
||||
.is_none_or(|byte| byte.is_ascii_whitespace() || matches!(byte, b'"' | b'\'' | b'`' | b'(' | b'=' | b'}' | b'>'));
|
||||
let ends_prefix = after.is_none_or(|byte| {
|
||||
byte.is_ascii_whitespace() || matches!(byte, b'/' | b'\\' | b'"' | b'\'' | b'`' | b'?' | b'#' | b')' | b'<')
|
||||
});
|
||||
if starts_path && ends_prefix {
|
||||
output.push_str(&rewritten[copied..offset]);
|
||||
output.push_str(target);
|
||||
copied = end;
|
||||
}
|
||||
}
|
||||
if copied != 0 {
|
||||
output.push_str(&rewritten[copied..]);
|
||||
rewritten = Cow::Owned(output);
|
||||
}
|
||||
}
|
||||
match rewritten {
|
||||
Cow::Borrowed(_) => data,
|
||||
Cow::Owned(text) => Cow::Owned(text.into_bytes()),
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Serialize, Clone)]
|
||||
pub(crate) struct Config {
|
||||
#[serde(skip)]
|
||||
@@ -523,7 +468,7 @@ fn get_console_config_from_env() -> (bool, u32, u64, String) {
|
||||
/// - `true` if the path is for console access, `false` otherwise.
|
||||
pub fn is_console_path(path: &str) -> bool {
|
||||
matches!(path, FAVICON_PATH | APPLE_TOUCH_ICON_PATH | APPLE_TOUCH_ICON_PRECOMPOSED_PATH)
|
||||
|| has_path_prefix(path, console_prefix())
|
||||
|| has_path_prefix(path, CONSOLE_PREFIX)
|
||||
}
|
||||
|
||||
/// Setup comprehensive middleware stack with tower-http features
|
||||
@@ -542,35 +487,34 @@ fn setup_console_middleware_stack(
|
||||
rate_limit_rpm: u32,
|
||||
auth_timeout: u64,
|
||||
) -> Router {
|
||||
let console_prefix = console_prefix();
|
||||
let mut app = Router::new()
|
||||
.route(FAVICON_PATH, get(static_handler))
|
||||
.route(&format!("{console_prefix}{LICENSE}"), get(license_handler))
|
||||
.route(&format!("{console_prefix}{VERSION}"), get(version_handler))
|
||||
.nest(console_prefix, Router::new().fallback_service(get(static_handler)))
|
||||
.route(&format!("{CONSOLE_PREFIX}{LICENSE}"), get(license_handler))
|
||||
.route(&format!("{CONSOLE_PREFIX}{VERSION}"), get(version_handler))
|
||||
.nest(CONSOLE_PREFIX, Router::new().fallback_service(get(static_handler)))
|
||||
.fallback_service(get(static_handler));
|
||||
|
||||
if rustfs_utils::get_env_bool(rustfs_config::ENV_HEALTH_ENDPOINT_ENABLE, rustfs_config::DEFAULT_HEALTH_ENDPOINT_ENABLE) {
|
||||
app = app
|
||||
.route(&format!("{console_prefix}{HEALTH_PREFIX}"), get(health_check).head(health_check))
|
||||
.route(&format!("{CONSOLE_PREFIX}{HEALTH_PREFIX}"), get(health_check).head(health_check))
|
||||
.route(
|
||||
&format!("{console_prefix}{}", crate::server::HEALTH_COMPAT_LIVE_PATH),
|
||||
&format!("{CONSOLE_PREFIX}{}", crate::server::HEALTH_COMPAT_LIVE_PATH),
|
||||
get(health_check).head(health_check),
|
||||
)
|
||||
.route(&format!("{console_prefix}{HEALTH_READY_PATH}"), get(health_check).head(health_check));
|
||||
.route(&format!("{CONSOLE_PREFIX}{HEALTH_READY_PATH}"), get(health_check).head(health_check));
|
||||
} else {
|
||||
// Keep disabled health probes from falling through to the SPA fallback.
|
||||
app = app
|
||||
.route(
|
||||
&format!("{console_prefix}{HEALTH_PREFIX}"),
|
||||
&format!("{CONSOLE_PREFIX}{HEALTH_PREFIX}"),
|
||||
get(health_route_disabled).head(health_route_disabled),
|
||||
)
|
||||
.route(
|
||||
&format!("{console_prefix}{}", crate::server::HEALTH_COMPAT_LIVE_PATH),
|
||||
&format!("{CONSOLE_PREFIX}{}", crate::server::HEALTH_COMPAT_LIVE_PATH),
|
||||
get(health_route_disabled).head(health_route_disabled),
|
||||
)
|
||||
.route(
|
||||
&format!("{console_prefix}{HEALTH_READY_PATH}"),
|
||||
&format!("{CONSOLE_PREFIX}{HEALTH_READY_PATH}"),
|
||||
get(health_route_disabled).head(health_route_disabled),
|
||||
);
|
||||
}
|
||||
@@ -680,7 +624,7 @@ async fn health_check(
|
||||
uri: Uri,
|
||||
server_ctx: Option<Extension<Arc<crate::runtime_sources::ServerContextSlot>>>,
|
||||
) -> Response {
|
||||
let probe = if uri.path().strip_prefix(console_prefix()) == Some(HEALTH_READY_PATH) {
|
||||
let probe = if uri.path().strip_prefix(CONSOLE_PREFIX) == Some(HEALTH_READY_PATH) {
|
||||
HealthProbe::Readiness
|
||||
} else {
|
||||
HealthProbe::Liveness
|
||||
@@ -842,7 +786,6 @@ pub(crate) fn make_console_server() -> Router {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::server::CONSOLE_PREFIX;
|
||||
use axum::body::Body;
|
||||
use axum::routing::get;
|
||||
use http::{Request, StatusCode};
|
||||
@@ -924,14 +867,9 @@ mod tests {
|
||||
async fn console_config_handler_serializes_admin_discovery_paths() {
|
||||
init_console_cfg(IpAddr::V4(Ipv4Addr::LOCALHOST), 9001);
|
||||
|
||||
let response = config_handler(
|
||||
format!("http://127.0.0.1:9001{CONSOLE_PREFIX}/api/v1/config")
|
||||
.parse()
|
||||
.expect("console URI"),
|
||||
HeaderMap::new(),
|
||||
)
|
||||
.await
|
||||
.into_response();
|
||||
let response = config_handler(Uri::from_static("http://127.0.0.1:9001/rustfs/console/api/v1/config"), HeaderMap::new())
|
||||
.await
|
||||
.into_response();
|
||||
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let body = response.into_body();
|
||||
@@ -951,8 +889,7 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn external_admin_paths_are_not_console_paths() {
|
||||
assert!(is_console_path(&format!("{CONSOLE_PREFIX}/")));
|
||||
assert!(!is_console_path(&format!("{CONSOLE_PREFIX}-other/index.html")));
|
||||
assert!(is_console_path("/rustfs/console/"));
|
||||
assert!(is_console_path("/apple-touch-icon.png"));
|
||||
assert!(is_console_path("/apple-touch-icon-precomposed.png"));
|
||||
assert!(!is_console_path("/minio/admin/v3/info"));
|
||||
@@ -1336,87 +1273,3 @@ mod tests {
|
||||
assert!(value.get("expired").is_none());
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod console_asset_prefix_tests {
|
||||
use super::rewrite_console_asset;
|
||||
use crate::server::CONSOLE_PREFIX;
|
||||
use std::borrow::Cow;
|
||||
|
||||
#[test]
|
||||
fn rewrites_exported_assets_and_client_routes() {
|
||||
let fixtures = [
|
||||
(
|
||||
"index.html",
|
||||
r#"<script src="/rustfs/console/_next/app.js"></script>"#,
|
||||
r#"<script src="/console/_next/app.js"></script>"#,
|
||||
),
|
||||
(
|
||||
"app.js",
|
||||
r#"let base="/rustfs/console";fetch(`${host}/rustfs/console/version`)"#,
|
||||
r#"let base="/console";fetch(`${host}/console/version`)"#,
|
||||
),
|
||||
("app.css", "url(/rustfs/console/logo.svg)", "url(/console/logo.svg)"),
|
||||
(
|
||||
"route.txt",
|
||||
r#"2:I[1,["/rustfs/console/_next/app.js"],"default"]"#,
|
||||
r#"2:I[1,["/console/_next/app.js"],"default"]"#,
|
||||
),
|
||||
("config.json", r#"{"url":"\/rustfs\/console\/login"}"#, r#"{"url":"\/console\/login"}"#),
|
||||
("site.webmanifest", r#"{"start_url":"/rustfs/console/"}"#, r#"{"start_url":"/console/"}"#),
|
||||
(
|
||||
"logo.svg",
|
||||
r#"<image href="/rustfs/console/logo.png"/>"#,
|
||||
r#"<image href="/console/logo.png"/>"#,
|
||||
),
|
||||
];
|
||||
for (path, input, expected) in fixtures {
|
||||
let input = input
|
||||
.replace("/rustfs/console", CONSOLE_PREFIX)
|
||||
.replace("\\/rustfs\\/console", &CONSOLE_PREFIX.replace('/', "\\/"));
|
||||
assert_eq!(
|
||||
rewrite_console_asset(path, Cow::Borrowed(input.as_bytes()), "/console").as_ref(),
|
||||
expected.as_bytes(),
|
||||
"{path}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn restores_the_standard_path_from_an_oem_build() {
|
||||
let input = format!(r#"<script src="{CONSOLE_PREFIX}/_next/app.js"></script>"#);
|
||||
let output = rewrite_console_asset("index.html", Cow::Borrowed(input.as_bytes()), rustfs_config::DEFAULT_CONSOLE_PREFIX);
|
||||
assert_eq!(output.as_ref(), br#"<script src="/rustfs/console/_next/app.js"></script>"#);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn preserves_unrelated_urls_paths_and_default_bytes() {
|
||||
let input = format!(
|
||||
r#"["https://github.com/rustfs/console","/other{CONSOLE_PREFIX}","{CONSOLE_PREFIX}-extra","{CONSOLE_PREFIX}/index.html"]"#
|
||||
);
|
||||
let expected = format!(
|
||||
r#"["https://github.com/rustfs/console","/other{CONSOLE_PREFIX}","{CONSOLE_PREFIX}-extra","/console/index.html"]"#
|
||||
);
|
||||
assert_eq!(
|
||||
rewrite_console_asset("app.js", Cow::Borrowed(input.as_bytes()), "/console").as_ref(),
|
||||
expected.as_bytes()
|
||||
);
|
||||
let unchanged = rewrite_console_asset("app.js", Cow::Borrowed(input.as_bytes()), CONSOLE_PREFIX);
|
||||
assert!(matches!(unchanged, Cow::Borrowed(_)));
|
||||
assert_eq!(unchanged.as_ref(), input.as_bytes());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn preserves_binary_invalid_utf8_and_text_without_paths() {
|
||||
let invalid_utf8 = [b"\xff".as_slice(), CONSOLE_PREFIX.as_bytes()].concat();
|
||||
for (path, bytes) in [
|
||||
("image.png", CONSOLE_PREFIX.as_bytes()),
|
||||
("app.js", invalid_utf8.as_slice()),
|
||||
("app.js", b"https://github.com/rustfs/console".as_slice()),
|
||||
] {
|
||||
let output = rewrite_console_asset(path, Cow::Borrowed(bytes), "/console");
|
||||
assert!(matches!(output, Cow::Borrowed(_)), "{path}");
|
||||
assert_eq!(output.as_ref(), bytes);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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(),
|
||||
|
||||
@@ -24,7 +24,7 @@ use crate::admin::storage_api::config::{
|
||||
read_admin_config_without_migrate, read_admin_server_config_snapshot, save_admin_server_config_snapshot,
|
||||
};
|
||||
use crate::admin::utils::json_response;
|
||||
use crate::server::{ADMIN_PREFIX, MINIO_ADMIN_PREFIX, console_prefix};
|
||||
use crate::server::{ADMIN_PREFIX, CONSOLE_PREFIX, MINIO_ADMIN_PREFIX};
|
||||
use http::StatusCode;
|
||||
use hyper::Method;
|
||||
use matchit::Params;
|
||||
@@ -824,8 +824,7 @@ fn build_console_redirect(
|
||||
let fragment =
|
||||
build_console_callback_fragment(access_key, secret_key, session_token, expiration, redirect_after, logout_token);
|
||||
|
||||
let console_prefix = console_prefix();
|
||||
let callback_path = format!("{console_prefix}{CONSOLE_OIDC_CALLBACK_SUFFIX}");
|
||||
let callback_path = format!("{CONSOLE_PREFIX}{CONSOLE_OIDC_CALLBACK_SUFFIX}");
|
||||
if let Some(base_url) = browser_redirect_url(&callback_path)? {
|
||||
return Ok(format!("{base_url}#{fragment}"));
|
||||
}
|
||||
@@ -837,8 +836,7 @@ fn build_console_redirect(
|
||||
}
|
||||
|
||||
fn build_console_login_redirect(req: &S3Request<Body>) -> S3Result<String> {
|
||||
let console_prefix = console_prefix();
|
||||
let login_path = format!("{console_prefix}{CONSOLE_LOGIN_SUFFIX}");
|
||||
let login_path = format!("{CONSOLE_PREFIX}{CONSOLE_LOGIN_SUFFIX}");
|
||||
if let Some(url) = browser_redirect_url(&login_path)? {
|
||||
return Ok(url);
|
||||
}
|
||||
@@ -1674,26 +1672,6 @@ mod tests {
|
||||
assert_eq!(callback, "https://internal:9000/rustfs/admin/v3/oidc/callback/default");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn console_prefix_process_case_oidc() {
|
||||
if std::env::var_os("RUSTFS_TEST_CONSOLE_PREFIX_PROCESS").is_none() {
|
||||
return;
|
||||
}
|
||||
crate::server::init_console_prefix().expect("initialize console prefix");
|
||||
let prefix = console_prefix();
|
||||
let req = build_oidc_request("http://internal/rustfs/admin/v3/oidc/callback/default", Some("internal:9000"), None);
|
||||
assert_eq!(
|
||||
build_console_login_redirect(&req).expect("login URL"),
|
||||
format!("https://console.example.com{prefix}/auth/login")
|
||||
);
|
||||
let redirect = build_console_redirect(&req, "access", "secret", "token", None, None, None).expect("console callback URL");
|
||||
assert!(redirect.starts_with(&format!("https://console.example.com{prefix}/auth/oidc-callback/#")));
|
||||
assert_eq!(
|
||||
derive_callback_uri(&req, "default").expect("admin callback URL"),
|
||||
"https://console.example.com/rustfs/admin/v3/oidc/callback/default"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_build_console_redirect_uses_browser_redirect_url() {
|
||||
let req = build_oidc_request("http://internal/rustfs/admin/v3/oidc/callback/default", Some("internal:9000"), None);
|
||||
@@ -1703,7 +1681,7 @@ mod tests {
|
||||
.expect("console redirect should use browser redirect URL")
|
||||
});
|
||||
|
||||
assert!(redirect.starts_with(&format!("https://console.example.com{}/auth/oidc-callback/#", console_prefix())));
|
||||
assert!(redirect.starts_with("https://console.example.com/rustfs/console/auth/oidc-callback/#"));
|
||||
assert!(redirect.contains("redirect=%2Fbuckets"));
|
||||
assert!(redirect.contains("logoutToken=logout-token"));
|
||||
}
|
||||
@@ -1716,7 +1694,7 @@ mod tests {
|
||||
build_console_login_redirect(&req).expect("login redirect should use browser redirect URL")
|
||||
});
|
||||
|
||||
assert_eq!(redirect, format!("https://console.example.com{}/auth/login", console_prefix()));
|
||||
assert_eq!(redirect, "https://console.example.com/rustfs/console/auth/login");
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -3853,30 +3853,6 @@ fn pending_remote_peer_ids(peers: &BTreeMap<String, PeerInfo>, local_peer: &Peer
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// The peers a pending remove / rotation still has to notify: every remote
|
||||
/// peer that has not acked, with the local site excluded by the same
|
||||
/// deployment-id-or-endpoint identity [`pending_remote_peer_ids`] finalizes
|
||||
/// on. The tick-driven `local_peer` carries the node's own listen address
|
||||
/// rather than the registered site endpoint (and a handler's carries the
|
||||
/// request `Host`, which behind a load balancer differs too), so an
|
||||
/// endpoint-only check dialed the site itself, timed out against the
|
||||
/// lifecycle lock this very request holds, and reported the operation as
|
||||
/// `Partial` (backlog#2367 A-4).
|
||||
fn pending_peers_awaiting_notification<'a>(
|
||||
peers: &'a BTreeMap<String, PeerInfo>,
|
||||
local_peer: &PeerInfo,
|
||||
acked_deployment_ids: &BTreeSet<String>,
|
||||
) -> Vec<&'a PeerInfo> {
|
||||
peers
|
||||
.values()
|
||||
.filter(|peer| {
|
||||
peer.deployment_id != local_peer.deployment_id
|
||||
&& !same_identity_endpoint(&peer.endpoint, &local_peer.endpoint)
|
||||
&& !acked_deployment_ids.contains(&peer.deployment_id)
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn pending_all_remote_peers_acked(
|
||||
peers: &BTreeMap<String, PeerInfo>,
|
||||
local_peer: &PeerInfo,
|
||||
@@ -4092,7 +4068,12 @@ async fn drive_pending_rotation(pending: &PendingRotation, local_peer: &PeerInfo
|
||||
};
|
||||
|
||||
let mut peer_errors = Vec::new();
|
||||
for peer in pending_peers_awaiting_notification(&pending.peers, local_peer, &pending.acked_deployment_ids) {
|
||||
for peer in pending.peers.values() {
|
||||
if same_identity_endpoint(&peer.endpoint, &local_peer.endpoint)
|
||||
|| pending.acked_deployment_ids.contains(&peer.deployment_id)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
// A superseded join returns BEFORE `apply_iam`, so a no-op answer
|
||||
// means the peer never installed the new secret. Acking it would
|
||||
// finalize a rotation half the mesh cannot authenticate against
|
||||
@@ -4395,9 +4376,12 @@ async fn drive_pending_remove(pending_remove: &PendingRemove, local_peer: &PeerI
|
||||
if secret_candidates.is_empty() {
|
||||
peer_errors.push("site replication service account secret unavailable".to_string());
|
||||
} else {
|
||||
for peer in
|
||||
pending_peers_awaiting_notification(&pending_remove.original_peers, local_peer, &pending_remove.acked_deployment_ids)
|
||||
{
|
||||
for peer in pending_remove.original_peers.values() {
|
||||
if same_identity_endpoint(&peer.endpoint, &local_peer.endpoint)
|
||||
|| pending_remove.acked_deployment_ids.contains(&peer.deployment_id)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
if let Err(err) = PeerAdminRequest::put(
|
||||
&runtime_peer_connection(peer)?,
|
||||
SITE_REPLICATION_PEER_REMOVE_PATH,
|
||||
@@ -4947,7 +4931,7 @@ async fn ensure_site_replication_bucket_targets(bucket: &str) -> S3Result<()> {
|
||||
return Ok(());
|
||||
};
|
||||
let config = bucket_replication_config_for_target_refresh(bucket).await?;
|
||||
let written = ensure_site_replication_bucket_targets_with_runtime(
|
||||
ensure_site_replication_bucket_targets_with_runtime(
|
||||
bucket,
|
||||
&runtime.state,
|
||||
&runtime.local_peer,
|
||||
@@ -4955,11 +4939,7 @@ async fn ensure_site_replication_bucket_targets(bucket: &str) -> S3Result<()> {
|
||||
&runtime.service_account_secret_key,
|
||||
expected_incarnation_id,
|
||||
)
|
||||
.await?;
|
||||
if written {
|
||||
reload_bucket_metadata_on_peers(bucket, "site_replication_bucket_targets", false).await;
|
||||
}
|
||||
Ok(())
|
||||
.await
|
||||
}
|
||||
|
||||
async fn ensure_site_replication_bucket_setup(bucket: &str) -> S3Result<bool> {
|
||||
@@ -5029,9 +5009,6 @@ async fn cleanup_removed_site_replication_bucket(bucket: &str, removed_deploymen
|
||||
Err(err) => return Err(ApiError::from(err).into()),
|
||||
}
|
||||
|
||||
if removed > 0 {
|
||||
reload_bucket_metadata_on_peers(bucket, "site_replication_bucket_cleanup", true).await;
|
||||
}
|
||||
Ok(removed)
|
||||
}
|
||||
|
||||
@@ -5348,7 +5325,7 @@ async fn refresh_bucket_targets_after_endpoint_edit(pending_id: &str, service_ac
|
||||
let local_peer = current_local_runtime_peer(&target_state);
|
||||
let _targets_guard = lock_bucket_targets_metadata(&bucket.name).await;
|
||||
let replication_config = bucket_replication_config_for_target_refresh(&bucket.name).await?;
|
||||
let written = ensure_site_replication_bucket_targets_with_runtime(
|
||||
ensure_site_replication_bucket_targets_with_runtime(
|
||||
&bucket.name,
|
||||
&target_state,
|
||||
&local_peer,
|
||||
@@ -5357,9 +5334,6 @@ async fn refresh_bucket_targets_after_endpoint_edit(pending_id: &str, service_ac
|
||||
expected_incarnation_id,
|
||||
)
|
||||
.await?;
|
||||
if written {
|
||||
reload_bucket_metadata_on_peers(&bucket.name, "site_replication_endpoint_refresh", false).await;
|
||||
}
|
||||
|
||||
rewritten.push(bucket.name.clone());
|
||||
|
||||
@@ -5403,12 +5377,16 @@ async fn site_bucket_resync_manifest_entry(bucket: &str, peer: &PeerInfo, now: O
|
||||
..Default::default()
|
||||
};
|
||||
let _targets_guard = lock_bucket_targets_metadata(bucket).await;
|
||||
// Read what is persisted, not this node's cache: the wiring may have
|
||||
// been written by another node moments ago (`start_site_bucket_resync`
|
||||
// already reads its targets from disk), and an operator resync must see
|
||||
// the same records the drive will use.
|
||||
let (config, targets) = match site_bucket_resync_persisted_wiring(bucket).await {
|
||||
Ok(wiring) => wiring,
|
||||
let (config, _) = match metadata_sys::get_replication_config(bucket).await {
|
||||
Ok(config) => config,
|
||||
Err(err) => {
|
||||
entry.status = "failed".to_string();
|
||||
entry.err_detail = summarize_peer_error_detail(&err.to_string());
|
||||
return entry;
|
||||
}
|
||||
};
|
||||
let targets = match metadata_sys::list_bucket_targets(bucket).await {
|
||||
Ok(targets) => targets,
|
||||
Err(err) => {
|
||||
entry.status = "failed".to_string();
|
||||
entry.err_detail = summarize_peer_error_detail(&err.to_string());
|
||||
@@ -5439,15 +5417,6 @@ async fn site_bucket_resync_manifest_entry(bucket: &str, peer: &PeerInfo, now: O
|
||||
entry
|
||||
}
|
||||
|
||||
/// The persisted replication configuration and bucket targets, bypassing the
|
||||
/// node-local metadata cache. `ConfigNotFound` surfaces for a bucket without
|
||||
/// a replication configuration, matching the cached read's error.
|
||||
async fn site_bucket_resync_persisted_wiring(bucket: &str) -> Result<(ReplicationConfiguration, BucketTargets), StorageError> {
|
||||
let metadata = metadata_sys::get_config_from_disk(bucket).await?;
|
||||
let config = metadata.replication_config.ok_or(StorageError::ConfigNotFound)?;
|
||||
Ok((config, metadata.bucket_target_config.unwrap_or_default()))
|
||||
}
|
||||
|
||||
async fn start_site_bucket_resync(bucket: &str, target_arn: &str, resync_id: &str) -> ResyncBucketStatus {
|
||||
let mut bucket_status = ResyncBucketStatus {
|
||||
bucket: bucket.to_string(),
|
||||
@@ -5470,8 +5439,17 @@ async fn start_site_bucket_resync(bucket: &str, target_arn: &str, resync_id: &st
|
||||
}
|
||||
};
|
||||
|
||||
let (config, targets) = match site_bucket_resync_persisted_wiring(bucket).await {
|
||||
Ok(wiring) => wiring,
|
||||
let (config, _) = match metadata_sys::get_replication_config(bucket).await {
|
||||
Ok(config) => config,
|
||||
Err(err) => {
|
||||
bucket_status.status = "failed".to_string();
|
||||
bucket_status.err_detail = err.to_string();
|
||||
return bucket_status;
|
||||
}
|
||||
};
|
||||
|
||||
let targets = match metadata_sys::list_bucket_targets_from_disk(bucket).await {
|
||||
Ok(targets) => targets,
|
||||
Err(err) => {
|
||||
bucket_status.status = "failed".to_string();
|
||||
bucket_status.err_detail = err.to_string();
|
||||
@@ -6068,10 +6046,6 @@ async fn apply_bucket_meta_item(item: SRBucketMeta) -> S3Result<()> {
|
||||
drop(lifecycle_guard);
|
||||
drop(targets_guard);
|
||||
|
||||
if !skip_config_write {
|
||||
reload_bucket_metadata_on_peers(&item.bucket, "site_replication_bucket_meta", item.r#type == "lc-config").await;
|
||||
}
|
||||
|
||||
if item.r#type == "replication-config" {
|
||||
// Rebuild the local outbound rules too: a site that joined an already-replicated
|
||||
// bucket receives this item before it has any `site-repl-*` rule of its own.
|
||||
@@ -7458,7 +7432,6 @@ impl Operation for SRPeerBucketOpsHandler {
|
||||
)
|
||||
.await
|
||||
.map_err(ApiError::from)?;
|
||||
reload_bucket_metadata_on_peers(&bucket, "site_replication_make_bucket", false).await;
|
||||
}
|
||||
"configure-replication" => {
|
||||
store
|
||||
@@ -15461,54 +15434,4 @@ mod tests {
|
||||
"no replicated config write may bypass the source stamp"
|
||||
);
|
||||
}
|
||||
|
||||
/// backlog#2367 A-4: `remove --all` notified "the peer" at the site's own
|
||||
/// registered endpoint. The tick-driven local peer carries the node's
|
||||
/// listen address, so an endpoint-only self check let the loop dial the
|
||||
/// site itself and report `Partial: failed to notify 1 peer(s)`.
|
||||
#[test]
|
||||
fn pending_notifications_skip_the_local_site_by_deployment_id() {
|
||||
let local_registered = PeerInfo {
|
||||
deployment_id: "site-b".to_string(),
|
||||
..peer("site-b", "http://site-b.example.com:9000")
|
||||
};
|
||||
let remote = PeerInfo {
|
||||
deployment_id: "site-a".to_string(),
|
||||
..peer("site-a", "http://site-a.example.com:9000")
|
||||
};
|
||||
let acked = PeerInfo {
|
||||
deployment_id: "site-c".to_string(),
|
||||
..peer("site-c", "http://site-c.example.com:9000")
|
||||
};
|
||||
let peers = BTreeMap::from([
|
||||
(local_registered.deployment_id.clone(), local_registered.clone()),
|
||||
(remote.deployment_id.clone(), remote),
|
||||
(acked.deployment_id.clone(), acked.clone()),
|
||||
]);
|
||||
let acked_ids = BTreeSet::from([acked.deployment_id]);
|
||||
|
||||
// The tick resolves the local peer from its own listen address.
|
||||
let local_from_tick = PeerInfo {
|
||||
deployment_id: "site-b".to_string(),
|
||||
..peer("site-b", "http://127.0.0.1:9000")
|
||||
};
|
||||
let to_notify: Vec<&str> = pending_peers_awaiting_notification(&peers, &local_from_tick, &acked_ids)
|
||||
.iter()
|
||||
.map(|peer| peer.deployment_id.as_str())
|
||||
.collect();
|
||||
assert_eq!(to_notify, vec!["site-a"], "the local site and the acked peer are never dialed");
|
||||
|
||||
// Identity stays consistent with what finalization waits for.
|
||||
assert_eq!(
|
||||
pending_remote_peer_ids(&peers, &local_from_tick),
|
||||
BTreeSet::from(["site-a".to_string(), "site-c".to_string()])
|
||||
);
|
||||
|
||||
// A handler-resolved local peer (registered endpoint) agrees.
|
||||
let to_notify: Vec<&str> = pending_peers_awaiting_notification(&peers, &local_registered, &acked_ids)
|
||||
.iter()
|
||||
.map(|peer| peer.deployment_id.as_str())
|
||||
.collect();
|
||||
assert_eq!(to_notify, vec!["site-a"]);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -73,7 +73,6 @@ const SITE_REPLICATION_RESYNC_ROUTE: &str = "/rustfs/admin/v3/site-replication/r
|
||||
const SITE_REPLICATION_REPAIR_ROUTE: &str = "/rustfs/admin/v3/site-replication/repair";
|
||||
const SITE_REPLICATION_REPAIR_STATUS_ROUTE: &str = "/rustfs/admin/v3/site-replication/repair/status";
|
||||
const IAM_POLICY_ATTACH_ROUTE: &str = "/rustfs/admin/v3/idp/builtin/policy/attach";
|
||||
const DATA_USAGE_INFO_ROUTE: &str = "/rustfs/admin/v3/datausageinfo";
|
||||
const IAM_POLICY_DETACH_ROUTE: &str = "/rustfs/admin/v3/idp/builtin/policy/detach";
|
||||
const IAM_POLICY_ENTITIES_ROUTE: &str = "/rustfs/admin/v3/idp/builtin/policy-entities";
|
||||
const IAM_ACCESS_KEYS_BULK_ROUTE: &str = "/rustfs/admin/v3/list-access-keys-bulk";
|
||||
@@ -1078,25 +1077,12 @@ fn advertised_admin_capabilities() -> Vec<AdvertisedAdminCapability> {
|
||||
("admin.account.mfa", HttpMethod::Get, ACCOUNT_MFA_ROUTE),
|
||||
("admin.mfa.challenge", HttpMethod::Get, MFA_CHALLENGE_ROUTE),
|
||||
("admin.user.mfa", HttpMethod::Get, USER_MFA_ROUTE),
|
||||
// `rc du` is gated on this name. Before it was advertised the client
|
||||
// inferred it from a `1.0.0-rc.` version prefix, which no longer
|
||||
// matches once the server reports `1.0.0` (backlog#2367 E-2).
|
||||
("admin.data-usage", HttpMethod::Get, DATA_USAGE_INFO_ROUTE),
|
||||
]
|
||||
.into_iter()
|
||||
.map(|(name, method, route)| AdvertisedAdminCapability {
|
||||
name,
|
||||
status: admin_route_capability(method, route),
|
||||
})
|
||||
.chain(std::iter::once(AdvertisedAdminCapability {
|
||||
// `rc watch` streams `GET /{bucket}?events=`, a misc extension route
|
||||
// dispatched by `admin::router` rather than an admin policy route,
|
||||
// so its status is not an inventory lookup (same version-prefix
|
||||
// inference on the client as `admin.data-usage`).
|
||||
name: "listen_notification",
|
||||
status: CapabilityStatus::supported()
|
||||
.with_reason("bucket listen notification (?events=) is dispatched by the admin router"),
|
||||
}))
|
||||
.collect()
|
||||
}
|
||||
|
||||
@@ -1272,9 +1258,6 @@ mod tests {
|
||||
"admin.iam.access-keys-bulk",
|
||||
"admin.iam.access-keys-bulk.ldap",
|
||||
"admin.iam.access-keys-bulk.openid",
|
||||
// rc pinned these two by version prefix until 1.0.0 (backlog#2367 E-2).
|
||||
"admin.data-usage",
|
||||
"listen_notification",
|
||||
];
|
||||
for name in expected_supported {
|
||||
let entry = response
|
||||
|
||||
@@ -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"),
|
||||
|
||||
@@ -69,7 +69,7 @@ use super::storage_api::multipart_usecase::{
|
||||
};
|
||||
use crate::app::object::{
|
||||
ConcurrencyManager, ForegroundWriteAdmission, get_concurrency_manager, guard_put_object_body_read_timeout,
|
||||
put_object_body_read_timeout, reject_oversize_single_upload,
|
||||
put_object_body_read_timeout,
|
||||
};
|
||||
use crate::app::object_data_cache::{
|
||||
ObjectDataCacheAdapter, invalidate_object_data_cache_after_complete_multipart_success,
|
||||
@@ -1169,9 +1169,6 @@ impl DefaultMultipartUsecase {
|
||||
validate_table_catalog_object_mutation(&bucket, &key).await?;
|
||||
|
||||
let mut size = resolve_upload_part_size(&req.headers, content_length)?;
|
||||
if let Some(size) = size {
|
||||
reject_oversize_single_upload(size)?;
|
||||
}
|
||||
let mut body_stream = body.ok_or_else(|| s3_error!(IncompleteBody))?;
|
||||
let Some(store) = self.object_store() else {
|
||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
||||
@@ -3212,36 +3209,6 @@ mod tests {
|
||||
assert_eq!(err.code(), &S3ErrorCode::IncompleteBody);
|
||||
}
|
||||
|
||||
/// issue #7596: a part whose declared length exceeds the 5 GiB
|
||||
/// single-request ceiling is rejected before the body is polled or the
|
||||
/// store is consulted. Exact-cap and zero-length parts pass admission.
|
||||
#[tokio::test]
|
||||
async fn execute_upload_part_rejects_oversize_declared_part_before_reading_the_body() {
|
||||
let ceiling = i64::try_from(rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE).expect("ceiling fits i64");
|
||||
|
||||
for (declared, expect_too_large) in [(ceiling + 1, true), (ceiling, false), (0, false)] {
|
||||
let (body, polls) = crate::app::object::PollCountingBody::streaming_blob();
|
||||
let input = UploadPartInput::builder()
|
||||
.bucket("bucket".to_string())
|
||||
.key("object".to_string())
|
||||
.upload_id("upload-id".to_string())
|
||||
.part_number(1)
|
||||
.body(Some(body))
|
||||
.content_length(Some(declared))
|
||||
.build()
|
||||
.unwrap();
|
||||
let req = build_request(input, Method::PUT);
|
||||
|
||||
let err = make_usecase().execute_upload_part(req).await.unwrap_err();
|
||||
if expect_too_large {
|
||||
assert_eq!(err.code(), &S3ErrorCode::EntityTooLarge, "declared {declared}");
|
||||
assert_eq!(polls.load(std::sync::atomic::Ordering::SeqCst), 0, "body must not be polled");
|
||||
} else {
|
||||
assert_ne!(err.code(), &S3ErrorCode::EntityTooLarge, "declared {declared} must pass admission");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn execute_upload_part_rejects_invalid_part_number_before_body_lookup() {
|
||||
for part_number in [-1, 0, 10001] {
|
||||
|
||||
@@ -223,10 +223,8 @@ pub(crate) use self::extract::*;
|
||||
pub(crate) use self::get::*;
|
||||
pub(crate) use self::internal_put::*;
|
||||
pub(crate) use self::on_demand_migration_put::*;
|
||||
#[cfg(test)]
|
||||
pub(crate) use self::put::PollCountingBody;
|
||||
use self::put::*;
|
||||
pub(crate) use self::put::{guard_put_object_body_read_timeout, put_object_body_read_timeout, reject_oversize_single_upload};
|
||||
pub(crate) use self::put::{guard_put_object_body_read_timeout, put_object_body_read_timeout};
|
||||
#[cfg(test)]
|
||||
pub(crate) use self::restore::RestoreStatusCommitBarrier;
|
||||
pub(crate) use self::shared::*;
|
||||
|
||||
@@ -109,21 +109,6 @@ fn resolve_put_object_authoritative_size(headers: &HeaderMap, content_length: Op
|
||||
Ok(size)
|
||||
}
|
||||
|
||||
/// Reject a declared upload length above the single-request ceiling
|
||||
/// ([`rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE`]) with `EntityTooLarge`.
|
||||
///
|
||||
/// Applies to `PutObject` and `UploadPart`. A negative or unknown length is
|
||||
/// left to the caller's existing validation.
|
||||
pub(crate) fn reject_oversize_single_upload(size: i64) -> S3Result<()> {
|
||||
if u64::try_from(size).is_ok_and(|size| size > rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE) {
|
||||
return Err(S3Error::with_message(
|
||||
S3ErrorCode::EntityTooLarge,
|
||||
ApiError::error_code_to_message(&S3ErrorCode::EntityTooLarge),
|
||||
));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Resolve the S3 request-body inter-chunk read timeout from the environment.
|
||||
///
|
||||
/// Returns `Duration::ZERO` when disabled (`RUSTFS_HTTP_REQUEST_BODY_READ_TIMEOUT=0`),
|
||||
@@ -1302,12 +1287,6 @@ impl DefaultObjectUsecase {
|
||||
// Resolve the authoritative decoded/plain object length (rejecting negative/unknown) before anything else consumes it.
|
||||
let size = resolve_put_object_authoritative_size(&req.headers, content_length)?;
|
||||
|
||||
// The streaming-body limit (s3s `put_object_max_size`) only fires once the
|
||||
// client has already streamed 5 GiB. The declared length is authoritative,
|
||||
// so reject an oversize single PUT here, before any body byte is read
|
||||
// (issue #7596).
|
||||
reject_oversize_single_upload(size)?;
|
||||
|
||||
if let Some(limit) = max_content_length
|
||||
&& u64::try_from(size).is_ok_and(|size| size > limit)
|
||||
{
|
||||
@@ -3339,77 +3318,6 @@ mod tests {
|
||||
assert_eq!(err.code(), &S3ErrorCode::InvalidStorageClass);
|
||||
}
|
||||
|
||||
/// issue #7596: a single PUT whose declared length exceeds the 5 GiB
|
||||
/// ceiling must be rejected from the headers, before any body byte is
|
||||
/// requested.
|
||||
#[tokio::test]
|
||||
async fn execute_put_object_rejects_oversize_content_length_before_reading_the_body() {
|
||||
let ceiling = i64::try_from(rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE).expect("ceiling fits i64");
|
||||
let (body, polls) = PollCountingBody::streaming_blob();
|
||||
let input = PutObjectInput::builder()
|
||||
.bucket("test-bucket".to_string())
|
||||
.key("huge.bin".to_string())
|
||||
.body(Some(body))
|
||||
.content_length(Some(ceiling + 1))
|
||||
.build()
|
||||
.unwrap();
|
||||
|
||||
let req = build_request(input, Method::PUT);
|
||||
let usecase = DefaultObjectUsecase::without_context();
|
||||
let fs = FS::new();
|
||||
|
||||
let err = Box::pin(usecase.execute_put_object(&fs, req)).await.unwrap_err();
|
||||
assert_eq!(err.code(), &S3ErrorCode::EntityTooLarge);
|
||||
assert_eq!(polls.load(std::sync::atomic::Ordering::SeqCst), 0, "body must not be polled");
|
||||
}
|
||||
|
||||
/// Admission uses the logical object size, not the wire length: a signed
|
||||
/// aws-chunked request whose framed `Content-Length` exceeds the cap but
|
||||
/// whose decoded length is within it must not be rejected as oversize,
|
||||
/// while a decoded length above the cap must be.
|
||||
#[tokio::test]
|
||||
async fn execute_put_object_oversize_admission_uses_decoded_length_for_aws_chunked() {
|
||||
let ceiling = i64::try_from(rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE).expect("ceiling fits i64");
|
||||
let framing_overhead = 1_000_000;
|
||||
|
||||
for (decoded, expect_too_large) in [(ceiling, false), (ceiling + 1, true)] {
|
||||
let (body, polls) = PollCountingBody::streaming_blob();
|
||||
let input = PutObjectInput::builder()
|
||||
.bucket("test-bucket".to_string())
|
||||
.key("huge.bin".to_string())
|
||||
.body(Some(body))
|
||||
.content_length(Some(decoded + framing_overhead))
|
||||
.build()
|
||||
.unwrap();
|
||||
|
||||
let mut req = build_request(input, Method::PUT);
|
||||
req.headers
|
||||
.insert(http::header::CONTENT_ENCODING, HeaderValue::from_static("aws-chunked"));
|
||||
req.headers.insert(
|
||||
HeaderName::from_static("x-amz-content-sha256"),
|
||||
HeaderValue::from_static("STREAMING-AWS4-HMAC-SHA256-PAYLOAD"),
|
||||
);
|
||||
req.headers.insert(
|
||||
HeaderName::from_static("x-amz-decoded-content-length"),
|
||||
HeaderValue::from_str(&decoded.to_string()).unwrap(),
|
||||
);
|
||||
let usecase = DefaultObjectUsecase::without_context();
|
||||
let fs = FS::new();
|
||||
|
||||
let err = Box::pin(usecase.execute_put_object(&fs, req)).await.unwrap_err();
|
||||
if expect_too_large {
|
||||
assert_eq!(err.code(), &S3ErrorCode::EntityTooLarge, "decoded {decoded}");
|
||||
assert_eq!(polls.load(std::sync::atomic::Ordering::SeqCst), 0, "body must not be polled");
|
||||
} else {
|
||||
assert_ne!(
|
||||
err.code(),
|
||||
&S3ErrorCode::EntityTooLarge,
|
||||
"framed wire length above the cap must not reject a decoded length at the cap"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn execute_put_object_rejects_post_object_sse_kms_from_headers() {
|
||||
let input = PutObjectInput::builder()
|
||||
@@ -4276,55 +4184,3 @@ mod tests {
|
||||
assert!(is_err_object_not_found(&lookup_err), "{lookup_err}");
|
||||
}
|
||||
}
|
||||
|
||||
/// Test-only request body that records how often it is polled, so admission
|
||||
/// tests can prove a rejection happened before any body byte was requested.
|
||||
#[cfg(test)]
|
||||
pub(crate) struct PollCountingBody {
|
||||
pub(crate) polls: std::sync::Arc<std::sync::atomic::AtomicUsize>,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl PollCountingBody {
|
||||
pub(crate) fn streaming_blob() -> (StreamingBlob, std::sync::Arc<std::sync::atomic::AtomicUsize>) {
|
||||
let polls = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let body = StreamingBlob::new(Self {
|
||||
polls: std::sync::Arc::clone(&polls),
|
||||
});
|
||||
(body, polls)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl Stream for PollCountingBody {
|
||||
type Item = Result<Bytes, StdError>;
|
||||
|
||||
fn poll_next(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
|
||||
self.polls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
Poll::Ready(Some(Ok(Bytes::from_static(b"x"))))
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl ByteStream for PollCountingBody {}
|
||||
|
||||
#[cfg(test)]
|
||||
mod oversize_single_upload_tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn reject_oversize_single_upload_enforces_the_single_request_ceiling() {
|
||||
let ceiling = i64::try_from(rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE).expect("ceiling fits i64");
|
||||
|
||||
assert!(reject_oversize_single_upload(0).is_ok());
|
||||
assert!(reject_oversize_single_upload(ceiling).is_ok(), "exact ceiling is allowed");
|
||||
assert!(reject_oversize_single_upload(-1).is_ok(), "unknown length is left to later validation");
|
||||
|
||||
let err = reject_oversize_single_upload(ceiling + 1).expect_err("one byte over must be rejected");
|
||||
assert_eq!(*err.code(), S3ErrorCode::EntityTooLarge);
|
||||
assert_eq!(
|
||||
err.message(),
|
||||
Some(ApiError::error_code_to_message(&S3ErrorCode::EntityTooLarge).as_str())
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -478,57 +478,6 @@ fn error_chain_s3s_body_stream_error(err: &(dyn std::error::Error + 'static)) ->
|
||||
None
|
||||
}
|
||||
|
||||
/// Walk an error chain (including `io::Error` custom payloads) and return
|
||||
/// whether any link satisfies `pred`.
|
||||
fn error_chain_any(err: &(dyn std::error::Error + 'static), pred: &dyn Fn(&(dyn std::error::Error + 'static)) -> bool) -> bool {
|
||||
if pred(err) {
|
||||
return true;
|
||||
}
|
||||
if let Some(io_err) = err.downcast_ref::<std::io::Error>()
|
||||
&& let Some(inner) = io_err.get_ref()
|
||||
&& error_chain_any(inner, pred)
|
||||
{
|
||||
return true;
|
||||
}
|
||||
let mut current = err.source();
|
||||
while let Some(err) = current {
|
||||
if error_chain_any(err, pred) {
|
||||
return true;
|
||||
}
|
||||
current = err.source();
|
||||
}
|
||||
false
|
||||
}
|
||||
|
||||
/// s3s raises `BodySizeLimitExceeded` when the streaming-body budget
|
||||
/// (`put_object_max_size`) runs out mid-stream. The type lives in s3s's
|
||||
/// private `http` module, so it is recognised by its `Display` form
|
||||
/// (`body size {size} exceeds limit {limit}`), like the other s3s body-stream
|
||||
/// errors above. Switch to a typed downcast once s3s re-exports the type.
|
||||
fn is_body_size_limit_exceeded_display(err: &(dyn std::error::Error + 'static)) -> bool {
|
||||
let text = err.to_string();
|
||||
text.starts_with("body size ") && text.contains(" exceeds limit ")
|
||||
}
|
||||
|
||||
fn error_chain_has_body_size_limit_exceeded(err: &(dyn std::error::Error + 'static)) -> bool {
|
||||
error_chain_any(err, &is_body_size_limit_exceeded_display)
|
||||
}
|
||||
|
||||
/// hyper reports a request body whose connection hit EOF before
|
||||
/// `Content-Length` bytes arrived as a `Kind::Body` error carrying an
|
||||
/// `UnexpectedEof` `io::Error` (its `IncompleteBody` marker is private).
|
||||
/// That is a client-side short body, not a server fault.
|
||||
fn is_hyper_body_eof(err: &(dyn std::error::Error + 'static)) -> bool {
|
||||
err.downcast_ref::<hyper::Error>()
|
||||
.and_then(|hyper_err| std::error::Error::source(hyper_err))
|
||||
.and_then(|cause| cause.downcast_ref::<std::io::Error>())
|
||||
.is_some_and(|io_err| io_err.kind() == std::io::ErrorKind::UnexpectedEof)
|
||||
}
|
||||
|
||||
fn error_chain_has_hyper_body_eof(err: &(dyn std::error::Error + 'static)) -> bool {
|
||||
error_chain_any(err, &is_hyper_body_eof)
|
||||
}
|
||||
|
||||
impl From<ApiError> for S3Error {
|
||||
fn from(err: ApiError) -> Self {
|
||||
let status = custom_error_status(&err.code);
|
||||
@@ -586,22 +535,6 @@ impl From<StorageError> for ApiError {
|
||||
};
|
||||
}
|
||||
|
||||
if error_chain_has_body_size_limit_exceeded(inner) {
|
||||
return ApiError {
|
||||
code: S3ErrorCode::EntityTooLarge,
|
||||
message: ApiError::error_code_to_message(&S3ErrorCode::EntityTooLarge),
|
||||
source: Some(Box::new(err)),
|
||||
};
|
||||
}
|
||||
|
||||
if error_chain_has_hyper_body_eof(inner) {
|
||||
return ApiError {
|
||||
code: S3ErrorCode::IncompleteBody,
|
||||
message: ApiError::error_code_to_message(&S3ErrorCode::IncompleteBody),
|
||||
source: Some(Box::new(err)),
|
||||
};
|
||||
}
|
||||
|
||||
if matches!(s3s_body_stream_error, Some(S3sBodyStreamError::IncompleteBody)) {
|
||||
return ApiError {
|
||||
code: S3ErrorCode::IncompleteBody,
|
||||
@@ -745,22 +678,6 @@ impl From<std::io::Error> for ApiError {
|
||||
source: Some(Box::new(err)),
|
||||
};
|
||||
}
|
||||
if error_chain_has_body_size_limit_exceeded(inner) {
|
||||
return ApiError {
|
||||
code: S3ErrorCode::EntityTooLarge,
|
||||
message: ApiError::error_code_to_message(&S3ErrorCode::EntityTooLarge),
|
||||
source: Some(Box::new(err)),
|
||||
};
|
||||
}
|
||||
|
||||
if error_chain_has_hyper_body_eof(inner) {
|
||||
return ApiError {
|
||||
code: S3ErrorCode::IncompleteBody,
|
||||
message: ApiError::error_code_to_message(&S3ErrorCode::IncompleteBody),
|
||||
source: Some(Box::new(err)),
|
||||
};
|
||||
}
|
||||
|
||||
if matches!(s3s_body_stream_error, Some(S3sBodyStreamError::IncompleteBody)) {
|
||||
return ApiError {
|
||||
code: S3ErrorCode::IncompleteBody,
|
||||
@@ -1033,117 +950,6 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn body_size_limit_exceeded_maps_to_entity_too_large_across_io_boundaries() {
|
||||
// Shape observed in production (issue #7596):
|
||||
// Custom { UnexpectedEof, Custom { Other, BodySizeLimitExceeded { size, limit } } }
|
||||
let nested = || {
|
||||
IoError::new(
|
||||
ErrorKind::UnexpectedEof,
|
||||
IoError::other(MockS3sBodyStreamError("body size 16384 exceeds limit 6389")),
|
||||
)
|
||||
};
|
||||
|
||||
let direct: ApiError = nested().into();
|
||||
assert_eq!(direct.code, S3ErrorCode::EntityTooLarge);
|
||||
assert_eq!(direct.message, ApiError::error_code_to_message(&S3ErrorCode::EntityTooLarge));
|
||||
|
||||
let storage: ApiError = StorageError::Io(nested()).into();
|
||||
assert_eq!(storage.code, S3ErrorCode::EntityTooLarge);
|
||||
assert!(storage.source.is_some());
|
||||
|
||||
// An unrelated message that merely mentions a limit stays internal.
|
||||
let other: ApiError = IoError::other(MockS3sBodyStreamError("limit exceeded for something else")).into();
|
||||
assert_eq!(other.code, S3ErrorCode::InternalError);
|
||||
}
|
||||
|
||||
/// Trip s3s's real streaming-body budget with a tiny limit so the
|
||||
/// display-based matcher is checked against the pinned dependency's
|
||||
/// actual error, not only the mocked string.
|
||||
#[tokio::test]
|
||||
async fn real_s3s_body_size_limit_error_maps_to_entity_too_large() {
|
||||
use futures::StreamExt;
|
||||
|
||||
let real_error = || async {
|
||||
let mut body = s3s::Body::from(bytes::Bytes::from_static(b"hello"));
|
||||
body.set_limit(Some(4));
|
||||
body.next()
|
||||
.await
|
||||
.expect("one frame")
|
||||
.expect_err("five bytes must exceed a four-byte budget")
|
||||
};
|
||||
|
||||
let err = real_error().await;
|
||||
assert!(is_body_size_limit_exceeded_display(err.as_ref()), "unexpected display: {err}");
|
||||
|
||||
let err = real_error().await;
|
||||
let storage: ApiError = StorageError::Io(IoError::new(ErrorKind::UnexpectedEof, IoError::other(err))).into();
|
||||
assert_eq!(storage.code, S3ErrorCode::EntityTooLarge);
|
||||
assert_eq!(storage.message, ApiError::error_code_to_message(&S3ErrorCode::EntityTooLarge));
|
||||
|
||||
let err = real_error().await;
|
||||
let direct: ApiError = IoError::other(err).into();
|
||||
assert_eq!(direct.code, S3ErrorCode::EntityTooLarge);
|
||||
}
|
||||
|
||||
/// Drive a real hyper HTTP/1 server so the test sees hyper's own body EOF
|
||||
/// error (`hyper::Error(Body, UnexpectedEof, IncompleteBody)`), which has no
|
||||
/// public constructor.
|
||||
async fn capture_hyper_body_eof_error() -> hyper::Error {
|
||||
use http_body_util::BodyExt;
|
||||
use hyper::service::service_fn;
|
||||
use hyper_util::rt::TokioIo;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use tokio::io::AsyncWriteExt;
|
||||
|
||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.expect("bind");
|
||||
let addr = listener.local_addr().expect("local addr");
|
||||
let captured: Arc<Mutex<Option<hyper::Error>>> = Arc::new(Mutex::new(None));
|
||||
let server_slot = Arc::clone(&captured);
|
||||
let server = tokio::spawn(async move {
|
||||
let (stream, _) = listener.accept().await.expect("accept");
|
||||
let slot = server_slot;
|
||||
let service = service_fn(move |req: hyper::Request<hyper::body::Incoming>| {
|
||||
let slot = Arc::clone(&slot);
|
||||
async move {
|
||||
let err = req.into_body().collect().await.expect_err("short body must fail");
|
||||
*slot.lock().expect("slot") = Some(err);
|
||||
Ok::<_, std::convert::Infallible>(hyper::Response::new(String::new()))
|
||||
}
|
||||
});
|
||||
let _ = hyper::server::conn::http1::Builder::new()
|
||||
.serve_connection(TokioIo::new(stream), service)
|
||||
.await;
|
||||
});
|
||||
|
||||
let mut client = tokio::net::TcpStream::connect(addr).await.expect("connect");
|
||||
client
|
||||
.write_all(b"PUT /bucket/key HTTP/1.1\r\nHost: localhost\r\nContent-Length: 100\r\n\r\nabc")
|
||||
.await
|
||||
.expect("write partial body");
|
||||
client.shutdown().await.expect("shutdown write side");
|
||||
let _ = tokio::time::timeout(std::time::Duration::from_secs(10), server).await;
|
||||
let captured = captured.lock().expect("slot").take();
|
||||
captured.expect("hyper body error captured")
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn hyper_body_eof_maps_to_incomplete_body_across_io_boundaries() {
|
||||
let hyper_err = capture_hyper_body_eof_error().await;
|
||||
assert!(is_hyper_body_eof(&hyper_err), "unexpected hyper error shape: {hyper_err:?}");
|
||||
|
||||
// Shape observed in production (issue #7596):
|
||||
// Custom { UnexpectedEof, Custom { Other, hyper::Error(Body, UnexpectedEof, IncompleteBody) } }
|
||||
let nested = IoError::new(ErrorKind::UnexpectedEof, IoError::other(hyper_err));
|
||||
let storage: ApiError = StorageError::Io(nested).into();
|
||||
assert_eq!(storage.code, S3ErrorCode::IncompleteBody);
|
||||
assert_eq!(storage.message, ApiError::error_code_to_message(&S3ErrorCode::IncompleteBody));
|
||||
|
||||
let hyper_err = capture_hyper_body_eof_error().await;
|
||||
let direct: ApiError = IoError::other(hyper_err).into();
|
||||
assert_eq!(direct.code, S3ErrorCode::IncompleteBody);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn server_side_source_read_error_maps_to_service_unavailable_before_incomplete_body() {
|
||||
let short_source = IoError::new(ErrorKind::UnexpectedEof, rustfs_rio::IncompleteBody { remaining: 17 });
|
||||
|
||||
@@ -418,7 +418,7 @@ impl PathCategory {
|
||||
PathCategory::InternodeRpc
|
||||
} else if path.starts_with("/rustfs/admin/") || path.starts_with("/minio/admin/") {
|
||||
PathCategory::AdminApi
|
||||
} else if crate::server::has_path_prefix(path, crate::server::console_prefix()) {
|
||||
} else if path.starts_with("/rustfs/console") {
|
||||
PathCategory::Console
|
||||
} else if path == "/health"
|
||||
|| path.starts_with("/health/")
|
||||
@@ -766,10 +766,8 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn test_path_category_classify_console() {
|
||||
let prefix = crate::server::console_prefix();
|
||||
assert_eq!(PathCategory::classify(&format!("{prefix}/index.html")), PathCategory::Console);
|
||||
assert_eq!(PathCategory::classify(prefix), PathCategory::Console);
|
||||
assert_eq!(PathCategory::classify(&format!("{prefix}-other/index.html")), PathCategory::S3DataPlane);
|
||||
assert_eq!(PathCategory::classify("/rustfs/console/index.html"), PathCategory::Console);
|
||||
assert_eq!(PathCategory::classify("/rustfs/console"), PathCategory::Console);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -158,11 +158,13 @@ static HTTP_STATUS_CLASS_METRICS: std::sync::LazyLock<[HttpStatusClassMetrics; 6
|
||||
static HTTP_TRANSPORT_FAILURES_COUNTER: std::sync::LazyLock<metrics::Counter> =
|
||||
std::sync::LazyLock::new(|| counter!(METRIC_HTTP_SERVER_FAILURES_TOTAL, LABEL_HTTP_STATUS_CLASS => "transport"));
|
||||
|
||||
const RUSTFS_S3_PUT_OBJECT_MAX_SIZE: u64 = 5 * 1024 * 1024 * 1024;
|
||||
|
||||
fn rustfs_s3_config() -> S3Config {
|
||||
let mut s3_config = S3Config::default();
|
||||
s3_config.normalize_forward_slash_path = true;
|
||||
s3_config.enable_sig_v2 = true;
|
||||
s3_config.put_object_max_size = Some(rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE);
|
||||
s3_config.put_object_max_size = Some(RUSTFS_S3_PUT_OBJECT_MAX_SIZE);
|
||||
s3_config.sig_v4_allowed_services.push("s3tables".to_string());
|
||||
s3_config
|
||||
}
|
||||
@@ -964,7 +966,6 @@ pub async fn start_http_server(
|
||||
readiness: Arc<GlobalReadiness>,
|
||||
server_ctx: Arc<ServerContextSlot>,
|
||||
) -> Result<(ShutdownHandle, SocketAddr)> {
|
||||
crate::server::init_console_prefix()?;
|
||||
let server_addr = parse_and_resolve_address(config.address.as_str()).map_err(Error::other)?;
|
||||
|
||||
// The listening address and port are obtained from the parameters
|
||||
@@ -1212,7 +1213,6 @@ pub async fn start_http_server(
|
||||
let now_time = jiff::Zoned::now().strftime("%Y-%m-%d %H:%M:%S").to_string();
|
||||
if config.console_enable {
|
||||
admin::console::init_console_cfg(local_ip, local_port);
|
||||
let console_prefix = crate::server::console_prefix();
|
||||
|
||||
info!(
|
||||
target: "rustfs::console::startup",
|
||||
@@ -1220,7 +1220,7 @@ pub async fn start_http_server(
|
||||
component = LOG_COMPONENT_SERVER,
|
||||
subsystem = LOG_SUBSYSTEM_STARTUP,
|
||||
service = "console",
|
||||
endpoint = %format!("{protocol}://{local_ip_str}:{local_port}{console_prefix}/index.html"),
|
||||
endpoint = %format!("{protocol}://{local_ip_str}:{local_port}/rustfs/console/index.html"),
|
||||
"Startup endpoint available"
|
||||
);
|
||||
info!(
|
||||
@@ -1229,7 +1229,7 @@ pub async fn start_http_server(
|
||||
component = LOG_COMPONENT_SERVER,
|
||||
subsystem = LOG_SUBSYSTEM_STARTUP,
|
||||
service = "console_localhost",
|
||||
endpoint = %format!("{protocol}://127.0.0.1:{local_port}{console_prefix}/index.html"),
|
||||
endpoint = %format!("{protocol}://127.0.0.1:{local_port}/rustfs/console/index.html"),
|
||||
"Startup endpoint available"
|
||||
);
|
||||
} else {
|
||||
@@ -3049,7 +3049,7 @@ mod tests {
|
||||
assert!(s3_config.normalize_forward_slash_path);
|
||||
assert!(s3_config.normalize_content_length);
|
||||
assert!(s3_config.enable_sig_v2);
|
||||
assert_eq!(s3_config.put_object_max_size, Some(rustfs_config::MAX_SINGLE_PUT_OBJECT_SIZE));
|
||||
assert_eq!(s3_config.put_object_max_size, Some(RUSTFS_S3_PUT_OBJECT_MAX_SIZE));
|
||||
assert!(s3_config.sig_v4_allowed_services.iter().any(|service| service == "s3"));
|
||||
assert!(s3_config.sig_v4_allowed_services.iter().any(|service| service == "sts"));
|
||||
assert!(s3_config.sig_v4_allowed_services.iter().any(|service| service == "s3tables"));
|
||||
|
||||
@@ -20,10 +20,10 @@ use crate::server::RemoteAddr;
|
||||
use crate::server::cors;
|
||||
use crate::server::hybrid::{HybridBody, is_grpc_request};
|
||||
use crate::server::{
|
||||
ADMIN_PREFIX, HEALTH_COMPAT_LIVE_PATH, HEALTH_PREFIX, HEALTH_READY_PATH, HealthProbe, MINIO_ADMIN_PREFIX,
|
||||
ADMIN_PREFIX, CONSOLE_PREFIX, HEALTH_COMPAT_LIVE_PATH, HEALTH_PREFIX, HEALTH_READY_PATH, HealthProbe, MINIO_ADMIN_PREFIX,
|
||||
MINIO_ADMIN_V3_PREFIX, MINIO_HEALTH_CLUSTER_PATH, MINIO_HEALTH_CLUSTER_READ_PATH, MINIO_HEALTH_LIVE_PATH,
|
||||
MINIO_HEALTH_READY_PATH, PROFILE_CPU_PATH, PROFILE_MEMORY_PATH, RPC_PREFIX, RUSTFS_ADMIN_PREFIX, active_http_requests,
|
||||
build_health_response_parts, collect_probe_readiness, console_prefix, has_path_prefix, is_admin_path, is_table_catalog_path,
|
||||
build_health_response_parts, collect_probe_readiness, has_path_prefix, is_admin_path, is_table_catalog_path,
|
||||
kms_probe_staleness_limit, kms_ready_from_probe,
|
||||
};
|
||||
use crate::shared_types::ReadinessDegradedReason;
|
||||
@@ -625,7 +625,7 @@ where
|
||||
// Create redirect response
|
||||
let redirect_response = Response::builder()
|
||||
.status(StatusCode::FOUND)
|
||||
.header(http::header::LOCATION, format!("{}/", console_prefix()))
|
||||
.header(http::header::LOCATION, "/rustfs/console/")
|
||||
.body(HybridBody::Rest {
|
||||
rest_body: RestBody::default(),
|
||||
})
|
||||
@@ -1861,7 +1861,7 @@ fn is_object_attributes_request<B>(req: &HttpRequest<B>) -> bool {
|
||||
|| has_path_prefix(path, RUSTFS_ADMIN_PREFIX)
|
||||
|| has_path_prefix(path, MINIO_ADMIN_V3_PREFIX)
|
||||
|| is_table_catalog_path(path)
|
||||
|| has_path_prefix(path, console_prefix())
|
||||
|| has_path_prefix(path, CONSOLE_PREFIX)
|
||||
|| has_path_prefix(path, RPC_PREFIX)
|
||||
{
|
||||
return false;
|
||||
@@ -2242,74 +2242,7 @@ fn rewrite_double_slash_root(uri: &Uri) -> Option<Uri> {
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
#[tokio::test]
|
||||
async fn console_prefix_process_case_browser_redirect() {
|
||||
if std::env::var_os("RUSTFS_TEST_CONSOLE_PREFIX_PROCESS").is_none() {
|
||||
return;
|
||||
}
|
||||
crate::server::init_console_prefix().expect("initialize console prefix");
|
||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.expect("redirect listener");
|
||||
let addr = listener.local_addr().expect("redirect listener address");
|
||||
let server = tokio::spawn(async move {
|
||||
let (stream, _) = listener.accept().await.expect("redirect client");
|
||||
let inner = tower::service_fn(|_request: Request<Incoming>| async {
|
||||
Ok::<_, Infallible>(Response::new(HybridBody::<Empty<Bytes>, Empty<Bytes>>::Rest { rest_body: Empty::new() }))
|
||||
});
|
||||
let service = RedirectLayer.layer(inner);
|
||||
hyper::server::conn::http1::Builder::new()
|
||||
.serve_connection(
|
||||
hyper_util::rt::TokioIo::new(stream),
|
||||
hyper_util::service::TowerToHyperService::new(service),
|
||||
)
|
||||
.await
|
||||
.expect("redirect connection");
|
||||
});
|
||||
let client = reqwest::Client::builder()
|
||||
.no_proxy()
|
||||
.http1_only()
|
||||
.redirect(reqwest::redirect::Policy::none())
|
||||
.timeout(Duration::from_secs(5))
|
||||
.build()
|
||||
.expect("redirect client");
|
||||
let response = client
|
||||
.get(format!("http://{addr}/"))
|
||||
.header(http::header::USER_AGENT, "Mozilla/5.0")
|
||||
.header(http::header::CONNECTION, "close")
|
||||
.send()
|
||||
.await
|
||||
.expect("browser response");
|
||||
assert_eq!(response.status(), StatusCode::FOUND);
|
||||
assert_eq!(response.headers()[http::header::LOCATION], format!("{}/", console_prefix()));
|
||||
response.bytes().await.expect("redirect body");
|
||||
tokio::time::timeout(Duration::from_secs(5), server)
|
||||
.await
|
||||
.expect("bounded redirect server shutdown")
|
||||
.expect("redirect task");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn console_prefix_process_case_classification() {
|
||||
if std::env::var_os("RUSTFS_TEST_CONSOLE_PREFIX_PROCESS").is_none() {
|
||||
return;
|
||||
}
|
||||
crate::server::init_console_prefix().expect("initialize console prefix");
|
||||
let prefix = crate::server::console_prefix();
|
||||
let console_uri = format!("{prefix}/index.html").parse().expect("console URI");
|
||||
assert!(is_empty_body_console_path(&Method::GET, &console_uri));
|
||||
let request = HttpRequest::builder()
|
||||
.uri(format!("{prefix}/index.html?attributes"))
|
||||
.body(())
|
||||
.expect("console attributes request");
|
||||
assert!(!is_object_attributes_request(&request));
|
||||
let s3_request = HttpRequest::builder()
|
||||
.uri("/bucket/object?attributes")
|
||||
.body(())
|
||||
.expect("S3 attributes request");
|
||||
assert!(is_object_attributes_request(&s3_request));
|
||||
}
|
||||
|
||||
use super::*;
|
||||
use crate::server::CONSOLE_PREFIX;
|
||||
use crate::server::compress::{HttpCompressionConfig, PathAwareHttpCompressionPredicate, PathCategoryInjectionLayer};
|
||||
use crate::server::{FAVICON_PATH, LICENSE, RemoteAddr, VERSION};
|
||||
use futures::future::{Ready, ready};
|
||||
@@ -2400,7 +2333,7 @@ mod tests {
|
||||
for path in [
|
||||
"/rustfs/admin/v3/metrics",
|
||||
"/minio/admin/v3/storageinfo",
|
||||
CONSOLE_PREFIX,
|
||||
"/rustfs/console/",
|
||||
"/rustfs/rpc/test",
|
||||
"/health/ready",
|
||||
"/_iceberg/v1/config",
|
||||
@@ -2691,7 +2624,7 @@ mod tests {
|
||||
for path in [
|
||||
"/rustfs/admin/v3/info",
|
||||
"/minio/admin/v3/info",
|
||||
CONSOLE_PREFIX,
|
||||
"/rustfs/console/",
|
||||
HEALTH_PREFIX,
|
||||
"/iceberg/v1/config",
|
||||
"/rustfs/rpc/v1/read-file",
|
||||
@@ -4050,7 +3983,7 @@ mod tests {
|
||||
"/minio/admin/v3/pools/cancel?versionId=unused",
|
||||
"/rustfs/admin/v3/pools/cancel?versionId=unused",
|
||||
"/rustfs/rpc/read_file_stream?versionId=unused",
|
||||
&format!("{CONSOLE_PREFIX}/index.html?versionId=unused"),
|
||||
"/rustfs/console/index.html?versionId=unused",
|
||||
"/health?versionId=unused",
|
||||
"/health/ready?versionId=unused",
|
||||
"/profile/cpu?versionId=unused",
|
||||
|
||||
@@ -72,7 +72,7 @@ pub(crate) use prefix::{
|
||||
HEALTH_COMPAT_LIVE_PATH, HEALTH_PREFIX, HEALTH_READY_PATH, LICENSE, MINIO_ADMIN_PREFIX, MINIO_ADMIN_V3_PREFIX,
|
||||
MINIO_HEALTH_CLUSTER_PATH, MINIO_HEALTH_CLUSTER_READ_PATH, MINIO_HEALTH_LIVE_PATH, MINIO_HEALTH_READY_PATH, PROFILE_CPU_PATH,
|
||||
PROFILE_MEMORY_PATH, RPC_PREFIX, RUSTFS_ADMIN_PREFIX, TABLE_CATALOG_COMPAT_PREFIX, TABLE_CATALOG_PREFIX, TONIC_PREFIX,
|
||||
VERSION, console_prefix, has_path_prefix, init_console_prefix, is_admin_path, is_table_catalog_path,
|
||||
VERSION, has_path_prefix, is_admin_path, is_table_catalog_path,
|
||||
};
|
||||
pub(crate) use readiness::ReadinessDegradedReason;
|
||||
pub(crate) use readiness::ReadinessGateLayer;
|
||||
|
||||
+4
-219
@@ -83,81 +83,10 @@ pub(crate) const RUSTFS_ADMIN_PREFIX: &str = "/rustfs/admin/v3";
|
||||
/// MinIO-compatible admin API prefix accepted by RustFS.
|
||||
pub(crate) const MINIO_ADMIN_V3_PREFIX: &str = "/minio/admin/v3";
|
||||
|
||||
/// Console asset base path embedded at build time and used as the startup default.
|
||||
/// It must match NEXT_PUBLIC_BASE_PATH when building the bundled frontend.
|
||||
pub(crate) const CONSOLE_PREFIX: &str = match option_env!("RUSTFS_CONSOLE_BASE_PATH") {
|
||||
Some(path) if !path.is_empty() => path,
|
||||
_ => rustfs_config::DEFAULT_CONSOLE_PREFIX,
|
||||
};
|
||||
|
||||
static CONFIGURED_CONSOLE_PREFIX: std::sync::OnceLock<String> = std::sync::OnceLock::new();
|
||||
|
||||
/// The prefix is fixed before listeners start; request handling never reads the environment.
|
||||
pub(crate) fn console_prefix() -> &'static str {
|
||||
CONFIGURED_CONSOLE_PREFIX.get().map(String::as_str).unwrap_or(CONSOLE_PREFIX)
|
||||
}
|
||||
|
||||
pub(crate) fn init_console_prefix() -> std::io::Result<()> {
|
||||
let raw = match std::env::var(rustfs_config::ENV_RUSTFS_CONSOLE_PREFIX) {
|
||||
Ok(value) => value,
|
||||
Err(std::env::VarError::NotPresent) => CONSOLE_PREFIX.to_string(),
|
||||
Err(err) => return Err(std::io::Error::new(std::io::ErrorKind::InvalidInput, err)),
|
||||
};
|
||||
let prefix = validate_console_prefix(&raw)?;
|
||||
if CONFIGURED_CONSOLE_PREFIX.get_or_init(|| prefix.clone()) != &prefix {
|
||||
return Err(std::io::Error::new(
|
||||
std::io::ErrorKind::InvalidInput,
|
||||
"RUSTFS_CONSOLE_PREFIX cannot change after server initialization",
|
||||
));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn validate_console_prefix(raw: &str) -> std::io::Result<String> {
|
||||
let prefix = raw.strip_suffix('/').unwrap_or(raw);
|
||||
// Keep the value safe in HTTP headers, Axum routes, and embedded HTML/JS.
|
||||
if !prefix.starts_with('/')
|
||||
|| prefix.len() > 256
|
||||
|| prefix[1..].split('/').any(|segment| {
|
||||
segment.is_empty()
|
||||
|| matches!(segment, "." | "..")
|
||||
|| !segment
|
||||
.bytes()
|
||||
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b'~'))
|
||||
})
|
||||
{
|
||||
return Err(std::io::Error::new(
|
||||
std::io::ErrorKind::InvalidInput,
|
||||
"RUSTFS_CONSOLE_PREFIX must be a non-root absolute path of at most 256 bytes with nonempty URL-safe segments",
|
||||
));
|
||||
}
|
||||
let reserved = [
|
||||
ADMIN_PREFIX,
|
||||
MINIO_ADMIN_PREFIX,
|
||||
TABLE_CATALOG_PREFIX,
|
||||
TABLE_CATALOG_COMPAT_PREFIX,
|
||||
RPC_PREFIX,
|
||||
TONIC_PREFIX,
|
||||
"/rustfs/peer",
|
||||
HEALTH_PREFIX,
|
||||
"/minio/health",
|
||||
"/profile",
|
||||
"/index.html",
|
||||
FAVICON_PATH,
|
||||
APPLE_TOUCH_ICON_PATH,
|
||||
APPLE_TOUCH_ICON_PRECOMPOSED_PATH,
|
||||
];
|
||||
if reserved
|
||||
.iter()
|
||||
.any(|path| has_path_prefix(prefix, path) || has_path_prefix(path, prefix))
|
||||
{
|
||||
return Err(std::io::Error::new(
|
||||
std::io::ErrorKind::InvalidInput,
|
||||
"RUSTFS_CONSOLE_PREFIX overlaps a reserved server route",
|
||||
));
|
||||
}
|
||||
Ok(prefix.to_string())
|
||||
}
|
||||
/// Predefined console prefix for RustFS server routes.
|
||||
/// This prefix is used for endpoints that handle console-related tasks
|
||||
/// such as user interface and management.
|
||||
pub(crate) const CONSOLE_PREFIX: &str = "/rustfs/console";
|
||||
|
||||
/// Predefined RPC prefix for RustFS server routes.
|
||||
/// This prefix is used for endpoints that handle remote procedure calls (RPC).
|
||||
@@ -182,147 +111,3 @@ pub const LOGO: &str = r#"
|
||||
░▀░▀░▀▀▀░▀▀▀░░▀░░▀░░░▀▀▀
|
||||
|
||||
"#;
|
||||
|
||||
#[cfg(test)]
|
||||
mod console_prefix_tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn console_prefix_validation() {
|
||||
for (raw, expected) in [
|
||||
(CONSOLE_PREFIX, CONSOLE_PREFIX),
|
||||
("/console", "/console"),
|
||||
("/management/console/", "/management/console"),
|
||||
("/health-dashboard", "/health-dashboard"),
|
||||
] {
|
||||
assert_eq!(validate_console_prefix(raw).expect("valid console prefix"), expected);
|
||||
}
|
||||
for raw in [
|
||||
"",
|
||||
"/",
|
||||
"console",
|
||||
"//console",
|
||||
"/console//",
|
||||
"/a//b",
|
||||
"/a/../b",
|
||||
"/a/./b",
|
||||
"/%2e%2e",
|
||||
"/console?x=1",
|
||||
"/console#x",
|
||||
"/console\\x",
|
||||
"/a\n",
|
||||
"/{param}",
|
||||
"/<script>",
|
||||
"/控制台",
|
||||
"/rustfs",
|
||||
"/rustfs/admin",
|
||||
"/rustfs/admin/v3/ui",
|
||||
"/minio",
|
||||
"/minio/admin",
|
||||
"/health",
|
||||
"/health/ui",
|
||||
"/iceberg",
|
||||
"/_iceberg/v1",
|
||||
"/rustfs/rpc",
|
||||
"/rustfs/peer",
|
||||
"/node_service.NodeService",
|
||||
"/profile",
|
||||
"/index.html",
|
||||
"/favicon.ico",
|
||||
"/index.html",
|
||||
] {
|
||||
assert_eq!(validate_console_prefix(raw).expect_err(raw).kind(), std::io::ErrorKind::InvalidInput);
|
||||
}
|
||||
assert!(validate_console_prefix(&format!("/{}", "a".repeat(255))).is_ok());
|
||||
assert!(validate_console_prefix(&format!("/{}", "a".repeat(256))).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn configured_console_prefix_subprocesses() {
|
||||
// Startup configuration is process-wide; isolate each value from other unit tests.
|
||||
for prefix in [None, Some("/console"), Some("/management/console/")] {
|
||||
let mut command = std::process::Command::new(std::env::current_exe().expect("test executable"));
|
||||
command
|
||||
.args(["console_prefix_process_case", "--test-threads=1"])
|
||||
.env("RUSTFS_TEST_CONSOLE_PREFIX_PROCESS", "1")
|
||||
.env("RUSTFS_CONSOLE_BASE_PATH", "/runtime-ignored/console")
|
||||
.env("RUSTFS_BROWSER_REDIRECT_URL", "https://console.example.com")
|
||||
.env("RUSTFS_HEALTH_ENDPOINT_ENABLE", "true")
|
||||
.env("RUSTFS_CONSOLE_RATE_LIMIT_ENABLE", "false");
|
||||
if let Some(prefix) = prefix {
|
||||
command.env(rustfs_config::ENV_RUSTFS_CONSOLE_PREFIX, prefix);
|
||||
} else {
|
||||
command.env_remove(rustfs_config::ENV_RUSTFS_CONSOLE_PREFIX);
|
||||
}
|
||||
let output = command.output().expect("run isolated console tests");
|
||||
assert!(
|
||||
output.status.success(),
|
||||
"prefix {prefix:?}: {}{}",
|
||||
String::from_utf8_lossy(&output.stdout),
|
||||
String::from_utf8_lossy(&output.stderr)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn console_prefix_process_case_routes() {
|
||||
if std::env::var_os("RUSTFS_TEST_CONSOLE_PREFIX_PROCESS").is_none() {
|
||||
return;
|
||||
}
|
||||
use axum::body::Body;
|
||||
use http::{Request, StatusCode};
|
||||
use tower::ServiceExt;
|
||||
init_console_prefix().expect("initialize configured console prefix");
|
||||
let expected = std::env::var(rustfs_config::ENV_RUSTFS_CONSOLE_PREFIX).unwrap_or_else(|_| CONSOLE_PREFIX.to_string());
|
||||
let prefix = expected.trim_end_matches('/');
|
||||
assert_eq!(console_prefix(), prefix);
|
||||
assert!(crate::admin::console::is_console_path(&format!("{prefix}/index.html")));
|
||||
assert!(!crate::admin::console::is_console_path(&format!("{prefix}-other/index.html")));
|
||||
for path in ["/rustfs/admin/v3/info", "/minio/admin/v3/info", "/health", "/bucket/object"] {
|
||||
assert!(!crate::admin::console::is_console_path(path), "reserved or S3 path {path}");
|
||||
}
|
||||
assert_eq!(
|
||||
crate::server::compress::PathCategory::classify(&format!("{prefix}/asset.js")),
|
||||
crate::server::compress::PathCategory::Console
|
||||
);
|
||||
if prefix != CONSOLE_PREFIX {
|
||||
assert!(!crate::admin::console::is_console_path(CONSOLE_PREFIX));
|
||||
}
|
||||
crate::admin::console::init_console_cfg(std::net::Ipv4Addr::LOCALHOST.into(), 9001);
|
||||
let router = crate::admin::console::make_console_server();
|
||||
for (suffix, expected_status, expected_ready) in [
|
||||
("/health", StatusCode::OK, None),
|
||||
("/health/live", StatusCode::OK, None),
|
||||
("/health/ready", StatusCode::SERVICE_UNAVAILABLE, Some(false)),
|
||||
] {
|
||||
let response = router
|
||||
.clone()
|
||||
.oneshot(
|
||||
Request::builder()
|
||||
.uri(format!("{prefix}{suffix}"))
|
||||
.body(Body::empty())
|
||||
.expect("health request"),
|
||||
)
|
||||
.await
|
||||
.expect("health response");
|
||||
assert_eq!(response.status(), expected_status, "{suffix}");
|
||||
let body = axum::body::to_bytes(response.into_body(), 65536).await.expect("health body");
|
||||
let payload: serde_json::Value = serde_json::from_slice(&body).expect("health JSON");
|
||||
assert_eq!(payload.get("ready").and_then(serde_json::Value::as_bool), expected_ready, "{suffix}");
|
||||
}
|
||||
for suffix in ["/version", "/license"] {
|
||||
let response = router
|
||||
.clone()
|
||||
.oneshot(
|
||||
Request::builder()
|
||||
.uri(format!("{prefix}{suffix}"))
|
||||
.body(Body::empty())
|
||||
.expect("console request"),
|
||||
)
|
||||
.await
|
||||
.expect("console response");
|
||||
assert_eq!(response.status(), StatusCode::OK, "{suffix}");
|
||||
assert_eq!(response.headers()[http::header::CONTENT_TYPE], "application/json");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -54,10 +54,9 @@
|
||||
//! dimension, whose key space (bucket names) is attacker-chosen.
|
||||
|
||||
use crate::server::{
|
||||
FAVICON_PATH, HEALTH_COMPAT_LIVE_PATH, HEALTH_PREFIX, HEALTH_READY_PATH, MINIO_HEALTH_CLUSTER_PATH,
|
||||
CONSOLE_PREFIX, FAVICON_PATH, HEALTH_COMPAT_LIVE_PATH, HEALTH_PREFIX, HEALTH_READY_PATH, MINIO_HEALTH_CLUSTER_PATH,
|
||||
MINIO_HEALTH_CLUSTER_READ_PATH, MINIO_HEALTH_LIVE_PATH, MINIO_HEALTH_READY_PATH, PROFILE_CPU_PATH, PROFILE_MEMORY_PATH,
|
||||
RPC_PREFIX, RemoteAddr, TONIC_PREFIX, console_prefix, has_path_prefix, is_admin_path, is_table_catalog_path,
|
||||
strip_valid_port_suffix,
|
||||
RPC_PREFIX, RemoteAddr, TONIC_PREFIX, has_path_prefix, is_admin_path, is_table_catalog_path, strip_valid_port_suffix,
|
||||
};
|
||||
use crate::storage_api::server::layer::request_context::RequestContext;
|
||||
use bytes::Bytes;
|
||||
@@ -406,7 +405,7 @@ fn is_rate_limit_exempt_path(path: &str) -> bool {
|
||||
| FAVICON_PATH
|
||||
) || has_path_prefix(path, RPC_PREFIX)
|
||||
|| has_path_prefix(path, TONIC_PREFIX)
|
||||
|| has_path_prefix(path, console_prefix())
|
||||
|| has_path_prefix(path, CONSOLE_PREFIX)
|
||||
}
|
||||
|
||||
/// Apply the standard throttling headers shared by every rate-limited scope.
|
||||
@@ -631,18 +630,6 @@ where
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
#[test]
|
||||
fn console_prefix_process_case_classification() {
|
||||
if std::env::var_os("RUSTFS_TEST_CONSOLE_PREFIX_PROCESS").is_none() {
|
||||
return;
|
||||
}
|
||||
crate::server::init_console_prefix().expect("initialize console prefix");
|
||||
let prefix = crate::server::console_prefix();
|
||||
assert!(is_rate_limit_exempt_path(&format!("{prefix}/version")));
|
||||
assert!(!is_rate_limit_exempt_path(&format!("{prefix}-other/version")));
|
||||
assert!(!is_rate_limit_exempt_path("/bucket/object"));
|
||||
}
|
||||
|
||||
use super::*;
|
||||
use http_body_util::BodyExt;
|
||||
use serial_test::serial;
|
||||
@@ -792,7 +779,7 @@ mod tests {
|
||||
"/favicon.ico",
|
||||
"/rustfs/rpc/anything",
|
||||
"/node_service.NodeService/Ping",
|
||||
&format!("{}/index.html", console_prefix()),
|
||||
"/rustfs/console/index.html",
|
||||
] {
|
||||
assert!(is_rate_limit_exempt_path(path), "{path} must be exempt");
|
||||
}
|
||||
|
||||
@@ -137,7 +137,7 @@ fn is_probe_path(path: &str) -> bool {
|
||||
let is_prefix_probe = has_path_prefix(path, crate::server::RUSTFS_ADMIN_PREFIX)
|
||||
|| has_path_prefix(path, crate::server::MINIO_ADMIN_V3_PREFIX)
|
||||
|| is_table_catalog_path(path)
|
||||
|| has_path_prefix(path, crate::server::console_prefix())
|
||||
|| has_path_prefix(path, crate::server::CONSOLE_PREFIX)
|
||||
|| has_path_prefix(path, crate::server::RPC_PREFIX)
|
||||
|| has_path_prefix(path, crate::server::ADMIN_PREFIX)
|
||||
|| has_path_prefix(path, crate::server::MINIO_ADMIN_PREFIX)
|
||||
@@ -1158,18 +1158,6 @@ where
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
#[test]
|
||||
fn console_prefix_process_case_classification() {
|
||||
if std::env::var_os("RUSTFS_TEST_CONSOLE_PREFIX_PROCESS").is_none() {
|
||||
return;
|
||||
}
|
||||
crate::server::init_console_prefix().expect("initialize console prefix");
|
||||
let prefix = crate::server::console_prefix();
|
||||
assert!(is_probe_path(&format!("{prefix}/index.html")));
|
||||
assert!(!is_probe_path(&format!("{prefix}-other/index.html")));
|
||||
assert!(!is_probe_path("/bucket/object"));
|
||||
}
|
||||
|
||||
use super::*;
|
||||
use crate::storage_api::server::readiness::{DiskOption, new_disk};
|
||||
use rustfs_madmin::{BackendInfo, Disk};
|
||||
@@ -1636,7 +1624,7 @@ mod tests {
|
||||
assert!(is_probe_path("/rustfs/admin/v3/info"));
|
||||
assert!(is_probe_path(&format!("{}/config", crate::server::TABLE_CATALOG_PREFIX)));
|
||||
assert!(is_probe_path("/_iceberg/v1/config"));
|
||||
assert!(is_probe_path(&format!("{}/", crate::server::console_prefix())));
|
||||
assert!(is_probe_path("/rustfs/console/"));
|
||||
assert!(!is_probe_path("/minio/adminx/object"));
|
||||
assert!(!is_probe_path("/rustfs/adminx/object"));
|
||||
assert!(!is_probe_path("/bucket/object"));
|
||||
|
||||
@@ -1630,44 +1630,6 @@ pub(crate) fn build_site_replication_config(
|
||||
}
|
||||
}
|
||||
|
||||
/// Reload `bucket`'s metadata on every other node of this site after a
|
||||
/// site-replication write. Every S3 bucket-config write does this
|
||||
/// (`app::bucket_usecase::notify_bucket_metadata_reload`); the
|
||||
/// site-replication writers did not, so on a multi-node site a node other
|
||||
/// than the one that applied the write served the previous targets and
|
||||
/// rules for up to the 15-minute refresh — a `resync start` routed to such a
|
||||
/// node reported every freshly wired bucket as `Config not found` or
|
||||
/// `recorded remote target no longer exists` (backlog#2367 A-5, backlog#2195
|
||||
/// item 2). Best effort like the S3 path: the write is durable and the
|
||||
/// refresh loop is the fallback, so an unreachable node must not fail the
|
||||
/// operation that already committed.
|
||||
pub(crate) async fn reload_bucket_metadata_on_peers(bucket: &str, operation: &'static str, scanner_maintenance_change: bool) {
|
||||
if scanner_maintenance_change {
|
||||
rustfs_scanner::record_scanner_maintenance_change(bucket);
|
||||
}
|
||||
let Some(notification_sys) = crate::admin::runtime_sources::current_notification_system() else {
|
||||
return;
|
||||
};
|
||||
let result = if scanner_maintenance_change {
|
||||
notification_sys.load_bucket_metadata_for_scanner_maintenance(bucket).await
|
||||
} else {
|
||||
notification_sys.load_bucket_metadata(bucket).await
|
||||
};
|
||||
if let Err(err) = result {
|
||||
warn!(
|
||||
event = EVENT_ADMIN_SITE_REPLICATION_STATE,
|
||||
component = LOG_COMPONENT_ADMIN,
|
||||
subsystem = LOG_SUBSYSTEM_SITE_REPLICATION,
|
||||
bucket = %bucket,
|
||||
operation,
|
||||
result = "peer_metadata_reload_failed",
|
||||
error = %err,
|
||||
"admin site replication state"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// Returns whether the bucket targets were rewritten.
|
||||
pub(crate) async fn ensure_site_replication_bucket_targets_with_runtime(
|
||||
bucket: &str,
|
||||
state: &SiteReplicationState,
|
||||
@@ -1675,7 +1637,7 @@ pub(crate) async fn ensure_site_replication_bucket_targets_with_runtime(
|
||||
config: Option<&ReplicationConfiguration>,
|
||||
service_account_secret_key: &str,
|
||||
expected_incarnation_id: Uuid,
|
||||
) -> S3Result<bool> {
|
||||
) -> S3Result<()> {
|
||||
let existing = match metadata_sys::list_bucket_targets(bucket).await {
|
||||
Ok(targets) => targets,
|
||||
Err(StorageError::ConfigNotFound) => BucketTargets::default(),
|
||||
@@ -1687,7 +1649,7 @@ pub(crate) async fn ensure_site_replication_bucket_targets_with_runtime(
|
||||
let updated =
|
||||
reconcile_site_replication_bucket_targets(existing, bucket, state, local_peer, config, service_account_secret_key)?;
|
||||
if updated.targets.is_empty() {
|
||||
return Ok(false);
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let json_targets = serde_json::to_vec(&updated)
|
||||
@@ -1696,12 +1658,12 @@ pub(crate) async fn ensure_site_replication_bucket_targets_with_runtime(
|
||||
// client — noticeable now that startup reconciles all buckets, not just the one bucket
|
||||
// an operation touched.
|
||||
if json_targets == existing_json {
|
||||
return Ok(false);
|
||||
return Ok(());
|
||||
}
|
||||
metadata_sys::update_if_incarnation(bucket, BUCKET_TARGETS_FILE, json_targets, expected_incarnation_id)
|
||||
.await
|
||||
.map_err(ApiError::from)?;
|
||||
Ok(true)
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn bucket_replication_config_for_target_refresh(bucket: &str) -> S3Result<Option<ReplicationConfiguration>> {
|
||||
@@ -1712,14 +1674,13 @@ pub(crate) async fn bucket_replication_config_for_target_refresh(bucket: &str) -
|
||||
}
|
||||
}
|
||||
|
||||
/// Returns whether the replication configuration was rewritten.
|
||||
pub(crate) async fn ensure_site_replication_bucket_replication_config_with_runtime(
|
||||
bucket: &str,
|
||||
state: &SiteReplicationState,
|
||||
local_peer: &PeerInfo,
|
||||
service_account_secret_key: &str,
|
||||
expected_incarnation_id: Uuid,
|
||||
) -> S3Result<bool> {
|
||||
) -> S3Result<()> {
|
||||
let existing = match metadata_sys::get_replication_config(bucket).await {
|
||||
Ok((existing, _)) => Some(existing),
|
||||
Err(StorageError::ConfigNotFound) => None,
|
||||
@@ -1728,7 +1689,7 @@ pub(crate) async fn ensure_site_replication_bucket_replication_config_with_runti
|
||||
|
||||
let Some(desired) = build_site_replication_config(bucket, state, local_peer, service_account_secret_key, existing.as_ref())?
|
||||
else {
|
||||
return Ok(false);
|
||||
return Ok(());
|
||||
};
|
||||
|
||||
// Derived rules are state owned by this site: rebuild them from the current peer
|
||||
@@ -1760,7 +1721,7 @@ pub(crate) async fn ensure_site_replication_bucket_replication_config_with_runti
|
||||
};
|
||||
|
||||
if rules == existing_rules && role == existing_role {
|
||||
return Ok(false);
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let data = serialize(&ReplicationConfiguration { role, rules })
|
||||
@@ -1769,7 +1730,7 @@ pub(crate) async fn ensure_site_replication_bucket_replication_config_with_runti
|
||||
.await
|
||||
.map_err(ApiError::from)?;
|
||||
|
||||
Ok(true)
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn ensure_site_replication_bucket_setup_with_runtime(
|
||||
@@ -1787,9 +1748,9 @@ pub(crate) async fn ensure_site_replication_bucket_setup_with_runtime_for_incarn
|
||||
runtime: &SiteReplicationRuntime,
|
||||
expected_incarnation_id: Uuid,
|
||||
) -> S3Result<()> {
|
||||
let targets_guard = lock_bucket_targets_metadata(bucket).await;
|
||||
let _targets_guard = lock_bucket_targets_metadata(bucket).await;
|
||||
let config = bucket_replication_config_for_target_refresh(bucket).await?;
|
||||
let targets_written = ensure_site_replication_bucket_targets_with_runtime(
|
||||
ensure_site_replication_bucket_targets_with_runtime(
|
||||
bucket,
|
||||
&runtime.state,
|
||||
&runtime.local_peer,
|
||||
@@ -1798,7 +1759,7 @@ pub(crate) async fn ensure_site_replication_bucket_setup_with_runtime_for_incarn
|
||||
expected_incarnation_id,
|
||||
)
|
||||
.await?;
|
||||
let config_written = ensure_site_replication_bucket_replication_config_with_runtime(
|
||||
ensure_site_replication_bucket_replication_config_with_runtime(
|
||||
bucket,
|
||||
&runtime.state,
|
||||
&runtime.local_peer,
|
||||
@@ -1806,10 +1767,6 @@ pub(crate) async fn ensure_site_replication_bucket_setup_with_runtime_for_incarn
|
||||
expected_incarnation_id,
|
||||
)
|
||||
.await?;
|
||||
drop(targets_guard);
|
||||
if targets_written || config_written {
|
||||
reload_bucket_metadata_on_peers(bucket, "site_replication_bucket_setup", config_written).await;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -1834,7 +1791,6 @@ pub(crate) async fn ensure_site_replication_bucket_versioning(bucket: &str) -> S
|
||||
metadata_sys::update_if_incarnation(bucket, BUCKET_VERSIONING_CONFIG, bucket_versioning_xml()?, expected_incarnation_id)
|
||||
.await
|
||||
.map_err(ApiError::from)?;
|
||||
reload_bucket_metadata_on_peers(bucket, "site_replication_bucket_versioning", false).await;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -52,11 +52,7 @@ pub(crate) struct SiteReplicationRetryEvent {
|
||||
/// deletion body (if it was a deletion) recorded in
|
||||
/// [`SiteReplicationState::iam_deletion_replays`]. Only then may a
|
||||
/// successful deletion replay plus a stable snapshot resend settle the
|
||||
/// entry. Every entry this binary creates starts recorded: the IAM
|
||||
/// change hook records deletion bodies, and the other creators (the add
|
||||
/// bootstrap's snapshot send, the drain's own replay) never carry a
|
||||
/// deletion. A legacy entry persisted by a binary that predates recording
|
||||
/// (serde default `false`), or one degraded by record overflow, keeps the
|
||||
/// entry; a legacy entry (or one degraded by record overflow) keeps the
|
||||
/// escalation semantics because an unrecorded deletion may hide in it.
|
||||
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
|
||||
pub(crate) deletions_recorded: bool,
|
||||
@@ -312,12 +308,7 @@ fn push_site_replication_retry_event(
|
||||
updated_at: Some(OffsetDateTime::now_utc()),
|
||||
edit_generation: generation,
|
||||
peer_unreachable,
|
||||
// See the field doc: only a row persisted by an older binary is
|
||||
// unrecorded. Stamping at creation is what lets an entry first
|
||||
// created by the bootstrap snapshot send settle after a later
|
||||
// deletion is replayed, instead of escalating forever
|
||||
// (backlog#2367 A-3).
|
||||
deletions_recorded: true,
|
||||
deletions_recorded: false,
|
||||
});
|
||||
Ok(evicted)
|
||||
}
|
||||
@@ -607,7 +598,22 @@ pub(crate) fn record_failed_iam_delivery(
|
||||
item: &SRIAMItem,
|
||||
error: &str,
|
||||
) -> S3Result<()> {
|
||||
let existed = state
|
||||
.retry_queue
|
||||
.iter()
|
||||
.any(|event| retry_event_matches(event, peer, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH));
|
||||
upsert_site_replication_retry_event(&mut state.retry_queue, peer, SITE_REPLICATION_PEER_IAM_ITEM_WIRE_PATH, error, None)?;
|
||||
if !existed
|
||||
&& let Some(event) = state
|
||||
.retry_queue
|
||||
.iter_mut()
|
||||
.find(|event| retry_event_matches(event, peer, SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH))
|
||||
{
|
||||
// Fresh entry: every failure it will ever collapse goes through this
|
||||
// recording path, so a deletion replay plus a stable snapshot resend
|
||||
// can later settle it instead of escalating.
|
||||
event.deletions_recorded = true;
|
||||
}
|
||||
|
||||
let Some(entity) = iam_item_deletion_entity(item) else {
|
||||
return Ok(());
|
||||
@@ -1412,33 +1418,6 @@ pub(crate) fn site_replication_retry_backoff_elapsed(event: &SiteReplicationRetr
|
||||
now.unix_timestamp().saturating_sub(updated_at.unix_timestamp()) >= delay
|
||||
}
|
||||
|
||||
/// Backoff evaluation time for the heavyweight tick: halfway to the next
|
||||
/// tick. Backoffs are multiples of the tick interval, so an entry stamped δ
|
||||
/// seconds after a tick is `600 − δ` old at the next one and slipped a whole
|
||||
/// extra interval for every δ > 0 — a first replay landed at T+1200 rather
|
||||
/// than T+600 (backlog#2367 A-1). Evaluating at the midpoint bounds the slip
|
||||
/// to half an interval either way; timestamps written back stay real time.
|
||||
pub(crate) fn heavyweight_retry_drain_horizon(now: OffsetDateTime) -> OffsetDateTime {
|
||||
let half_interval = crate::site_replication_reconcile::RECONCILE_INTERVAL / 2;
|
||||
now + time::Duration::seconds(i64::try_from(half_interval.as_secs()).unwrap_or(i64::MAX))
|
||||
}
|
||||
|
||||
/// What the lightweight 30-second pass may act on. It replays bounded bucket
|
||||
/// ops only, but probes every backed-off class: promotion is a state flip
|
||||
/// the heavyweight tick then replays, so an IAM or bucket-metadata snapshot
|
||||
/// owed to a peer that came back is resent at the next tick instead of
|
||||
/// after its own backoff has fully elapsed (backlog#2367 A-1).
|
||||
pub(crate) fn lightweight_retry_drain_partition(
|
||||
state: &SiteReplicationState,
|
||||
now: OffsetDateTime,
|
||||
) -> (Vec<SiteReplicationRetryEvent>, Vec<SiteReplicationRetryEvent>) {
|
||||
let mut actionable = actionable_site_replication_retry_events(state, now);
|
||||
actionable.retain(|event| {
|
||||
classify_site_replication_retry_event(event).is_some_and(|action| is_lightweight_retry_drain_action(&action))
|
||||
});
|
||||
(actionable, deferred_site_replication_retry_events(state, now))
|
||||
}
|
||||
|
||||
/// The subset of the retry queue the background drain is allowed to touch.
|
||||
pub(crate) fn actionable_site_replication_retry_events(
|
||||
state: &SiteReplicationState,
|
||||
@@ -1687,7 +1666,14 @@ async fn drain_site_replication_retry_queue_lightweight_inner() -> S3Result<()>
|
||||
return Ok(());
|
||||
}
|
||||
let now = OffsetDateTime::now_utc();
|
||||
let (actionable, deferred) = lightweight_retry_drain_partition(&runtime.state, now);
|
||||
let mut actionable = actionable_site_replication_retry_events(&runtime.state, now);
|
||||
let mut deferred = deferred_site_replication_retry_events(&runtime.state, now);
|
||||
actionable.retain(|event| {
|
||||
classify_site_replication_retry_event(event).is_some_and(|action| is_lightweight_retry_drain_action(&action))
|
||||
});
|
||||
deferred.retain(|event| {
|
||||
classify_site_replication_retry_event(event).is_some_and(|action| is_lightweight_retry_drain_action(&action))
|
||||
});
|
||||
if actionable.is_empty() && deferred.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
@@ -1711,7 +1697,10 @@ async fn drain_site_replication_retry_queue_lightweight_inner() -> S3Result<()>
|
||||
return Ok(());
|
||||
}
|
||||
let now = OffsetDateTime::now_utc();
|
||||
let (actionable, _) = lightweight_retry_drain_partition(&runtime.state, now);
|
||||
let mut actionable = actionable_site_replication_retry_events(&runtime.state, now);
|
||||
actionable.retain(|event| {
|
||||
classify_site_replication_retry_event(event).is_some_and(|action| is_lightweight_retry_drain_action(&action))
|
||||
});
|
||||
if actionable.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
@@ -1728,9 +1717,9 @@ pub(crate) async fn drain_site_replication_retry_queue_inner() -> S3Result<()> {
|
||||
// The alert must fire even when nothing is drainable this tick —
|
||||
// escalated markers are exactly the entries the drain skips.
|
||||
log_site_replication_retry_liabilities(&runtime.state);
|
||||
let horizon = heavyweight_retry_drain_horizon(OffsetDateTime::now_utc());
|
||||
let actionable = actionable_site_replication_retry_events(&runtime.state, horizon);
|
||||
let deferred = deferred_site_replication_retry_events(&runtime.state, horizon);
|
||||
let now = OffsetDateTime::now_utc();
|
||||
let actionable = actionable_site_replication_retry_events(&runtime.state, now);
|
||||
let deferred = deferred_site_replication_retry_events(&runtime.state, now);
|
||||
if actionable.is_empty() && deferred.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
@@ -1775,8 +1764,8 @@ pub(crate) async fn drain_site_replication_retry_queue_inner() -> S3Result<()> {
|
||||
{
|
||||
return Ok(());
|
||||
}
|
||||
let horizon = heavyweight_retry_drain_horizon(OffsetDateTime::now_utc());
|
||||
let actionable = actionable_site_replication_retry_events(&runtime.state, horizon);
|
||||
let now = OffsetDateTime::now_utc();
|
||||
let actionable = actionable_site_replication_retry_events(&runtime.state, now);
|
||||
if actionable.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
@@ -845,8 +845,7 @@ fn test_record_failed_iam_delivery_records_deletions_and_flags_entry() {
|
||||
record_failed_iam_delivery(&mut state, &target, &policy_delete_item("readonly"), "peer offline").expect("record failure");
|
||||
assert_eq!(state.iam_deletion_replays.len(), 2);
|
||||
|
||||
// A legacy entry (persisted by a binary that predates recording, so it
|
||||
// deserialized with the `false` default) is never stamped.
|
||||
// A legacy entry (created without recording) is never stamped.
|
||||
let legacy = PeerInfo {
|
||||
deployment_id: "legacy-dep".to_string(),
|
||||
..peer("legacy", "https://legacy.example.com")
|
||||
@@ -860,12 +859,6 @@ fn test_record_failed_iam_delivery_records_deletions_and_flags_entry() {
|
||||
None,
|
||||
)
|
||||
.expect("upsert retry event");
|
||||
state
|
||||
.retry_queue
|
||||
.iter_mut()
|
||||
.find(|event| event.peer_deployment_id == legacy.deployment_id)
|
||||
.expect("legacy entry")
|
||||
.deletions_recorded = false;
|
||||
record_failed_iam_delivery(&mut state, &legacy, &user_delete_item("bob"), "peer offline").expect("record failure");
|
||||
let legacy_event = state
|
||||
.retry_queue
|
||||
@@ -878,43 +871,6 @@ fn test_record_failed_iam_delivery_records_deletions_and_flags_entry() {
|
||||
);
|
||||
}
|
||||
|
||||
/// backlog#2367 A-3: an entry first created by a non-deletion failure — the
|
||||
/// add bootstrap's snapshot send, or the drain's own replay — hides no
|
||||
/// unrecorded deletion, so a deletion recorded later plus a stable snapshot
|
||||
/// resend must settle it instead of escalating it to the permanent marker
|
||||
/// that only `replicate repair` clears.
|
||||
#[test]
|
||||
fn test_bootstrap_created_iam_entry_settles_after_deletion_replay() {
|
||||
let target = PeerInfo {
|
||||
deployment_id: "remote-dep".to_string(),
|
||||
..peer("remote", "https://remote.example.com")
|
||||
};
|
||||
let mut state = deletion_replay_state(&target);
|
||||
upsert_site_replication_retry_event(
|
||||
&mut state.retry_queue,
|
||||
&target,
|
||||
SITE_REPLICATION_PEER_IAM_ITEM_WIRE_PATH,
|
||||
"peer request to https://remote.example.com failed (connect): connection refused",
|
||||
None,
|
||||
)
|
||||
.expect("bootstrap send failure");
|
||||
assert!(state.retry_queue[0].deletions_recorded, "a fresh entry carries no unrecorded deletion");
|
||||
|
||||
record_failed_iam_delivery(&mut state, &target, &user_delete_item("alice"), "peer offline").expect("record failure");
|
||||
assert_eq!(state.retry_queue.len(), 1, "the hook failure collapses into the bootstrap entry");
|
||||
assert!(state.retry_queue[0].deletions_recorded);
|
||||
assert_eq!(state.iam_deletion_replays.len(), 1);
|
||||
|
||||
let observed = state.retry_queue[0].clone();
|
||||
let replayed: Vec<String> = state.iam_deletion_replays.iter().map(|record| record.id.clone()).collect();
|
||||
assert!(
|
||||
settle_replayed_iam_retry_events(&mut state, &target, &observed, &replayed),
|
||||
"the replayed deletion plus the snapshot resend settle the entry"
|
||||
);
|
||||
assert!(state.retry_queue.is_empty(), "no escalation marker may remain: {:?}", state.retry_queue);
|
||||
assert!(state.iam_deletion_replays.is_empty());
|
||||
}
|
||||
|
||||
/// Overflowing the per-peer record cap degrades the entry back to the
|
||||
/// escalation semantics: the record set is no longer complete, so a replay
|
||||
/// can no longer prove the peer converged.
|
||||
@@ -1670,74 +1626,6 @@ fn test_deferred_retry_events_do_not_probe_fresh_application_failures() {
|
||||
assert!(actionable_site_replication_retry_events(&state, now).is_empty());
|
||||
}
|
||||
|
||||
/// backlog#2367 A-1: the lightweight pass replays bucket ops only, but
|
||||
/// probes every backed-off class so a recovered peer's IAM snapshot is
|
||||
/// promoted within 30 seconds instead of waiting for the heavyweight tick
|
||||
/// to notice it.
|
||||
#[test]
|
||||
fn test_lightweight_partition_probes_snapshot_entries_but_replays_bucket_ops_only() {
|
||||
let now = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
|
||||
let mut state = SiteReplicationState::default();
|
||||
state
|
||||
.peers
|
||||
.insert("remote".to_string(), peer("remote", "https://remote.example.com"));
|
||||
|
||||
let bucket_make = "/rustfs/admin/v3/site-replication/peer/bucket-ops?bucket=photos&operation=make-with-versioning";
|
||||
let mut iam_unreachable = drain_event(
|
||||
"remote",
|
||||
SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH,
|
||||
3,
|
||||
Some(now - time::Duration::seconds(30)),
|
||||
);
|
||||
iam_unreachable.peer_unreachable = true;
|
||||
let mut bucket_unreachable = drain_event("remote", bucket_make, 3, Some(now - time::Duration::seconds(30)));
|
||||
bucket_unreachable.peer_unreachable = true;
|
||||
state.retry_queue = vec![
|
||||
iam_unreachable,
|
||||
bucket_unreachable,
|
||||
// Already promoted (or never stamped): due now.
|
||||
drain_event("remote", SITE_REPLICATION_RETRY_BUCKET_METADATA_SNAPSHOT_PATH, 1, None),
|
||||
drain_event("remote", bucket_make, 1, None),
|
||||
];
|
||||
|
||||
let (actionable, deferred) = lightweight_retry_drain_partition(&state, now);
|
||||
let deferred_paths: Vec<&str> = deferred.iter().map(|event| event.path.as_str()).collect();
|
||||
assert!(
|
||||
deferred_paths.contains(&SITE_REPLICATION_RETRY_IAM_SNAPSHOT_PATH),
|
||||
"the backed-off IAM snapshot must be probed by the lightweight pass: {deferred_paths:?}"
|
||||
);
|
||||
assert!(deferred_paths.contains(&bucket_make));
|
||||
assert_eq!(
|
||||
actionable.iter().map(|event| event.path.as_str()).collect::<Vec<_>>(),
|
||||
vec![bucket_make],
|
||||
"only the bounded bucket op is replayed by the lightweight pass"
|
||||
);
|
||||
}
|
||||
|
||||
/// backlog#2367 A-1: the heavyweight tick evaluates backoff halfway to its
|
||||
/// next tick. A first failure stamped one second after a tick is 599 s old
|
||||
/// at the next tick; without the horizon it slipped to the tick after.
|
||||
#[test]
|
||||
fn test_heavyweight_horizon_absorbs_tick_phase() {
|
||||
let now = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("timestamp");
|
||||
let horizon = heavyweight_retry_drain_horizon(now);
|
||||
assert_eq!(horizon - now, time::Duration::seconds(300));
|
||||
|
||||
let elapsed_at_horizon = |secs_ago: i64| {
|
||||
site_replication_retry_backoff_elapsed(
|
||||
&drain_event("remote", "/p", 1, Some(now - time::Duration::seconds(secs_ago))),
|
||||
horizon,
|
||||
)
|
||||
};
|
||||
// Stamped just after the previous tick: due at this tick, not the next.
|
||||
assert!(elapsed_at_horizon(599));
|
||||
// Due before the next tick's midpoint: drained now rather than a whole
|
||||
// interval late.
|
||||
assert!(elapsed_at_horizon(301));
|
||||
// Due after the midpoint: waits for the next tick.
|
||||
assert!(!elapsed_at_horizon(299));
|
||||
}
|
||||
|
||||
/// The drain settles a peer-edit success under a freshly allocated
|
||||
/// generation; legacy queue entries carry `edit_generation: None` and
|
||||
/// must be cleared by that generation-scoped settlement (`(Some, None)`
|
||||
|
||||
@@ -32,7 +32,7 @@ use tokio::time::Instant;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tracing::warn;
|
||||
|
||||
pub(crate) const RECONCILE_INTERVAL: Duration = Duration::from_secs(600);
|
||||
const RECONCILE_INTERVAL: Duration = Duration::from_secs(600);
|
||||
pub(crate) const RETRY_DRAIN_INTERVAL: Duration = Duration::from_secs(30);
|
||||
|
||||
/// A reconciler reports its own failures; the outcome carries no value because neither
|
||||
|
||||
@@ -81,7 +81,6 @@ pub(crate) async fn init_startup_listen_context(
|
||||
config: &Config,
|
||||
instance_ctx: &Arc<InstanceContext>,
|
||||
) -> Result<StartupListenContext> {
|
||||
crate::server::init_console_prefix()?;
|
||||
log_sanitized_server_config(config);
|
||||
let readiness = Arc::new(GlobalReadiness::new());
|
||||
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -15,7 +15,6 @@ ROOT = Path(__file__).resolve().parent.parent
|
||||
ALWAYS_JOBS = ("classify-changes", "typos", "quick-checks")
|
||||
CODE_JOBS = (
|
||||
"test-and-lint", "test-ilm-integration-serial", "test-and-lint-rio-v2",
|
||||
"offline-enrollment-root-boundary",
|
||||
"connect-short-credential-boundary", "test-and-lint-protocols",
|
||||
"build-rustfs-debug-binary", "uring-integration", "e2e-tests",
|
||||
"s3-implemented-tests", "s3-lifecycle-behavior-tests",
|
||||
|
||||
@@ -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