Compare commits

...

19 Commits

Author SHA1 Message Date
xiaomage e55c9ffece feat(ci): add fault-tolerance degradation suite to the functional chain
Adds a fault-tolerance suite that verifies read/write behavior under
drive and node loss against the erasure-coding contract and snapshots
health-endpoint responses at every degradation tier. Scenarios are
derived from product source (default_parity_count, erasure set sizing):

- A: single-node 4 drives (EC:2, read quorum 2): hide 1/2/3 drives
- B: multi-node 4x1 (one set of 4, EC:2): stop 1/2/3 nodes
- C: multi-node 4x4 (one set of 16, EC:4, read quorum 12): 1 node down
  lands exactly on the read-quorum boundary; 2 nodes down breaks it
- C2: multi-node 4x4 with RUSTFS_STORAGE_CLASS_STANDARD=EC:8 (read
  quorum 8, write quorum 9, lock majority 9): 2 nodes down puts reads
  inside the reported divergence window (read quorum met while the
  lock majority is broken)

By default a reads-refused-despite-met-read-quorum observation is
reported as known-divergence without failing the suite; the strict
input escalates it. Chain order becomes:
upgrade -> s3 -> kms -> tier -> storage -> heal -> pool -> security ->
replication -> fault-tolerance -> performance.

Covers the findings of the 2026-09 external degradation report
(reads at read quorum, health-endpoint readiness truthfulness,
degradedReasons capture) as automated regression probes.
2026-09-11 22:45:56 +08:00
cxymds c478a392e7 fix(heal): select writable recovery intent owner (#7668) 2026-09-11 22:38:23 +08:00
cxymds c8ccc1e198 fix(heal): recover stale peer state after restart (#7665) 2026-09-11 22:37:48 +08:00
cxymds 74ba5c205c fix(ecstore): preserve multi-pool version histories (#7658)
* fix(ecstore): preserve multi-pool version histories

* fix(ecstore): satisfy clippy in rebalance test

* fix(ecstore): preserve decommission delete compatibility
2026-09-11 22:36:15 +08:00
hector 9fcd54734b Revert "feat(ci): add fault-tolerance degradation suite to the functional chain" (#7666)
This reverts commit 6eb200022e.
2026-09-11 22:30:16 +08:00
xiaomage 6eb200022e feat(ci): add fault-tolerance degradation suite to the functional chain
Adds a fault-tolerance suite that verifies read/write behavior under
drive and node loss against the erasure-coding contract and snapshots
health-endpoint responses at every degradation tier. Scenarios are
derived from product source (default_parity_count, erasure set sizing):

- A: single-node 4 drives (EC:2, read quorum 2): hide 1/2/3 drives
- B: multi-node 4x1 (one set of 4, EC:2): stop 1/2/3 nodes
- C: multi-node 4x4 (one set of 16, EC:4, read quorum 12): 1 node down
  lands exactly on the read-quorum boundary; 2 nodes down breaks it
- C2: multi-node 4x4 with RUSTFS_STORAGE_CLASS_STANDARD=EC:8 (read
  quorum 8, write quorum 9, lock majority 9): 2 nodes down puts reads
  inside the reported divergence window (read quorum met while the
  lock majority is broken)

By default a reads-refused-despite-met-read-quorum observation is
reported as known-divergence without failing the suite; the strict
input escalates it. Chain order becomes:
upgrade -> s3 -> kms -> tier -> storage -> heal -> pool -> security ->
replication -> fault-tolerance -> performance.

Covers the findings of the 2026-09 external degradation report
(reads at read quorum, health-endpoint readiness truthfulness,
degradedReasons capture) as automated regression probes.
2026-09-11 21:53:22 +08:00
houseme 92ce782999 fix(scanner): wait for ABBA telemetry window (#7660)
Size the scanner collector wait budget from the requested sample count and interval instead of using a fixed 120 second timeout. This preserves the full duration/60+1 telemetry sample set for measured ABBA runs.

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-11 21:44:37 +08:00
cxymds bf425ca32c fix(heal): reuse API error boundary for decoding failures (#7647)
* fix(heal): reuse API error boundary for decoding failures

* Update heal.rs

Signed-off-by: houseme <housemecn@gmail.com>

* test(ecstore): synchronize batch cleanup metadata reads

Acquire object read locks while observing unversioned and explicit-version batch cleanup, then release them before waiting for progress. This prevents snapshots from spanning per-disk marker removal while retaining the existing quorum, timeout, and remote delete count assertions.

Validation: cargo fmt --all --check and git diff --check passed. The focused nextest test passed 20 stress iterations each with test-util and test-util,rio-v2, with retries disabled.

---------

Signed-off-by: houseme <housemecn@gmail.com>
Co-authored-by: houseme <housemecn@gmail.com>
2026-09-11 13:46:35 +08:00
houseme 0ae38aabbc test(e2e): target multi-set outage heal candidate (#7653)
* test(e2e): target multi-set outage heal candidate

Require the outage write used by EC8+4 multi-set root-heal evidence to miss the same erasure index owned by the selected replacement drive. This avoids accepting a candidate from a different set and turning a valid heal into a false negative.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

* fix(e2e): satisfy G14 heal lint gates

Remove clippy-only noise from the G14 multi-set heal evidence test and align the admin route policy inventory with the registered heal catch-all route.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

* update

---------

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-11 08:22:10 +08:00
houseme 77ba4f1b2e fix(e2e): record EC8+4 drive restart set size (#7648)
* fix(e2e): record EC8+4 drive restart set size

Include the erasure set drive count in the distributed EC8+4 drive restart oracle so Scanner/Heal release evidence matches the registry contract.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

* fix(s3): tighten s3s footprint ratchet

Replace release-merge s3_error! macro calls with equivalent S3Error constructors so the s3gate migration ratchet does not grow on the PR merge tree.

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

---------

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 22:52:06 +08:00
cxymds 7b6e45d372 fix(heal): preserve null-version tags and encoded object paths (#7644) 2026-09-10 21:10:32 +08:00
cxymds 09c85a5f29 test(ilm): restore noncurrent compensation tests to serial CI (#7645) 2026-09-10 20:36:32 +08:00
cxymds 00aeb12914 fix(ilm): preserve cleanup ownership on tiered overwrites (#7639)
* fix(ilm): preserve cleanup ownership on tiered overwrites

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

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

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

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

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

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

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

Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 16:42:21 +08:00
houseme d286f3d06c chore(deps): refresh release dependencies (#7632)
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
2026-09-10 16:41:31 +08:00
41 changed files with 3563 additions and 265 deletions
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-darwin=874c881d7b45f12378a5817c7f42c95c4981960a2ec9ce12dcf4af239ae1f9d5
sha256-linux=9351e25b45bf7dfce18b951a5e3740225f457cacc53b8bf9f500f6947763ec0e
sha256-darwin=15cb0cf9909bfbfc5a835fb08bd3675516c641db1b92ecc1e170a1cfa0fa2fd5
sha256-linux=3163fdd29df5def86cf511ca7db05880d5c0caaa608a7031a71394327f7c217a
+1 -1
View File
@@ -1 +1 @@
sha256=6d18f9cce820c51d5589de944e8cc185f73eeca0ea9a9916651943e3759169d0
sha256=5fbb230b89212b7c3d7229d6cef3e7e2d16f0ecfec62237ebc770785706f67d9
+7
View File
@@ -539,6 +539,13 @@ path = "junit.xml"
filter = 'package(e2e_test)'
test-group = 'e2e-cluster-nightly'
# The EC8+4 multi-set heal proof deliberately uploads a larger workload so the
# background root-heal pass can be interrupted after targeting the replacement
# drive's erasure slots. Keep the extended budget scoped to this proof case.
[[profile.e2e-nightly.overrides]]
filter = 'package(e2e_test) & test(=heal_erasure_disk_rebuild_test::tests::test_cluster_root_heal_recovers_ec84_shards_across_multi_set_after_background_target_restart)'
slow-timeout = { period = "120s", terminate-after = 12, grace-period = "10s" }
# ---------------------------------------------------------------------------
# e2e-distributed profile — 4-node 4-disk Actions suite
# ---------------------------------------------------------------------------
+2 -6
View File
@@ -335,11 +335,7 @@ jobs:
# re-enabled by backlog#1304 (restore accepts serialize on a short CAS
# guard; the copy-back no longer holds the #4877 whole-copy-back lock,
# so the mid-restore ongoing read and fast 409 rejection it asserts are
# the implemented contract). The remaining exclusions each hit a
# DIFFERENT, independent issue (all tracked under rustfs/backlog#1148;
# they keep #[ignore] with a backlog reference):
# - test_noncurrent_{expiry,transition}_still_works_after_immediate_compensation_transition:
# noncurrent transition/expiry after an immediate compensation transition.
# the implemented contract).
- name: Run ignored ILM integration tests serially
env:
# Match the measured Test and Lint link budget. The default exposed
@@ -352,7 +348,7 @@ jobs:
NEXTEST_HIDE_PROGRESS_BAR=1 timeout --verbose --signal=TERM --kill-after=30s 80m \
cargo nextest run -j1 --run-ignored ignored-only \
-p rustfs-scanner -p rustfs \
-E '(binary(lifecycle_integration_test) or (package(rustfs) and test(lifecycle_transition_api_test))) and not (test(test_noncurrent_expiry_still_works_after_immediate_compensation_transition) or test(test_noncurrent_transition_still_works_after_immediate_compensation_transition))' \
-E 'binary(lifecycle_integration_test) or (package(rustfs) and test(lifecycle_transition_api_test))' \
--status-level all --final-status-level all \
2>&1 | tee artifacts/ilm-integration/nextest.log
status=${PIPESTATUS[0]}
@@ -0,0 +1,276 @@
# RustFS Fault-Tolerance (degradation) Test
#
# Scenario suite for the 2026-09 degradation report: verifies read/write
# behavior under drive and node loss against the erasure-coding contract and
# snapshots health-endpoint responses at every tier.
#
# A single-node 4 drives (SNMD): hide 1/2/3 drives, restore
# B multi-node 4x1 (one drive per node): stop 1/2/3 nodes, restore
# C multi-node 4x4 (16 drives, EC:4): stop 1 node (read-quorum boundary),
# stop 2 nodes, restore
# C2 multi-node 4x4 with EC:8: 2 nodes down puts 8 drives online -- reads
# satisfy the EC read quorum while the lock majority is broken (the
# reported divergence window: reads 503 with lock_quorum_unavailable)
#
# Expectations come from product source (default_parity_count, erasure set
# sizing). By default a "reads refused although the read quorum is met"
# observation is reported as known-divergence without failing the suite; the
# strict input turns those into failures once the product behavior changes.
name: RustFS Fault-Tolerance Test
on:
workflow_dispatch:
inputs:
package_url:
description: 'Direct .deb URL. Required unless the nightly default is wanted.'
required: false
type: string
strict:
description: 'Fail the suite when reads are refused despite a met read quorum'
type: boolean
default: false
cleanup_before:
description: 'Reset the nodes before the test (DESTROYS existing data/config)'
type: boolean
default: true
cleanup_after:
description: 'Reset the nodes after the test (DESTROYS test data/config)'
type: boolean
default: true
repository_dispatch:
# Chain handoff: dispatched when the replication suite finishes, ahead of
# the performance suite.
types: [rustfs-chain-fault-tolerance]
permissions:
contents: read
# The suite stops services and hides drive dirs on the shared fleet; only one
# functional suite may touch the environment at a time.
concurrency:
group: rustfs-shared-functional-tests
cancel-in-progress: false
defaults:
run:
shell: bash
env:
RUSTFS_ACCESS_KEY: ${{ secrets.RUSTFS_ACCESS_KEY }}
RUSTFS_SECRET_KEY: ${{ secrets.RUSTFS_SECRET_KEY }}
RUSTFS_NODES: ${{ secrets.RUSTFS_NODES || vars.RUSTFS_NODES }}
RUSTFS_SSH_USER: ${{ secrets.RUSTFS_SSH_USER || vars.RUSTFS_SSH_USER }}
PF_TESTING_GH_TOKEN: ${{ secrets.PF_TESTING_GH_TOKEN }}
RUSTFS_NIGHTLY_PACKAGE_URL: ${{ vars.RUSTFS_NIGHTLY_PACKAGE_URL || 'https://dl.rustfs.com/artifacts/rustfs/packages/nightly/rustfs-nightly-latest.deb' }}
jobs:
fault-tolerance-test:
runs-on: smoke-testing
timeout-minutes: 480
if: ${{ github.event_name == 'workflow_dispatch' || github.event_name == 'repository_dispatch' }}
steps:
- name: Initialize functional evidence
id: evidence
run: |
set -euo pipefail
umask 077
FUNCTIONAL_ARTIFACTS_DIR="${RUNNER_TEMP}/rustfs-ft-${GITHUB_RUN_ID}-${GITHUB_RUN_ATTEMPT}"
mkdir -- "${FUNCTIONAL_ARTIFACTS_DIR}" "${FUNCTIONAL_ARTIFACTS_DIR}/evidence"
{
printf 'FUNCTIONAL_ARTIFACTS_DIR=%s\n' "${FUNCTIONAL_ARTIFACTS_DIR}"
printf 'LOG_FILE=%s/suite.log\n' "${FUNCTIONAL_ARTIFACTS_DIR}"
printf 'REPORT_FILE=%s/report.md\n' "${FUNCTIONAL_ARTIFACTS_DIR}"
printf 'EVIDENCE_DIR=%s/evidence\n' "${FUNCTIONAL_ARTIFACTS_DIR}"
} >> "${GITHUB_ENV}"
# auto-testing is private: clone it with the dedicated PF token (not
# GITHUB_TOKEN) and retry transient GitHub/network failures.
- name: Checkout auto-testing scripts (with retry)
env:
GH_TOKEN: ${{ secrets.PF_TESTING_GH_TOKEN }}
run: |
set -euo pipefail
rm -rf auto-testing
for attempt in 1 2 3 4 5; do
if gh repo clone rustfs/auto-testing auto-testing -- --depth 1 --quiet; then
echo "auto-testing cloned (attempt ${attempt})"
exit 0
fi
rm -rf auto-testing
echo "clone attempt ${attempt} failed; retrying in $((attempt * 15))s" >&2
sleep $((attempt * 15))
done
echo "ERROR: unable to clone rustfs/auto-testing after 5 attempts" >&2
exit 1
- name: Show environment
run: |
uname -a
jq --version
aws --version
df -h /data | tail -1
- name: Cleanup environment (before)
if: ${{ inputs.cleanup_before != 'false' }}
run: |
./auto-testing/rustfs-fault-tolerance-test.sh --cleanup -y --log-file "${LOG_FILE}"
- name: Run fault-tolerance scenarios (A, B, C, C2)
id: test
run: |
ARGS=(--all -y --package-url "${{ inputs.package_url || env.RUSTFS_NIGHTLY_PACKAGE_URL }}" --log-file "${LOG_FILE}")
if [ "${{ inputs.strict }}" = "true" ]; then
ARGS+=(--strict)
fi
./auto-testing/rustfs-fault-tolerance-test.sh "${ARGS[@]}"
- name: Generate report
if: ${{ always() && steps.evidence.outcome == 'success' }}
run: |
set -euo pipefail
{
echo "# RustFS fault-tolerance test report"
echo ""
echo "- Run: ${GITHUB_SERVER_URL}/${GITHUB_REPOSITORY}/actions/runs/${GITHUB_RUN_ID}"
echo "- Package: ${{ inputs.package_url || 'nightly (R2 latest)' }}"
echo "- Strict mode: ${{ inputs.strict || 'false' }}"
echo ""
echo "## Per-probe results"
echo ""
echo '```'
grep -E '^FT-(CASE|SUMMARY|REPORT)' "${LOG_FILE}" || echo "(no FT-CASE lines found)"
echo '```'
echo ""
echo "## Health snapshots"
echo ""
for f in "${FUNCTIONAL_ARTIFACTS_DIR}"/evidence/*.code; do
[ -e "${f}" ] || continue
printf '%s -> %s\n' "$(basename "${f}" .code)" "$(cat "${f}")"
done
} > "${REPORT_FILE}"
- name: File failure issue in rustfs/backlog
if: ${{ failure() && steps.evidence.outcome == 'success' }}
env:
GH_TOKEN: ${{ secrets.PF_TESTING_GH_TOKEN }}
SUITE: fault-tolerance
SUITE_LABEL: Fault-Tolerance
run: |
set -euo pipefail
TITLE="[functional][${SUITE}] ${SUITE_LABEL} suite failed (run ${GITHUB_RUN_ID})"
RUN_URL="${GITHUB_SERVER_URL}/${GITHUB_REPOSITORY}/actions/runs/${GITHUB_RUN_ID}"
EXISTING="$(gh issue list -R rustfs/backlog \
--search "in:title \"run ${GITHUB_RUN_ID}\"" \
--json number --jq '.[].number' || true)"
if [ -n "${EXISTING}" ]; then
echo "backlog issue already exists for run ${GITHUB_RUN_ID}; skipping"
exit 0
fi
redact() {
sed -E \
-e 's/(RUSTFS_(ACCESS_KEY|SECRET_KEY)[=: ]+)[^[:space:]]+/\1[REDACTED]/Ig' \
-e 's/(Authorization:).*/\1 [REDACTED]/Ig' \
-e 's/(X-Amz-Signature=)[^&[:space:]]+/\1[REDACTED]/Ig' \
-e 's/^.*(password|secret|token)[=: ].*/[REDACTED SENSITIVE LINE]/Ig'
}
BODY_FILE="$(mktemp)"
{
echo "The **${SUITE_LABEL}** functional suite failed."
echo ""
echo "- Suite: \`${SUITE}\`"
echo "- Run: ${RUN_URL}"
echo "- Attempt: ${GITHUB_RUN_ATTEMPT}"
echo "- Workflow Commit: ${GITHUB_SHA}"
echo "- Trigger: ${GITHUB_EVENT_NAME}"
echo "- Date: $(date -u +%Y-%m-%d)"
echo ""
echo "## Report (errors and symptoms)"
echo ""
if [ -s "${REPORT_FILE:-}" ]; then
redact < "${REPORT_FILE}"
elif [ -s "${LOG_FILE:-}" ]; then
echo "(report file missing; log tail below)"
echo ""
tail -n 200 "${LOG_FILE}" | redact
else
echo "(no report or log file was produced)"
fi
} | head -c 55000 > "${BODY_FILE}"
gh label create functional-test -R rustfs/backlog --color d73a4a 2>/dev/null || true
if ! gh issue create -R rustfs/backlog --title "${TITLE}" \
--body-file "${BODY_FILE}" --label functional-test; then
gh issue create -R rustfs/backlog --title "${TITLE}" --body-file "${BODY_FILE}"
fi
echo "filed backlog issue for suite ${SUITE}"
- name: Upload test logs & evidence
if: ${{ always() && steps.evidence.outcome == 'success' }}
uses: actions/upload-artifact@b7c566a772e6b6bfb58ed0dc250532a479d7789f # v6
with:
name: rustfs-fault-tolerance-${{ github.run_id }}-${{ github.run_attempt }}
path: |
${{ env.FUNCTIONAL_ARTIFACTS_DIR }}/report.md
${{ env.FUNCTIONAL_ARTIFACTS_DIR }}/suite.log
${{ env.FUNCTIONAL_ARTIFACTS_DIR }}/evidence/
if-no-files-found: warn
- name: Cleanup environment (after)
if: ${{ always() && inputs.cleanup_after != 'false' }}
run: |
./auto-testing/rustfs-fault-tolerance-test.sh --cleanup -y --log-file "${LOG_FILE}" || true
- name: "Continue functional chain (next: Performance)"
# Only chain-triggered runs forward to the next suite; standalone
# workflow_dispatch runs stop after their own cleanup. A failed
# handoff retries, then files an alert issue in rustfs/backlog.
if: ${{ always() && github.event_name == 'repository_dispatch' }}
continue-on-error: true
env:
GH_TOKEN: ${{ secrets.PF_TESTING_GH_TOKEN }}
run: |
set -uo pipefail
if [ -z "${GH_TOKEN:-}" ]; then
echo "PF_TESTING_GH_TOKEN is not configured; cannot dispatch the next suite" >&2
exit 1
fi
DISPATCHED=0
for attempt in 1 2 3; do
if gh api --method POST repos/rustfs/rustfs/dispatches \
-f event_type='rustfs-chain-performance' \
-F 'client_payload[from_suite]=fault-tolerance'; then
echo "dispatched next suite Performance (attempt ${attempt})"
DISPATCHED=1
break
fi
echo "dispatch attempt ${attempt} failed; retrying in ${attempt}0s" >&2
sleep "${attempt}0"
done
if [ "${DISPATCHED:-0}" -ne 1 ]; then
echo "ERROR: functional chain stalled: could not dispatch Performance after 3 attempts" >&2
TITLE="[functional][chain] stalled after fault-tolerance (run ${GITHUB_RUN_ID})"
BODY_FILE="$(mktemp)"
{
echo "The functional chain could not hand off from **fault-tolerance** to **Performance** after 3 attempts."
echo ""
echo "- Failed suite job: ${GITHUB_SERVER_URL}/${GITHUB_REPOSITORY}/actions/runs/${GITHUB_RUN_ID}"
echo "- Expected next event: 'rustfs-chain-performance'"
echo "- Likely cause: PF_TESTING_GH_TOKEN lacks contents:write on rustfs/rustfs, or the GitHub API was unavailable."
echo "- Recovery: re-dispatch manually with"
FENCE="$(printf "\x60\x60\x60")"; echo " ${FENCE}"
echo " gh api --method POST repos/rustfs/rustfs/dispatches -f event_type='rustfs-chain-performance'"
FENCE="$(printf "\x60\x60\x60")"; echo " ${FENCE}"
} > "${BODY_FILE}"
gh issue create -R rustfs/backlog --title "${TITLE}" \
--body-file "${BODY_FILE}" --label functional-test \
|| gh issue create -R rustfs/backlog --title "${TITLE}" --body-file "${BODY_FILE}" \
|| echo "could not file the stall alert issue either; check the token" >&2
exit 1
fi
- name: Notify on failure
if: failure()
run: |
echo "RustFS fault-tolerance test failed"
echo "Package source: ${{ inputs.package_url || 'nightly (R2 latest)' }}"
echo "See the uploaded log artifact and FT-CASE lines for details."
@@ -14,7 +14,7 @@
# Functional chain driver: runs the ten functional suites in a fixed order
# (upgrade -> s3 -> kms -> tier -> storage -> heal -> pool -> security ->
# replication -> performance). Each suite attempts the next handoff even
# replication -> fault-tolerance -> performance). Each suite attempts the next handoff even
# when its tests fail.
#
# Each suite workflow can still be dispatched standalone (workflow_dispatch);
@@ -329,7 +329,7 @@ jobs:
'
done
- name: "Continue functional chain (next: Performance)"
- name: "Continue functional chain (next: Fault tolerance)"
if: ${{ always() && github.event_name == 'repository_dispatch' }}
env:
GH_TOKEN: ${{ secrets.PF_TESTING_GH_TOKEN }}
@@ -342,9 +342,9 @@ jobs:
DISPATCHED=0
for attempt in 1 2 3; do
if gh api --method POST repos/rustfs/rustfs/dispatches \
-f event_type='rustfs-chain-performance' \
-f event_type='rustfs-chain-fault-tolerance' \
-F 'client_payload[from_suite]=replication'; then
echo "dispatched next suite Performance (attempt ${attempt})"
echo "dispatched next suite Fault tolerance (attempt ${attempt})"
DISPATCHED=1
break
fi
@@ -352,19 +352,19 @@ jobs:
sleep "${attempt}0"
done
if [ "${DISPATCHED:-0}" -ne 1 ]; then
echo "ERROR: functional chain stalled: could not dispatch Performance after 3 attempts" >&2
echo "ERROR: functional chain stalled: could not dispatch Fault tolerance after 3 attempts" >&2
TITLE="[functional][chain] stalled after replication (run ${GITHUB_RUN_ID})"
BODY_FILE="$(mktemp)"
trap 'rm -f "${BODY_FILE}"' EXIT
{
echo "The functional chain could not hand off from **replication** to **Performance** after 3 attempts."
echo "The functional chain could not hand off from **replication** to **Fault tolerance** after 3 attempts."
echo ""
echo "- Failed suite job: ${GITHUB_SERVER_URL}/${GITHUB_REPOSITORY}/actions/runs/${GITHUB_RUN_ID}"
echo "- Expected next event: 'rustfs-chain-performance'"
echo "- Expected next event: 'rustfs-chain-fault-tolerance'"
echo "- Likely cause: PF_TESTING_GH_TOKEN lacks contents:write on rustfs/rustfs, or the GitHub API was unavailable."
echo "- Recovery: re-dispatch manually with"
FENCE="$(printf "\x60\x60\x60")"; echo " ${FENCE}"
echo " gh api --method POST repos/rustfs/rustfs/dispatches -f event_type='rustfs-chain-performance'"
echo " gh api --method POST repos/rustfs/rustfs/dispatches -f event_type='rustfs-chain-fault-tolerance'"
FENCE="$(printf "\x60\x60\x60")"; echo " ${FENCE}"
} > "${BODY_FILE}"
gh issue create -R rustfs/backlog --title "${TITLE}" \
Generated
+42 -42
View File
@@ -688,9 +688,9 @@ dependencies = [
[[package]]
name = "async-compat"
version = "0.2.5"
version = "0.2.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a1ba85bc55464dcbf728b56d97e119d673f4cf9062be330a9a26f3acf504a590"
checksum = "4c97d7ff3c25d6c10d64170c12acaf5d4245e76dece3779c1d92b153a64f11df"
dependencies = [
"futures-core",
"futures-io",
@@ -1641,9 +1641,9 @@ checksum = "bef38d45163c2f1dde094a7dfd33ccf595c92905c8f8f4fdc18d06fb1037718a"
[[package]]
name = "bitflags"
version = "2.13.1"
version = "2.13.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b588b76d00fde79687d7646a9b5bdf3cc0f655e0bbd080335a95d7e96f3587da"
checksum = "3ded4057c258ba199e2d26386d3af3780957ecaee6c4ef4041c6b4b8b97c0b06"
dependencies = [
"serde_core",
]
@@ -2270,9 +2270,9 @@ dependencies = [
[[package]]
name = "console"
version = "0.16.4"
version = "0.16.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4fe5f465a4f6fee88fad41b85d990f84c835335e85b5d9e6e63e0d06d28cba7c"
checksum = "e96a4956774c13c126a8b5af4daa79384f4d826534c95a02d76afb39e2ab64e3"
dependencies = [
"encode_unicode",
"libc",
@@ -3951,7 +3951,7 @@ version = "0.3.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1e0e367e4e7da84520dedcac1901e4da967309406d1e51017ae1abfb97adbd38"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"objc2",
]
@@ -4416,7 +4416,7 @@ version = "25.12.19"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "35f6839d7b3b98adde531effaf34f0c2badc6f4735d26fe74709d8e513a96ef3"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"rustc_version",
]
@@ -5155,9 +5155,9 @@ dependencies = [
[[package]]
name = "hickory-net"
version = "0.26.2"
version = "0.26.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "084e7bd6a377435d568f652153e571b50970d7ccc1d1eeec0519f834632287e1"
checksum = "c480823ed7c2c5d0f09c41020cb6b7c28029ce60ec42dc942158dcf22f8e0a4d"
dependencies = [
"async-trait",
"cfg-if",
@@ -5179,9 +5179,9 @@ dependencies = [
[[package]]
name = "hickory-proto"
version = "0.26.2"
version = "0.26.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7e2da0694c15b44c6f68a6b05e0233617008c54080e31d6eb848d858a9c5b38d"
checksum = "12b92608f679a6fa515dd1d15c1ff89443026e391200a2c840c7afcba482893d"
dependencies = [
"data-encoding",
"idna",
@@ -5199,9 +5199,9 @@ dependencies = [
[[package]]
name = "hickory-resolver"
version = "0.26.2"
version = "0.26.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0e4f9f4603319422d482738f3f6fe5aac03157fdbfed1cd85a3ff45adb09072f"
checksum = "3f3da5255c95d5a716857d54b5b8f4e8d67c3484d3beaaaae2ce25063b3ba981"
dependencies = [
"cfg-if",
"futures-util",
@@ -5706,7 +5706,7 @@ version = "0.7.15"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ed3bd0ecfbb87805f538bb7b32e5239ca0763890c623e349860ecba69469f2bb"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"cfg-if",
"libc",
]
@@ -6212,7 +6212,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c9f8ff371890db2cf65a0758dba9a79f9cd965de369f6dbdc6581a22780af45e"
dependencies = [
"async-trait",
"bitflags 2.13.1",
"bitflags 2.13.2",
"bytes",
"chrono",
"dashmap",
@@ -6794,7 +6794,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0f27695f286b461da077b8c2f72f47feaa04ce3c3f9c0976257410e90e21208a"
dependencies = [
"base64 0.22.1",
"bitflags 2.13.1",
"bitflags 2.13.2",
"btoi",
"byteorder",
"bytes",
@@ -6826,7 +6826,7 @@ version = "0.7.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "22f9786d56d972959e1408b6a93be6af13b9c1392036c5c1fafa08a1b0c6ee87"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"byteorder",
"derive_builder",
"getset",
@@ -6874,7 +6874,7 @@ version = "0.29.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "71e2746dc3a24dd78b3cfcb7be93368c6de9963d30f43a6a73998a9cf4b17b46"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"cfg-if",
"cfg_aliases",
"libc",
@@ -6887,7 +6887,7 @@ version = "0.30.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "74523f3a35e05aba87a1d978330aef40f67b0304ac79c1c00b294c9830543db6"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"cfg-if",
"cfg_aliases",
"libc",
@@ -6899,7 +6899,7 @@ version = "0.31.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "cf20d2fde8ff38632c426f1165ed7436270b44f199fc55284c38276f9db47c3d"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"cfg-if",
"cfg_aliases",
"libc",
@@ -7100,7 +7100,7 @@ version = "0.13.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d164abbde0b3c03edb9edb9cb8d31a7f5b79015c692b7c771f6e0840e9106b9f"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"libloading",
"nvml-wrapper-sys",
"static_assertions",
@@ -7151,7 +7151,7 @@ version = "0.3.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2a180dd8642fa45cdb7dd721cd4c11b1cadd4929ce112ebd8b9f5803cc79d536"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"dispatch2",
"objc2",
]
@@ -7168,7 +7168,7 @@ version = "0.3.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e3e0adef53c21f888deb4fa59fc59f7eb17404926ee8a6f59f5df0fd7f9f3272"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"objc2",
]
@@ -7721,7 +7721,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e33c6dbf1a8fb7f71742cd70f5e9f0986e60b2d19dc0b28d9ca0d1323259274a"
dependencies = [
"arrayvec",
"bitflags 2.13.1",
"bitflags 2.13.2",
"thiserror 2.0.20",
"zerocopy",
"zerocopy-derive",
@@ -7777,7 +7777,7 @@ version = "0.1.8"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "575828d9d7d205188048eb1508560607a03d21eafdbba47b8cade1736c1c28e1"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"c-enum",
"perf-event-open-sys2",
]
@@ -8294,7 +8294,7 @@ checksum = "4b45fcc2344c680f5025fe57779faef368840d0bd1f42f216291f0dc4ace4744"
dependencies = [
"bit-set",
"bit-vec 0.8.0",
"bitflags 2.13.1",
"bitflags 2.13.2",
"num-traits",
"rand 0.9.5",
"rand_chacha 0.9.0",
@@ -8436,7 +8436,7 @@ version = "0.13.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e9f068eba8e7071c5f9511831b44f32c740d5adf574e990f946ddb53db2f314e"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"memchr",
"unicase",
]
@@ -8874,7 +8874,7 @@ version = "11.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "498cd0dc59d73224351ee52a95fee0f1a617a2eae0e7d9d720cc622c73a54186"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
]
[[package]]
@@ -8983,7 +8983,7 @@ version = "0.5.18"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ed2bf2547551a7053d6fdfafda3f938979645c44812fbfcda098faae3f1a362d"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
]
[[package]]
@@ -9310,7 +9310,7 @@ checksum = "036204edbd199552a5b3832f63c60dcdf395dc44c7f06b4af1c0e8139cc11bce"
dependencies = [
"aes 0.9.3",
"aws-lc-rs",
"bitflags 2.13.1",
"bitflags 2.13.2",
"block-padding 0.4.2",
"byteorder",
"bytes",
@@ -9391,7 +9391,7 @@ version = "3.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "093197e526668d92bba562e2bbbe98d1af9831bf080b619c736316ca1fa35101"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"bytes",
"chrono",
"dashmap",
@@ -11061,7 +11061,7 @@ version = "1.1.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b6fe4565b9518b83ef4f91bb47ce29620ca828bd32cb7e408f0062e9930ba190"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"errno",
"libc",
"linux-raw-sys",
@@ -11431,7 +11431,7 @@ version = "3.7.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b7f4bc775c73d9a02cde8bf7b2ec4c9d12743edf609006c7facc23998404cd1d"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"core-foundation 0.10.1",
"core-foundation-sys",
"libc",
@@ -12321,7 +12321,7 @@ version = "0.7.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a13f3d0daba03132c0aa9767f98351b3488edc2c100cda2d2ec2b04f3d8d3c8b"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"core-foundation 0.9.4",
"system-configuration-sys",
]
@@ -12761,9 +12761,9 @@ dependencies = [
[[package]]
name = "toml_edit"
version = "0.25.13+spec-1.1.0"
version = "0.25.14+spec-1.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "6975367e4d2ef766d86af01ffad14b622fecc8d4357a998fbc4deb6e9bacaf9b"
checksum = "d2195eec204e2764644a4ea619704f9fbe5e0673038eded55ad9956f24fca0cc"
dependencies = [
"indexmap 2.14.2",
"toml_datetime",
@@ -12876,7 +12876,7 @@ version = "0.6.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4cfcf7e2740e6fc6d4d688b4ef00650406bb94adf4731e43c096c3a19fe40840"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"bytes",
"futures-util",
"http 1.5.0",
@@ -12895,7 +12895,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "08a05a66a4fdd61cbbe0a1d755ffe0ca6aba159dd4820936a0ff8a8278245b9c"
dependencies = [
"async-compression",
"bitflags 2.13.1",
"bitflags 2.13.2",
"bytes",
"futures-core",
"futures-util",
@@ -13244,9 +13244,9 @@ checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821"
[[package]]
name = "uuid"
version = "1.26.0"
version = "1.26.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b5772d71c9be8a8a6ac2117d949c5b224c1b72241bb611d9a3012edcf8af7812"
checksum = "2ef6dac1e96601b4fb3acccccff2139741fcb757cb9a36089bf5be91cfb285ce"
dependencies = [
"getrandom 0.4.3",
"js-sys",
@@ -13820,7 +13820,7 @@ version = "0.36.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3f3fd376f71958b862e7afb20cfe5a22830e1963462f3a17f49d82a6c1d1f42d"
dependencies = [
"bitflags 2.13.1",
"bitflags 2.13.2",
"windows-sys 0.59.0",
]
+1 -1
View File
@@ -335,7 +335,7 @@ tracing-subscriber = { version = "0.3.23" }
transform-stream = "0.3.1"
url = "2.5.8"
urlencoding = "2.1.3"
uuid = { version = "1.26.0" }
uuid = { version = "1.26.1" }
vaultrs = { version = "0.8.0" }
tar = "0.4.46"
walkdir = "2.5.0"
+6
View File
@@ -61,6 +61,7 @@ pub(crate) struct VersionShardCensus {
pub has_xl_meta: bool,
pub data_dir: Option<String>,
pub erasure_index: Option<usize>,
pub erasure_distribution: Option<Vec<usize>>,
pub data_blocks: Option<usize>,
pub parity_blocks: Option<usize>,
pub expected_part_numbers: BTreeSet<usize>,
@@ -90,6 +91,7 @@ impl VersionShardCensus {
&& manifest.is_complete()
&& self.data_dir == manifest.data_dir
&& self.erasure_index == manifest.erasure_index
&& self.erasure_distribution == manifest.erasure_distribution
&& self.data_blocks == manifest.data_blocks
&& self.parity_blocks == manifest.parity_blocks
&& self.expected_part_numbers == manifest.expected_part_numbers
@@ -317,6 +319,7 @@ pub(crate) fn census_object_version_on_disk(
has_xl_meta: false,
data_dir: None,
erasure_index: None,
erasure_distribution: None,
data_blocks: None,
parity_blocks: None,
expected_part_numbers: BTreeSet::new(),
@@ -334,6 +337,7 @@ pub(crate) fn census_object_version_on_disk(
};
let data_dir = file_info.data_dir.map(|id| id.to_string());
let erasure_index = Some(file_info.erasure.index);
let erasure_distribution = Some(file_info.erasure.distribution.clone());
let inline_data_fingerprint = file_info.data.as_deref().map(shard_fingerprint).transpose()?;
let part_dir = data_dir.as_ref().map_or_else(|| object_dir.clone(), |id| object_dir.join(id));
let present_part_fingerprints = match std::fs::read_dir(&part_dir) {
@@ -366,6 +370,7 @@ pub(crate) fn census_object_version_on_disk(
has_xl_meta: true,
data_dir,
erasure_index,
erasure_distribution,
data_blocks: Some(file_info.erasure.data_blocks),
parity_blocks: Some(file_info.erasure.parity_blocks),
expected_part_numbers,
@@ -421,6 +426,7 @@ mod tests {
has_xl_meta: true,
data_dir: Some("data-dir".to_string()),
erasure_index: Some(3),
erasure_distribution: Some(vec![1, 2, 3, 4]),
data_blocks: Some(2),
parity_blocks: Some(2),
expected_part_numbers: BTreeSet::from([1]),
@@ -32,6 +32,7 @@ const EC84_NODE_COUNT: usize = 3;
const EC84_DRIVES_PER_NODE: usize = 4;
const EC84_DATA_BLOCKS: usize = 8;
const EC84_PARITY_BLOCKS: usize = 4;
const EC84_ERASURE_SET_DRIVE_COUNT: usize = EC84_DATA_BLOCKS + EC84_PARITY_BLOCKS;
const EC84_TARGET_DRIVE_RESTART_CASE: &str = "ec84-target-drive-restart";
const EC84_TARGET_DRIVE_RESTART_ORACLE: &str = "ec84-target-drive-restart.json";
const EC84_HEAL_CONTROL_READY_TIMEOUT: Duration = Duration::from_secs(45);
@@ -185,6 +186,7 @@ async fn write_scanner_heal_evidence(context: ScannerHealEvidenceContext, payloa
"binary_sha256": string_field(&context.run, "binary.sha256")?,
"test_binary_sha256": string_field(&context.run, "test_binary.sha256")?,
"topology": {"nodes": EC84_NODE_COUNT, "drives_per_node": EC84_DRIVES_PER_NODE},
"erasure_set_drive_count": EC84_ERASURE_SET_DRIVE_COUNT,
"pid_before": payload.pid_before,
"pid_after": payload.pid_after,
"unclean_shutdown_marker": false,
@@ -251,6 +253,7 @@ async fn put_large_inventory(client: &Client, bucket: &str) -> TestResult<Vec<Ex
has_xl_meta: false,
data_dir: None,
erasure_index: None,
erasure_distribution: None,
data_blocks: None,
parity_blocks: None,
expected_part_numbers: Default::default(),
@@ -410,4 +413,16 @@ mod tests {
"admin POST failed: 503 Service Unavailable cluster heal coordination unavailable".into();
assert!(!is_cluster_heal_coordination_unavailable(wrong_status.as_ref()));
}
#[test]
fn ec84_drive_restart_evidence_shape_matches_registry() {
assert_eq!(EC84_ERASURE_SET_DRIVE_COUNT, EC84_NODE_COUNT * EC84_DRIVES_PER_NODE);
let registry: Value = serde_json::from_str(include_str!("../../../../.config/scanner-heal-required-tests.json"))
.expect("scanner/heal registry is valid JSON");
let case = &registry["cases"][EC84_TARGET_DRIVE_RESTART_CASE];
assert_eq!(case["erasure_set_drive_count"], EC84_ERASURE_SET_DRIVE_COUNT);
assert_eq!(case["topology"]["nodes"], EC84_NODE_COUNT);
assert_eq!(case["topology"]["drives_per_node"], EC84_DRIVES_PER_NODE);
}
}
@@ -24,7 +24,7 @@ mod tests {
use crate::storage_api::RUSTFS_META_BUCKET;
use aws_sdk_s3::{
error::{ProvideErrorMetadata, SdkError},
operation::put_object::PutObjectError,
operation::{delete_object::DeleteObjectError, put_object::PutObjectError},
primitives::ByteStream,
};
use http::Method;
@@ -45,6 +45,7 @@ mod tests {
struct ReplacementDriveSelection {
replaced_disk: PathBuf,
drive_index: usize,
replacement_format_path: PathBuf,
replacement_format: Vec<u8>,
expected_pool_metadata: Option<VersionShardCensus>,
@@ -461,6 +462,12 @@ mod tests {
shard_census: VersionShardCensus,
}
#[derive(Debug, Default)]
struct OutagePeerManifest {
erasure_indices: HashSet<usize>,
erasure_distribution: Option<Vec<usize>>,
}
fn deterministic_object_body(len: usize, seed: u8) -> Vec<u8> {
let mut value = seed;
std::iter::repeat_with(|| {
@@ -486,6 +493,81 @@ mod tests {
Ok(matching)
}
fn collect_outage_peer_manifest(
cluster: &RustFSTestClusterEnvironment,
offline_node_index: usize,
bucket: &str,
key: &str,
erasure_set_drive_count: usize,
) -> Result<OutagePeerManifest, Box<dyn Error + Send + Sync>> {
let mut manifest = OutagePeerManifest::default();
for (node_index, node) in cluster.nodes.iter().enumerate() {
if node_index == offline_node_index {
continue;
}
for (drive_index, drive) in node.data_dirs.iter().enumerate() {
let census = census_object_version_on_disk(Path::new(drive), bucket, key, None)?;
if !census.has_xl_meta {
continue;
}
assert!(
census.is_complete(),
"online node {node_index} drive {drive_index} must hold a complete outage-object shard: {census:?}"
);
let erasure_index = census.erasure_index.ok_or_else(|| {
format!("online node {node_index} drive {drive_index} outage-object shard has no erasure index: {census:?}")
})?;
assert!(
(1..=erasure_set_drive_count).contains(&erasure_index),
"online node {node_index} drive {drive_index} outage-object erasure index is out of range: {census:?}"
);
let distribution = census.erasure_distribution.as_ref().ok_or_else(|| {
format!(
"online node {node_index} drive {drive_index} outage-object shard has no erasure distribution: {census:?}"
)
})?;
assert_eq!(
distribution.len(),
erasure_set_drive_count,
"online node {node_index} drive {drive_index} outage-object distribution must match the erasure set: {census:?}"
);
match &manifest.erasure_distribution {
Some(existing) => {
assert_eq!(existing, distribution, "outage-object shards must agree on one erasure distribution")
}
None => manifest.erasure_distribution = Some(distribution.clone()),
}
assert!(
manifest.erasure_indices.insert(erasure_index),
"outage-object erasure index {erasure_index} is duplicated across online drives"
);
}
}
Ok(manifest)
}
fn outage_candidate_replacement_erasure_index(
peer_manifest: &OutagePeerManifest,
erasure_set_drive_count: usize,
replacement_set_slot: usize,
) -> Option<usize> {
let distribution = peer_manifest.erasure_distribution.as_ref()?;
distribution.get(replacement_set_slot).copied().filter(|replacement_index| {
(1..=erasure_set_drive_count).contains(replacement_index)
&& !peer_manifest.erasure_indices.contains(replacement_index)
})
}
fn outage_candidate_targets_replacement(
peer_manifest: &OutagePeerManifest,
erasure_set_drive_count: usize,
replacement_set_slot: usize,
) -> bool {
let min_online_data_shards = erasure_set_drive_count.saturating_sub(4);
peer_manifest.erasure_indices.len() >= min_online_data_shards
&& outage_candidate_replacement_erasure_index(peer_manifest, erasure_set_drive_count, replacement_set_slot).is_some()
}
fn metadata_count(disk: &Path, bucket: &str, expected_manifests: &[PhysicalObjectManifest]) -> usize {
expected_manifests
.iter()
@@ -601,6 +683,10 @@ mod tests {
error.as_service_error().and_then(ProvideErrorMetadata::code) == Some("ServiceUnavailable")
}
fn is_service_unavailable_delete(error: &SdkError<DeleteObjectError>) -> bool {
error.as_service_error().and_then(ProvideErrorMetadata::code) == Some("ServiceUnavailable")
}
fn select_replacement_drive(
cluster: &RustFSTestClusterEnvironment,
node_index: usize,
@@ -612,7 +698,7 @@ mod tests {
.ok_or_else(|| format!("replacement node {node_index} is absent"))?;
let mut incomplete_pool_metadata = Vec::new();
for drive in &node.data_dirs {
for (drive_index, drive) in node.data_dirs.iter().enumerate() {
let replaced_disk = PathBuf::from(drive);
let replacement_format_path = replaced_disk.join(".rustfs.sys").join("format.json");
let replacement_format = std::fs::read(&replacement_format_path).map_err(|err| {
@@ -621,6 +707,7 @@ mod tests {
if !require_pool_metadata {
return Ok(ReplacementDriveSelection {
replaced_disk,
drive_index,
replacement_format_path,
replacement_format,
expected_pool_metadata: None,
@@ -631,6 +718,7 @@ mod tests {
if census.is_complete() {
return Ok(ReplacementDriveSelection {
replaced_disk,
drive_index,
replacement_format_path,
replacement_format,
expected_pool_metadata: Some(census),
@@ -1180,7 +1268,7 @@ mod tests {
async fn test_cluster_root_heal_recovers_ec84_shards_across_multi_set_after_background_target_restart()
-> Result<(), Box<dyn Error + Send + Sync>> {
timeout(
Duration::from_secs(600),
Duration::from_secs(900),
run_cluster_root_heal_interruption(InterruptionScenario::BackgroundTargetRestartEc84MultiSet),
)
.await?
@@ -1332,11 +1420,14 @@ mod tests {
let ReplacementDriveSelection {
replaced_disk,
drive_index: replacement_drive_index,
replacement_format_path,
replacement_format,
expected_pool_metadata,
} = select_replacement_drive(&cluster, 1, background_enabled)?;
let default_online_object_count = if !outage_target_manifest_required { 96 } else { 24 };
let replacement_global_drive_index = topology.drives_per_node + replacement_drive_index;
let replacement_set_slot = replacement_global_drive_index % erasure_set_drive_count;
let default_online_object_count = if !outage_target_manifest_required { 64 } else { 24 };
let online_object_count = std::env::var("RUSTFS_HEAL_CHAOS_OBJECT_COUNT")
.ok()
.and_then(|value| value.parse::<usize>().ok())
@@ -1391,6 +1482,23 @@ mod tests {
expected_manifests.push(PhysicalObjectManifest { key, shard_census });
attempt_count += 1;
}
for manifest in &expected_manifests {
let distribution = manifest
.shard_census
.erasure_distribution
.as_ref()
.ok_or_else(|| format!("replacement baseline manifest has no erasure distribution: {manifest:?}"))?;
assert_eq!(
distribution.len(),
erasure_set_drive_count,
"replacement baseline distribution must match the erasure set: {manifest:?}"
);
assert_eq!(
manifest.shard_census.erasure_index,
distribution.get(replacement_set_slot).copied(),
"replacement baseline shard must match the selected drive's erasure-set slot"
);
}
if background_enabled {
wait_for_scanner_cycle_after(&cluster, 0).await?;
@@ -1415,6 +1523,9 @@ mod tests {
let mut outage_write_deferred_until_rejoin = false;
let mut service_unavailable_outage_writes = 0usize;
let mut last_service_unavailable = None;
let mut outage_peer_manifest = OutagePeerManifest::default();
let mut replacement_outage_erasure_index = None;
let mut rejected_outage_keys = Vec::new();
for attempt in 0..max_outage_write_attempts {
let candidate_key = format!("cluster/written-while-node-down-{attempt:04}.bin");
let put_result = timeout(
@@ -1429,6 +1540,24 @@ mod tests {
.await;
match put_result {
Ok(Ok(_)) => {
if outage_target_manifest_required {
let candidate_peer_manifest =
collect_outage_peer_manifest(&cluster, 1, bucket, &candidate_key, erasure_set_drive_count)?;
if !outage_candidate_targets_replacement(
&candidate_peer_manifest,
erasure_set_drive_count,
replacement_set_slot,
) {
rejected_outage_keys.push(candidate_key);
continue;
}
replacement_outage_erasure_index = outage_candidate_replacement_erasure_index(
&candidate_peer_manifest,
erasure_set_drive_count,
replacement_set_slot,
);
outage_peer_manifest = candidate_peer_manifest;
}
outage_key = Some(candidate_key);
break;
}
@@ -1456,57 +1585,37 @@ mod tests {
}
};
let mut outage_peer_erasure_indices = HashSet::new();
if !outage_write_deferred_until_rejoin {
for (node_index, node) in cluster.nodes.iter().enumerate() {
if node_index == 1 {
continue;
}
for (drive_index, drive) in node.data_dirs.iter().enumerate() {
let census = census_object_version_on_disk(Path::new(drive), bucket, &outage_key, None)?;
if !census.has_xl_meta {
continue;
}
assert!(
census.is_complete(),
"online node {node_index} drive {drive_index} must hold a complete outage-object shard: {census:?}"
);
let erasure_index = census.erasure_index.ok_or_else(|| {
format!(
"online node {node_index} drive {drive_index} outage-object shard has no erasure index: {census:?}"
)
})?;
assert!(
(1..=erasure_set_drive_count).contains(&erasure_index),
"online node {node_index} drive {drive_index} outage-object erasure index is out of range: {census:?}"
);
assert!(
outage_peer_erasure_indices.insert(erasure_index),
"outage-object erasure index {erasure_index} is duplicated across online drives"
);
}
}
if !outage_write_deferred_until_rejoin && outage_peer_manifest.erasure_indices.is_empty() {
outage_peer_manifest = collect_outage_peer_manifest(&cluster, 1, bucket, &outage_key, erasure_set_drive_count)?;
replacement_outage_erasure_index =
outage_candidate_replacement_erasure_index(&outage_peer_manifest, erasure_set_drive_count, replacement_set_slot);
}
assert!(
outage_write_deferred_until_rejoin
|| (!outage_peer_erasure_indices.is_empty() && outage_peer_erasure_indices.len() <= erasure_set_drive_count),
|| (!outage_peer_manifest.erasure_indices.is_empty()
&& outage_peer_manifest.erasure_indices.len() <= erasure_set_drive_count),
"outage-object must occupy one non-empty erasure set"
);
if outage_target_manifest_required {
let min_online_data_shards = erasure_set_drive_count.saturating_sub(4);
assert!(
outage_peer_erasure_indices.len() >= min_online_data_shards,
outage_peer_manifest.erasure_indices.len() >= min_online_data_shards,
"online drives in the selected erasure set must retain at least the EC data quorum"
);
}
let missing_outage_erasure_indices = (1..=erasure_set_drive_count)
.filter(|index| !outage_peer_erasure_indices.contains(index))
.filter(|index| !outage_peer_manifest.erasure_indices.contains(index))
.collect::<HashSet<_>>();
if outage_target_manifest_required {
assert!(
!missing_outage_erasure_indices.is_empty(),
"the stopped target must account for at least one missing outage-object erasure index"
);
assert_eq!(
replacement_outage_erasure_index.filter(|index| missing_outage_erasure_indices.contains(index)),
replacement_outage_erasure_index,
"the outage object must target the selected replacement drive's erasure-set slot"
);
}
let heal_body = r#"{"recursive":true,"dryRun":false,"remove":false,"recreate":true,"scanMode":2,"updateParity":false,"nolock":false}"#;
@@ -1531,6 +1640,24 @@ mod tests {
}
cluster.start_node_from_binary(1, &server_binary).await?;
for rejected_key in rejected_outage_keys {
let delete_deadline = Instant::now() + Duration::from_secs(60);
loop {
let delete_result = timeout(
Duration::from_secs(30),
clients[0].delete_object().bucket(bucket).key(&rejected_key).send(),
)
.await;
match delete_result {
Ok(Ok(_)) => break,
Ok(Err(error)) if is_service_unavailable_delete(&error) && Instant::now() < delete_deadline => {
sleep(Duration::from_secs(1)).await;
}
Ok(Err(error)) => return Err(error.into()),
Err(error) => return Err(error.into()),
}
}
}
let status_url = format!("{}/rustfs/admin/v3/background-heal/status", cluster.nodes[0].url);
let recovery_deadline = Instant::now() + Duration::from_secs(60);
@@ -1823,6 +1950,7 @@ mod tests {
scenario,
InterruptionScenario::BackgroundTargetRestart
| InterruptionScenario::BackgroundTargetRestartEc84
| InterruptionScenario::BackgroundTargetRestartEc84MultiSet
| InterruptionScenario::BackgroundCoordinatorRestart
) {
cluster.stop_node_gracefully(interruption_node).await?;
@@ -1994,9 +2122,9 @@ mod tests {
assert_eq!(
outage_census
.erasure_index
.filter(|index| missing_outage_erasure_indices.contains(index)),
.filter(|index| replacement_outage_erasure_index == Some(*index)),
outage_census.erasure_index,
"the outage object must be rebuilt into one of the stopped node's missing erasure slots"
"the outage object must be rebuilt into the selected replacement drive's erasure-set slot"
);
}
@@ -2123,7 +2251,21 @@ mod tests {
evidence_context.run.binary.sha256,
"server build changed during restart"
);
let evidence = serde_json::json!({
let outage_write_diagnostic = (!outage_target_manifest_required).then(|| {
serde_json::json!({
"attempted": true,
"required": false,
"accepted": true,
"attempts": if outage_write_deferred_until_rejoin {
max_outage_write_attempts
} else {
service_unavailable_outage_writes + 1
},
"service_unavailable": service_unavailable_outage_writes,
"deferred_until_rejoin": outage_write_deferred_until_rejoin,
})
});
let mut evidence = serde_json::json!({
"schema": 1, "case": evidence_context.case.id, "evidence": evidence_context.case.evidence,
"run_id": evidence_context.run.run_id, "source_revision": evidence_context.run.source_revision,
"test_build": compiled_test_identity(),
@@ -2146,6 +2288,9 @@ mod tests {
"unclean_shutdown_marker": unclean_shutdown_marker_observed.unwrap_or(false),
"objects": evidence_objects, "node_listings": node_listings,
});
if let Some(outage_write) = outage_write_diagnostic {
evidence["outage_write"] = outage_write;
}
let data = serde_json::to_vec(&evidence)?;
if data.len() > 1024 * 1024 {
return Err("scanner/heal oracle exceeds the 1 MiB artifact budget".into());
@@ -2259,4 +2404,28 @@ mod tests {
Ok(())
}
#[test]
fn outage_candidate_must_target_replacement_erasure_index() {
let distribution = vec![4, 7, 10, 1, 5, 8, 11, 2, 6, 9, 12, 3];
let replacement_set_slot = 8;
let replacement_index = distribution[replacement_set_slot];
let peers_missing_replacement = OutagePeerManifest {
erasure_indices: (1..=12).filter(|index| *index != replacement_index).collect(),
erasure_distribution: Some(distribution.clone()),
};
assert!(outage_candidate_targets_replacement(&peers_missing_replacement, 12, replacement_set_slot));
let peers_missing_other_slot = OutagePeerManifest {
erasure_indices: (1..=12).filter(|index| *index != 9).collect(),
erasure_distribution: Some(distribution),
};
assert!(!outage_candidate_targets_replacement(&peers_missing_other_slot, 12, replacement_set_slot));
let insufficient_peer_shards = OutagePeerManifest {
erasure_indices: [1, 3, 4, 5, 6, 7, 8].into_iter().collect(),
erasure_distribution: Some(vec![1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12]),
};
assert!(!outage_candidate_targets_replacement(&insufficient_peer_shards, 12, replacement_set_slot));
}
}
+77 -2
View File
@@ -52,8 +52,8 @@ use aws_sdk_s3::error::ProvideErrorMetadata;
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{
BucketLifecycleConfiguration, BucketVersioningStatus, CompletedMultipartUpload, CompletedPart, ExpirationStatus,
LifecycleRule, LifecycleRuleFilter, NoncurrentVersionTransition, RestoreRequest, Transition, TransitionStorageClass,
VersioningConfiguration,
LifecycleRule, LifecycleRuleFilter, MetadataDirective, NoncurrentVersionTransition, RestoreRequest, Transition,
TransitionStorageClass, VersioningConfiguration,
};
use http::Method;
use serde::Deserialize;
@@ -963,6 +963,81 @@ async fn test_hermetic_transition_main_path() -> TestResult {
Ok(())
}
/// PUT and materialized self-copy must retain cleanup ownership of a replaced
/// transitioned null version while publishing the new bytes and metadata.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_hermetic_transition_overwrite_and_self_copy() -> TestResult {
let mut cold = RustFSTestEnvironment::new().await?;
cold.access_key = "coldtieradmin".to_string();
cold.secret_key = "coldtiersecret".to_string();
cold.start_rustfs_server_without_cleanup(vec![]).await?;
let cold_client = cold.create_s3_client();
cold_client.create_bucket().bucket(TIER_BUCKET).send().await?;
let mut hot = RustFSTestEnvironment::new().await?;
start_tier_source(&mut hot, crate::common::FAST_DATA_USAGE_SCANNER_ENV).await?;
let hot_client = hot.create_s3_client();
add_rustfs_tier(&hot, &cold).await?;
hot_client.create_bucket().bucket(SOURCE_BUCKET).send().await?;
let data = payload();
for self_copy in [false, true] {
hot_client
.put_bucket_lifecycle_configuration()
.bucket(SOURCE_BUCKET)
.lifecycle_configuration(BucketLifecycleConfiguration::builder().rules(transition_rule()?).build()?)
.send()
.await?;
put_multipart_object(&hot_client, SOURCE_BUCKET, OBJECT_KEY, &data).await?;
wait_for_transition(&hot_client, SOURCE_BUCKET, OBJECT_KEY, StdDuration::from_secs(90)).await?;
assert_eq!(cold_tier_object_count(&cold_client).await?, 1);
// Keep the replacement local so disappearance of the old remote
// object cannot be confused with another automatic transition.
hot_client.delete_bucket_lifecycle().bucket(SOURCE_BUCKET).send().await?;
let expected = if self_copy { data.clone() } else { vec![0x73; 513] };
if self_copy {
hot_client
.copy_object()
.bucket(SOURCE_BUCKET)
.key(OBJECT_KEY)
.copy_source(format!("{SOURCE_BUCKET}/{}", urlencoding::encode(OBJECT_KEY)))
.metadata_directive(MetadataDirective::Replace)
.content_type("text/plain")
.metadata("replacement", "kept")
.send()
.await?;
} else {
hot_client
.put_object()
.bucket(SOURCE_BUCKET)
.key(OBJECT_KEY)
.body(ByteStream::from(expected.clone()))
.content_type("text/plain")
.metadata("replacement", "kept")
.send()
.await?;
}
wait_for_cold_tier_empty(&cold_client, StdDuration::from_secs(90)).await?;
let current = hot_client.get_object().bucket(SOURCE_BUCKET).key(OBJECT_KEY).send().await?;
assert_eq!(current.content_type(), Some("text/plain"));
assert_eq!(current.metadata().and_then(|m| m.get("replacement")).map(String::as_str), Some("kept"));
assert!(
current
.metadata()
.is_none_or(|metadata| !metadata.contains_key(USER_META_KEY))
);
assert_eq!(current.body.collect().await?.into_bytes().as_ref(), expected.as_slice());
hot_client
.delete_object()
.bucket(SOURCE_BUCKET)
.key(OBJECT_KEY)
.send()
.await?;
}
Ok(())
}
/// Restore a transitioned object through a real RustFS remote tier.
///
/// The test covers the externally visible copy-back contract that a mock tier
@@ -27,7 +27,9 @@ mod tests {
use aws_sdk_s3::Client;
use aws_sdk_s3::error::ProvideErrorMetadata;
use aws_sdk_s3::primitives::ByteStream;
use aws_sdk_s3::types::{BucketVersioningStatus, CompletedMultipartUpload, CompletedPart, VersioningConfiguration};
use aws_sdk_s3::types::{
BucketVersioningStatus, CompletedMultipartUpload, CompletedPart, Tag, Tagging, VersioningConfiguration,
};
use tracing::info;
fn create_s3_client(env: &RustFSTestEnvironment) -> Client {
@@ -83,6 +85,156 @@ mod tests {
Ok(())
}
async fn assert_version_tags(client: &Client, bucket: &str, key: &str, version: Option<&str>, value: Option<&str>) {
let tags = client
.get_object_tagging()
.bucket(bucket)
.key(key)
.set_version_id(version.map(str::to_owned))
.send()
.await
.expect("GetObjectTagging must accept the exact version selector");
let expected = value
.map(|value| vec![Tag::builder().key("generation").value(value).build().expect("valid tag")])
.unwrap_or_default();
assert_eq!(tags.tag_set(), expected, "version selector: {version:?}");
}
async fn assert_null_tagging_across_versioning_changes(client: &Client, bucket: &str, key: &str) {
assert_version_tags(client, bucket, key, None, None).await;
assert_version_tags(client, bucket, key, Some("null"), None).await;
client
.put_object_tagging()
.bucket(bucket)
.key(key)
.version_id("null")
.tagging(
Tagging::builder()
.tag_set(
Tag::builder()
.key("generation")
.value("original-null")
.build()
.expect("valid tag"),
)
.build()
.expect("valid tagging"),
)
.send()
.await
.expect("tag the original null version");
enable_versioning(client, bucket)
.await
.expect("enable versioning over a null version");
let versioned = client
.put_object()
.bucket(bucket)
.key(key)
.tagging("generation=versioned")
.body(ByteStream::from_static(b"new version"))
.send()
.await
.expect("write a newer UUID version");
let version = versioned.version_id().expect("versioned PUT must return a UUID");
assert_ne!(version, "null");
assert_version_tags(client, bucket, key, None, Some("versioned")).await;
assert_version_tags(client, bucket, key, Some(version), Some("versioned")).await;
assert_version_tags(client, bucket, key, Some("null"), Some("original-null")).await;
client
.put_object_tagging()
.bucket(bucket)
.key(key)
.version_id("null")
.tagging(
Tagging::builder()
.tag_set(
Tag::builder()
.key("generation")
.value("updated-null")
.build()
.expect("valid tag"),
)
.build()
.expect("valid tagging"),
)
.send()
.await
.expect("update tags on the noncurrent null version");
assert_version_tags(client, bucket, key, Some("null"), Some("updated-null")).await;
assert_version_tags(client, bucket, key, None, Some("versioned")).await;
client
.delete_object_tagging()
.bucket(bucket)
.key(key)
.version_id("null")
.send()
.await
.expect("delete only the noncurrent null version tags");
assert_version_tags(client, bucket, key, Some("null"), None).await;
assert_version_tags(client, bucket, key, Some(version), Some("versioned")).await;
suspend_versioning(client, bucket).await.expect("suspend versioning");
client
.put_object()
.bucket(bucket)
.key(key)
.tagging("generation=suspended-null")
.body(ByteStream::from_static(b"replacement null version"))
.send()
.await
.expect("replace the null version while suspended");
assert_version_tags(client, bucket, key, None, Some("suspended-null")).await;
assert_version_tags(client, bucket, key, Some("null"), Some("suspended-null")).await;
assert_version_tags(client, bucket, key, Some(version), Some("versioned")).await;
enable_versioning(client, bucket).await.expect("re-enable versioning");
client
.put_object()
.bucket(bucket)
.key(key)
.tagging("generation=latest")
.body(ByteStream::from_static(b"latest version"))
.send()
.await
.expect("write a new latest version");
assert_version_tags(client, bucket, key, Some("null"), Some("suspended-null")).await;
assert_version_tags(client, bucket, key, None, Some("latest")).await;
client
.delete_object()
.bucket(bucket)
.key(key)
.version_id("null")
.send()
.await
.expect("remove only the null version");
assert_version_tags(client, bucket, key, None, Some("latest")).await;
let absent_version = uuid::Uuid::new_v4().to_string();
for (missing_key, selector, expected_code) in [
(key, Some("null"), "NoSuchVersion"),
(key, Some(absent_version.as_str()), "NoSuchVersion"),
("never-created", None, "NoSuchKey"),
("never-created", Some("null"), "NoSuchVersion"),
("never-created", Some(absent_version.as_str()), "NoSuchVersion"),
] {
let missing = client
.get_object_tagging()
.bucket(bucket)
.key(missing_key)
.set_version_id(selector.map(str::to_owned))
.send()
.await
.expect_err("a missing version must not fall back to latest");
assert_eq!(
missing.as_service_error().and_then(ProvideErrorMetadata::code),
Some(expected_code),
"key: {missing_key}, version selector: {selector:?}"
);
}
assert_version_tags(client, bucket, key, None, Some("latest")).await;
}
/// Test 1: PutObject should return version_id when versioning is enabled
/// This directly addresses the Veeam issue from #1066
#[tokio::test]
@@ -262,7 +414,9 @@ mod tests {
info!("🧪 TEST: PutObject behavior without versioning (no regression)");
let mut env = RustFSTestEnvironment::new().await.expect("Failed to create test environment");
env.start_rustfs_server(vec![]).await.expect("Failed to start RustFS");
env.start_rustfs_server_without_cleanup(vec![])
.await
.expect("Failed to start isolated RustFS");
let client = create_s3_client(&env);
let bucket = "test-no-versioning";
@@ -290,6 +444,9 @@ mod tests {
output.version_id().is_none() || output.version_id() == Some("null"),
"non-versioned PUT must omit version ID or return the S3 null version"
);
// Reuse this unversioned fixture to prove explicit null never becomes
// an implicit latest-version read after enable/suspend transitions.
assert_null_tagging_across_versioning_changes(&client, bucket, key).await;
info!("✅ PASSED: PutObject works correctly without versioning");
}
@@ -779,7 +779,7 @@ fn free_version_physical_topology_generation(api: &ECStore) -> String {
rustfs_utils::crypto::hex(hasher.finalize().as_slice())
}
fn free_version_remote_tuple_matches(candidate: &ObjectInfo, expected: &ObjectInfo) -> std::io::Result<bool> {
pub(crate) fn free_version_remote_tuple_matches(candidate: &ObjectInfo, expected: &ObjectInfo) -> std::io::Result<bool> {
if candidate.transitioned_object.tier != expected.transitioned_object.tier
|| candidate.transitioned_object.name != expected.transitioned_object.name
{
@@ -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]
@@ -577,6 +577,26 @@ fn heal_control_auth_may_need_replay_scope_refresh(err: &Error) -> bool {
)
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum HealControlRetryAction {
Reconnect,
RefreshReplayScope,
}
fn heal_control_retry_action(
err: &Error,
reconnect_attempted: bool,
replay_scope_refresh_attempted: bool,
) -> Option<HealControlRetryAction> {
if !replay_scope_refresh_attempted && heal_control_auth_may_need_replay_scope_refresh(err) {
return Some(HealControlRetryAction::RefreshReplayScope);
}
if !reconnect_attempted && PeerRestClient::is_network_like_error(err) {
return Some(HealControlRetryAction::Reconnect);
}
None
}
fn decode_remote_version_state_capability(expected_member: &str, result: &[u8]) -> Result<Uuid> {
let (topology_member, process_epoch) = rustfs_protos::decode_remote_version_state_capability(result).map_err(Error::other)?;
if topology_member != expected_member {
@@ -1753,34 +1773,41 @@ impl PeerRestClient {
return Err(Error::other("heal control command exceeds size limit"));
}
let capability_probe = rustfs_protos::is_heal_control_capability_probe(&command);
let result = self
.heal_control_once(version, &topology_fingerprint, &command, capability_probe)
.await;
if result
.as_ref()
.err()
.is_some_and(heal_control_auth_may_need_replay_scope_refresh)
{
self.prepare_heal_control_auth_retry().await;
return self
.finalize_result(
self.heal_control_once(version, &topology_fingerprint, &command, capability_probe)
.await,
)
let mut reconnect_attempted = false;
let mut replay_scope_refresh_attempted = false;
loop {
let result = self
.heal_control_once(version, &topology_fingerprint, &command, capability_probe)
.await;
let Some(action) = result
.as_ref()
.err()
.and_then(|err| heal_control_retry_action(err, reconnect_attempted, replay_scope_refresh_attempted))
else {
return self.finalize_result(result).await;
};
match action {
HealControlRetryAction::Reconnect => reconnect_attempted = true,
HealControlRetryAction::RefreshReplayScope => replay_scope_refresh_attempted = true,
}
self.prepare_heal_control_retry(action).await;
}
self.finalize_result(result).await
}
async fn prepare_heal_control_auth_retry(&self) {
if let Err(err) = clear_peer_replay_state_for_addr(&self.grid_host) {
async fn prepare_heal_control_retry(&self, action: HealControlRetryAction) {
if action == HealControlRetryAction::RefreshReplayScope
&& let Err(err) = clear_peer_replay_state_for_addr(&self.grid_host)
{
debug!(
peer = %self.grid_host,
error = %err,
"could not clear heal control replay state before retry"
);
}
self.evict_connection().await;
// A restart can leave both the local offline gate and the peer replay
// epoch stale. Clear the gate on either recovery step so the next
// bounded attempt reaches a fresh channel instead of fast-failing.
self.prepare_retry().await;
}
async fn heal_control_once(
@@ -3867,6 +3894,51 @@ mod tests {
)));
}
#[test]
fn heal_control_retry_plan_allows_one_reconnect_and_one_epoch_refresh() {
let offline = Error::RemoteClientUnavailable("peer http://127.0.0.1:9000 is temporarily offline".to_string());
let stale_epoch = Error::from(tonic::Status::unauthenticated("No valid auth token"));
assert_eq!(heal_control_retry_action(&offline, false, false), Some(HealControlRetryAction::Reconnect));
assert_eq!(heal_control_retry_action(&offline, true, false), None);
assert_eq!(
heal_control_retry_action(&stale_epoch, true, false),
Some(HealControlRetryAction::RefreshReplayScope),
"a reconnect may expose the restarted peer's stale replay epoch"
);
assert_eq!(heal_control_retry_action(&stale_epoch, true, true), None);
assert_eq!(
heal_control_retry_action(&stale_epoch, false, false),
Some(HealControlRetryAction::RefreshReplayScope)
);
assert_eq!(
heal_control_retry_action(&offline, false, true),
Some(HealControlRetryAction::Reconnect),
"an epoch refresh may be followed by one bounded reconnect"
);
assert_eq!(heal_control_retry_action(&offline, true, true), None);
assert_eq!(
heal_control_retry_action(&Error::from(tonic::Status::permission_denied("bad signature")), false, false),
None,
"authorization failures must never be retried"
);
}
#[tokio::test]
async fn heal_control_epoch_refresh_clears_offline_gate() {
let client = test_peer_client();
client.offline.store(true, Ordering::Release);
client
.prepare_heal_control_retry(HealControlRetryAction::RefreshReplayScope)
.await;
assert!(
!client.offline.load(Ordering::Acquire),
"epoch refresh must not leave the following attempt behind the offline gate"
);
}
#[test]
fn peer_rest_client_network_classifier_keeps_slow_peers_online() {
// The per-RPC channel deadline (RUSTFS_INTERNODE_RPC_TIMEOUT, 30s)
+14
View File
@@ -7470,6 +7470,20 @@ impl PoolMeta {
.is_some_and(is_decommission_suspended)
}
pub(crate) fn has_active_decommission_capacity_reservation(&self, idx: usize) -> bool {
self.pools
.get(idx)
.and_then(|pool| pool.decommission.as_ref())
.is_some_and(|info| {
info.has_decommission_state()
&& is_decommission_active(info.complete, info.failed, info.canceled)
&& info
.capacity_reservation
.as_ref()
.is_some_and(DecommissionCapacityReservation::active)
})
}
pub(crate) fn scanner_pause_backlog_pool_writable(&self, idx: usize) -> bool {
self.pools.get(idx).is_some_and(|pool| {
!pool
+5 -6
View File
@@ -219,7 +219,10 @@ impl Sets {
let mut disk_set = Vec::with_capacity(set_count);
let lock_registry = runtime_sources::lock_registry();
let pool_lockers = runtime_sources::lock_registry()
.as_ref()
.map(|registry| registry.clients_for_endpoints(endpoints.endpoints.as_ref()))
.unwrap_or_default();
for i in 0..set_count {
let mut set_drive = Vec::with_capacity(set_drive_count);
@@ -270,10 +273,6 @@ impl Sets {
}
}
let lockers = lock_registry
.as_ref()
.map(|registry| registry.clients_for_endpoints(&set_endpoints))
.unwrap_or_default();
let set_disks = SetDisks::new_with_instance_ctx(
runtime_sources::local_node_name().await,
Arc::new(RwLock::new(set_drive)),
@@ -283,7 +282,7 @@ impl Sets {
pool_idx,
set_endpoints,
fm.clone(),
lockers,
pool_lockers.clone(),
instance_ctx.clone(),
)
.await;
+54
View File
@@ -41,6 +41,14 @@ struct DanglingDeleteGraceError {
grace_secs: i64,
}
/// Marks a conditional-file write that failed before its publication rename.
/// Callers may choose another owner only while this marker is preserved; every
/// unmarked error remains commit-ambiguous and must fail closed.
#[derive(Debug)]
struct ConditionalFileNotCommittedError {
source: io::Error,
}
// DiskError == StorageErr
#[derive(Debug, thiserror::Error)]
pub enum DiskError {
@@ -220,6 +228,18 @@ impl std::fmt::Display for DanglingDeleteGraceError {
impl StdError for DanglingDeleteGraceError {}
impl std::fmt::Display for ConditionalFileNotCommittedError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
self.source.fmt(f)
}
}
impl StdError for ConditionalFileNotCommittedError {
fn source(&self) -> Option<&(dyn StdError + 'static)> {
Some(&self.source)
}
}
fn classify_internode_missing_error(error: &InternodeHttpError) -> Option<DiskError> {
if error.is_remote_file_not_found() {
return Some(DiskError::FileNotFound);
@@ -293,6 +313,22 @@ impl DiskError {
})
}
pub(crate) fn conditional_file_not_committed(source: io::Error) -> io::Error {
io::Error::new(source.kind(), ConditionalFileNotCommittedError { source })
}
/// Whether a local conditional-file replacement failed before the target
/// publication rename and therefore cannot have committed new owner bytes.
pub fn is_conditional_file_not_committed(&self) -> bool {
matches!(
self,
DiskError::Io(io_error)
if io_error
.get_ref()
.is_some_and(|source| source.downcast_ref::<ConditionalFileNotCommittedError>().is_some())
)
}
pub fn is_dangling_delete_grace(&self) -> bool {
matches!(self, DiskError::Io(io_error) if Self::io_error_is_dangling_delete_grace(io_error))
}
@@ -627,6 +663,9 @@ impl From<tokio::task::JoinError> for DiskError {
impl Clone for DiskError {
fn clone(&self) -> Self {
match self {
DiskError::Io(io_error) if self.is_conditional_file_not_committed() => DiskError::Io(
DiskError::conditional_file_not_committed(io::Error::new(io_error.kind(), io_error.to_string())),
),
DiskError::Io(io_error) => DiskError::Io(
rustfs_rio::clone_internode_http_io_error(io_error)
.and_then(std::io::Error::into_inner)
@@ -820,6 +859,21 @@ mod tests {
use super::*;
use std::collections::HashMap;
#[test]
fn conditional_file_not_committed_marker_is_explicit_and_clone_safe() {
let marked = DiskError::from(DiskError::conditional_file_not_committed(io::Error::new(
io::ErrorKind::PermissionDenied,
"staging rejected",
)));
assert!(marked.is_conditional_file_not_committed());
assert!(marked.clone().is_conditional_file_not_committed());
assert!(!DiskError::Timeout.is_conditional_file_not_committed());
assert!(
!DiskError::Io(io::Error::new(io::ErrorKind::PermissionDenied, "rename rejected"))
.is_conditional_file_not_committed()
);
}
#[test]
fn terminal_read_error_preserves_kind_and_disk_classification() {
let timeout = terminal_read_error_to_io(DiskError::Timeout);
+14 -4
View File
@@ -8957,7 +8957,8 @@ impl DiskAPI for LocalDisk {
.truncate(false)
.read(true)
.write(true)
.open(&lock_path)?;
.open(&lock_path)
.map_err(DiskError::conditional_file_not_committed)?;
flock(&lock, FlockOperation::NonBlockingLockExclusive).map_err(std::io::Error::from)?;
let result = (|| {
let current = match std::fs::read(&file_path) {
@@ -9012,10 +9013,15 @@ impl DiskAPI for LocalDisk {
.ok_or_else(|| std::io::Error::new(ErrorKind::InvalidInput, "conditional file has no parent"))?;
let temporary = parent.join(format!(".{}.{}.tmp", path.replace('/', "_"), Uuid::new_v4()));
let write_result = (|| -> std::io::Result<()> {
let mut staged = std::fs::OpenOptions::new().create_new(true).write(true).open(&temporary)?;
staged.write_all(&replacement)?;
let not_committed = DiskError::conditional_file_not_committed;
let mut staged = std::fs::OpenOptions::new()
.create_new(true)
.write(true)
.open(&temporary)
.map_err(not_committed)?;
staged.write_all(&replacement).map_err(not_committed)?;
if sync_metadata {
staged.sync_all()?;
staged.sync_all().map_err(not_committed)?;
}
std::fs::rename(&temporary, &file_path)?;
Ok(())
@@ -22650,6 +22656,10 @@ mod test {
.await
.expect_err("directory fsync failure must fail the CAS update");
assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::Other));
assert!(
!err.is_conditional_file_not_committed(),
"an error after publication rename must remain commit-ambiguous"
);
assert_eq!(
disk.read_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
.await
@@ -3076,6 +3076,23 @@ impl SetDisks {
&self,
bucket: &str,
object: &str,
) -> Result<Option<rustfs_filemeta::FileInfoVersions>> {
self.load_file_info_versions_for_cleanup(bucket, object, false).await
}
pub(crate) async fn load_file_info_versions_for_tier_cleanup(
&self,
bucket: &str,
object: &str,
) -> Result<Option<rustfs_filemeta::FileInfoVersions>> {
self.load_file_info_versions_for_cleanup(bucket, object, true).await
}
async fn load_file_info_versions_for_cleanup(
&self,
bucket: &str,
object: &str,
retain_unconfirmed_tier_references: bool,
) -> Result<Option<rustfs_filemeta::FileInfoVersions>> {
let disk_object = rustfs_utils::path::encode_dir_object(object);
let disks = self.get_disks_internal().await;
@@ -3152,12 +3169,24 @@ impl SetDisks {
)));
}
let file_info_versions = FileMeta {
let mut file_info_versions = FileMeta {
versions,
..Default::default()
}
.get_all_file_info_versions(bucket, object, true)
.map_err(decode_error)?;
if retain_unconfirmed_tier_references {
// A failed overwrite may leave its live source on a minority
// of disks. Preserve that reference even if quorum merging
// selects only the replacement and its cleanup owner.
file_info_versions.versions.extend(
transition_copies
.into_values()
.flatten()
.map(|(version, _)| version)
.filter(|version| !version.tier_free_version()),
);
}
for file_info in file_info_versions
.versions
@@ -12214,6 +12243,64 @@ mod tests {
assert!(result.is_err(), "missing disks must prevent metadata write quorum");
}
#[tokio::test]
async fn tier_overwrite_cleanup_rejects_unreadable_disk_despite_metadata_quorum() {
let bucket = "tier-unreadable-disk";
let object = "object";
let mut dirs = Vec::new();
let mut disks = Vec::new();
let mut fi = metadata_test_fileinfo(object);
fi.mod_time = Some(OffsetDateTime::now_utc());
for index in 1..=3 {
let (dir, disk) = read_multiple_test_disk(bucket, &[]).await;
fi.erasure.index = index;
disk.write_metadata(bucket, bucket, object, fi.clone())
.await
.expect("seed metadata quorum");
dirs.push(dir);
disks.push(Some(disk));
}
disks.push(None);
let set = io_primitives_test_set(disks, 2).await;
assert!(
set.load_file_info_versions_exact(bucket, object).await.is_err(),
"exact reads must preserve release's unreadable-replica fence"
);
assert!(
set.load_file_info_versions_for_tier_cleanup(bucket, object).await.is_err(),
"unreadable replica may still reference the old remote object"
);
}
#[tokio::test]
async fn tier_overwrite_cleanup_rejects_minority_metadata_in_an_absent_set() {
let bucket = "tier-minority-metadata";
let object = "object";
let mut dirs = Vec::new();
let mut disks = Vec::new();
for index in 1..=4 {
let (dir, disk) = read_multiple_test_disk(bucket, &[]).await;
if index == 1 {
let mut fi = metadata_test_fileinfo(object);
fi.mod_time = Some(OffsetDateTime::now_utc());
disk.write_metadata(bucket, bucket, object, fi)
.await
.expect("seed minority metadata");
}
dirs.push(dir);
disks.push(Some(disk));
}
let set = io_primitives_test_set(disks, 2).await;
assert!(
set.load_file_info_versions_exact(bucket, object).await.is_err(),
"exact reads must preserve release's minority-ownership fence"
);
assert!(
set.load_file_info_versions_for_tier_cleanup(bucket, object).await.is_err(),
"absence on a majority cannot prove this physical set has no remote reference"
);
}
#[tokio::test]
async fn load_file_info_versions_exact_returns_versions_from_read_quorum() {
let bucket = "exact-versions-bucket";
+78 -4
View File
@@ -3852,7 +3852,7 @@ pub struct SetDisks {
pub default_parity_count: usize,
pub set_index: usize,
pub pool_index: usize,
/// Stable namespace shared by every object lock created for this set.
/// Stable namespace shared by every object lock created for this pool.
set_lock_namespace: Arc<str>,
pub format: FormatV3,
#[allow(dead_code, reason = "asserted by this file's tests (backlog#1823)")]
@@ -4491,7 +4491,7 @@ impl SetDisks {
instance_ctx: Arc<InstanceContext>,
) -> Arc<Self> {
let ctx = instance_ctx;
let set_lock_namespace: Arc<str> = format!("set-{pool_index}-{set_index}").into();
let set_lock_namespace: Arc<str> = format!("pool-{pool_index}").into();
let shared_lockers = Arc::from(lockers.to_vec());
Arc::new(SetDisks {
locker_owner,
@@ -4605,7 +4605,9 @@ impl SetDisks {
pub(crate) async fn shares_namespace_lock_domain(&self, other: &Self) -> bool {
match (self.ctx.is_dist_erasure().await, other.ctx.is_dist_erasure().await) {
(false, false) => Arc::ptr_eq(&self.local_lock_manager, &other.local_lock_manager),
(true, true) => same_distributed_lock_domain(&self.lockers, &other.lockers),
(true, true) => {
self.set_lock_namespace == other.set_lock_namespace && same_distributed_lock_domain(&self.lockers, &other.lockers)
}
_ => false,
}
}
@@ -7123,7 +7125,7 @@ mod tests {
ctx.update_erasure_type(SetupType::Erasure).await;
let set = make_test_set_disks_with_ctx(Vec::new(), ctx).await;
assert_eq!(&*set.set_lock_namespace, "set-0-0");
assert_eq!(&*set.set_lock_namespace, "pool-0");
let before = Arc::strong_count(&set.set_lock_namespace);
let lock = set
.new_ns_lock("bucket", "object")
@@ -8348,6 +8350,78 @@ mod tests {
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn test_new_ns_lock_distributed_write_succeeds_with_three_lockers_one_offline() {
let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await;
let manager_a = Arc::new(rustfs_lock::GlobalLockManager::new());
let manager_b = Arc::new(rustfs_lock::GlobalLockManager::new());
let healthy_a: Arc<dyn LockClient> = Arc::new(LocalClient::with_manager(manager_a));
let healthy_b: Arc<dyn LockClient> = Arc::new(LocalClient::with_manager(manager_b));
let failing_client: Arc<dyn LockClient> = Arc::new(FailingClient);
let set_disks = make_test_set_disks(vec![healthy_a, failing_client, healthy_b]).await;
let guard = set_disks
.new_ns_lock("bucket", "object")
.await
.expect("namespace lock should be created")
.get_write_lock(Duration::from_millis(500))
.await
.expect("two healthy lockers should satisfy the three-locker write quorum");
match guard {
NamespaceLockGuard::Standard(_) => {}
NamespaceLockGuard::Fast(_) => panic!("Expected distributed guard for dist-erasure"),
}
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn namespace_lock_domain_includes_pool_namespace() {
let _setup_type_guard = SetupTypeGuard::switch_to(SetupType::DistErasure).await;
let first: Arc<dyn LockClient> = Arc::new(LocalClient::with_manager(Arc::new(rustfs_lock::GlobalLockManager::new())));
let second: Arc<dyn LockClient> = Arc::new(LocalClient::with_manager(Arc::new(rustfs_lock::GlobalLockManager::new())));
let lockers = vec![first, second];
let same_pool_first_set = make_test_set_disks_with_ctx(lockers.clone(), bootstrap_ctx()).await;
let same_pool_second_set = SetDisks::new_with_instance_ctx(
"test-owner".to_string(),
Arc::new(RwLock::new(vec![None, None])),
2,
1,
1,
0,
same_pool_first_set.set_endpoints.clone(),
FormatV3::new(2, 2),
lockers.clone(),
bootstrap_ctx(),
)
.await;
let other_pool_set = SetDisks::new_with_instance_ctx(
"test-owner".to_string(),
Arc::new(RwLock::new(vec![None, None])),
2,
1,
0,
1,
same_pool_first_set.set_endpoints.clone(),
FormatV3::new(1, 2),
lockers,
bootstrap_ctx(),
)
.await;
assert!(
same_pool_first_set.shares_namespace_lock_domain(&same_pool_second_set).await,
"sets in the same pool share the object namespace lock domain"
);
assert!(
!same_pool_first_set.shares_namespace_lock_domain(&other_pool_set).await,
"different pool namespaces must not be deduplicated solely by identical clients"
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial]
async fn streaming_reader_holds_read_lock_until_eof() {
+59
View File
@@ -3280,6 +3280,65 @@ mod heal_result_report_tests {
);
}
#[tokio::test]
async fn deep_heal_rebuilds_missing_part_when_metadata_remains_current() {
let (temp_dirs, disks, set) = hermetic_set_disks_isolated(4).await;
let bucket = "deep-heal-missing-part-current-meta";
let object = "object.bin";
for disk in &disks {
disk.make_volume(bucket).await.expect("bucket volume should be created");
}
let payload = vec![0x7b; 1024 * 1024];
let mut reader = PutObjReader::from_vec(payload);
set.put_object(
bucket,
object,
&mut reader,
&ObjectOptions {
no_lock: true,
..Default::default()
},
)
.await
.expect("source object should be written before shard loss");
let source = disks[2]
.read_version("", bucket, object, "", &ReadOptions::default())
.await
.expect("source metadata should be readable");
let data_dir = source.data_dir.expect("non-inline source should have a data directory");
let missing_part = temp_dirs[1]
.path()
.join(bucket)
.join(object)
.join(data_dir.to_string())
.join("part.1");
tokio::fs::remove_file(&missing_part)
.await
.expect("target shard should be removed while xl.meta remains current");
let (result, error) = set
.heal_object(
bucket,
object,
"",
&HealOpts {
no_lock: true,
scan_mode: HealScanMode::Deep,
..Default::default()
},
)
.await
.expect("deep heal should finish after a single shard is removed");
assert!(error.is_none(), "deep heal should recover the missing shard: {error:?}");
assert_eq!(result.after.drives[1].state, DriveState::Ok.to_string());
assert!(
missing_part.exists(),
"deep heal must reconstruct the missing shard on the original disk slot"
);
}
#[tokio::test]
async fn replacement_target_readback_checks_the_requested_historical_version() {
let (temp_dirs, disks, set) = hermetic_set_disks_isolated(4).await;
+103
View File
@@ -3821,6 +3821,13 @@ impl SetDisks {
}
fi.metadata = user_defined;
if fi.version_id.is_none_or(|id| id.is_nil()) && !opts.data_movement && expected_restore_operation_id.is_none() {
// Every disk must publish the same cleanup owner alongside a
// replaced null version. This transient key is not persisted
// on the new object; recovery discovers the free-version in
// the committed xl.meta even if this request is cancelled.
fi.set_tier_free_version_id(&Uuid::new_v4().to_string());
}
fi.mod_time = mod_time;
fi.size = w_size as i64;
fi.versioned = opts.versioned || opts.version_suspended;
@@ -18174,6 +18181,102 @@ mod put_object_tmp_cleanup_tests {
drop(temp_dirs);
}
#[tokio::test]
#[serial_test::serial(capacity_dirty_scope)]
async fn tier_overwrite_failed_quorum_and_cancellation_preserve_live_source() {
for cancel_before_rename in [false, true] {
let (dirs, disks, set) = hermetic_set_disks(4).await;
let bucket = "tier-overwrite-failure";
let object = "still-live";
make_completion_test_bucket(&disks, bucket).await;
let old_body = vec![0x31; TEST_OBJECT_SIZE];
let mut metadata = HashMap::from([(
"x-amz-restore".to_string(),
"ongoing-request=\"false\", expiry-date=\"2099-01-01T00:00:00Z\"".to_string(),
)]);
for (suffix, value) in [
(rustfs_utils::http::SUFFIX_TRANSITION_STATUS, "complete".to_string()),
(rustfs_utils::http::SUFFIX_TRANSITION_TIER, "WARM".to_string()),
(rustfs_utils::http::SUFFIX_TRANSITIONED_OBJECTNAME, "remote/still-live".to_string()),
(rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID, "exact-live-version".to_string()),
(rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_STATE, "exact".to_string()),
(rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID, "ab".repeat(32)),
] {
rustfs_utils::http::insert_str(&mut metadata, suffix, value);
}
set.put_object(
bucket,
object,
&mut PutObjReader::from_vec(old_body.clone()),
&ObjectOptions {
user_defined: metadata,
write_completion: WriteCompletion::TailDrained,
..Default::default()
},
)
.await
.expect("seed live transitioned source");
wait_for_tmp_workspace_to_drain(&dirs, "seed write must drain").await;
let before = set
.load_file_info_versions_exact(bucket, object)
.await
.expect("read original metadata")
.expect("original exists");
if cancel_before_rename {
let barrier = PutObjectCommitBarrier::install(bucket, object, PutObjectCommitPause::AfterQuotaReservation);
let writer = Arc::clone(&set);
let put = tokio::spawn(async move {
writer
.put_object(
bucket,
object,
&mut PutObjReader::from_vec(vec![0x32; TEST_OBJECT_SIZE]),
&ObjectOptions::default(),
)
.await
});
barrier.wait_until_paused().await;
put.abort();
assert!(put.await.expect_err("cancel paused replacement").is_cancelled());
wait_for_tmp_workspace_to_drain(&dirs, "cancelled replacement must roll back").await;
drop(barrier);
} else {
let _fault = rename_fault_injection::fail_rename_on(object, &[2, 3]);
let err = set
.put_object(
bucket,
object,
&mut PutObjReader::from_vec(vec![0x32; TEST_OBJECT_SIZE]),
&ObjectOptions {
write_completion: WriteCompletion::TailDrained,
..Default::default()
},
)
.await
.expect_err("two disk commits cannot satisfy write quorum three");
assert!(matches!(err, Error::ErasureWriteQuorum | Error::InsufficientWriteQuorum(_, _)), "{err}");
}
let after = set
.load_file_info_versions_exact(bucket, object)
.await
.expect("read rolled-back metadata")
.expect("live source must survive");
assert_eq!(after.versions, before.versions, "failed replacement must preserve the live version");
assert_eq!(
after.free_versions, before.free_versions,
"failed replacement must not publish a cleanup owner"
);
let mut reader = set
.get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
.await
.expect("live source remains readable");
let mut actual = Vec::new();
reader.stream.read_to_end(&mut actual).await.expect("read original bytes");
assert_eq!(actual, old_body);
}
}
#[tokio::test]
async fn cooperative_cancellation_while_waiting_for_namespace_lock_cleans_tmp_workspace() {
let (temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
+250 -12
View File
@@ -2832,6 +2832,224 @@ mod tests {
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
async fn tier_overwrite_put_and_self_copy_recover_persisted_cleanup_owners() {
use crate::bucket::lifecycle::bucket_lifecycle_ops::ExpiryState;
use crate::bucket::lifecycle::tier_free_version_recovery::recover_tier_free_versions;
use rustfs_filemeta::TransitionVersionState::{Exact, KnownDisabled, SuspendedNull};
use rustfs_s3_client::transition_api::ReaderImpl;
use rustfs_utils::http::{
SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID,
SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, insert_str,
};
let temp_dir = tempfile::tempdir().expect("create tier overwrite store");
let (mut ctx, mut store, mut shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-overwrite", &[4])).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
let tier = "OVERWRITE-TIER";
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier).await;
let lease = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier)
.await
.expect("tier identity");
let identity = rustfs_utils::crypto::hex(lease.backend_identity());
drop(lease);
for state in [Exact, KnownDisabled, SuspendedNull] {
for suspended in [false, true] {
for self_copy in [false, true] {
let bucket = format!("tier-overwrite-{}", Uuid::new_v4());
let object = "object";
let remote = format!("remote/{bucket}");
let version = match state {
Exact => "opaque-overwrite-version",
SuspendedNull => "null",
_ => "",
};
let payload = vec![0x5b; if suspended { 512 * 1024 } else { 257 }];
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("create bucket");
backend.set_put_remote_version(Some(version.to_string())).await;
let lease = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier)
.await
.expect("seed tier lease");
lease
.put(
&remote,
ReaderImpl::Body(bytes::Bytes::from(payload.clone())),
payload.len().try_into().expect("payload size"),
)
.await
.expect("seed remote bytes");
drop(lease);
let mut metadata = HashMap::from([
("content-type".to_string(), "application/octet-stream".to_string()),
(
"x-amz-restore".to_string(),
"ongoing-request=\"false\", expiry-date=\"2099-01-01T00:00:00Z\"".to_string(),
),
]);
for (suffix, value) in [
(SUFFIX_TRANSITION_STATUS, "complete"),
(SUFFIX_TRANSITION_TIER, tier),
(SUFFIX_TRANSITION_TIER_DESTINATION_ID, identity.as_str()),
(SUFFIX_TRANSITIONED_OBJECTNAME, remote.as_str()),
(SUFFIX_TRANSITIONED_VERSION_STATE, state.as_str()),
] {
insert_str(&mut metadata, suffix, value.to_string());
}
if !version.is_empty() {
insert_str(&mut metadata, SUFFIX_TRANSITIONED_VERSION_ID, version.to_string());
}
let options = ObjectOptions {
version_suspended: suspended,
..Default::default()
};
store
.put_object(
&bucket,
object,
&mut PutObjReader::from_vec(payload.clone()),
&ObjectOptions {
user_defined: metadata,
..options.clone()
},
)
.await
.expect("seed transitioned source with locally restored bytes");
let expected = if self_copy {
payload.clone()
} else {
vec![0x73; payload.len()]
};
let new_metadata = HashMap::from([
("content-type".to_string(), "text/plain".to_string()),
("x-amz-meta-replacement".to_string(), "kept".to_string()),
]);
if self_copy {
let mut source = store
.get_object_info(&bucket, object, &options)
.await
.expect("self-copy source");
source.metadata_only = false;
source.user_defined = Arc::new(new_metadata);
source.put_object_reader = Some(PutObjReader::from_vec(expected.clone()));
store
.copy_object(&bucket, object, &bucket, object, &mut source, &options, &options)
.await
.expect("materialized self-copy");
} else {
store
.put_object(
&bucket,
object,
&mut PutObjReader::from_vec(expected.clone()),
&ObjectOptions {
user_defined: new_metadata,
..options.clone()
},
)
.await
.expect("overwrite transitioned null version");
}
let set = store.pools[0].get_disks_by_key(object);
let versions = set
.load_file_info_versions_exact(&bucket, object)
.await
.expect("read committed disk metadata")
.expect("replacement metadata exists");
let free: Vec<_> = versions
.versions
.iter()
.chain(versions.free_versions.iter())
.filter(|fi| fi.tier_free_version())
.collect();
assert_eq!(free.len(), 1, "{state:?}, suspended={suspended}, copy={self_copy}");
assert_eq!(free[0].transitioned_objname, remote);
assert_eq!(free[0].transition_version_state, state);
assert!(backend.contains(&remote).await, "commit must not delete remote bytes before cleanup");
let removed_before = backend.remove_count().await;
// Restart before queue delivery. The new runtime must
// reconstruct ownership solely from the committed xl.meta.
let tier_config = ctx
.tier_config_mgr()
.read()
.await
.tiers
.get(tier)
.expect("tier configuration survives restart")
.clone_with_credentials();
drop(set);
shutdown.cancel();
drop(store);
drop(ctx);
(ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-overwrite-restart", &[4]))
.await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
{
let manager = ctx.tier_config_mgr();
let mut manager = manager.write().await;
manager.tiers.insert(tier.to_string(), tier_config);
manager
.install_test_driver(tier, Box::new(backend.clone()))
.expect("rebind the same remote destination after restart");
}
let set = store.pools[0].get_disks_by_key(object);
ExpiryState::resize_workers(1, Arc::clone(&store)).await;
let recovered = recover_tier_free_versions(Arc::clone(&store), 100, None, None)
.await
.expect("recover persisted cleanup owner");
assert!(recovered.enqueued >= 1);
tokio::time::timeout(Duration::from_secs(30), async {
loop {
let versions = set
.load_file_info_versions_exact(&bucket, object)
.await
.expect("read cleanup progress")
.expect("new object must survive cleanup");
if versions
.versions
.iter()
.chain(versions.free_versions.iter())
.all(|fi| !fi.tier_free_version())
{
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("cleanup must converge");
assert!(!backend.contains(&remote).await);
assert_eq!(backend.remove_count().await, removed_before + 1, "one remote DELETE per owner");
assert_eq!(backend.remove_versions().await.last(), Some(&(remote.clone(), version.to_string())));
let mut reader = store
.get_object_reader(&bucket, object, None, HeaderMap::new(), &options)
.await
.expect("replacement remains readable");
let mut actual = Vec::new();
reader.stream.read_to_end(&mut actual).await.expect("read replacement bytes");
assert_eq!(actual, expected);
let current = store
.get_object_info(&bucket, object, &options)
.await
.expect("replacement metadata");
assert_eq!(current.user_defined.get("content-type").map(String::as_str), Some("text/plain"));
assert_eq!(current.user_defined.get("x-amz-meta-replacement").map(String::as_str), Some("kept"));
assert!(current.transitioned_object.status.is_empty());
}
}
}
shutdown.cancel();
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]
@@ -15503,12 +15721,20 @@ mod tests {
);
tokio::time::timeout(Duration::from_secs(30), async {
loop {
let metadata_absent = store.pools[0]
.get_disks_by_key(causal)
.load_file_info_versions_exact(bucket, causal)
.await
.expect("causal batch cleanup metadata should remain readable")
.is_none();
let metadata_absent = {
// Synchronize with cleanup so the snapshot cannot span per-disk marker removal.
let mut read_opts = ObjectOptions::default();
let _guards = store
.acquire_all_physical_object_read_locks("batch_transitioned_delete_test", bucket, causal, &mut read_opts)
.await
.expect("causal batch cleanup observation should acquire object read locks");
store.pools[0]
.get_disks_by_key(causal)
.load_file_info_versions_exact(bucket, causal)
.await
.expect("causal batch cleanup metadata should remain readable")
.is_none()
};
if metadata_absent && backend.remove_versions().await.len() >= 2 {
return;
}
@@ -15587,12 +15813,24 @@ mod tests {
);
tokio::time::timeout(Duration::from_secs(30), async {
loop {
let metadata_absent = store.pools[0]
.get_disks_by_key(versioned_causal)
.load_file_info_versions_exact(bucket, versioned_causal)
.await
.expect("versioned causal batch cleanup metadata should remain readable")
.is_none();
let metadata_absent = {
let mut read_opts = ObjectOptions::default();
let _guards = store
.acquire_all_physical_object_read_locks(
"batch_transitioned_delete_test",
bucket,
versioned_causal,
&mut read_opts,
)
.await
.expect("versioned causal batch cleanup observation should acquire object read locks");
store.pools[0]
.get_disks_by_key(versioned_causal)
.load_file_info_versions_exact(bucket, versioned_causal)
.await
.expect("versioned causal batch cleanup metadata should remain readable")
.is_none()
};
if metadata_absent && backend.remove_versions().await.len() == 3 {
return;
}
+475 -38
View File
@@ -20,6 +20,7 @@ fn to_filemeta_err(err: Error) -> rustfs_filemeta::Error {
err.narrow_to_filemeta().unwrap_or_else(rustfs_filemeta::Error::other)
}
use crate::bucket::lifecycle::bucket_lifecycle_ops::free_version_remote_tuple_matches;
use crate::bucket::metadata_sys::{
get_versioning_config, has_authoritative_never_versioned_state, has_authoritative_never_versioned_state_in,
};
@@ -70,7 +71,7 @@ use tokio::io::duplex;
use tokio::sync::broadcast::{self};
use tokio::sync::mpsc::{self, Receiver, Sender};
use tokio::sync::{OnceCell, RwLock};
use tokio::task::JoinSet;
use tokio::task::{JoinHandle, JoinSet};
use tokio_util::sync::CancellationToken;
use tracing::{Instrument, debug, error, info, warn};
use uuid::Uuid;
@@ -4331,6 +4332,7 @@ impl ECStore {
"store list_merged started"
);
let rx = rx.child_token();
let mut futures = Vec::new();
let mut inputs = Vec::new();
@@ -4346,16 +4348,10 @@ impl ECStore {
}
}
tokio::spawn(
async move {
if let Err(err) = merge_entry_channels(rx, inputs, sender.clone(), 1).await {
error!("merge_entry_channels err {:?}", err)
}
}
.instrument(tracing::Span::current()),
);
let merge_task = spawn_listing_merge(rx, inputs, sender);
let results = join_all(futures).await;
merge_task.await.map_err(Error::from)??;
let mut all_at_eof = true;
@@ -4422,6 +4418,7 @@ impl ECStore {
) -> Result<()> {
check_list_objs_args(bucket, prefix, &None)?;
let rx = rx.child_token();
let mut futures = Vec::new();
let mut inputs = Vec::new();
@@ -4783,17 +4780,11 @@ impl ECStore {
.instrument(tracing::Span::current()),
);
tokio::spawn(
async move {
if let Err(err) = merge_entry_channels(rx, inputs, merge_tx, 1).await {
error!("merge_entry_channels err {:?}", err)
}
}
.instrument(tracing::Span::current()),
);
let merge_task = spawn_listing_merge(rx, inputs, merge_tx);
let walk_started = std::time::Instant::now();
let walk_results = join_all(futures).await;
merge_task.await.map_err(Error::from)??;
let mut errs = Vec::new();
for walk_result in walk_results {
match walk_result {
@@ -5068,6 +5059,130 @@ async fn send_or_cancel(rx: &CancellationToken, out_channel: &Sender<MetaCacheEn
}
}
/// Each input has already been resolved inside its own erasure set. This is a
/// union of version histories, never a quorum vote between unrelated pools.
fn merge_object_entry_versions(first: &mut MetaCacheEntry, others: impl Iterator<Item = MetaCacheEntry>) -> Result<()> {
let name = first.name.clone();
let mut versions: HashMap<(Option<Uuid>, bool), (FileMetaShallowVersion, ObjectInfo)> = HashMap::new();
for mut entry in std::iter::once(std::mem::take(first)).chain(others) {
let meta = match entry.cached.take() {
Some(meta) => meta,
None => FileMeta::load(&entry.metadata).map_err(|_| Error::FileCorrupt)?,
};
if meta.versions.is_empty() {
return Err(Error::FileCorrupt);
}
for version in meta.versions {
let parsed = version.parse_version_meta().map_err(|_| Error::FileCorrupt)?;
if !parsed.valid() || parsed.version_type != version.header.version_type {
return Err(Error::FileCorrupt);
}
let fi = parsed.into_fileinfo("", &name, true).map_err(|_| Error::FileCorrupt)?;
let version_id = fi.version_id.filter(|id| !id.is_nil());
if version_id != version.header.version_id.filter(|id| !id.is_nil())
|| fi.mod_time != version.header.mod_time
|| fi.tier_free_version() != version.header.free_version()
{
return Err(Error::FileCorrupt);
}
let info = ObjectInfo::from_file_info(&fi, "", &name, true);
let identity = (version_id, version.header.free_version());
match versions.entry(identity) {
std::collections::hash_map::Entry::Vacant(slot) => {
slot.insert((version, info));
}
std::collections::hash_map::Entry::Occupied(mut slot) => {
let (previous, previous_info) = slot.get();
// Suspended and unversioned writes replace the one null
// slot. Distinct UUID versions never supersede each other.
if version_id.is_none() && info.mod_time != previous_info.mod_time {
if info.mod_time > previous_info.mod_time {
slot.insert((version, info));
}
continue;
}
let equivalent = if info.delete_marker && previous_info.delete_marker {
super::object::is_equivalent_data_movement_delete_marker(&info, previous_info)
} else {
crate::data_movement::is_equivalent_data_movement_object_identity(&info, previous_info, true, true)
};
if !equivalent {
return Err(Error::FileCorrupt);
}
// Equivalent migrated copies can have different coding or
// data directories. Choose a stable representation without
// making input order part of the S3 version order.
if version.meta < previous.meta {
slot.insert((version, info));
}
}
}
}
}
let mut live_remote_references = HashMap::<String, HashMap<String, Vec<ObjectInfo>>>::new();
for (_, info) in versions.values() {
if !info.transitioned_object.free_version && info.transitioned_object.status == rustfs_filemeta::TRANSITION_COMPLETE {
live_remote_references
.entry(info.transitioned_object.tier.clone())
.or_default()
.entry(info.transitioned_object.name.clone())
.or_default()
.push(info.clone());
}
}
// Keep cleanup durable in its source xl.meta, but do not expose it to a
// merged recovery walk while another physical pool still owns the tuple.
versions.retain(|_, (_, info)| {
!info.transitioned_object.free_version
|| !live_remote_references
.get(info.transitioned_object.tier.as_str())
.and_then(|by_name| by_name.get(info.transitioned_object.name.as_str()))
.is_some_and(|candidates| {
candidates
.iter()
.any(|live| free_version_remote_tuple_matches(info, live).unwrap_or(false))
})
});
let mut merged = FileMeta::new();
merged.versions = versions.into_values().map(|(version, _)| version).collect();
merged.versions.sort_by(|a, b| {
if a.header.sorts_before(&b.header) {
std::cmp::Ordering::Less
} else if b.header.sorts_before(&a.header) {
std::cmp::Ordering::Greater
} else {
std::cmp::Ordering::Equal
}
});
let metadata = merged.marshal_msg()?;
*first = MetaCacheEntry {
name,
metadata,
cached: Some(merged),
reusable: true,
};
Ok(())
}
/// `rx` is private to the producers. Cancelling it on a merge error must not
/// cancel the request token, which would suppress that error at the API edge.
fn spawn_listing_merge(
rx: CancellationToken,
inputs: Vec<Receiver<MetaCacheEntry>>,
sender: Sender<MetaCacheEntry>,
) -> JoinHandle<Result<()>> {
tokio::spawn(
async move {
let result = merge_entry_channels(rx.clone(), inputs, sender, 1).await;
if result.is_err() {
rx.cancel();
}
result
}
.instrument(tracing::Span::current()),
)
}
async fn merge_entry_channels(
rx: CancellationToken,
in_channels: Vec<Receiver<MetaCacheEntry>>,
@@ -5133,6 +5248,7 @@ async fn merge_entry_channels(
// after anything greater has been emitted).
let mut last_emitted = String::new();
let mut group: Vec<Box<MergeHead>> = Vec::new();
let mut object_entries: Vec<MetaCacheEntry> = Vec::new();
let mut refill: Vec<usize> = Vec::with_capacity(in_channels.len());
while let Some(Reverse(first)) = heap.pop() {
@@ -5150,7 +5266,7 @@ async fn merge_entry_channels(
// Resolve the same-name group to one winner (heads arrive in ascending
// channel order):
// - prefix dir vs prefix dir: the first (lowest channel) wins;
// - object vs object: the later channel wins (legacy authority rule);
// - object vs object: merge the independently resolved version stacks;
// - object vs prefix dir: same-name means both end with the separator,
// i.e. the object is an explicit "directory marker" for the same S3
// key — it shadows the prefix dir so the key does not surface as
@@ -5168,11 +5284,27 @@ async fn merge_entry_channels(
if dir_winner.is_none() {
dir_winner = Some(head);
}
} else if let Some(winner) = object_winner.as_ref() {
// Key-only candidates carry no version metadata and cannot
// replace a resolved stack or contribute a quorum vote.
if head.entry.is_object() {
if winner.entry.is_object() {
object_entries.push(head.entry);
} else {
object_winner = Some(head);
}
}
} else {
object_winner = Some(head);
}
}
if !object_entries.is_empty()
&& let Some(winner) = object_winner.as_mut()
{
merge_object_entry_versions(&mut winner.entry, object_entries.drain(..))?;
}
if let Some(head) = object_winner.or(dir_winner)
&& head.entry.name != last_emitted
{
@@ -5605,6 +5737,7 @@ impl Sets {
"sets list_merged started"
);
let rx = rx.child_token();
let mut futures = Vec::new();
let mut inputs = Vec::new();
@@ -5617,16 +5750,10 @@ impl Sets {
futures.push(async move { set.list_path(rx_clone, opts, send).await });
}
tokio::spawn(
async move {
if let Err(err) = merge_entry_channels(rx, inputs, sender.clone(), 1).await {
error!("merge_entry_channels err {:?}", err);
}
}
.instrument(tracing::Span::current()),
);
let merge_task = spawn_listing_merge(rx, inputs, sender);
let results = join_all(futures).await;
merge_task.await.map_err(Error::from)??;
let mut all_at_eof = true;
let mut errs = Vec::new();
for result in results {
@@ -5677,6 +5804,7 @@ impl Sets {
) -> Result<()> {
check_list_objs_args(bucket, prefix, &None)?;
let rx = rx.child_token();
let mut futures = Vec::new();
let mut inputs = Vec::new();
@@ -6007,17 +6135,11 @@ impl Sets {
.instrument(tracing::Span::current()),
);
tokio::spawn(
async move {
if let Err(err) = merge_entry_channels(rx, inputs, merge_tx, 1).await {
error!("merge_entry_channels err {:?}", err)
}
}
.instrument(tracing::Span::current()),
);
let merge_task = spawn_listing_merge(rx, inputs, merge_tx);
let walk_started = std::time::Instant::now();
let walk_results = join_all(futures).await;
merge_task.await.map_err(Error::from)??;
let mut errs = Vec::new();
for walk_result in walk_results {
match walk_result {
@@ -7016,7 +7138,7 @@ mod test {
};
use crate::cache_value::metacache_set::{FallbackClaimTracker, TestReaderBehavior, list_path_raw};
use crate::disk::{DiskAPI, DiskOption, STORAGE_FORMAT_FILE, endpoint::Endpoint, error::DiskError, new_disk};
use crate::error::StorageError;
use crate::error::{Result, StorageError};
use crate::object_api::ObjectInfo;
use rustfs_filemeta::{
FileInfo, FileMeta, FileMetaVersion, MetaCacheEntries, MetaCacheEntriesSorted, MetaCacheEntry, MetaDeleteMarker,
@@ -7338,6 +7460,47 @@ mod test {
}
}
fn test_transitioned_meta_entry(name: &str, remote_object: &str, delete_source: bool) -> MetaCacheEntry {
let mut source = FileInfo::new(name, 2, 2);
source.volume = "bucket".to_string();
source.name = name.to_string();
source.version_id = Some(Uuid::from_u128(1));
source.versioned = true;
source.size = 1;
source.mod_time = Some(time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp"));
source.transition_status = rustfs_filemeta::TRANSITION_COMPLETE.to_string();
source.transition_tier = "WARM".to_string();
source.transitioned_objname = remote_object.to_string();
source.transition_version = Some("remote-version".to_string());
source.transition_version_state = rustfs_filemeta::TransitionVersionState::Exact;
rustfs_utils::http::metadata_compat::insert_str(
&mut source.metadata,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
"00".repeat(32),
);
let mut meta = FileMeta::new();
meta.add_version(source.clone())
.expect("test metadata should accept transitioned source");
if delete_source {
let mut delete = FileInfo {
name: name.to_string(),
version_id: source.version_id,
..Default::default()
};
delete.set_tier_free_version_id(&Uuid::from_u128(2).to_string());
meta.delete_version(&delete)
.expect("transitioned delete should create a free-version owner");
}
let metadata = meta.marshal_msg().expect("test transitioned metadata should marshal");
MetaCacheEntry {
name: name.to_string(),
metadata,
cached: Some(meta),
reusable: false,
}
}
fn test_object_with_delete_marker_meta_entry(
name: &str,
object_mod_time: time::OffsetDateTime,
@@ -10399,7 +10562,7 @@ mod test {
}
#[tokio::test]
async fn merge_entry_channels_documents_candidate_metadata_authority_risk() {
async fn merge_entry_channels_preserves_cross_pool_delete_marker_versions() {
let (tx_a, rx_a) = mpsc::channel(4);
let (tx_b, rx_b) = mpsc::channel(4);
let (tx_c, rx_c) = mpsc::channel(4);
@@ -10428,9 +10591,13 @@ mod test {
.expect("merged entry should be present");
assert_eq!(merged.name, "obj-a");
assert!(
!merged.is_latest_delete_marker(),
"current merge consumes candidate metadata bytes; future index-backed strong modes must live-verify metadata instead"
merged.is_latest_delete_marker(),
"a newer marker must remain current across independently resolved pools"
);
let versions = merged.file_info_versions("bucket").expect("merged versions should decode");
assert_eq!(versions.versions.len(), 2, "retain the historical object and deduplicate the marker");
assert!(versions.versions[0].deleted && versions.versions[0].is_latest);
assert!(!versions.versions[1].deleted && !versions.versions[1].is_latest);
assert!(
matches!(timeout(Duration::from_secs(1), out_rx.recv()).await, Ok(None)),
"merge should not emit a duplicate entry for the same key"
@@ -10442,6 +10609,276 @@ mod test {
.expect("merge task should succeed");
}
fn rewrite_test_version(mut entry: MetaCacheEntry, change: impl FnOnce(&mut FileMetaVersion)) -> MetaCacheEntry {
let meta = entry.cached.as_mut().expect("test metadata should be decoded");
assert_eq!(meta.versions.len(), 1);
let mut version = meta.versions[0].parse_version_meta().expect("test version should decode");
change(&mut version);
meta.versions[0] = version.try_into().expect("test version should encode");
entry.metadata = meta.marshal_msg().expect("test metadata should encode");
entry
}
async fn merge_test_object_entries(entries: Vec<MetaCacheEntry>) -> Result<MetaCacheEntry> {
let mut inputs = Vec::with_capacity(entries.len());
for entry in entries {
let (sender, receiver) = mpsc::channel(1);
sender.send(entry).await.expect("fixture entry should queue");
inputs.push(receiver);
}
let (sender, mut receiver) = mpsc::channel(1);
let task = tokio::spawn(merge_entry_channels(CancellationToken::new(), inputs, sender, 1));
let entry = receiver.recv().await;
task.await.expect("merge must not panic")?;
Ok(entry.expect("a valid same-key group must produce an entry"))
}
#[tokio::test]
async fn merge_entry_channels_orders_complete_histories_independently_of_pool_order() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let first = test_object_meta_entry_with_erasure_versions("key", &[(time, "first", 4, 2)]);
let second = rewrite_test_version(
test_object_meta_entry_with_erasure_versions("key", &[(time, "second", 4, 2)]),
|version| version.object.as_mut().expect("object version").version_id = Some(Uuid::from_u128(2)),
);
let marker = test_delete_marker_meta_entry("key", time + time::Duration::seconds(1));
let inputs = [first, second, marker];
let mut expected = None;
for order in [[0, 1, 2], [0, 2, 1], [1, 0, 2], [1, 2, 0], [2, 0, 1], [2, 1, 0]] {
let entry = merge_test_object_entries(order.map(|index| inputs[index].clone()).to_vec())
.await
.expect("disjoint version chains should merge");
let versions = entry.file_info_versions("bucket").expect("merged versions should decode");
assert_eq!(versions.versions.len(), 3);
assert!(versions.versions[0].deleted && versions.versions[0].is_latest);
assert!(
versions.versions[1..]
.iter()
.all(|version| !version.deleted && !version.is_latest)
);
assert!(versions.versions.iter().all(|version| version.num_versions == 3));
let identities = versions.versions.iter().map(|version| version.version_id).collect::<Vec<_>>();
assert!(identities.contains(&Some(Uuid::from_u128(1))));
assert!(identities.contains(&Some(Uuid::from_u128(2))));
if let Some(expected) = &expected {
assert_eq!(&identities, expected, "equal-time versions must have stable pagination order");
} else {
expected = Some(identities);
}
}
}
#[tokio::test]
async fn merge_entry_channels_key_only_candidates_do_not_override_version_metadata() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let marker = test_delete_marker_meta_entry("key", time);
for entries in [
vec![test_meta_entry("key"), marker.clone()],
vec![marker.clone(), test_meta_entry("key")],
] {
let mut merged = merge_test_object_entries(entries)
.await
.expect("merge a name with resolved metadata");
assert!(merged.is_latest_delete_marker());
assert_eq!(merged.file_info_versions("bucket").expect("decode marker").versions.len(), 1);
}
}
#[tokio::test]
async fn merge_entry_channels_accepts_equivalent_migrated_coding_and_data_dirs() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let first = test_object_meta_entry_with_erasure_versions("key", &[(time, "same-etag", 4, 2)]);
let second = rewrite_test_version(
test_object_meta_entry_with_erasure_versions("key", &[(time, "same-etag", 6, 2)]),
|version| version.object.as_mut().expect("object version").data_dir = Some(Uuid::from_u128(42)),
);
let forward = merge_test_object_entries(vec![first.clone(), second.clone()])
.await
.expect("valid migration copies");
let reverse = merge_test_object_entries(vec![second, first])
.await
.expect("reversed migration copies");
assert_eq!(forward.metadata, reverse.metadata, "representation must not depend on channel order");
let versions = forward.file_info_versions("bucket").expect("merged metadata should decode");
assert_eq!(versions.versions.len(), 1);
assert_eq!(versions.versions[0].metadata.get("etag").map(String::as_str), Some("same-etag"));
}
#[tokio::test]
async fn merge_entry_channels_defers_free_version_while_same_remote_source_is_live() {
let live = test_transitioned_meta_entry("key", "remote/shared", false);
let free = test_transitioned_meta_entry("key", "remote/shared", true);
for inputs in [vec![live.clone(), free.clone()], vec![free.clone(), live.clone()]] {
let merged = merge_test_object_entries(inputs)
.await
.expect("same remote source and cleanup owner should merge");
let versions = merged
.file_info_versions_with_free_versions("bucket")
.expect("merged transition history should decode");
assert_eq!(versions.versions.len(), 1);
assert!(versions.free_versions.is_empty(), "a live remote reference must defer cleanup discovery");
}
let unrelated = test_transitioned_meta_entry("key", "remote/other", false);
let merged = merge_test_object_entries(vec![free, unrelated])
.await
.expect("unrelated remote references should merge");
let versions = merged
.file_info_versions_with_free_versions("bucket")
.expect("merged transition history should decode");
assert_eq!(versions.versions.len(), 1);
assert_eq!(versions.free_versions.len(), 1, "an unrelated source must not suppress cleanup");
}
#[tokio::test]
async fn merge_entry_channels_rejects_conflicting_version_identity_and_metadata() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let original = test_object_meta_entry_with_erasure_versions("key", &[(time, "original", 4, 2)]);
for key in [
"etag",
"x-amz-tagging",
"x-amz-object-lock-mode",
"x-amz-object-lock-retain-until-date",
] {
let changed = rewrite_test_version(original.clone(), |version| {
version
.object
.as_mut()
.expect("object version")
.meta_user
.insert(key.to_string(), "changed".to_string());
});
for pair in [[original.clone(), changed.clone()], [changed, original.clone()]] {
let err = merge_test_object_entries(pair.to_vec())
.await
.expect_err("conflicting copies must fail");
assert_eq!(err, StorageError::FileCorrupt, "conflict in {key} must not become arbitrary metadata");
}
}
let marker = rewrite_test_version(test_delete_marker_meta_entry("key", time), |version| {
version.delete_marker.as_mut().expect("delete marker").version_id = Some(Uuid::from_u128(1));
});
assert_eq!(
merge_test_object_entries(vec![original, marker])
.await
.expect_err("UUID type conflict"),
StorageError::FileCorrupt
);
}
#[tokio::test]
async fn merge_entry_channels_reconciles_null_overwrite_without_losing_uuid_history() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let history = test_object_meta_entry_with_erasure_versions("key", &[(time, "history", 4, 2)]);
let old_null = rewrite_test_version(history.clone(), |version| {
version.object.as_mut().expect("null object").version_id = None;
});
let marker = rewrite_test_version(test_delete_marker_meta_entry("key", time + time::Duration::seconds(1)), |version| {
version.delete_marker.as_mut().expect("null marker").version_id = Some(Uuid::nil());
});
for inputs in [
vec![old_null.clone(), marker.clone(), history.clone()],
vec![history, marker, old_null],
] {
let entry = merge_test_object_entries(inputs)
.await
.expect("new null slot should replace old null slot");
let versions = entry.file_info_versions("bucket").expect("null versions should decode");
assert_eq!(versions.versions.len(), 2);
assert!(versions.versions[0].deleted && versions.versions[0].is_latest);
assert!(versions.versions[0].version_id.is_none_or(|id| id.is_nil()));
assert_eq!(versions.versions[1].version_id, Some(Uuid::from_u128(1)));
}
}
#[tokio::test]
async fn merge_entry_channels_rejects_corrupt_version_headers_and_empty_stacks() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let original = test_object_meta_entry_with_erasure_versions("key", &[(time, "etag", 4, 2)]);
for empty in [false, true] {
let mut corrupt = original.clone();
let meta = corrupt.cached.as_mut().expect("fixture metadata");
if empty {
meta.versions.clear();
} else {
meta.versions[0].header.version_id = Some(Uuid::from_u128(99));
}
corrupt.metadata = meta.marshal_msg().expect("encode corrupt fixture");
assert_eq!(
merge_test_object_entries(vec![original.clone(), corrupt])
.await
.expect_err("corrupt candidate must fail"),
StorageError::FileCorrupt
);
}
let mut malformed = original.clone();
malformed.cached = None;
malformed.metadata = vec![0xff];
assert_eq!(
merge_test_object_entries(vec![original, malformed])
.await
.expect_err("malformed metadata must fail"),
StorageError::FileCorrupt
);
}
#[tokio::test]
async fn merge_entry_channels_does_not_combine_subquorum_markers_across_erasure_sets() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let old = test_object_meta_entry_with_erasure_versions("key", &[(time, "history", 4, 2)]);
let marked = test_object_with_delete_marker_meta_entry("key", time, time + time::Duration::seconds(1));
let mut inputs = Vec::new();
for marker_copies in [1, 2] {
let resolver = list_metadata_resolution_params("bucket".to_string(), 3, 3, true, 0);
let copies = (0..3)
.map(|index| Some(if index < marker_copies { marked.clone() } else { old.clone() }))
.collect();
let entry = resolve_listing_entries(MetaCacheEntries(copies), resolver, false)
.expect("each set independently retains its quorum-backed history");
inputs.push(entry);
}
let merged = merge_test_object_entries(inputs).await.expect("merge resolved histories");
let versions = merged.file_info_versions("bucket").expect("decode merged history");
assert_eq!(
versions.versions.len(),
1,
"three marker copies across two EC domains do not form a quorum"
);
assert!(!versions.versions[0].deleted);
}
#[tokio::test]
async fn listing_merge_preserves_error_after_partial_output_without_cancelling_request() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let (first_tx, first_rx) = mpsc::channel(2);
let (second_tx, second_rx) = mpsc::channel(1);
first_tx.send(test_meta_entry("a/")).await.expect("queue preceding prefix");
first_tx
.send(test_object_meta_entry_with_erasure_versions("b", &[(time, "one", 4, 2)]))
.await
.expect("queue first copy");
second_tx
.send(test_object_meta_entry_with_erasure_versions("b", &[(time, "two", 4, 2)]))
.await
.expect("queue conflicting copy");
drop(first_tx);
drop(second_tx);
let request = CancellationToken::new();
let workers = request.child_token();
let (sender, mut receiver) = mpsc::channel(1);
let task = super::spawn_listing_merge(workers.clone(), vec![first_rx, second_rx], sender);
assert_eq!(receiver.recv().await.expect("preceding result should arrive").name, "a/");
assert!(receiver.recv().await.is_none());
assert_eq!(
task.await
.expect("merge task must not panic")
.expect_err("conflict must propagate"),
StorageError::FileCorrupt
);
assert!(workers.is_cancelled(), "failed merge must stop the disk producers");
assert!(!request.is_cancelled(), "the API must still observe the merge error");
}
#[tokio::test]
async fn merge_entry_channels_handles_single_channel() {
let (tx, rx) = mpsc::channel(4);
+611 -23
View File
@@ -50,7 +50,7 @@ use crate::services::notification_sys::{
use crate::services::tier::tier::{TierConfigMgr, TierDestinationId, TierOperationLease, tier_destination_id_from_metadata};
use crate::set_disk::{
SetDisks, get_lock_acquire_timeout, get_object_lock_diag_slow_acquire_threshold, get_object_lock_diag_slow_hold_threshold,
is_lock_optimization_enabled, is_object_lock_diag_enabled, same_distributed_lock_domain,
is_lock_optimization_enabled, is_object_lock_diag_enabled,
};
use crate::storage_api_contracts::{
list::ListOperations as _,
@@ -2013,7 +2013,7 @@ fn effective_object_actual_size(info: &ObjectInfo) -> Option<i64> {
info.get_actual_size().ok()
}
fn is_equivalent_data_movement_delete_marker(source: &ObjectInfo, target: &ObjectInfo) -> bool {
pub(super) fn is_equivalent_data_movement_delete_marker(source: &ObjectInfo, target: &ObjectInfo) -> bool {
is_data_movement_delete_marker(source)
&& is_data_movement_delete_marker(target)
&& source.version_id == target.version_id
@@ -3440,10 +3440,15 @@ impl ECStore {
for pool in &self.pools {
let hashed_set = pool.get_disks_by_key(object);
let lock_domain_already_held = !distributed
|| locked_sets
.iter()
.any(|locked_set| same_distributed_lock_domain(&locked_set.lockers, &hashed_set.lockers));
let mut lock_domain_already_held = !distributed;
if !lock_domain_already_held {
for locked_set in &locked_sets {
if locked_set.shares_namespace_lock_domain(&hashed_set).await {
lock_domain_already_held = true;
break;
}
}
}
if lock_domain_already_held {
continue;
}
@@ -3500,10 +3505,15 @@ impl ECStore {
let mut locked_sets = vec![fixed_set];
for pool in &self.pools {
for set in &pool.disk_set {
let lock_domain_already_held = !distributed
|| locked_sets
.iter()
.any(|locked_set| same_distributed_lock_domain(&locked_set.lockers, &set.lockers));
let mut lock_domain_already_held = !distributed;
if !lock_domain_already_held {
for locked_set in &locked_sets {
if locked_set.shares_namespace_lock_domain(set).await {
lock_domain_already_held = true;
break;
}
}
}
if lock_domain_already_held {
continue;
}
@@ -3581,10 +3591,15 @@ impl ECStore {
let mut locked_sets = vec![fixed_set];
for pool in &self.pools {
for set in &pool.disk_set {
let lock_domain_already_held = !distributed
|| locked_sets
.iter()
.any(|locked_set| same_distributed_lock_domain(&locked_set.lockers, &set.lockers));
let mut lock_domain_already_held = !distributed;
if !lock_domain_already_held {
for locked_set in &locked_sets {
if locked_set.shares_namespace_lock_domain(set).await {
lock_domain_already_held = true;
break;
}
}
}
if lock_domain_already_held {
continue;
}
@@ -3667,11 +3682,15 @@ impl ECStore {
.get(pool_idx)
.ok_or_else(|| Error::other(format!("invalid data movement publication pool {pool_idx}")))?;
let set = pool.get_disks_by_key(object);
let lock_domain_already_held = !locked_sets.is_empty()
&& (!distributed
|| locked_sets.iter().any(|locked_set: &Arc<crate::set_disk::SetDisks>| {
same_distributed_lock_domain(&locked_set.lockers, &set.lockers)
}));
let mut lock_domain_already_held = !locked_sets.is_empty() && !distributed;
if !lock_domain_already_held {
for locked_set in &locked_sets {
if locked_set.shares_namespace_lock_domain(&set).await {
lock_domain_already_held = true;
break;
}
}
}
if lock_domain_already_held {
continue;
}
@@ -4772,7 +4791,13 @@ impl ECStore {
return Ok(ObjectInfo::default());
}
let gopts = delete_pool_lookup_opts(&opts, true);
let creates_latest_marker = should_create_delete_marker_for_missing_object(&opts);
let mut gopts = delete_pool_lookup_opts(&opts, true);
if creates_latest_marker {
// An unwritable source still owns its current version. Hiding it
// during lookup would turn a rejected write into a new-pool marker.
gopts.skip_rebalancing = false;
}
if opts.data_movement {
let existing_pool_info = self.get_pool_info_existing_with_opts(bucket, object, &gopts).await;
@@ -4898,7 +4923,12 @@ impl ECStore {
}
// Determine which pool contains it
let (mut pinfo, errs) = match self.get_pool_info_existing_with_opts(bucket, object, &gopts).await {
let existing_pool_info = if creates_latest_marker {
self.get_pool_info_for_delete_marker(bucket, object, &gopts).await
} else {
self.get_pool_info_existing_with_opts(bucket, object, &gopts).await
};
let (mut pinfo, errs) = match existing_pool_info {
Ok(res) => res,
Err(err) if is_err_read_quorum(&err) => return Err(StorageError::ErasureWriteQuorum),
Err(err) if is_err_object_not_found(&err) && should_create_delete_marker_for_missing_object(&opts) => {
@@ -4935,7 +4965,18 @@ impl ECStore {
}
};
if pinfo.object_info.delete_marker && opts.version_id.is_none() {
if creates_latest_marker && self.is_suspended(pinfo.index).await {
let has_active_reservation = self
.pool_meta
.read()
.await
.has_active_decommission_capacity_reservation(pinfo.index);
if has_active_reservation {
pinfo.index = self.get_pool_idx_no_lock(bucket, object, 0).await?;
}
}
if pinfo.object_info.delete_marker && opts.version_id.is_none() && !creates_latest_marker {
pinfo.object_info.name = decode_dir_object(object);
return Ok(pinfo.object_info);
}
@@ -4957,7 +4998,13 @@ impl ECStore {
}
for pool in self.pools.iter() {
if creates_latest_marker && pool.pool_idx != pinfo.index {
continue;
}
if self.is_suspended(pool.pool_idx).await || self.is_pool_rebalancing(pool.pool_idx).await {
if creates_latest_marker {
return Err(StorageError::SlowDown);
}
continue;
}
@@ -4982,7 +5029,7 @@ impl ECStore {
return Ok(obj);
}
Err(err) => {
if !is_err_object_not_found(&err) && !is_err_version_not_found(&err) {
if creates_latest_marker || (!is_err_object_not_found(&err) && !is_err_version_not_found(&err)) {
return Err(err);
}
}
@@ -5717,6 +5764,7 @@ mod tests {
GetObjectBodyCacheHook, GetObjectBodyCacheHookLookup, GetObjectBodySource, clear_get_object_body_cache_hook,
lookup_get_object_body_cache_hook, register_get_object_body_cache_hook,
};
use crate::set_disk::same_distributed_lock_domain;
use crate::set_disk::{SetDisks, disk_call_counters};
use crate::storage_api_contracts::bucket::MakeBucketOptions;
use crate::storage_api_contracts::lifecycle::TransitionedObject;
@@ -8271,6 +8319,546 @@ mod tests {
}
}
async fn multipool_version_test_store(bucket: &str) -> (Vec<tempfile::TempDir>, Arc<ECStore>) {
let ctx = Arc::new(crate::runtime::instance::InstanceContext::new());
let (mut dirs, first_set) = make_local_set_disks_with_ctx(4, 2, Arc::clone(&ctx)).await;
let (second_dirs, second_set) = make_local_set_disks_with_ctx(4, 2, Arc::clone(&ctx)).await;
dirs.extend(second_dirs);
let store = Arc::new(new_prepared_reader_test_store_with_ctx(&[first_set, second_set], ctx).await);
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
store
.handle_make_bucket(
bucket,
&MakeBucketOptions {
versioning_enabled: true,
..Default::default()
},
)
.await
.expect("create the versioned bucket in both pools");
(dirs, store)
}
#[tokio::test]
async fn multipool_delete_marker_stays_with_existing_versions() {
let bucket = "multipool-marker-routing";
let object = "history.bin";
let (_dirs, store) = multipool_version_test_store(bucket).await;
let mut expected_versions = Vec::new();
for value in 1..=3_u8 {
let written = store.pools[1]
.put_object(
bucket,
object,
&mut PutObjReader::from_vec(vec![value; 4097]),
&ObjectOptions {
versioned: true,
user_defined: HashMap::from([
(rustfs_utils::http::AMZ_OBJECT_TAGGING.to_string(), format!("generation={value}")),
("x-amz-meta-generation".to_string(), value.to_string()),
]),
..Default::default()
},
)
.await
.expect("write a historical version deterministically to pool 1");
expected_versions.push(written.version_id.expect("versioned PUT must acknowledge a UUID"));
}
let marker = store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("delete the current version");
assert!(marker.delete_marker);
let marker_id = marker.version_id.expect("DELETE must acknowledge a marker UUID");
assert!(!expected_versions.contains(&marker_id));
let local = store.pools[1]
.clone()
.inner_list_object_versions(bucket, object, None, None, None, 10)
.await
.expect("list the object-owning pool");
assert_eq!(local.objects.len(), 4, "the marker must be committed beside the three existing versions");
assert_eq!(local.objects[0].version_id, Some(marker_id));
assert!(local.objects[0].delete_marker && local.objects[0].is_latest);
let versions = store
.clone()
.inner_list_object_versions(bucket, object, None, None, None, 10)
.await
.expect("list all pools");
assert_eq!(versions.objects.len(), 4);
assert_eq!(versions.objects.iter().filter(|version| version.is_latest).count(), 1);
for (index, version_id) in expected_versions.iter().copied().enumerate() {
assert!(
versions
.objects
.iter()
.any(|version| version.version_id == Some(version_id) && !version.is_latest)
);
let mut reader = store
.handle_get_object_reader(
bucket,
object,
None,
HeaderMap::new(),
&ObjectOptions {
version_id: Some(version_id.to_string()),
..Default::default()
},
)
.await
.expect("historical version should remain readable");
let mut payload = Vec::new();
reader
.stream
.read_to_end(&mut payload)
.await
.expect("read all historical bytes");
let value = u8::try_from(index + 1).expect("fixture generation fits u8");
assert_eq!(payload, vec![value; 4097]);
let listed = versions
.objects
.iter()
.find(|version| version.version_id == Some(version_id))
.expect("listed historical version");
assert_eq!(listed.user_tags.as_str(), format!("generation={value}"));
assert_eq!(listed.user_defined.get("x-amz-meta-generation"), Some(&value.to_string()));
}
let repeated = store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("repeat simple DELETE");
assert!(repeated.delete_marker);
assert_ne!(
repeated.version_id,
Some(marker_id),
"each enabled-versioning DELETE creates a new marker"
);
store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
version_id: repeated.version_id.map(|id| id.to_string()),
..Default::default()
},
)
.await
.expect("remove the newest marker by identity");
store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
version_id: Some(marker_id.to_string()),
..Default::default()
},
)
.await
.expect("remove the original marker by identity");
let current = store
.get_object_info(bucket, object, &ObjectOptions::default())
.await
.expect("previous version becomes current");
assert_eq!(current.version_id, expected_versions.last().copied());
assert!(!current.delete_marker);
}
#[tokio::test]
async fn multipool_existing_split_history_lists_and_paginates_without_metadata_writes() {
for marker_pool in 0..2 {
let bucket = format!("multipool-split-history-{marker_pool}");
let object = "history.bin";
let (dirs, store) = multipool_version_test_store(&bucket).await;
let mut expected = Vec::new();
for value in 1..=3_u8 {
let version = store.pools[1 - marker_pool]
.put_object(
&bucket,
object,
&mut PutObjReader::from_vec(vec![value; 4097]),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed history through the owning pool's normal write path");
expected.push(version.version_id);
}
let marker = store.pools[marker_pool]
.delete_object(
&bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("reproduce the previously committed split marker with a normal pool DELETE");
assert!(marker.delete_marker);
expected.push(marker.version_id);
expected.reverse();
let mut before = Vec::new();
for dir in &dirs {
before.push(
tokio::fs::read(dir.path().join(&bucket).join(object).join("xl.meta"))
.await
.expect("snapshot persisted version metadata"),
);
}
for max_keys in [1, 2, 4, 10] {
let mut key_marker = None;
let mut version_marker = None;
let mut listed = Vec::new();
let mut completed = false;
for _ in 0..6 {
let page = store
.clone()
.inner_list_object_versions(&bucket, object, key_marker.clone(), version_marker.clone(), None, max_keys)
.await
.expect("read a complete merged version page");
listed.extend(page.objects);
if !page.is_truncated {
completed = true;
break;
}
assert_ne!(
(&page.next_marker, &page.next_version_idmarker),
(&key_marker, &version_marker),
"version cursor must advance"
);
key_marker = page.next_marker;
version_marker = page.next_version_idmarker;
}
assert!(completed, "pagination must terminate");
assert_eq!(listed.iter().map(|version| version.version_id).collect::<Vec<_>>(), expected);
assert!(listed[0].delete_marker && listed[0].is_latest);
assert!(listed[1..].iter().all(|version| !version.delete_marker && !version.is_latest));
}
let visible = store
.clone()
.list_objects_generic(&bucket, "", None, None, 10, false)
.await
.expect("list current objects");
assert!(visible.objects.is_empty(), "the global current marker hides the object");
for (dir, before) in dirs.iter().zip(before) {
assert_eq!(
tokio::fs::read(dir.path().join(&bucket).join(object).join("xl.meta"))
.await
.expect("read unchanged metadata"),
before
);
}
}
}
#[tokio::test]
async fn multipool_conflicting_version_metadata_fails_the_listing_request() {
let bucket = "multipool-version-conflict";
let (_dirs, store) = multipool_version_test_store(bucket).await;
store.pools[0]
.put_object(
bucket,
"a.bin",
&mut PutObjReader::from_vec(b"preceding result".to_vec()),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed an entry before the conflict");
let version_id = Uuid::new_v4();
let mod_time = OffsetDateTime::now_utc();
let mut etags = Vec::new();
for (pool_idx, value) in [(0, 1), (1, 2)] {
let written = store.pools[pool_idx]
.put_object(
bucket,
"z.bin",
&mut PutObjReader::from_vec(vec![value; 4097]),
&ObjectOptions {
versioned: true,
version_id: Some(version_id.to_string()),
mod_time: Some(mod_time),
..Default::default()
},
)
.await
.expect("persist independent conflicting copies");
assert_eq!(written.version_id, Some(version_id));
etags.push(written.etag);
}
assert_ne!(etags[0], etags[1], "fixture must contain a semantic conflict");
let err = store
.clone()
.inner_list_object_versions(bucket, "", None, None, None, 10)
.await
.expect_err("partial output must not hide the merge error");
assert_eq!(err, StorageError::FileCorrupt);
let cancellation = tokio_util::sync::CancellationToken::new();
let (sender, mut receiver) = tokio::sync::mpsc::channel(1);
let walk = store
.clone()
.walk(cancellation.clone(), bucket, "", sender, WalkOptions::default());
let drain = async { while receiver.recv().await.is_some() {} };
let (walk_result, ()) = tokio::time::timeout(Duration::from_secs(5), async { tokio::join!(walk, drain) })
.await
.expect("bounded walk output must drain and terminate on a merge error");
assert_eq!(
walk_result.expect_err("walk must report the same metadata conflict"),
StorageError::FileCorrupt
);
assert!(!cancellation.is_cancelled(), "worker failure must not cancel the caller's request");
}
#[tokio::test]
async fn multipool_marker_rejects_unwritable_owner_without_falling_back() {
let bucket = "multipool-unwritable-owner";
let object = "history.bin";
let (_dirs, store) = multipool_version_test_store(bucket).await;
store.pools[1]
.put_object(
bucket,
object,
&mut PutObjReader::from_vec(vec![1; 4097]),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed the nonzero owner");
*store.pool_meta.write().await = PoolMeta {
pools: vec![prepared_pool_test_status(0, false), prepared_pool_test_status(1, true)],
..Default::default()
};
let err = store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect_err("suspended owner must reject marker creation");
assert_eq!(err, StorageError::SlowDown);
*store.pool_meta.write().await = PoolMeta::default();
let mut rebalancing = crate::services::rebalance::RebalanceStats {
participating: true,
..Default::default()
};
rebalancing.info.status = crate::services::rebalance::RebalStatus::Started;
*store.rebalance_meta.write().await = Some(crate::services::rebalance::RebalanceMeta {
pool_stats: vec![crate::services::rebalance::RebalanceStats::default(), rebalancing],
..Default::default()
});
let error = store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await;
*store.rebalance_meta.write().await = None;
assert_eq!(
error.expect_err("rebalance must not hide the owner during marker lookup"),
StorageError::SlowDown
);
// Isolate object quorum from the bucket metadata preflight: disabling
// a whole pool can otherwise fail before object ownership is looked up.
let quorum_bucket = RUSTFS_META_BUCKET;
store.pools[1]
.put_object(
quorum_bucket,
object,
&mut PutObjReader::from_vec(vec![1; 4097]),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed the object-quorum fixture");
for pool in &store.pools {
pool.put_object(
quorum_bucket,
"split.bin",
&mut PutObjReader::from_vec(vec![2; 4097]),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed another key with history in both pools");
}
let owner = &store.pools[1].disk_set[0];
let healthy = owner.disks.read().await.clone();
for disk in owner.disks.write().await.iter_mut().skip(1) {
*disk = None;
}
let error = store
.delete_object(
quorum_bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await;
let split_error = store
.delete_object(
quorum_bucket,
"split.bin",
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await;
*owner.disks.write().await = healthy;
assert_eq!(
error.expect_err("subquorum owner must not become a new-pool marker"),
StorageError::ErasureWriteQuorum
);
assert_eq!(
split_error.expect_err("a readable older pool does not prove the global current version"),
StorageError::ErasureWriteQuorum
);
let split = store.pools[0]
.clone()
.inner_list_object_versions(quorum_bucket, "split.bin", None, None, None, 10)
.await
.expect("inspect the readable older pool");
assert_eq!(split.objects.len(), 1);
assert!(!split.objects[0].delete_marker);
let other = store.pools[0]
.clone()
.inner_list_object_versions(quorum_bucket, object, None, None, None, 10)
.await
.expect("inspect the other pool");
assert!(other.objects.is_empty());
let history = store.pools[1]
.clone()
.inner_list_object_versions(quorum_bucket, object, None, None, None, 10)
.await
.expect("inspect the restored owner");
assert_eq!(history.objects.len(), 1);
assert!(!history.objects[0].delete_marker);
}
#[tokio::test]
async fn multipool_suspended_and_batch_markers_keep_the_null_slot_on_the_owner() {
let bucket = "multipool-suspended-markers";
let object = "history.bin";
let (_dirs, store) = multipool_version_test_store(bucket).await;
let original = store.pools[1]
.put_object(
bucket,
object,
&mut PutObjReader::from_vec(vec![1; 4097]),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed the UUID history in pool 1");
store
.update_bucket_metadata_config(
bucket,
crate::bucket::metadata::BUCKET_VERSIONING_CONFIG,
b"<VersioningConfiguration><Status>Suspended</Status></VersioningConfiguration>".to_vec(),
)
.await
.expect("persist suspended bucket versioning");
let null = store
.put_object(
bucket,
object,
&mut PutObjReader::from_vec(vec![2; 4097]),
&ObjectOptions {
version_suspended: true,
..Default::default()
},
)
.await
.expect("write the null version beside the history");
assert!(null.version_id.is_none_or(|id| id.is_nil()));
for _ in 0..2 {
let marker = store
.delete_object(
bucket,
object,
ObjectOptions {
version_suspended: true,
..Default::default()
},
)
.await
.expect("replace the null slot with a marker");
assert!(marker.delete_marker);
assert!(marker.version_id.is_none_or(|id| id.is_nil()));
}
let (deleted, errors) = store
.delete_objects(
bucket,
vec![ObjectToDelete {
object_name: object.to_string(),
..Default::default()
}],
ObjectOptions::default(),
)
.await;
assert!(errors.iter().all(Option::is_none), "batch DELETE should succeed: {errors:?}");
assert_eq!(deleted.len(), 1);
assert!(deleted[0].delete_marker);
let versions = store.pools[1]
.clone()
.inner_list_object_versions(bucket, object, None, None, None, 10)
.await
.expect("inspect the owner after batch DELETE");
assert_eq!(versions.objects.len(), 2, "only one null marker and the UUID history remain");
assert!(versions.objects[0].delete_marker && versions.objects[0].is_latest);
assert_eq!(versions.objects[1].version_id, original.version_id);
let other = store.pools[0]
.clone()
.inner_list_object_versions(bucket, object, None, None, None, 10)
.await
.expect("inspect the unused pool");
assert!(other.objects.is_empty());
}
async fn assert_prepared_reader_blocks_writer(store: &ECStore, bucket: &str, object: &str) {
assert_pool_writer_is_blocked(store, 0, bucket, object).await;
}
+50 -1
View File
@@ -611,7 +611,18 @@ impl ECStore {
object: &str,
opts: &ObjectOptions,
) -> Result<(PoolObjInfo, Vec<PoolErr>)> {
self.internal_get_pool_info_existing_with_opts(bucket, object, opts).await
self.internal_get_pool_info_existing_with_opts(bucket, object, opts, false)
.await
}
pub(super) async fn get_pool_info_for_delete_marker(
&self,
bucket: &str,
object: &str,
opts: &ObjectOptions,
) -> Result<(PoolObjInfo, Vec<PoolErr>)> {
self.internal_get_pool_info_existing_with_opts(bucket, object, opts, true)
.await
}
async fn internal_get_pool_info_existing_with_opts(
@@ -619,6 +630,7 @@ impl ECStore {
bucket: &str,
object: &str,
opts: &ObjectOptions,
require_all_pool_reads: bool,
) -> Result<(PoolObjInfo, Vec<PoolErr>)> {
let mut futures = Vec::new();
for pool in self.pools.iter() {
@@ -647,6 +659,15 @@ impl ECStore {
});
}
Err(e) => {
// A readable older pool cannot prove ownership of the
// current version while another pool is unreadable. Check
// both raw and object-scoped quorum errors before sorting.
if require_all_pool_reads && !is_err_object_not_found(&e) && !is_err_version_not_found(&e) {
return Err(match e {
Error::ErasureReadQuorum | Error::InsufficientReadQuorum(_, _) => Error::ErasureWriteQuorum,
err => err,
});
}
ress.push(PoolObjInfo {
index,
err: Some(e),
@@ -656,6 +677,34 @@ impl ECStore {
}
}
if require_all_pool_reads {
let suspended_pools = {
let pool_meta = self.pool_meta.read().await;
(0..self.pools.len())
.map(|idx| pool_meta.is_suspended(idx))
.collect::<Vec<_>>()
};
let candidates = ress
.iter()
.map(|pinfo| LatestObjectInfoCandidate {
info: pinfo.err.is_none().then(|| pinfo.object_info.clone()),
idx: pinfo.index,
err: pinfo.err.clone(),
})
.collect();
let (object_info, index) =
resolve_latest_object_info_candidates_with_pool_state(candidates, &suspended_pools, bucket, object, opts)?;
let pools_with_object = self.pools_with_object(&ress, opts).await;
return Ok((
PoolObjInfo {
index,
object_info,
err: None,
},
pools_with_object,
));
}
ress.sort_by(|a, b| {
let at = a.object_info.mod_time.unwrap_or(OffsetDateTime::UNIX_EPOCH);
let bt = b.object_info.mod_time.unwrap_or(OffsetDateTime::UNIX_EPOCH);
+333 -1
View File
@@ -480,7 +480,110 @@ impl FileMeta {
Ok(())
}
pub fn add_version(&mut self, mut fi: FileInfo) -> Result<()> {
pub fn add_version(&mut self, fi: FileInfo) -> Result<()> {
if let Some(free_version) = self.overwritten_tier_free_version(&fi)? {
// The replacement and its cleanup owner must share one xl.meta
// commit. Keep the original intact if either insertion fails.
let mut next = self.clone();
next.add_version_inner(fi)?;
next.add_version_filemata(free_version)?;
*self = next;
return Ok(());
}
self.add_version_inner(fi)
}
fn overwritten_tier_free_version(&self, fi: &FileInfo) -> Result<Option<FileMetaVersion>> {
use rustfs_utils::http::{
SUFFIX_TIER_FV_ID, SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID,
SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE,
get_consistent_bytes, get_consistent_str, has_internal_suffix, strip_internal_prefix_preserving_case,
};
if fi.version_id.is_some_and(|id| !id.is_nil()) || !contains_key_str(&fi.metadata, SUFFIX_TIER_FV_ID) {
return Ok(None);
}
let Some(existing) = self
.versions
.iter()
.find(|v| v.header.version_id.is_none_or(|id| id.is_nil()))
else {
return Ok(None);
};
let old = existing.parse_version_meta()?;
let Some(mut object) = old.object else {
return Ok(None);
};
let status = get_consistent_bytes(&object.meta_sys, SUFFIX_TRANSITION_STATUS);
if status.is_none()
&& object
.meta_sys
.keys()
.any(|key| has_internal_suffix(key, SUFFIX_TRANSITION_STATUS))
{
// Empty status is a valid local object. The reader distinguishes
// it from conflicting aliases before the ordinary overwrite.
object.into_fileinfo(&fi.volume, &fi.name, false)?;
return Ok(None);
}
if status != Some(TRANSITION_COMPLETE.as_bytes()) {
return Ok(None);
}
// Reuse the reader's alias/state validation. A legacy empty remote
// version is valid and must not be mistaken for conflicting aliases.
object.into_fileinfo(&fi.volume, &fi.name, false)?;
if object
.meta_sys
.keys()
.any(|key| has_internal_suffix(key, SUFFIX_TRANSITION_TIER_DESTINATION_ID))
&& get_consistent_bytes(&object.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID).is_none()
{
return Err(Error::FileCorrupt);
}
let transition_suffixes = [
SUFFIX_TRANSITION_STATUS,
SUFFIX_TRANSITION_TIER,
SUFFIX_TRANSITION_TIER_DESTINATION_ID,
SUFFIX_TRANSITIONED_OBJECTNAME,
SUFFIX_TRANSITIONED_VERSION_ID,
SUFFIX_TRANSITIONED_VERSION_STATE,
];
let replacement = MetaObject::from(fi.clone());
if transition_suffixes
.iter()
.all(|suffix| get_consistent_bytes(&object.meta_sys, suffix) == get_consistent_bytes(&replacement.meta_sys, suffix))
{
return Ok(None);
}
let id = get_consistent_str(&fi.metadata, SUFFIX_TIER_FV_ID).ok_or(Error::FileCorrupt)?;
let id = Uuid::parse_str(id)?;
if id.is_nil() || self.versions.iter().any(|version| version.header.version_id == Some(id)) {
return Err(Error::FileCorrupt);
}
// The reader also accepts legacy key casing. Canonicalize only this
// cleanup source so init_free_version preserves every accepted field,
// including an explicitly empty unversioned remote version.
for suffix in transition_suffixes {
let value = object
.meta_sys
.iter()
.find(|(key, _)| {
strip_internal_prefix_preserving_case(key).is_some_and(|found| found.eq_ignore_ascii_case(suffix))
})
.map(|(_, value)| value.clone());
if let Some(value) = value {
rustfs_utils::http::insert_bytes(&mut object.meta_sys, suffix, value);
}
}
let (free_version, created) = object.init_free_version(fi)?;
if !created {
return Err(Error::FileCorrupt);
}
Ok(Some(free_version))
}
fn add_version_inner(&mut self, mut fi: FileInfo) -> Result<()> {
rustfs_utils::http::remove_str(&mut fi.metadata, rustfs_utils::http::SUFFIX_TIER_FV_ID);
// empty version_id means "null" (versioning disabled/suspended)
if fi.version_id.is_none() {
fi.version_id = Some(Uuid::nil());
@@ -1463,6 +1566,235 @@ mod test {
});
}
fn tier_overwrite_fixture(state: crate::TransitionVersionState) -> (FileMeta, FileInfo) {
let mut source = FileInfo::new("object", 2, 2);
source.erasure.index = 1;
source.mod_time = Some(OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("fixture timestamp"));
source.data_dir = Some(Uuid::new_v4());
source.transition_status = TRANSITION_COMPLETE.to_string();
source.transition_tier = "WARM".to_string();
source.transitioned_objname = "remote/old-object".to_string();
source.transition_version_state = state;
source.transition_version = match state {
crate::TransitionVersionState::Exact => Some("opaque-provider-version".to_string()),
crate::TransitionVersionState::SuspendedNull => Some("null".to_string()),
_ => None,
};
rustfs_utils::http::insert_str(
&mut source.metadata,
rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
"ab".repeat(32),
);
if state == crate::TransitionVersionState::KnownDisabled {
rustfs_utils::http::insert_str(
&mut source.metadata,
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID,
String::new(),
);
}
let mut meta = FileMeta::new();
meta.add_version(source.clone()).expect("seed transitioned null version");
(meta, source)
}
#[test]
fn tier_overwrite_preserves_exact_cleanup_owner_across_reload() {
use crate::TransitionVersionState::{Exact, KnownDisabled, SuspendedNull, Unknown};
use rustfs_utils::http::{MINIO_INTERNAL_PREFIX, RUSTFS_INTERNAL_PREFIX, SUFFIX_TIER_FV_ID};
for state in [Exact, KnownDisabled, SuspendedNull, Unknown] {
for inline in [false, true] {
let (mut meta, _) = tier_overwrite_fixture(state);
let old = meta.versions[0].parse_version_meta().expect("old metadata");
let old = old.object.expect("old object");
let id = Uuid::new_v4();
let mut replacement = FileInfo::new("object", 2, 2);
replacement.version_id = inline.then_some(Uuid::nil());
replacement.mod_time = Some(OffsetDateTime::from_unix_timestamp(1_700_000_001).expect("fixture timestamp"));
replacement.data_dir = Some(Uuid::new_v4());
replacement.size = 3;
if inline {
replacement.data = Some(Bytes::from_static(b"new"));
}
replacement.set_tier_free_version_id(&id.to_string());
meta.add_version(replacement.clone())
.expect("replace transitioned null version");
let bytes = meta.marshal_msg().expect("persist replacement and cleanup owner");
let mut reopened = FileMeta::load(&bytes).expect("reopen committed metadata");
assert_eq!(reopened.versions.len(), 2);
let (_, free) = reopened.find_version(Some(id)).expect("durable cleanup owner");
assert!(free.free_version());
let marker = free.delete_marker.expect("cleanup marker");
for suffix in [
rustfs_utils::http::SUFFIX_TRANSITION_TIER,
rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
rustfs_utils::http::SUFFIX_TRANSITIONED_OBJECTNAME,
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID,
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_STATE,
] {
for prefix in [RUSTFS_INTERNAL_PREFIX, MINIO_INTERNAL_PREFIX] {
let key = format!("{prefix}{suffix}");
assert_eq!(marker.meta_sys.get(&key), old.meta_sys.get(&key), "{state:?}: {key}");
}
}
let (_, current) = reopened.find_version(None).expect("replacement survives restart");
let current = current.object.expect("replacement object");
assert_eq!(current.size, 3);
assert!(!rustfs_utils::http::contains_key_bytes(&current.meta_sys, SUFFIX_TIER_FV_ID));
reopened
.add_version(replacement)
.expect("replaying replacement is idempotent");
assert_eq!(reopened.versions.len(), 2);
}
}
}
#[test]
fn tier_overwrite_rejects_cleanup_failure_without_mutating_source() {
for id in ["not-a-uuid".to_string(), Uuid::nil().to_string()] {
let (mut meta, _) = tier_overwrite_fixture(crate::TransitionVersionState::Exact);
let before = meta.clone();
let mut replacement = FileInfo::new("object", 2, 2);
replacement.mod_time = Some(OffsetDateTime::now_utc());
replacement.set_tier_free_version_id(&id);
assert!(meta.add_version(replacement).is_err());
assert_eq!(meta, before, "invalid cleanup identity must preserve source");
}
with_object_max_versions_for_test(1, || {
let (mut meta, _) = tier_overwrite_fixture(crate::TransitionVersionState::KnownDisabled);
let before = meta.clone();
let mut replacement = FileInfo::new("object", 2, 2);
replacement.mod_time = Some(OffsetDateTime::now_utc());
replacement.data = Some(Bytes::from_static(b"new"));
replacement.set_tier_free_version_id(&Uuid::new_v4().to_string());
assert_eq!(
meta.add_version(replacement)
.expect_err("cleanup owner exceeds version limit"),
Error::MaxVersionsExceeded
);
assert_eq!(meta, before, "failed cleanup insertion must also preserve inline bytes");
});
}
#[test]
fn tier_overwrite_allows_empty_transition_status_on_local_source() {
let mut source = FileInfo::new("object", 2, 2);
source.mod_time = Some(OffsetDateTime::now_utc());
source.data = Some(Bytes::from_static(b"old"));
rustfs_utils::http::insert_str(&mut source.metadata, rustfs_utils::http::SUFFIX_TRANSITION_STATUS, String::new());
let mut meta = FileMeta::new();
meta.add_version(source)
.expect("seed readable local metadata with empty status");
let mut replacement = FileInfo::new("object", 2, 2);
replacement.mod_time = Some(OffsetDateTime::now_utc());
replacement.size = 3;
replacement.data = Some(Bytes::from_static(b"new"));
replacement.set_tier_free_version_id(&Uuid::new_v4().to_string());
meta.add_version(replacement)
.expect("an ordinary overwrite must still succeed");
assert_eq!(meta.versions.len(), 1);
let (_, current) = meta.find_version(None).expect("replacement remains visible");
assert_eq!(current.object.expect("ordinary object").size, 3);
assert!(!meta.versions[0].header.free_version());
}
#[test]
fn tier_overwrite_preserves_legacy_metadata_casing() {
use rustfs_utils::http::{MINIO_INTERNAL_PREFIX, RUSTFS_INTERNAL_PREFIX};
for state in [
crate::TransitionVersionState::Exact,
crate::TransitionVersionState::KnownDisabled,
] {
let (mut meta, _) = tier_overwrite_fixture(state);
let mut source = meta.versions[0].parse_version_meta().expect("seeded source");
let object = source.object.as_mut().expect("transitioned source");
let expected = object.meta_sys.clone();
object.meta_sys = object
.meta_sys
.drain()
.map(|(key, value)| (key.to_ascii_uppercase(), value))
.collect();
meta.versions[0] = FileMetaShallowVersion::try_from(source).expect("legacy key casing");
let id = Uuid::new_v4();
let mut replacement = FileInfo::new("object", 2, 2);
replacement.mod_time = Some(OffsetDateTime::now_utc());
replacement.set_tier_free_version_id(&id.to_string());
meta.add_version(replacement).expect("overwrite readable legacy source");
let reopened = FileMeta::load(&meta.marshal_msg().expect("persist overwrite")).expect("reopen overwrite");
let (_, owner) = reopened
.find_version(Some(id))
.expect("legacy source must retain cleanup ownership");
let marker = owner.delete_marker.expect("cleanup marker");
for suffix in [
rustfs_utils::http::SUFFIX_TRANSITION_TIER,
rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
rustfs_utils::http::SUFFIX_TRANSITIONED_OBJECTNAME,
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID,
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_STATE,
] {
for prefix in [RUSTFS_INTERNAL_PREFIX, MINIO_INTERNAL_PREFIX] {
let key = format!("{prefix}{suffix}");
assert_eq!(marker.meta_sys.get(&key), expected.get(&key), "legacy {state:?}: {key}");
}
}
}
}
#[test]
fn tier_overwrite_rejects_conflicting_remote_metadata_aliases() {
use rustfs_utils::http::{
MINIO_INTERNAL_PREFIX, SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID,
SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE,
};
for suffix in [
SUFFIX_TRANSITION_STATUS,
SUFFIX_TRANSITION_TIER,
SUFFIX_TRANSITION_TIER_DESTINATION_ID,
SUFFIX_TRANSITIONED_OBJECTNAME,
SUFFIX_TRANSITIONED_VERSION_ID,
SUFFIX_TRANSITIONED_VERSION_STATE,
] {
let (mut meta, _) = tier_overwrite_fixture(crate::TransitionVersionState::Exact);
let mut old = meta.versions[0].parse_version_meta().expect("seeded source metadata");
old.object
.as_mut()
.expect("transitioned source")
.meta_sys
.insert(format!("{MINIO_INTERNAL_PREFIX}{suffix}"), b"conflicting-value".to_vec());
meta.versions[0] = FileMetaShallowVersion::try_from(old).expect("encode conflicting aliases");
let before = meta.clone();
let mut replacement = FileInfo::new("object", 2, 2);
replacement.mod_time = Some(OffsetDateTime::now_utc());
replacement.set_tier_free_version_id(&Uuid::new_v4().to_string());
assert_eq!(
meta.add_version(replacement)
.expect_err("ambiguous ownership must fail closed"),
Error::FileCorrupt
);
assert_eq!(meta, before, "conflicting {suffix} must not erase the old remote tuple");
}
}
#[test]
fn tier_overwrite_keeps_retained_remote_and_versioned_copy_ownership() {
let (mut meta, mut source) = tier_overwrite_fixture(crate::TransitionVersionState::Exact);
source.set_tier_free_version_id(&Uuid::new_v4().to_string());
meta.add_version(source.clone())
.expect("restore retains the same remote owner");
assert_eq!(meta.versions.len(), 1);
source.version_id = Some(Uuid::new_v4());
source.transition_status.clear();
source.transition_tier.clear();
source.transitioned_objname.clear();
source.transition_version = None;
source.transition_version_state = crate::TransitionVersionState::Unknown;
meta.add_version(source).expect("versioned write retains historical source");
assert_eq!(meta.versions.len(), 2);
assert!(meta.versions.iter().all(|version| !version.header.free_version()));
}
#[test]
fn add_version_filemata_uses_canonical_equal_time_order() {
let mod_time = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("valid test timestamp");
+39 -21
View File
@@ -520,32 +520,50 @@ impl RootHealRecovery {
let _guard = self.mutation.lock().await;
let disks = self.disks().await?;
let existing = Self::find(&disks, &request.id).await?;
let (disk, expected) = match existing {
Some((disk, bytes)) => (disk, Some(bytes)),
None => {
let disk = disks
.first()
.cloned()
.ok_or_else(|| Error::Other("No local disk available for root heal shutdown recovery".to_string()))?;
(disk, None)
}
};
if request.options.no_lock {
return Err(Error::Other("Administrator root heal cannot skip namespace locking".to_string()));
}
let bytes = serde_json::to_vec(&RootHealIntent::from_request(request))
.map_err(|error| Error::Other(format!("Serialize root heal recovery record: {error}")))?;
match EcstoreDiskAPI::compare_and_update_file(
disk.as_ref(),
RUSTFS_META_BUCKET,
&intent_path(&request.id)?,
expected,
Some(bytes.into()),
)
.await?
{
EcstoreConditionalFileUpdate::Updated => Ok(()),
_ => Err(Error::Other(format!("Root heal recovery record changed for {}", request.id))),
let path = intent_path(&request.id)?;
if let Some((disk, expected)) = existing {
return match EcstoreDiskAPI::compare_and_update_file(
disk.as_ref(),
RUSTFS_META_BUCKET,
&path,
Some(expected),
Some(bytes.into()),
)
.await?
{
EcstoreConditionalFileUpdate::Updated => Ok(()),
_ => Err(Error::Other(format!("Root heal recovery record changed for {}", request.id))),
};
}
if disks.is_empty() {
return Err(Error::Other("No local disk available for root heal shutdown recovery".to_string()));
}
let mut last_not_committed = None;
for disk in &disks {
match EcstoreDiskAPI::compare_and_update_file(
disk.as_ref(),
RUSTFS_META_BUCKET,
&path,
None,
Some(bytes.clone().into()),
)
.await
{
Ok(EcstoreConditionalFileUpdate::Updated) => return Ok(()),
Ok(_) => return Err(Error::Other(format!("Root heal recovery record changed for {}", request.id))),
Err(error) if error.is_conditional_file_not_committed() => last_not_committed = Some(error),
Err(error) => return Err(Error::Disk(error)),
}
}
match last_not_committed {
Some(error) => Err(Error::Disk(error)),
None => Err(Error::Other("No local disk accepted the root heal recovery record".to_string())),
}
}
@@ -17,6 +17,44 @@ use super::*;
use crate::heal::RUSTFS_META_BUCKET;
use std::collections::HashSet;
#[cfg(unix)]
struct RestoreDirectoryMode {
path: std::path::PathBuf,
mode: u32,
}
#[cfg(unix)]
impl RestoreDirectoryMode {
fn read_only(path: std::path::PathBuf) -> Self {
use std::os::unix::fs::PermissionsExt as _;
let mode = std::fs::metadata(&path)
.expect("metadata directory mode")
.permissions()
.mode();
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o555)).expect("make metadata directory read-only");
Self { path, mode }
}
}
#[cfg(unix)]
impl Drop for RestoreDirectoryMode {
fn drop(&mut self) {
use std::os::unix::fs::PermissionsExt as _;
let _ = std::fs::set_permissions(&self.path, std::fs::Permissions::from_mode(self.mode));
}
}
#[cfg(unix)]
fn ordered_recovery_disks(first: DiskStore, second: DiskStore) -> (DiskStore, DiskStore) {
if first.endpoint().to_string() <= second.endpoint().to_string() {
(first, second)
} else {
(second, first)
}
}
async fn recovery_disk() -> (TempDir, DiskStore) {
let temp = TempDir::new().expect("temporary root recovery disk");
let endpoint = Endpoint::try_from(temp.path().to_string_lossy().as_ref()).expect("disk endpoint");
@@ -78,6 +116,94 @@ fn completed_admin_status(heal_type: &HealType, completed_at: SystemTime) -> Com
}
}
#[cfg(unix)]
#[tokio::test]
async fn root_recovery_new_intent_skips_prepublication_read_only_owner() {
let (first_temp, first_disk) = recovery_disk().await;
let (second_temp, second_disk) = recovery_disk().await;
let first_endpoint = first_disk.endpoint().to_string();
let (read_only_disk, writable_disk) = ordered_recovery_disks(first_disk, second_disk);
let read_only_root = if read_only_disk.endpoint().to_string() == first_endpoint {
first_temp.path()
} else {
second_temp.path()
};
let _restore = RestoreDirectoryMode::read_only(read_only_root.join(RUSTFS_META_BUCKET));
let manager = recovery_manager(vec![read_only_disk.clone(), writable_disk.clone()]);
let mut request = admin_request(HealType::Object {
bucket: "bucket".to_string(),
object: "object".to_string(),
version_id: None,
});
let receipt = manager
.submit_heal_request_with_receipt(request.clone())
.await
.expect("a writable local disk should own the admin heal intent");
assert_eq!(receipt.result, HealAdmissionResult::Accepted);
let path = format!("root-heal-{}.json", request.id);
assert!(matches!(
read_only_disk.read_all(RUSTFS_META_BUCKET, &path).await,
Err(DiskError::FileNotFound)
));
assert!(writable_disk.read_all(RUSTFS_META_BUCKET, &path).await.is_ok());
request.retry_attempts = 1;
manager
.root_recovery
.persist(&request)
.await
.expect("an existing fallback owner should remain updateable");
let pending = manager.root_recovery.pending().await.expect("read the single durable owner");
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].id, request.id);
assert_eq!(pending[0].retry_attempts, 1);
}
#[cfg(unix)]
#[tokio::test]
async fn root_recovery_existing_owner_never_migrates_after_write_rejection() {
let (first_temp, first_disk) = recovery_disk().await;
let (second_temp, second_disk) = recovery_disk().await;
let first_endpoint = first_disk.endpoint().to_string();
let (owner_disk, alternate_disk) = ordered_recovery_disks(first_disk, second_disk);
let owner_root = if owner_disk.endpoint().to_string() == first_endpoint {
first_temp.path()
} else {
second_temp.path()
};
let manager = recovery_manager(vec![owner_disk.clone(), alternate_disk.clone()]);
let mut request = root_request();
manager
.root_recovery
.persist(&request)
.await
.expect("create the canonical owner");
let path = format!("root-heal-{}.json", request.id);
let committed = owner_disk
.read_all(RUSTFS_META_BUCKET, &path)
.await
.expect("canonical owner bytes");
let _restore = RestoreDirectoryMode::read_only(owner_root.join(RUSTFS_META_BUCKET));
request.retry_attempts = 1;
assert!(
manager.root_recovery.persist(&request).await.is_err(),
"an existing owner write rejection must fail closed"
);
assert_eq!(
owner_disk
.read_all(RUSTFS_META_BUCKET, &path)
.await
.expect("original owner remains"),
committed
);
assert!(matches!(
alternate_disk.read_all(RUSTFS_META_BUCKET, &path).await,
Err(DiskError::FileNotFound)
));
}
async fn active_root(manager: &HealManager, request: HealRequest) -> Arc<HealTask> {
let task = Arc::new(HealTask::from_request(request, manager.storage.clone()));
*task.status.write().await = HealTaskStatus::Running;
@@ -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.
+147 -10
View File
@@ -17,6 +17,7 @@ use crate::admin::router::{AdminOperation, Operation, S3Router};
use crate::admin::runtime_sources::app_context_from_req;
use crate::admin::storage_api::bucket::is_reserved_or_invalid_bucket;
use crate::admin::storage_api::bucket::utils::is_valid_object_prefix;
use crate::error::ApiError;
use crate::server::ADMIN_PREFIX;
use crate::server::RemoteAddr;
use crate::storage::rpc::node_service::heal::{
@@ -28,6 +29,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 +37,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 +71,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(|_| ApiError::invalid_request("invalid bucket name encoding"))?
.into_owned(),
obj_prefix: percent_decode_str(params.get("prefix").unwrap_or_default())
.decode_utf8()
.map_err(|_| ApiError::invalid_request("invalid object name encoding"))?
.into_owned(),
..Default::default()
};
validate_heal_target(&hip.bucket, &hip.obj_prefix)?;
@@ -164,13 +172,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 +208,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 +1624,135 @@ mod tests {
use tokio::sync::mpsc;
use tokio::time::Duration;
fn parse_registered_heal_request(uri: &Uri) -> s3s::S3Result<HealInitParams> {
let mut registered = super::S3Router::new(false);
super::register_heal_route(&mut registered).expect("register production Heal routes");
let mut router = Router::new();
for route in registered.registered_routes() {
router.insert(route.clone(), ()).expect("replay production route");
}
let path = format!("POST|{}", uri.path());
let matched = router.at(&path).expect("request must match a production Heal route");
let body = Bytes::from_static(
br#"{"recursive":false,"dryRun":true,"remove":false,"recreate":false,"scanMode":2,"updateParity":false,"nolock":false,"readRepair":false,"pool":0,"set":0}"#,
);
extract_heal_init_params(&body, uri, matched.params)
}
#[test]
fn test_heal_routes_accept_nested_and_encoded_object_paths() {
let mut router = super::S3Router::new(false);
super::register_heal_route(&mut router).expect("register production Heal routes");
for prefix in ["/rustfs/admin", "/minio/admin"] {
for target in [
"",
"test-bucket",
"test-bucket/object.bin",
"test-bucket/dir/sub/object.bin",
"test-bucket/dir%2Fobject.bin",
] {
let path = format!("{prefix}/v3/heal/{target}");
assert!(router.contains_compatible_route(http::Method::POST, &path), "{path}");
assert!(!router.contains_compatible_route(http::Method::GET, &path), "{path}");
}
}
}
#[test]
fn test_heal_target_decodes_once_and_keeps_start_status_stop_identity() {
for (wire, object) in [
("object.bin", "object.bin"),
("dir/sub/object.bin", "dir/sub/object.bin"),
("dir%2Fsub%2Fobject.bin", "dir/sub/object.bin"),
("dir%2fsub/object.bin", "dir/sub/object.bin"),
("%2Fobject.bin", "/object.bin"),
("dir/", "dir/"),
("dir%2F", "dir/"),
("literal%252Fslash", "literal%2Fslash"),
("space%20key%2Bplus", "space key+plus"),
("literal+plus", "literal+plus"),
("%E4%B8%AD%E6%96%87%2F%E6%96%87%E4%BB%B6", "中文/文件"),
("query%3Fhash%23percent%25", "query?hash#percent%"),
] {
for query in ["", "?clientToken=task", "?clientToken=task&forceStop=true"] {
let uri = format!("/rustfs/admin/v3/heal/test%2Dbucket/{wire}{query}")
.parse()
.expect("valid encoded URI");
let parsed = parse_registered_heal_request(&uri).expect("valid Heal target");
assert_eq!(parsed.bucket, "test-bucket");
assert_eq!(parsed.obj_prefix, object, "wire target: {wire}");
assert_eq!(
encode_heal_control_path(&parsed.bucket, &parsed.obj_prefix),
format!("test-bucket/{object}")
);
assert_eq!(parsed.client_token, if query.is_empty() { "" } else { "task" });
assert_eq!(parsed.force_stop, query.ends_with("forceStop=true"));
if query.is_empty() {
let request = build_heal_channel_request(&parsed);
assert_eq!(request.bucket, "test-bucket");
assert_eq!(request.object_prefix.as_deref(), Some(object));
assert_eq!(request.pool_index, Some(0));
assert_eq!(request.set_index, Some(0));
assert_eq!(request.dry_run, Some(true));
assert_eq!(request.scan_mode, Some(HealScanMode::Deep));
}
}
}
}
#[test]
fn test_heal_target_validates_decoded_paths_before_admission() {
for target in [
"test%2Fbucket/object",
"test%00bucket/object",
"test%FFbucket/object",
"test-bucket/dir%2F..%2Fobject",
"test-bucket/dir/%2e/object",
"test-bucket/dir%5C..%5Cobject",
"test-bucket/dir%2F%2Fobject",
"test-bucket/object%00",
"test-bucket/object%FF",
] {
let uri = format!("/rustfs/admin/v3/heal/{target}").parse().expect("encoded URI");
let err = parse_registered_heal_request(&uri).expect_err("decoded invalid target must fail closed");
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest, "target: {target}");
}
}
#[tokio::test]
async fn test_nested_heal_routes_still_require_authentication() {
use s3s::route::S3Route;
let mut router = super::S3Router::new(false);
super::register_heal_route(&mut router).expect("register production Heal routes");
for prefix in ["/rustfs/admin", "/minio/admin"] {
for object in ["dir/object.bin", "dir%2Fobject.bin", "literal%252Fslash"] {
let mut req = s3s::S3Request {
input: s3s::Body::empty(),
method: http::Method::POST,
uri: format!("{prefix}/v3/heal/test-bucket/{object}").parse().expect("Heal URI"),
headers: http::HeaderMap::new(),
extensions: http::Extensions::new(),
credentials: None,
region: None,
service: None,
trailing_headers: None,
};
let err = router
.check_access(&mut req)
.await
.expect_err("router must require a signature");
assert_eq!(err.code(), &S3ErrorCode::AccessDenied);
let err = router
.call(req)
.await
.expect_err("handler must independently require authentication");
assert_eq!(err.code(), &S3ErrorCode::InvalidRequest);
assert!(err.to_string().contains("authentication required"));
}
}
}
fn replacement_record(task_id: &str) -> rustfs_heal::ReplacementRecoveryRecord {
rustfs_heal::ReplacementRecoveryRecord {
task_id: task_id.to_string(),
+1 -1
View File
@@ -346,7 +346,7 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
admin(HttpMethod::Post, "/rustfs/admin/v3/rebalance/stop", REBALANCE, RouteRiskLevel::High),
admin(HttpMethod::Post, "/rustfs/admin/v3/heal/", HEAL, RouteRiskLevel::High),
admin(HttpMethod::Post, "/rustfs/admin/v3/heal/{bucket}", HEAL, RouteRiskLevel::High),
admin(HttpMethod::Post, "/rustfs/admin/v3/heal/{bucket}/{prefix}", HEAL, RouteRiskLevel::High),
admin(HttpMethod::Post, "/rustfs/admin/v3/heal/{bucket}/{*prefix}", HEAL, RouteRiskLevel::High),
admin(HttpMethod::Post, "/rustfs/admin/v3/background-heal/status", HEAL, RouteRiskLevel::High),
admin(
HttpMethod::Get,
+1 -1
View File
@@ -198,7 +198,7 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
admin_route(Method::POST, "/v3/rebalance/stop"),
admin_route(Method::POST, "/v3/heal/"),
admin_route_sample(Method::POST, "/v3/heal/{bucket}", "/v3/heal/test-bucket"),
admin_route_sample(Method::POST, "/v3/heal/{bucket}/{prefix}", "/v3/heal/test-bucket/prefix"),
admin_route_sample(Method::POST, "/v3/heal/{bucket}/{*prefix}", "/v3/heal/test-bucket/prefix"),
admin_route(Method::POST, "/v3/background-heal/status"),
admin_route(Method::GET, "/v4/heal/replacement-recovery"),
admin_route(Method::GET, "/v3/tier"),
+9 -1
View File
@@ -377,6 +377,10 @@ impl FS {
pub(crate) fn parse_object_version_id(version_id: Option<String>) -> S3Result<Option<Uuid>> {
if let Some(vid) = version_id {
if vid == "null" {
// A nil UUID selects the stored null version; None selects latest.
return Ok(Some(Uuid::nil()));
}
let uuid = Uuid::parse_str(&vid).map_err(|e| {
error!("Invalid version ID: {}", e);
s3_error!(InvalidArgument, "Invalid version ID")
@@ -1183,7 +1187,11 @@ impl S3 for FS {
error = %e,
"Object tags not found"
);
return Err(s3_error!(NoSuchKey));
return Err(S3Error::new(if opts.version_id.is_some() {
S3ErrorCode::NoSuchVersion
} else {
S3ErrorCode::NoSuchKey
}));
}
error!(
component = LOG_COMPONENT_STORAGE,
+33 -1
View File
@@ -17,7 +17,9 @@ mod tests {
use crate::config::WorkloadProfile;
use crate::server::cors;
use crate::storage::StorageError;
use crate::storage::ecfs::{FS, propagate_object_lock_peer_reload, validate_object_lock_configuration_input};
use crate::storage::ecfs::{
FS, parse_object_version_id, propagate_object_lock_peer_reload, validate_object_lock_configuration_input,
};
use crate::storage::ecfs_extend::{apply_bucket_default_lock_retention, map_bucket_object_lock_config_state};
use crate::storage::s3_api::common::{rustfs_initiator, rustfs_owner};
use crate::storage::storage_api::ecstore_bucket::metadata_sys::ObjectLockConfigState;
@@ -641,6 +643,36 @@ mod tests {
assert_eq!(metadata.get("content-type"), Some(&"application/octet-stream".to_string()));
}
#[test]
fn test_tagging_version_id_preserves_explicit_null_and_latest_selection() {
assert_eq!(parse_object_version_id(None).expect("latest version selector"), None);
assert_eq!(
parse_object_version_id(Some("null".to_owned())).expect("explicit null version selector"),
Some(uuid::Uuid::nil())
);
for version in [uuid::Uuid::nil(), uuid::Uuid::new_v4()] {
assert_eq!(
parse_object_version_id(Some(version.to_string())).expect("UUID version selector"),
Some(version)
);
}
}
#[test]
fn test_tagging_version_id_rejects_invalid_values_instead_of_selecting_latest() {
for version in [
"",
"NULL",
" null ",
"not-a-version",
"null/other",
"00000000-0000-0000-0000-00000000000g",
] {
let err = parse_object_version_id(Some(version.to_owned())).expect_err("invalid version must fail closed");
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument, "version: {version:?}");
}
}
#[tokio::test]
async fn test_get_object_tagging_returns_internal_error_when_store_uninitialized() {
if !store_uninitialized_premise_holds() {
+3 -1
View File
@@ -58,8 +58,10 @@ cd "$(dirname "$0")/.."
# 1589 -> 1588 on 2026-09-08: the GA blocker set (rustfs/backlog#2366) added
# three invocation lines to the endpoint-refresh paths and folded the five
# copies of the concurrent-change error into one constructor, netting -1.
# 1588 -> 1586 on 2026-09-10: the release merge no longer introduces direct
# s3_error! constructors for heal percent-decoding or tagging not-found errors.
S3S_IMPORT_FILES_BASELINE=213
S3_ERROR_LINES_BASELINE=1588
S3_ERROR_LINES_BASELINE=1586
# ecstore-scoped ratchet (rustfs/backlog#1842): the storage engine must not
# know S3 wire/DTO types (ARCHITECTURE.md invariant 4). The S3-*consuming*
# client was extracted to crates/s3-client, where s3s usage is legitimate;
+83 -10
View File
@@ -83,9 +83,12 @@ runtime_profile_for() {
background-target-crash|background-target-restart)
echo "background-4x1"
;;
background-target-crash-ec8-4|background-target-restart-ec8-4|background-target-restart-ec8-4-multi-set)
background-target-crash-ec8-4|background-target-restart-ec8-4)
echo "background-ec8-4"
;;
background-target-restart-ec8-4-multi-set)
echo "background-ec8-4-multi-set"
;;
background-target-crash-ec8-4-multi-pool)
echo "background-ec8-4-multi-pool"
;;
@@ -111,8 +114,13 @@ apply_runtime_profile() {
export RUSTFS_HEAL_CHAOS_OBJECT_SIZE_BYTES="${RUSTFS_HEAL_CHAOS_OBJECT_SIZE_BYTES:-8388608}"
export RUSTFS_HEAL_CHAOS_PARTIAL_TIMEOUT_SECS="${RUSTFS_HEAL_CHAOS_PARTIAL_TIMEOUT_SECS:-180}"
;;
background-ec8-4-multi-set)
export RUSTFS_HEAL_CHAOS_OBJECT_COUNT="${RUSTFS_HEAL_CHAOS_OBJECT_COUNT:-64}"
export RUSTFS_HEAL_CHAOS_OBJECT_SIZE_BYTES="${RUSTFS_HEAL_CHAOS_OBJECT_SIZE_BYTES:-16777216}"
export RUSTFS_HEAL_CHAOS_PARTIAL_TIMEOUT_SECS="${RUSTFS_HEAL_CHAOS_PARTIAL_TIMEOUT_SECS:-240}"
;;
background-ec8-4-multi-pool)
export RUSTFS_HEAL_CHAOS_OBJECT_COUNT="${RUSTFS_HEAL_CHAOS_OBJECT_COUNT:-96}"
export RUSTFS_HEAL_CHAOS_OBJECT_COUNT="${RUSTFS_HEAL_CHAOS_OBJECT_COUNT:-64}"
export RUSTFS_HEAL_CHAOS_OBJECT_SIZE_BYTES="${RUSTFS_HEAL_CHAOS_OBJECT_SIZE_BYTES:-4194304}"
export RUSTFS_HEAL_CHAOS_PARTIAL_TIMEOUT_SECS="${RUSTFS_HEAL_CHAOS_PARTIAL_TIMEOUT_SECS:-240}"
;;
@@ -165,6 +173,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 +228,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
@@ -267,9 +327,25 @@ if [[ "$NOFILE_SOFT" =~ ^[0-9]+$ && "$NOFILE_HARD" =~ ^[0-9]+$ && "$NOFILE_SOFT"
ulimit -n "$NOFILE_HARD" || true
fi
fi
mkdir -p "$(dirname "$RUN_DIR")"
TMP_DIR="$(mktemp -d "${TMPDIR:-/tmp}/rustfs-scanner-heal-evidence.XXXXXX")"
trap 'rm -rf "$TMP_DIR"' EXIT
RUN_PARENT="$(dirname "$RUN_DIR")"
mkdir -p "$RUN_PARENT"
RUN_TMP_ROOT_CREATED=0
if [[ -z "${TMPDIR:-}" ]]; then
RUN_TMP_ROOT="$RUN_PARENT/.tmp-$(basename "$RUN_DIR")"
mkdir -p "$RUN_TMP_ROOT"
export TMPDIR="$RUN_TMP_ROOT"
RUN_TMP_ROOT_CREATED=1
else
RUN_TMP_ROOT=""
fi
TMP_DIR="$(mktemp -d "${TMPDIR%/}/rustfs-scanner-heal-evidence.XXXXXX")"
cleanup_tmp() {
rm -rf "$TMP_DIR"
if [[ "$RUN_TMP_ROOT_CREATED" == 1 ]]; then
rm -rf "$RUN_TMP_ROOT"
fi
}
trap cleanup_tmp EXIT
BUILD_FEATURES="${RUSTFS_BUILD_FEATURES:-}"
TARGET_DIR="${CARGO_TARGET_DIR:-$ROOT/target}"
@@ -300,8 +376,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 +386,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
+7 -3
View File
@@ -654,9 +654,12 @@ def collect_live(prepared, request, request_path, adapter):
if request.get("evidence") == "measured":
expected_metrics_endpoints = request["release_evidence"]["distributed"]["metrics_endpoints"]
output = request_path.parent / "telemetry"
interval_seconds = 60
sample_count = request["duration_seconds"] // interval_seconds + 1
collector_timeout_seconds = (sample_count - 1) * interval_seconds + 300
args = ["bash", str(collector), "--alias", connection["alias"], "--endpoint", connection["endpoint"],
"--metrics-endpoints", connection["metrics_endpoints"], "--deployment", "distributed",
"--samples", str(request["duration_seconds"] // 60 + 1), "--interval-secs", "60",
"--samples", str(sample_count), "--interval-secs", str(interval_seconds),
"--out-dir", str(output)]
with (request_path.parent / "collector.log").open("wb") as log:
process = OwnedCommand(args, log)
@@ -664,10 +667,11 @@ def collect_live(prepared, request, request_path, adapter):
started = time.monotonic()
result = invoke(adapter, "measure", request_path, request["duration_seconds"] + 300)
require(time.monotonic() - started >= request["duration_seconds"], "measurement ended before required window")
require(process.wait(120) == 0, "scanner collector failed")
require(process.wait(collector_timeout_seconds) == 0,
f"scanner collector failed within {collector_timeout_seconds}s timeout")
require(output.joinpath("scanner-summary.csv").stat().st_size > 0, "missing collector samples")
samples = list((output / "status").glob("scanner-status.*.json"))
require(len(samples) == request["duration_seconds"] // 60 + 1, "missing scanner samples")
require(len(samples) == sample_count, "missing scanner samples")
for sample in samples:
status = read_json(sample)
require(isinstance(status.get("metrics"), dict) and status["metrics"], "invalid scanner status response")