mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-11 21:39:27 +00:00
Compare commits
19 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| e55c9ffece | |||
| c478a392e7 | |||
| c8ccc1e198 | |||
| 74ba5c205c | |||
| 9fcd54734b | |||
| 6eb200022e | |||
| 92ce782999 | |||
| bf425ca32c | |||
| 0ae38aabbc | |||
| 77ba4f1b2e | |||
| 7b6e45d372 | |||
| 09c85a5f29 | |||
| 00aeb12914 | |||
| ffb18979f8 | |||
| 97c7b451d2 | |||
| 9ba95cac37 | |||
| cec796328c | |||
| 1aa4fa145f | |||
| d286f3d06c |
@@ -1,2 +1,2 @@
|
||||
sha256-darwin=874c881d7b45f12378a5817c7f42c95c4981960a2ec9ce12dcf4af239ae1f9d5
|
||||
sha256-linux=9351e25b45bf7dfce18b951a5e3740225f457cacc53b8bf9f500f6947763ec0e
|
||||
sha256-darwin=15cb0cf9909bfbfc5a835fb08bd3675516c641db1b92ecc1e170a1cfa0fa2fd5
|
||||
sha256-linux=3163fdd29df5def86cf511ca7db05880d5c0caaa608a7031a71394327f7c217a
|
||||
|
||||
@@ -1 +1 @@
|
||||
sha256=6d18f9cce820c51d5589de944e8cc185f73eeca0ea9a9916651943e3759169d0
|
||||
sha256=5fbb230b89212b7c3d7229d6cef3e7e2d16f0ecfec62237ebc770785706f67d9
|
||||
|
||||
@@ -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
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
@@ -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
@@ -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
@@ -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"
|
||||
|
||||
@@ -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 = ®istry["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));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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";
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -3821,6 +3821,13 @@ impl SetDisks {
|
||||
}
|
||||
|
||||
fi.metadata = user_defined;
|
||||
if fi.version_id.is_none_or(|id| id.is_nil()) && !opts.data_movement && expected_restore_operation_id.is_none() {
|
||||
// Every disk must publish the same cleanup owner alongside a
|
||||
// replaced null version. This transient key is not persisted
|
||||
// on the new object; recovery discovers the free-version in
|
||||
// the committed xl.meta even if this request is cancelled.
|
||||
fi.set_tier_free_version_id(&Uuid::new_v4().to_string());
|
||||
}
|
||||
fi.mod_time = mod_time;
|
||||
fi.size = w_size as i64;
|
||||
fi.versioned = opts.versioned || opts.version_suspended;
|
||||
@@ -18174,6 +18181,102 @@ mod put_object_tmp_cleanup_tests {
|
||||
drop(temp_dirs);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
#[serial_test::serial(capacity_dirty_scope)]
|
||||
async fn tier_overwrite_failed_quorum_and_cancellation_preserve_live_source() {
|
||||
for cancel_before_rename in [false, true] {
|
||||
let (dirs, disks, set) = hermetic_set_disks(4).await;
|
||||
let bucket = "tier-overwrite-failure";
|
||||
let object = "still-live";
|
||||
make_completion_test_bucket(&disks, bucket).await;
|
||||
let old_body = vec![0x31; TEST_OBJECT_SIZE];
|
||||
let mut metadata = HashMap::from([(
|
||||
"x-amz-restore".to_string(),
|
||||
"ongoing-request=\"false\", expiry-date=\"2099-01-01T00:00:00Z\"".to_string(),
|
||||
)]);
|
||||
for (suffix, value) in [
|
||||
(rustfs_utils::http::SUFFIX_TRANSITION_STATUS, "complete".to_string()),
|
||||
(rustfs_utils::http::SUFFIX_TRANSITION_TIER, "WARM".to_string()),
|
||||
(rustfs_utils::http::SUFFIX_TRANSITIONED_OBJECTNAME, "remote/still-live".to_string()),
|
||||
(rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID, "exact-live-version".to_string()),
|
||||
(rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_STATE, "exact".to_string()),
|
||||
(rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID, "ab".repeat(32)),
|
||||
] {
|
||||
rustfs_utils::http::insert_str(&mut metadata, suffix, value);
|
||||
}
|
||||
set.put_object(
|
||||
bucket,
|
||||
object,
|
||||
&mut PutObjReader::from_vec(old_body.clone()),
|
||||
&ObjectOptions {
|
||||
user_defined: metadata,
|
||||
write_completion: WriteCompletion::TailDrained,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("seed live transitioned source");
|
||||
wait_for_tmp_workspace_to_drain(&dirs, "seed write must drain").await;
|
||||
let before = set
|
||||
.load_file_info_versions_exact(bucket, object)
|
||||
.await
|
||||
.expect("read original metadata")
|
||||
.expect("original exists");
|
||||
|
||||
if cancel_before_rename {
|
||||
let barrier = PutObjectCommitBarrier::install(bucket, object, PutObjectCommitPause::AfterQuotaReservation);
|
||||
let writer = Arc::clone(&set);
|
||||
let put = tokio::spawn(async move {
|
||||
writer
|
||||
.put_object(
|
||||
bucket,
|
||||
object,
|
||||
&mut PutObjReader::from_vec(vec![0x32; TEST_OBJECT_SIZE]),
|
||||
&ObjectOptions::default(),
|
||||
)
|
||||
.await
|
||||
});
|
||||
barrier.wait_until_paused().await;
|
||||
put.abort();
|
||||
assert!(put.await.expect_err("cancel paused replacement").is_cancelled());
|
||||
wait_for_tmp_workspace_to_drain(&dirs, "cancelled replacement must roll back").await;
|
||||
drop(barrier);
|
||||
} else {
|
||||
let _fault = rename_fault_injection::fail_rename_on(object, &[2, 3]);
|
||||
let err = set
|
||||
.put_object(
|
||||
bucket,
|
||||
object,
|
||||
&mut PutObjReader::from_vec(vec![0x32; TEST_OBJECT_SIZE]),
|
||||
&ObjectOptions {
|
||||
write_completion: WriteCompletion::TailDrained,
|
||||
..Default::default()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect_err("two disk commits cannot satisfy write quorum three");
|
||||
assert!(matches!(err, Error::ErasureWriteQuorum | Error::InsufficientWriteQuorum(_, _)), "{err}");
|
||||
}
|
||||
let after = set
|
||||
.load_file_info_versions_exact(bucket, object)
|
||||
.await
|
||||
.expect("read rolled-back metadata")
|
||||
.expect("live source must survive");
|
||||
assert_eq!(after.versions, before.versions, "failed replacement must preserve the live version");
|
||||
assert_eq!(
|
||||
after.free_versions, before.free_versions,
|
||||
"failed replacement must not publish a cleanup owner"
|
||||
);
|
||||
let mut reader = set
|
||||
.get_object_reader(bucket, object, None, HeaderMap::new(), &ObjectOptions::default())
|
||||
.await
|
||||
.expect("live source remains readable");
|
||||
let mut actual = Vec::new();
|
||||
reader.stream.read_to_end(&mut actual).await.expect("read original bytes");
|
||||
assert_eq!(actual, old_body);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn cooperative_cancellation_while_waiting_for_namespace_lock_cleans_tmp_workspace() {
|
||||
let (temp_dirs, disk_stores, set_disks) = hermetic_set_disks(4).await;
|
||||
|
||||
@@ -2832,6 +2832,224 @@ mod tests {
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
async fn tier_overwrite_put_and_self_copy_recover_persisted_cleanup_owners() {
|
||||
use crate::bucket::lifecycle::bucket_lifecycle_ops::ExpiryState;
|
||||
use crate::bucket::lifecycle::tier_free_version_recovery::recover_tier_free_versions;
|
||||
use rustfs_filemeta::TransitionVersionState::{Exact, KnownDisabled, SuspendedNull};
|
||||
use rustfs_s3_client::transition_api::ReaderImpl;
|
||||
use rustfs_utils::http::{
|
||||
SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||
SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE, insert_str,
|
||||
};
|
||||
|
||||
let temp_dir = tempfile::tempdir().expect("create tier overwrite store");
|
||||
let (mut ctx, mut store, mut shutdown) =
|
||||
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-overwrite", &[4])).await;
|
||||
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
|
||||
let tier = "OVERWRITE-TIER";
|
||||
let backend = register_mock_tier(&ctx.tier_config_mgr(), tier).await;
|
||||
let lease = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier)
|
||||
.await
|
||||
.expect("tier identity");
|
||||
let identity = rustfs_utils::crypto::hex(lease.backend_identity());
|
||||
drop(lease);
|
||||
|
||||
for state in [Exact, KnownDisabled, SuspendedNull] {
|
||||
for suspended in [false, true] {
|
||||
for self_copy in [false, true] {
|
||||
let bucket = format!("tier-overwrite-{}", Uuid::new_v4());
|
||||
let object = "object";
|
||||
let remote = format!("remote/{bucket}");
|
||||
let version = match state {
|
||||
Exact => "opaque-overwrite-version",
|
||||
SuspendedNull => "null",
|
||||
_ => "",
|
||||
};
|
||||
let payload = vec![0x5b; if suspended { 512 * 1024 } else { 257 }];
|
||||
store
|
||||
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||
.await
|
||||
.expect("create bucket");
|
||||
backend.set_put_remote_version(Some(version.to_string())).await;
|
||||
let lease = TierConfigMgr::acquire_operation_lease(&ctx.tier_config_mgr(), tier)
|
||||
.await
|
||||
.expect("seed tier lease");
|
||||
lease
|
||||
.put(
|
||||
&remote,
|
||||
ReaderImpl::Body(bytes::Bytes::from(payload.clone())),
|
||||
payload.len().try_into().expect("payload size"),
|
||||
)
|
||||
.await
|
||||
.expect("seed remote bytes");
|
||||
drop(lease);
|
||||
let mut metadata = HashMap::from([
|
||||
("content-type".to_string(), "application/octet-stream".to_string()),
|
||||
(
|
||||
"x-amz-restore".to_string(),
|
||||
"ongoing-request=\"false\", expiry-date=\"2099-01-01T00:00:00Z\"".to_string(),
|
||||
),
|
||||
]);
|
||||
for (suffix, value) in [
|
||||
(SUFFIX_TRANSITION_STATUS, "complete"),
|
||||
(SUFFIX_TRANSITION_TIER, tier),
|
||||
(SUFFIX_TRANSITION_TIER_DESTINATION_ID, identity.as_str()),
|
||||
(SUFFIX_TRANSITIONED_OBJECTNAME, remote.as_str()),
|
||||
(SUFFIX_TRANSITIONED_VERSION_STATE, state.as_str()),
|
||||
] {
|
||||
insert_str(&mut metadata, suffix, value.to_string());
|
||||
}
|
||||
if !version.is_empty() {
|
||||
insert_str(&mut metadata, SUFFIX_TRANSITIONED_VERSION_ID, version.to_string());
|
||||
}
|
||||
let options = ObjectOptions {
|
||||
version_suspended: suspended,
|
||||
..Default::default()
|
||||
};
|
||||
store
|
||||
.put_object(
|
||||
&bucket,
|
||||
object,
|
||||
&mut PutObjReader::from_vec(payload.clone()),
|
||||
&ObjectOptions {
|
||||
user_defined: metadata,
|
||||
..options.clone()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("seed transitioned source with locally restored bytes");
|
||||
let expected = if self_copy {
|
||||
payload.clone()
|
||||
} else {
|
||||
vec![0x73; payload.len()]
|
||||
};
|
||||
let new_metadata = HashMap::from([
|
||||
("content-type".to_string(), "text/plain".to_string()),
|
||||
("x-amz-meta-replacement".to_string(), "kept".to_string()),
|
||||
]);
|
||||
if self_copy {
|
||||
let mut source = store
|
||||
.get_object_info(&bucket, object, &options)
|
||||
.await
|
||||
.expect("self-copy source");
|
||||
source.metadata_only = false;
|
||||
source.user_defined = Arc::new(new_metadata);
|
||||
source.put_object_reader = Some(PutObjReader::from_vec(expected.clone()));
|
||||
store
|
||||
.copy_object(&bucket, object, &bucket, object, &mut source, &options, &options)
|
||||
.await
|
||||
.expect("materialized self-copy");
|
||||
} else {
|
||||
store
|
||||
.put_object(
|
||||
&bucket,
|
||||
object,
|
||||
&mut PutObjReader::from_vec(expected.clone()),
|
||||
&ObjectOptions {
|
||||
user_defined: new_metadata,
|
||||
..options.clone()
|
||||
},
|
||||
)
|
||||
.await
|
||||
.expect("overwrite transitioned null version");
|
||||
}
|
||||
|
||||
let set = store.pools[0].get_disks_by_key(object);
|
||||
let versions = set
|
||||
.load_file_info_versions_exact(&bucket, object)
|
||||
.await
|
||||
.expect("read committed disk metadata")
|
||||
.expect("replacement metadata exists");
|
||||
let free: Vec<_> = versions
|
||||
.versions
|
||||
.iter()
|
||||
.chain(versions.free_versions.iter())
|
||||
.filter(|fi| fi.tier_free_version())
|
||||
.collect();
|
||||
assert_eq!(free.len(), 1, "{state:?}, suspended={suspended}, copy={self_copy}");
|
||||
assert_eq!(free[0].transitioned_objname, remote);
|
||||
assert_eq!(free[0].transition_version_state, state);
|
||||
assert!(backend.contains(&remote).await, "commit must not delete remote bytes before cleanup");
|
||||
let removed_before = backend.remove_count().await;
|
||||
|
||||
// Restart before queue delivery. The new runtime must
|
||||
// reconstruct ownership solely from the committed xl.meta.
|
||||
let tier_config = ctx
|
||||
.tier_config_mgr()
|
||||
.read()
|
||||
.await
|
||||
.tiers
|
||||
.get(tier)
|
||||
.expect("tier configuration survives restart")
|
||||
.clone_with_credentials();
|
||||
drop(set);
|
||||
shutdown.cancel();
|
||||
drop(store);
|
||||
drop(ctx);
|
||||
(ctx, store, shutdown) =
|
||||
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "tier-overwrite-restart", &[4]))
|
||||
.await;
|
||||
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
|
||||
{
|
||||
let manager = ctx.tier_config_mgr();
|
||||
let mut manager = manager.write().await;
|
||||
manager.tiers.insert(tier.to_string(), tier_config);
|
||||
manager
|
||||
.install_test_driver(tier, Box::new(backend.clone()))
|
||||
.expect("rebind the same remote destination after restart");
|
||||
}
|
||||
let set = store.pools[0].get_disks_by_key(object);
|
||||
ExpiryState::resize_workers(1, Arc::clone(&store)).await;
|
||||
let recovered = recover_tier_free_versions(Arc::clone(&store), 100, None, None)
|
||||
.await
|
||||
.expect("recover persisted cleanup owner");
|
||||
assert!(recovered.enqueued >= 1);
|
||||
tokio::time::timeout(Duration::from_secs(30), async {
|
||||
loop {
|
||||
let versions = set
|
||||
.load_file_info_versions_exact(&bucket, object)
|
||||
.await
|
||||
.expect("read cleanup progress")
|
||||
.expect("new object must survive cleanup");
|
||||
if versions
|
||||
.versions
|
||||
.iter()
|
||||
.chain(versions.free_versions.iter())
|
||||
.all(|fi| !fi.tier_free_version())
|
||||
{
|
||||
break;
|
||||
}
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("cleanup must converge");
|
||||
assert!(!backend.contains(&remote).await);
|
||||
assert_eq!(backend.remove_count().await, removed_before + 1, "one remote DELETE per owner");
|
||||
assert_eq!(backend.remove_versions().await.last(), Some(&(remote.clone(), version.to_string())));
|
||||
let mut reader = store
|
||||
.get_object_reader(&bucket, object, None, HeaderMap::new(), &options)
|
||||
.await
|
||||
.expect("replacement remains readable");
|
||||
let mut actual = Vec::new();
|
||||
reader.stream.read_to_end(&mut actual).await.expect("read replacement bytes");
|
||||
assert_eq!(actual, expected);
|
||||
let current = store
|
||||
.get_object_info(&bucket, object, &options)
|
||||
.await
|
||||
.expect("replacement metadata");
|
||||
assert_eq!(current.user_defined.get("content-type").map(String::as_str), Some("text/plain"));
|
||||
assert_eq!(current.user_defined.get("x-amz-meta-replacement").map(String::as_str), Some("kept"));
|
||||
assert!(current.transitioned_object.status.is_empty());
|
||||
}
|
||||
}
|
||||
}
|
||||
shutdown.cancel();
|
||||
}
|
||||
|
||||
#[cfg(feature = "test-util")]
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
#[serial_test::serial(storage_class_env)]
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -480,7 +480,110 @@ impl FileMeta {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn add_version(&mut self, mut fi: FileInfo) -> Result<()> {
|
||||
pub fn add_version(&mut self, fi: FileInfo) -> Result<()> {
|
||||
if let Some(free_version) = self.overwritten_tier_free_version(&fi)? {
|
||||
// The replacement and its cleanup owner must share one xl.meta
|
||||
// commit. Keep the original intact if either insertion fails.
|
||||
let mut next = self.clone();
|
||||
next.add_version_inner(fi)?;
|
||||
next.add_version_filemata(free_version)?;
|
||||
*self = next;
|
||||
return Ok(());
|
||||
}
|
||||
self.add_version_inner(fi)
|
||||
}
|
||||
|
||||
fn overwritten_tier_free_version(&self, fi: &FileInfo) -> Result<Option<FileMetaVersion>> {
|
||||
use rustfs_utils::http::{
|
||||
SUFFIX_TIER_FV_ID, SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||
SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||
get_consistent_bytes, get_consistent_str, has_internal_suffix, strip_internal_prefix_preserving_case,
|
||||
};
|
||||
|
||||
if fi.version_id.is_some_and(|id| !id.is_nil()) || !contains_key_str(&fi.metadata, SUFFIX_TIER_FV_ID) {
|
||||
return Ok(None);
|
||||
}
|
||||
let Some(existing) = self
|
||||
.versions
|
||||
.iter()
|
||||
.find(|v| v.header.version_id.is_none_or(|id| id.is_nil()))
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
let old = existing.parse_version_meta()?;
|
||||
let Some(mut object) = old.object else {
|
||||
return Ok(None);
|
||||
};
|
||||
let status = get_consistent_bytes(&object.meta_sys, SUFFIX_TRANSITION_STATUS);
|
||||
if status.is_none()
|
||||
&& object
|
||||
.meta_sys
|
||||
.keys()
|
||||
.any(|key| has_internal_suffix(key, SUFFIX_TRANSITION_STATUS))
|
||||
{
|
||||
// Empty status is a valid local object. The reader distinguishes
|
||||
// it from conflicting aliases before the ordinary overwrite.
|
||||
object.into_fileinfo(&fi.volume, &fi.name, false)?;
|
||||
return Ok(None);
|
||||
}
|
||||
if status != Some(TRANSITION_COMPLETE.as_bytes()) {
|
||||
return Ok(None);
|
||||
}
|
||||
// Reuse the reader's alias/state validation. A legacy empty remote
|
||||
// version is valid and must not be mistaken for conflicting aliases.
|
||||
object.into_fileinfo(&fi.volume, &fi.name, false)?;
|
||||
if object
|
||||
.meta_sys
|
||||
.keys()
|
||||
.any(|key| has_internal_suffix(key, SUFFIX_TRANSITION_TIER_DESTINATION_ID))
|
||||
&& get_consistent_bytes(&object.meta_sys, SUFFIX_TRANSITION_TIER_DESTINATION_ID).is_none()
|
||||
{
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
let transition_suffixes = [
|
||||
SUFFIX_TRANSITION_STATUS,
|
||||
SUFFIX_TRANSITION_TIER,
|
||||
SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||
SUFFIX_TRANSITIONED_OBJECTNAME,
|
||||
SUFFIX_TRANSITIONED_VERSION_ID,
|
||||
SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||
];
|
||||
let replacement = MetaObject::from(fi.clone());
|
||||
if transition_suffixes
|
||||
.iter()
|
||||
.all(|suffix| get_consistent_bytes(&object.meta_sys, suffix) == get_consistent_bytes(&replacement.meta_sys, suffix))
|
||||
{
|
||||
return Ok(None);
|
||||
}
|
||||
let id = get_consistent_str(&fi.metadata, SUFFIX_TIER_FV_ID).ok_or(Error::FileCorrupt)?;
|
||||
let id = Uuid::parse_str(id)?;
|
||||
if id.is_nil() || self.versions.iter().any(|version| version.header.version_id == Some(id)) {
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
// The reader also accepts legacy key casing. Canonicalize only this
|
||||
// cleanup source so init_free_version preserves every accepted field,
|
||||
// including an explicitly empty unversioned remote version.
|
||||
for suffix in transition_suffixes {
|
||||
let value = object
|
||||
.meta_sys
|
||||
.iter()
|
||||
.find(|(key, _)| {
|
||||
strip_internal_prefix_preserving_case(key).is_some_and(|found| found.eq_ignore_ascii_case(suffix))
|
||||
})
|
||||
.map(|(_, value)| value.clone());
|
||||
if let Some(value) = value {
|
||||
rustfs_utils::http::insert_bytes(&mut object.meta_sys, suffix, value);
|
||||
}
|
||||
}
|
||||
let (free_version, created) = object.init_free_version(fi)?;
|
||||
if !created {
|
||||
return Err(Error::FileCorrupt);
|
||||
}
|
||||
Ok(Some(free_version))
|
||||
}
|
||||
|
||||
fn add_version_inner(&mut self, mut fi: FileInfo) -> Result<()> {
|
||||
rustfs_utils::http::remove_str(&mut fi.metadata, rustfs_utils::http::SUFFIX_TIER_FV_ID);
|
||||
// empty version_id means "null" (versioning disabled/suspended)
|
||||
if fi.version_id.is_none() {
|
||||
fi.version_id = Some(Uuid::nil());
|
||||
@@ -1463,6 +1566,235 @@ mod test {
|
||||
});
|
||||
}
|
||||
|
||||
fn tier_overwrite_fixture(state: crate::TransitionVersionState) -> (FileMeta, FileInfo) {
|
||||
let mut source = FileInfo::new("object", 2, 2);
|
||||
source.erasure.index = 1;
|
||||
source.mod_time = Some(OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("fixture timestamp"));
|
||||
source.data_dir = Some(Uuid::new_v4());
|
||||
source.transition_status = TRANSITION_COMPLETE.to_string();
|
||||
source.transition_tier = "WARM".to_string();
|
||||
source.transitioned_objname = "remote/old-object".to_string();
|
||||
source.transition_version_state = state;
|
||||
source.transition_version = match state {
|
||||
crate::TransitionVersionState::Exact => Some("opaque-provider-version".to_string()),
|
||||
crate::TransitionVersionState::SuspendedNull => Some("null".to_string()),
|
||||
_ => None,
|
||||
};
|
||||
rustfs_utils::http::insert_str(
|
||||
&mut source.metadata,
|
||||
rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||
"ab".repeat(32),
|
||||
);
|
||||
if state == crate::TransitionVersionState::KnownDisabled {
|
||||
rustfs_utils::http::insert_str(
|
||||
&mut source.metadata,
|
||||
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID,
|
||||
String::new(),
|
||||
);
|
||||
}
|
||||
let mut meta = FileMeta::new();
|
||||
meta.add_version(source.clone()).expect("seed transitioned null version");
|
||||
(meta, source)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tier_overwrite_preserves_exact_cleanup_owner_across_reload() {
|
||||
use crate::TransitionVersionState::{Exact, KnownDisabled, SuspendedNull, Unknown};
|
||||
use rustfs_utils::http::{MINIO_INTERNAL_PREFIX, RUSTFS_INTERNAL_PREFIX, SUFFIX_TIER_FV_ID};
|
||||
|
||||
for state in [Exact, KnownDisabled, SuspendedNull, Unknown] {
|
||||
for inline in [false, true] {
|
||||
let (mut meta, _) = tier_overwrite_fixture(state);
|
||||
let old = meta.versions[0].parse_version_meta().expect("old metadata");
|
||||
let old = old.object.expect("old object");
|
||||
let id = Uuid::new_v4();
|
||||
let mut replacement = FileInfo::new("object", 2, 2);
|
||||
replacement.version_id = inline.then_some(Uuid::nil());
|
||||
replacement.mod_time = Some(OffsetDateTime::from_unix_timestamp(1_700_000_001).expect("fixture timestamp"));
|
||||
replacement.data_dir = Some(Uuid::new_v4());
|
||||
replacement.size = 3;
|
||||
if inline {
|
||||
replacement.data = Some(Bytes::from_static(b"new"));
|
||||
}
|
||||
replacement.set_tier_free_version_id(&id.to_string());
|
||||
|
||||
meta.add_version(replacement.clone())
|
||||
.expect("replace transitioned null version");
|
||||
let bytes = meta.marshal_msg().expect("persist replacement and cleanup owner");
|
||||
let mut reopened = FileMeta::load(&bytes).expect("reopen committed metadata");
|
||||
assert_eq!(reopened.versions.len(), 2);
|
||||
let (_, free) = reopened.find_version(Some(id)).expect("durable cleanup owner");
|
||||
assert!(free.free_version());
|
||||
let marker = free.delete_marker.expect("cleanup marker");
|
||||
for suffix in [
|
||||
rustfs_utils::http::SUFFIX_TRANSITION_TIER,
|
||||
rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||
rustfs_utils::http::SUFFIX_TRANSITIONED_OBJECTNAME,
|
||||
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID,
|
||||
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||
] {
|
||||
for prefix in [RUSTFS_INTERNAL_PREFIX, MINIO_INTERNAL_PREFIX] {
|
||||
let key = format!("{prefix}{suffix}");
|
||||
assert_eq!(marker.meta_sys.get(&key), old.meta_sys.get(&key), "{state:?}: {key}");
|
||||
}
|
||||
}
|
||||
let (_, current) = reopened.find_version(None).expect("replacement survives restart");
|
||||
let current = current.object.expect("replacement object");
|
||||
assert_eq!(current.size, 3);
|
||||
assert!(!rustfs_utils::http::contains_key_bytes(¤t.meta_sys, SUFFIX_TIER_FV_ID));
|
||||
reopened
|
||||
.add_version(replacement)
|
||||
.expect("replaying replacement is idempotent");
|
||||
assert_eq!(reopened.versions.len(), 2);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tier_overwrite_rejects_cleanup_failure_without_mutating_source() {
|
||||
for id in ["not-a-uuid".to_string(), Uuid::nil().to_string()] {
|
||||
let (mut meta, _) = tier_overwrite_fixture(crate::TransitionVersionState::Exact);
|
||||
let before = meta.clone();
|
||||
let mut replacement = FileInfo::new("object", 2, 2);
|
||||
replacement.mod_time = Some(OffsetDateTime::now_utc());
|
||||
replacement.set_tier_free_version_id(&id);
|
||||
assert!(meta.add_version(replacement).is_err());
|
||||
assert_eq!(meta, before, "invalid cleanup identity must preserve source");
|
||||
}
|
||||
with_object_max_versions_for_test(1, || {
|
||||
let (mut meta, _) = tier_overwrite_fixture(crate::TransitionVersionState::KnownDisabled);
|
||||
let before = meta.clone();
|
||||
let mut replacement = FileInfo::new("object", 2, 2);
|
||||
replacement.mod_time = Some(OffsetDateTime::now_utc());
|
||||
replacement.data = Some(Bytes::from_static(b"new"));
|
||||
replacement.set_tier_free_version_id(&Uuid::new_v4().to_string());
|
||||
assert_eq!(
|
||||
meta.add_version(replacement)
|
||||
.expect_err("cleanup owner exceeds version limit"),
|
||||
Error::MaxVersionsExceeded
|
||||
);
|
||||
assert_eq!(meta, before, "failed cleanup insertion must also preserve inline bytes");
|
||||
});
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tier_overwrite_allows_empty_transition_status_on_local_source() {
|
||||
let mut source = FileInfo::new("object", 2, 2);
|
||||
source.mod_time = Some(OffsetDateTime::now_utc());
|
||||
source.data = Some(Bytes::from_static(b"old"));
|
||||
rustfs_utils::http::insert_str(&mut source.metadata, rustfs_utils::http::SUFFIX_TRANSITION_STATUS, String::new());
|
||||
let mut meta = FileMeta::new();
|
||||
meta.add_version(source)
|
||||
.expect("seed readable local metadata with empty status");
|
||||
let mut replacement = FileInfo::new("object", 2, 2);
|
||||
replacement.mod_time = Some(OffsetDateTime::now_utc());
|
||||
replacement.size = 3;
|
||||
replacement.data = Some(Bytes::from_static(b"new"));
|
||||
replacement.set_tier_free_version_id(&Uuid::new_v4().to_string());
|
||||
meta.add_version(replacement)
|
||||
.expect("an ordinary overwrite must still succeed");
|
||||
assert_eq!(meta.versions.len(), 1);
|
||||
let (_, current) = meta.find_version(None).expect("replacement remains visible");
|
||||
assert_eq!(current.object.expect("ordinary object").size, 3);
|
||||
assert!(!meta.versions[0].header.free_version());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tier_overwrite_preserves_legacy_metadata_casing() {
|
||||
use rustfs_utils::http::{MINIO_INTERNAL_PREFIX, RUSTFS_INTERNAL_PREFIX};
|
||||
|
||||
for state in [
|
||||
crate::TransitionVersionState::Exact,
|
||||
crate::TransitionVersionState::KnownDisabled,
|
||||
] {
|
||||
let (mut meta, _) = tier_overwrite_fixture(state);
|
||||
let mut source = meta.versions[0].parse_version_meta().expect("seeded source");
|
||||
let object = source.object.as_mut().expect("transitioned source");
|
||||
let expected = object.meta_sys.clone();
|
||||
object.meta_sys = object
|
||||
.meta_sys
|
||||
.drain()
|
||||
.map(|(key, value)| (key.to_ascii_uppercase(), value))
|
||||
.collect();
|
||||
meta.versions[0] = FileMetaShallowVersion::try_from(source).expect("legacy key casing");
|
||||
let id = Uuid::new_v4();
|
||||
let mut replacement = FileInfo::new("object", 2, 2);
|
||||
replacement.mod_time = Some(OffsetDateTime::now_utc());
|
||||
replacement.set_tier_free_version_id(&id.to_string());
|
||||
meta.add_version(replacement).expect("overwrite readable legacy source");
|
||||
let reopened = FileMeta::load(&meta.marshal_msg().expect("persist overwrite")).expect("reopen overwrite");
|
||||
let (_, owner) = reopened
|
||||
.find_version(Some(id))
|
||||
.expect("legacy source must retain cleanup ownership");
|
||||
let marker = owner.delete_marker.expect("cleanup marker");
|
||||
for suffix in [
|
||||
rustfs_utils::http::SUFFIX_TRANSITION_TIER,
|
||||
rustfs_utils::http::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||
rustfs_utils::http::SUFFIX_TRANSITIONED_OBJECTNAME,
|
||||
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_ID,
|
||||
rustfs_utils::http::SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||
] {
|
||||
for prefix in [RUSTFS_INTERNAL_PREFIX, MINIO_INTERNAL_PREFIX] {
|
||||
let key = format!("{prefix}{suffix}");
|
||||
assert_eq!(marker.meta_sys.get(&key), expected.get(&key), "legacy {state:?}: {key}");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tier_overwrite_rejects_conflicting_remote_metadata_aliases() {
|
||||
use rustfs_utils::http::{
|
||||
MINIO_INTERNAL_PREFIX, SUFFIX_TRANSITION_STATUS, SUFFIX_TRANSITION_TIER, SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||
SUFFIX_TRANSITIONED_OBJECTNAME, SUFFIX_TRANSITIONED_VERSION_ID, SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||
};
|
||||
for suffix in [
|
||||
SUFFIX_TRANSITION_STATUS,
|
||||
SUFFIX_TRANSITION_TIER,
|
||||
SUFFIX_TRANSITION_TIER_DESTINATION_ID,
|
||||
SUFFIX_TRANSITIONED_OBJECTNAME,
|
||||
SUFFIX_TRANSITIONED_VERSION_ID,
|
||||
SUFFIX_TRANSITIONED_VERSION_STATE,
|
||||
] {
|
||||
let (mut meta, _) = tier_overwrite_fixture(crate::TransitionVersionState::Exact);
|
||||
let mut old = meta.versions[0].parse_version_meta().expect("seeded source metadata");
|
||||
old.object
|
||||
.as_mut()
|
||||
.expect("transitioned source")
|
||||
.meta_sys
|
||||
.insert(format!("{MINIO_INTERNAL_PREFIX}{suffix}"), b"conflicting-value".to_vec());
|
||||
meta.versions[0] = FileMetaShallowVersion::try_from(old).expect("encode conflicting aliases");
|
||||
let before = meta.clone();
|
||||
let mut replacement = FileInfo::new("object", 2, 2);
|
||||
replacement.mod_time = Some(OffsetDateTime::now_utc());
|
||||
replacement.set_tier_free_version_id(&Uuid::new_v4().to_string());
|
||||
assert_eq!(
|
||||
meta.add_version(replacement)
|
||||
.expect_err("ambiguous ownership must fail closed"),
|
||||
Error::FileCorrupt
|
||||
);
|
||||
assert_eq!(meta, before, "conflicting {suffix} must not erase the old remote tuple");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tier_overwrite_keeps_retained_remote_and_versioned_copy_ownership() {
|
||||
let (mut meta, mut source) = tier_overwrite_fixture(crate::TransitionVersionState::Exact);
|
||||
source.set_tier_free_version_id(&Uuid::new_v4().to_string());
|
||||
meta.add_version(source.clone())
|
||||
.expect("restore retains the same remote owner");
|
||||
assert_eq!(meta.versions.len(), 1);
|
||||
source.version_id = Some(Uuid::new_v4());
|
||||
source.transition_status.clear();
|
||||
source.transition_tier.clear();
|
||||
source.transitioned_objname.clear();
|
||||
source.transition_version = None;
|
||||
source.transition_version_state = crate::TransitionVersionState::Unknown;
|
||||
meta.add_version(source).expect("versioned write retains historical source");
|
||||
assert_eq!(meta.versions.len(), 2);
|
||||
assert!(meta.versions.iter().all(|version| !version.header.free_version()));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn add_version_filemata_uses_canonical_equal_time_order() {
|
||||
let mod_time = OffsetDateTime::from_unix_timestamp(1_700_000_000).expect("valid test timestamp");
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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(),
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -198,7 +198,7 @@ fn expected_admin_route_matrix() -> Vec<RouteMatrixEntry> {
|
||||
admin_route(Method::POST, "/v3/rebalance/stop"),
|
||||
admin_route(Method::POST, "/v3/heal/"),
|
||||
admin_route_sample(Method::POST, "/v3/heal/{bucket}", "/v3/heal/test-bucket"),
|
||||
admin_route_sample(Method::POST, "/v3/heal/{bucket}/{prefix}", "/v3/heal/test-bucket/prefix"),
|
||||
admin_route_sample(Method::POST, "/v3/heal/{bucket}/{*prefix}", "/v3/heal/test-bucket/prefix"),
|
||||
admin_route(Method::POST, "/v3/background-heal/status"),
|
||||
admin_route(Method::GET, "/v4/heal/replacement-recovery"),
|
||||
admin_route(Method::GET, "/v3/tier"),
|
||||
|
||||
@@ -377,6 +377,10 @@ impl FS {
|
||||
|
||||
pub(crate) fn parse_object_version_id(version_id: Option<String>) -> S3Result<Option<Uuid>> {
|
||||
if let Some(vid) = version_id {
|
||||
if vid == "null" {
|
||||
// A nil UUID selects the stored null version; None selects latest.
|
||||
return Ok(Some(Uuid::nil()));
|
||||
}
|
||||
let uuid = Uuid::parse_str(&vid).map_err(|e| {
|
||||
error!("Invalid version ID: {}", e);
|
||||
s3_error!(InvalidArgument, "Invalid version ID")
|
||||
@@ -1183,7 +1187,11 @@ impl S3 for FS {
|
||||
error = %e,
|
||||
"Object tags not found"
|
||||
);
|
||||
return Err(s3_error!(NoSuchKey));
|
||||
return Err(S3Error::new(if opts.version_id.is_some() {
|
||||
S3ErrorCode::NoSuchVersion
|
||||
} else {
|
||||
S3ErrorCode::NoSuchKey
|
||||
}));
|
||||
}
|
||||
error!(
|
||||
component = LOG_COMPONENT_STORAGE,
|
||||
|
||||
@@ -17,7 +17,9 @@ mod tests {
|
||||
use crate::config::WorkloadProfile;
|
||||
use crate::server::cors;
|
||||
use crate::storage::StorageError;
|
||||
use crate::storage::ecfs::{FS, propagate_object_lock_peer_reload, validate_object_lock_configuration_input};
|
||||
use crate::storage::ecfs::{
|
||||
FS, parse_object_version_id, propagate_object_lock_peer_reload, validate_object_lock_configuration_input,
|
||||
};
|
||||
use crate::storage::ecfs_extend::{apply_bucket_default_lock_retention, map_bucket_object_lock_config_state};
|
||||
use crate::storage::s3_api::common::{rustfs_initiator, rustfs_owner};
|
||||
use crate::storage::storage_api::ecstore_bucket::metadata_sys::ObjectLockConfigState;
|
||||
@@ -641,6 +643,36 @@ mod tests {
|
||||
assert_eq!(metadata.get("content-type"), Some(&"application/octet-stream".to_string()));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_tagging_version_id_preserves_explicit_null_and_latest_selection() {
|
||||
assert_eq!(parse_object_version_id(None).expect("latest version selector"), None);
|
||||
assert_eq!(
|
||||
parse_object_version_id(Some("null".to_owned())).expect("explicit null version selector"),
|
||||
Some(uuid::Uuid::nil())
|
||||
);
|
||||
for version in [uuid::Uuid::nil(), uuid::Uuid::new_v4()] {
|
||||
assert_eq!(
|
||||
parse_object_version_id(Some(version.to_string())).expect("UUID version selector"),
|
||||
Some(version)
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_tagging_version_id_rejects_invalid_values_instead_of_selecting_latest() {
|
||||
for version in [
|
||||
"",
|
||||
"NULL",
|
||||
" null ",
|
||||
"not-a-version",
|
||||
"null/other",
|
||||
"00000000-0000-0000-0000-00000000000g",
|
||||
] {
|
||||
let err = parse_object_version_id(Some(version.to_owned())).expect_err("invalid version must fail closed");
|
||||
assert_eq!(err.code(), &S3ErrorCode::InvalidArgument, "version: {version:?}");
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_get_object_tagging_returns_internal_error_when_store_uninitialized() {
|
||||
if !store_uninitialized_premise_holds() {
|
||||
|
||||
@@ -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,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
|
||||
|
||||
@@ -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")
|
||||
|
||||
Reference in New Issue
Block a user