Compare commits

...

4 Commits

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

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

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

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

* fix(ecstore): satisfy clippy in rebalance test

* fix(ecstore): preserve decommission delete compatibility
2026-09-11 22:36:15 +08:00
13 changed files with 1720 additions and 96 deletions
@@ -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}" \
@@ -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
{
@@ -577,6 +577,26 @@ fn heal_control_auth_may_need_replay_scope_refresh(err: &Error) -> bool {
)
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum HealControlRetryAction {
Reconnect,
RefreshReplayScope,
}
fn heal_control_retry_action(
err: &Error,
reconnect_attempted: bool,
replay_scope_refresh_attempted: bool,
) -> Option<HealControlRetryAction> {
if !replay_scope_refresh_attempted && heal_control_auth_may_need_replay_scope_refresh(err) {
return Some(HealControlRetryAction::RefreshReplayScope);
}
if !reconnect_attempted && PeerRestClient::is_network_like_error(err) {
return Some(HealControlRetryAction::Reconnect);
}
None
}
fn decode_remote_version_state_capability(expected_member: &str, result: &[u8]) -> Result<Uuid> {
let (topology_member, process_epoch) = rustfs_protos::decode_remote_version_state_capability(result).map_err(Error::other)?;
if topology_member != expected_member {
@@ -1753,34 +1773,41 @@ impl PeerRestClient {
return Err(Error::other("heal control command exceeds size limit"));
}
let capability_probe = rustfs_protos::is_heal_control_capability_probe(&command);
let result = self
.heal_control_once(version, &topology_fingerprint, &command, capability_probe)
.await;
if result
.as_ref()
.err()
.is_some_and(heal_control_auth_may_need_replay_scope_refresh)
{
self.prepare_heal_control_auth_retry().await;
return self
.finalize_result(
self.heal_control_once(version, &topology_fingerprint, &command, capability_probe)
.await,
)
let mut reconnect_attempted = false;
let mut replay_scope_refresh_attempted = false;
loop {
let result = self
.heal_control_once(version, &topology_fingerprint, &command, capability_probe)
.await;
let Some(action) = result
.as_ref()
.err()
.and_then(|err| heal_control_retry_action(err, reconnect_attempted, replay_scope_refresh_attempted))
else {
return self.finalize_result(result).await;
};
match action {
HealControlRetryAction::Reconnect => reconnect_attempted = true,
HealControlRetryAction::RefreshReplayScope => replay_scope_refresh_attempted = true,
}
self.prepare_heal_control_retry(action).await;
}
self.finalize_result(result).await
}
async fn prepare_heal_control_auth_retry(&self) {
if let Err(err) = clear_peer_replay_state_for_addr(&self.grid_host) {
async fn prepare_heal_control_retry(&self, action: HealControlRetryAction) {
if action == HealControlRetryAction::RefreshReplayScope
&& let Err(err) = clear_peer_replay_state_for_addr(&self.grid_host)
{
debug!(
peer = %self.grid_host,
error = %err,
"could not clear heal control replay state before retry"
);
}
self.evict_connection().await;
// A restart can leave both the local offline gate and the peer replay
// epoch stale. Clear the gate on either recovery step so the next
// bounded attempt reaches a fresh channel instead of fast-failing.
self.prepare_retry().await;
}
async fn heal_control_once(
@@ -3867,6 +3894,51 @@ mod tests {
)));
}
#[test]
fn heal_control_retry_plan_allows_one_reconnect_and_one_epoch_refresh() {
let offline = Error::RemoteClientUnavailable("peer http://127.0.0.1:9000 is temporarily offline".to_string());
let stale_epoch = Error::from(tonic::Status::unauthenticated("No valid auth token"));
assert_eq!(heal_control_retry_action(&offline, false, false), Some(HealControlRetryAction::Reconnect));
assert_eq!(heal_control_retry_action(&offline, true, false), None);
assert_eq!(
heal_control_retry_action(&stale_epoch, true, false),
Some(HealControlRetryAction::RefreshReplayScope),
"a reconnect may expose the restarted peer's stale replay epoch"
);
assert_eq!(heal_control_retry_action(&stale_epoch, true, true), None);
assert_eq!(
heal_control_retry_action(&stale_epoch, false, false),
Some(HealControlRetryAction::RefreshReplayScope)
);
assert_eq!(
heal_control_retry_action(&offline, false, true),
Some(HealControlRetryAction::Reconnect),
"an epoch refresh may be followed by one bounded reconnect"
);
assert_eq!(heal_control_retry_action(&offline, true, true), None);
assert_eq!(
heal_control_retry_action(&Error::from(tonic::Status::permission_denied("bad signature")), false, false),
None,
"authorization failures must never be retried"
);
}
#[tokio::test]
async fn heal_control_epoch_refresh_clears_offline_gate() {
let client = test_peer_client();
client.offline.store(true, Ordering::Release);
client
.prepare_heal_control_retry(HealControlRetryAction::RefreshReplayScope)
.await;
assert!(
!client.offline.load(Ordering::Acquire),
"epoch refresh must not leave the following attempt behind the offline gate"
);
}
#[test]
fn peer_rest_client_network_classifier_keeps_slow_peers_online() {
// The per-RPC channel deadline (RUSTFS_INTERNODE_RPC_TIMEOUT, 30s)
+14
View File
@@ -7470,6 +7470,20 @@ impl PoolMeta {
.is_some_and(is_decommission_suspended)
}
pub(crate) fn has_active_decommission_capacity_reservation(&self, idx: usize) -> bool {
self.pools
.get(idx)
.and_then(|pool| pool.decommission.as_ref())
.is_some_and(|info| {
info.has_decommission_state()
&& is_decommission_active(info.complete, info.failed, info.canceled)
&& info
.capacity_reservation
.as_ref()
.is_some_and(DecommissionCapacityReservation::active)
})
}
pub(crate) fn scanner_pause_backlog_pool_writable(&self, idx: usize) -> bool {
self.pools.get(idx).is_some_and(|pool| {
!pool
+54
View File
@@ -41,6 +41,14 @@ struct DanglingDeleteGraceError {
grace_secs: i64,
}
/// Marks a conditional-file write that failed before its publication rename.
/// Callers may choose another owner only while this marker is preserved; every
/// unmarked error remains commit-ambiguous and must fail closed.
#[derive(Debug)]
struct ConditionalFileNotCommittedError {
source: io::Error,
}
// DiskError == StorageErr
#[derive(Debug, thiserror::Error)]
pub enum DiskError {
@@ -220,6 +228,18 @@ impl std::fmt::Display for DanglingDeleteGraceError {
impl StdError for DanglingDeleteGraceError {}
impl std::fmt::Display for ConditionalFileNotCommittedError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
self.source.fmt(f)
}
}
impl StdError for ConditionalFileNotCommittedError {
fn source(&self) -> Option<&(dyn StdError + 'static)> {
Some(&self.source)
}
}
fn classify_internode_missing_error(error: &InternodeHttpError) -> Option<DiskError> {
if error.is_remote_file_not_found() {
return Some(DiskError::FileNotFound);
@@ -293,6 +313,22 @@ impl DiskError {
})
}
pub(crate) fn conditional_file_not_committed(source: io::Error) -> io::Error {
io::Error::new(source.kind(), ConditionalFileNotCommittedError { source })
}
/// Whether a local conditional-file replacement failed before the target
/// publication rename and therefore cannot have committed new owner bytes.
pub fn is_conditional_file_not_committed(&self) -> bool {
matches!(
self,
DiskError::Io(io_error)
if io_error
.get_ref()
.is_some_and(|source| source.downcast_ref::<ConditionalFileNotCommittedError>().is_some())
)
}
pub fn is_dangling_delete_grace(&self) -> bool {
matches!(self, DiskError::Io(io_error) if Self::io_error_is_dangling_delete_grace(io_error))
}
@@ -627,6 +663,9 @@ impl From<tokio::task::JoinError> for DiskError {
impl Clone for DiskError {
fn clone(&self) -> Self {
match self {
DiskError::Io(io_error) if self.is_conditional_file_not_committed() => DiskError::Io(
DiskError::conditional_file_not_committed(io::Error::new(io_error.kind(), io_error.to_string())),
),
DiskError::Io(io_error) => DiskError::Io(
rustfs_rio::clone_internode_http_io_error(io_error)
.and_then(std::io::Error::into_inner)
@@ -820,6 +859,21 @@ mod tests {
use super::*;
use std::collections::HashMap;
#[test]
fn conditional_file_not_committed_marker_is_explicit_and_clone_safe() {
let marked = DiskError::from(DiskError::conditional_file_not_committed(io::Error::new(
io::ErrorKind::PermissionDenied,
"staging rejected",
)));
assert!(marked.is_conditional_file_not_committed());
assert!(marked.clone().is_conditional_file_not_committed());
assert!(!DiskError::Timeout.is_conditional_file_not_committed());
assert!(
!DiskError::Io(io::Error::new(io::ErrorKind::PermissionDenied, "rename rejected"))
.is_conditional_file_not_committed()
);
}
#[test]
fn terminal_read_error_preserves_kind_and_disk_classification() {
let timeout = terminal_read_error_to_io(DiskError::Timeout);
+14 -4
View File
@@ -8957,7 +8957,8 @@ impl DiskAPI for LocalDisk {
.truncate(false)
.read(true)
.write(true)
.open(&lock_path)?;
.open(&lock_path)
.map_err(DiskError::conditional_file_not_committed)?;
flock(&lock, FlockOperation::NonBlockingLockExclusive).map_err(std::io::Error::from)?;
let result = (|| {
let current = match std::fs::read(&file_path) {
@@ -9012,10 +9013,15 @@ impl DiskAPI for LocalDisk {
.ok_or_else(|| std::io::Error::new(ErrorKind::InvalidInput, "conditional file has no parent"))?;
let temporary = parent.join(format!(".{}.{}.tmp", path.replace('/', "_"), Uuid::new_v4()));
let write_result = (|| -> std::io::Result<()> {
let mut staged = std::fs::OpenOptions::new().create_new(true).write(true).open(&temporary)?;
staged.write_all(&replacement)?;
let not_committed = DiskError::conditional_file_not_committed;
let mut staged = std::fs::OpenOptions::new()
.create_new(true)
.write(true)
.open(&temporary)
.map_err(not_committed)?;
staged.write_all(&replacement).map_err(not_committed)?;
if sync_metadata {
staged.sync_all()?;
staged.sync_all().map_err(not_committed)?;
}
std::fs::rename(&temporary, &file_path)?;
Ok(())
@@ -22650,6 +22656,10 @@ mod test {
.await
.expect_err("directory fsync failure must fail the CAS update");
assert!(matches!(err, DiskError::Io(ref err) if err.kind() == ErrorKind::Other));
assert!(
!err.is_conditional_file_not_committed(),
"an error after publication rename must remain commit-ambiguous"
);
assert_eq!(
disk.read_all(RUSTFS_META_BUCKET, HEALING_MARKER_PATH)
.await
+475 -38
View File
@@ -20,6 +20,7 @@ fn to_filemeta_err(err: Error) -> rustfs_filemeta::Error {
err.narrow_to_filemeta().unwrap_or_else(rustfs_filemeta::Error::other)
}
use crate::bucket::lifecycle::bucket_lifecycle_ops::free_version_remote_tuple_matches;
use crate::bucket::metadata_sys::{
get_versioning_config, has_authoritative_never_versioned_state, has_authoritative_never_versioned_state_in,
};
@@ -70,7 +71,7 @@ use tokio::io::duplex;
use tokio::sync::broadcast::{self};
use tokio::sync::mpsc::{self, Receiver, Sender};
use tokio::sync::{OnceCell, RwLock};
use tokio::task::JoinSet;
use tokio::task::{JoinHandle, JoinSet};
use tokio_util::sync::CancellationToken;
use tracing::{Instrument, debug, error, info, warn};
use uuid::Uuid;
@@ -4331,6 +4332,7 @@ impl ECStore {
"store list_merged started"
);
let rx = rx.child_token();
let mut futures = Vec::new();
let mut inputs = Vec::new();
@@ -4346,16 +4348,10 @@ impl ECStore {
}
}
tokio::spawn(
async move {
if let Err(err) = merge_entry_channels(rx, inputs, sender.clone(), 1).await {
error!("merge_entry_channels err {:?}", err)
}
}
.instrument(tracing::Span::current()),
);
let merge_task = spawn_listing_merge(rx, inputs, sender);
let results = join_all(futures).await;
merge_task.await.map_err(Error::from)??;
let mut all_at_eof = true;
@@ -4422,6 +4418,7 @@ impl ECStore {
) -> Result<()> {
check_list_objs_args(bucket, prefix, &None)?;
let rx = rx.child_token();
let mut futures = Vec::new();
let mut inputs = Vec::new();
@@ -4783,17 +4780,11 @@ impl ECStore {
.instrument(tracing::Span::current()),
);
tokio::spawn(
async move {
if let Err(err) = merge_entry_channels(rx, inputs, merge_tx, 1).await {
error!("merge_entry_channels err {:?}", err)
}
}
.instrument(tracing::Span::current()),
);
let merge_task = spawn_listing_merge(rx, inputs, merge_tx);
let walk_started = std::time::Instant::now();
let walk_results = join_all(futures).await;
merge_task.await.map_err(Error::from)??;
let mut errs = Vec::new();
for walk_result in walk_results {
match walk_result {
@@ -5068,6 +5059,130 @@ async fn send_or_cancel(rx: &CancellationToken, out_channel: &Sender<MetaCacheEn
}
}
/// Each input has already been resolved inside its own erasure set. This is a
/// union of version histories, never a quorum vote between unrelated pools.
fn merge_object_entry_versions(first: &mut MetaCacheEntry, others: impl Iterator<Item = MetaCacheEntry>) -> Result<()> {
let name = first.name.clone();
let mut versions: HashMap<(Option<Uuid>, bool), (FileMetaShallowVersion, ObjectInfo)> = HashMap::new();
for mut entry in std::iter::once(std::mem::take(first)).chain(others) {
let meta = match entry.cached.take() {
Some(meta) => meta,
None => FileMeta::load(&entry.metadata).map_err(|_| Error::FileCorrupt)?,
};
if meta.versions.is_empty() {
return Err(Error::FileCorrupt);
}
for version in meta.versions {
let parsed = version.parse_version_meta().map_err(|_| Error::FileCorrupt)?;
if !parsed.valid() || parsed.version_type != version.header.version_type {
return Err(Error::FileCorrupt);
}
let fi = parsed.into_fileinfo("", &name, true).map_err(|_| Error::FileCorrupt)?;
let version_id = fi.version_id.filter(|id| !id.is_nil());
if version_id != version.header.version_id.filter(|id| !id.is_nil())
|| fi.mod_time != version.header.mod_time
|| fi.tier_free_version() != version.header.free_version()
{
return Err(Error::FileCorrupt);
}
let info = ObjectInfo::from_file_info(&fi, "", &name, true);
let identity = (version_id, version.header.free_version());
match versions.entry(identity) {
std::collections::hash_map::Entry::Vacant(slot) => {
slot.insert((version, info));
}
std::collections::hash_map::Entry::Occupied(mut slot) => {
let (previous, previous_info) = slot.get();
// Suspended and unversioned writes replace the one null
// slot. Distinct UUID versions never supersede each other.
if version_id.is_none() && info.mod_time != previous_info.mod_time {
if info.mod_time > previous_info.mod_time {
slot.insert((version, info));
}
continue;
}
let equivalent = if info.delete_marker && previous_info.delete_marker {
super::object::is_equivalent_data_movement_delete_marker(&info, previous_info)
} else {
crate::data_movement::is_equivalent_data_movement_object_identity(&info, previous_info, true, true)
};
if !equivalent {
return Err(Error::FileCorrupt);
}
// Equivalent migrated copies can have different coding or
// data directories. Choose a stable representation without
// making input order part of the S3 version order.
if version.meta < previous.meta {
slot.insert((version, info));
}
}
}
}
}
let mut live_remote_references = HashMap::<String, HashMap<String, Vec<ObjectInfo>>>::new();
for (_, info) in versions.values() {
if !info.transitioned_object.free_version && info.transitioned_object.status == rustfs_filemeta::TRANSITION_COMPLETE {
live_remote_references
.entry(info.transitioned_object.tier.clone())
.or_default()
.entry(info.transitioned_object.name.clone())
.or_default()
.push(info.clone());
}
}
// Keep cleanup durable in its source xl.meta, but do not expose it to a
// merged recovery walk while another physical pool still owns the tuple.
versions.retain(|_, (_, info)| {
!info.transitioned_object.free_version
|| !live_remote_references
.get(info.transitioned_object.tier.as_str())
.and_then(|by_name| by_name.get(info.transitioned_object.name.as_str()))
.is_some_and(|candidates| {
candidates
.iter()
.any(|live| free_version_remote_tuple_matches(info, live).unwrap_or(false))
})
});
let mut merged = FileMeta::new();
merged.versions = versions.into_values().map(|(version, _)| version).collect();
merged.versions.sort_by(|a, b| {
if a.header.sorts_before(&b.header) {
std::cmp::Ordering::Less
} else if b.header.sorts_before(&a.header) {
std::cmp::Ordering::Greater
} else {
std::cmp::Ordering::Equal
}
});
let metadata = merged.marshal_msg()?;
*first = MetaCacheEntry {
name,
metadata,
cached: Some(merged),
reusable: true,
};
Ok(())
}
/// `rx` is private to the producers. Cancelling it on a merge error must not
/// cancel the request token, which would suppress that error at the API edge.
fn spawn_listing_merge(
rx: CancellationToken,
inputs: Vec<Receiver<MetaCacheEntry>>,
sender: Sender<MetaCacheEntry>,
) -> JoinHandle<Result<()>> {
tokio::spawn(
async move {
let result = merge_entry_channels(rx.clone(), inputs, sender, 1).await;
if result.is_err() {
rx.cancel();
}
result
}
.instrument(tracing::Span::current()),
)
}
async fn merge_entry_channels(
rx: CancellationToken,
in_channels: Vec<Receiver<MetaCacheEntry>>,
@@ -5133,6 +5248,7 @@ async fn merge_entry_channels(
// after anything greater has been emitted).
let mut last_emitted = String::new();
let mut group: Vec<Box<MergeHead>> = Vec::new();
let mut object_entries: Vec<MetaCacheEntry> = Vec::new();
let mut refill: Vec<usize> = Vec::with_capacity(in_channels.len());
while let Some(Reverse(first)) = heap.pop() {
@@ -5150,7 +5266,7 @@ async fn merge_entry_channels(
// Resolve the same-name group to one winner (heads arrive in ascending
// channel order):
// - prefix dir vs prefix dir: the first (lowest channel) wins;
// - object vs object: the later channel wins (legacy authority rule);
// - object vs object: merge the independently resolved version stacks;
// - object vs prefix dir: same-name means both end with the separator,
// i.e. the object is an explicit "directory marker" for the same S3
// key — it shadows the prefix dir so the key does not surface as
@@ -5168,11 +5284,27 @@ async fn merge_entry_channels(
if dir_winner.is_none() {
dir_winner = Some(head);
}
} else if let Some(winner) = object_winner.as_ref() {
// Key-only candidates carry no version metadata and cannot
// replace a resolved stack or contribute a quorum vote.
if head.entry.is_object() {
if winner.entry.is_object() {
object_entries.push(head.entry);
} else {
object_winner = Some(head);
}
}
} else {
object_winner = Some(head);
}
}
if !object_entries.is_empty()
&& let Some(winner) = object_winner.as_mut()
{
merge_object_entry_versions(&mut winner.entry, object_entries.drain(..))?;
}
if let Some(head) = object_winner.or(dir_winner)
&& head.entry.name != last_emitted
{
@@ -5605,6 +5737,7 @@ impl Sets {
"sets list_merged started"
);
let rx = rx.child_token();
let mut futures = Vec::new();
let mut inputs = Vec::new();
@@ -5617,16 +5750,10 @@ impl Sets {
futures.push(async move { set.list_path(rx_clone, opts, send).await });
}
tokio::spawn(
async move {
if let Err(err) = merge_entry_channels(rx, inputs, sender.clone(), 1).await {
error!("merge_entry_channels err {:?}", err);
}
}
.instrument(tracing::Span::current()),
);
let merge_task = spawn_listing_merge(rx, inputs, sender);
let results = join_all(futures).await;
merge_task.await.map_err(Error::from)??;
let mut all_at_eof = true;
let mut errs = Vec::new();
for result in results {
@@ -5677,6 +5804,7 @@ impl Sets {
) -> Result<()> {
check_list_objs_args(bucket, prefix, &None)?;
let rx = rx.child_token();
let mut futures = Vec::new();
let mut inputs = Vec::new();
@@ -6007,17 +6135,11 @@ impl Sets {
.instrument(tracing::Span::current()),
);
tokio::spawn(
async move {
if let Err(err) = merge_entry_channels(rx, inputs, merge_tx, 1).await {
error!("merge_entry_channels err {:?}", err)
}
}
.instrument(tracing::Span::current()),
);
let merge_task = spawn_listing_merge(rx, inputs, merge_tx);
let walk_started = std::time::Instant::now();
let walk_results = join_all(futures).await;
merge_task.await.map_err(Error::from)??;
let mut errs = Vec::new();
for walk_result in walk_results {
match walk_result {
@@ -7016,7 +7138,7 @@ mod test {
};
use crate::cache_value::metacache_set::{FallbackClaimTracker, TestReaderBehavior, list_path_raw};
use crate::disk::{DiskAPI, DiskOption, STORAGE_FORMAT_FILE, endpoint::Endpoint, error::DiskError, new_disk};
use crate::error::StorageError;
use crate::error::{Result, StorageError};
use crate::object_api::ObjectInfo;
use rustfs_filemeta::{
FileInfo, FileMeta, FileMetaVersion, MetaCacheEntries, MetaCacheEntriesSorted, MetaCacheEntry, MetaDeleteMarker,
@@ -7338,6 +7460,47 @@ mod test {
}
}
fn test_transitioned_meta_entry(name: &str, remote_object: &str, delete_source: bool) -> MetaCacheEntry {
let mut source = FileInfo::new(name, 2, 2);
source.volume = "bucket".to_string();
source.name = name.to_string();
source.version_id = Some(Uuid::from_u128(1));
source.versioned = true;
source.size = 1;
source.mod_time = Some(time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp"));
source.transition_status = rustfs_filemeta::TRANSITION_COMPLETE.to_string();
source.transition_tier = "WARM".to_string();
source.transitioned_objname = remote_object.to_string();
source.transition_version = Some("remote-version".to_string());
source.transition_version_state = rustfs_filemeta::TransitionVersionState::Exact;
rustfs_utils::http::metadata_compat::insert_str(
&mut source.metadata,
rustfs_utils::http::metadata_compat::SUFFIX_TRANSITION_TIER_DESTINATION_ID,
"00".repeat(32),
);
let mut meta = FileMeta::new();
meta.add_version(source.clone())
.expect("test metadata should accept transitioned source");
if delete_source {
let mut delete = FileInfo {
name: name.to_string(),
version_id: source.version_id,
..Default::default()
};
delete.set_tier_free_version_id(&Uuid::from_u128(2).to_string());
meta.delete_version(&delete)
.expect("transitioned delete should create a free-version owner");
}
let metadata = meta.marshal_msg().expect("test transitioned metadata should marshal");
MetaCacheEntry {
name: name.to_string(),
metadata,
cached: Some(meta),
reusable: false,
}
}
fn test_object_with_delete_marker_meta_entry(
name: &str,
object_mod_time: time::OffsetDateTime,
@@ -10399,7 +10562,7 @@ mod test {
}
#[tokio::test]
async fn merge_entry_channels_documents_candidate_metadata_authority_risk() {
async fn merge_entry_channels_preserves_cross_pool_delete_marker_versions() {
let (tx_a, rx_a) = mpsc::channel(4);
let (tx_b, rx_b) = mpsc::channel(4);
let (tx_c, rx_c) = mpsc::channel(4);
@@ -10428,9 +10591,13 @@ mod test {
.expect("merged entry should be present");
assert_eq!(merged.name, "obj-a");
assert!(
!merged.is_latest_delete_marker(),
"current merge consumes candidate metadata bytes; future index-backed strong modes must live-verify metadata instead"
merged.is_latest_delete_marker(),
"a newer marker must remain current across independently resolved pools"
);
let versions = merged.file_info_versions("bucket").expect("merged versions should decode");
assert_eq!(versions.versions.len(), 2, "retain the historical object and deduplicate the marker");
assert!(versions.versions[0].deleted && versions.versions[0].is_latest);
assert!(!versions.versions[1].deleted && !versions.versions[1].is_latest);
assert!(
matches!(timeout(Duration::from_secs(1), out_rx.recv()).await, Ok(None)),
"merge should not emit a duplicate entry for the same key"
@@ -10442,6 +10609,276 @@ mod test {
.expect("merge task should succeed");
}
fn rewrite_test_version(mut entry: MetaCacheEntry, change: impl FnOnce(&mut FileMetaVersion)) -> MetaCacheEntry {
let meta = entry.cached.as_mut().expect("test metadata should be decoded");
assert_eq!(meta.versions.len(), 1);
let mut version = meta.versions[0].parse_version_meta().expect("test version should decode");
change(&mut version);
meta.versions[0] = version.try_into().expect("test version should encode");
entry.metadata = meta.marshal_msg().expect("test metadata should encode");
entry
}
async fn merge_test_object_entries(entries: Vec<MetaCacheEntry>) -> Result<MetaCacheEntry> {
let mut inputs = Vec::with_capacity(entries.len());
for entry in entries {
let (sender, receiver) = mpsc::channel(1);
sender.send(entry).await.expect("fixture entry should queue");
inputs.push(receiver);
}
let (sender, mut receiver) = mpsc::channel(1);
let task = tokio::spawn(merge_entry_channels(CancellationToken::new(), inputs, sender, 1));
let entry = receiver.recv().await;
task.await.expect("merge must not panic")?;
Ok(entry.expect("a valid same-key group must produce an entry"))
}
#[tokio::test]
async fn merge_entry_channels_orders_complete_histories_independently_of_pool_order() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let first = test_object_meta_entry_with_erasure_versions("key", &[(time, "first", 4, 2)]);
let second = rewrite_test_version(
test_object_meta_entry_with_erasure_versions("key", &[(time, "second", 4, 2)]),
|version| version.object.as_mut().expect("object version").version_id = Some(Uuid::from_u128(2)),
);
let marker = test_delete_marker_meta_entry("key", time + time::Duration::seconds(1));
let inputs = [first, second, marker];
let mut expected = None;
for order in [[0, 1, 2], [0, 2, 1], [1, 0, 2], [1, 2, 0], [2, 0, 1], [2, 1, 0]] {
let entry = merge_test_object_entries(order.map(|index| inputs[index].clone()).to_vec())
.await
.expect("disjoint version chains should merge");
let versions = entry.file_info_versions("bucket").expect("merged versions should decode");
assert_eq!(versions.versions.len(), 3);
assert!(versions.versions[0].deleted && versions.versions[0].is_latest);
assert!(
versions.versions[1..]
.iter()
.all(|version| !version.deleted && !version.is_latest)
);
assert!(versions.versions.iter().all(|version| version.num_versions == 3));
let identities = versions.versions.iter().map(|version| version.version_id).collect::<Vec<_>>();
assert!(identities.contains(&Some(Uuid::from_u128(1))));
assert!(identities.contains(&Some(Uuid::from_u128(2))));
if let Some(expected) = &expected {
assert_eq!(&identities, expected, "equal-time versions must have stable pagination order");
} else {
expected = Some(identities);
}
}
}
#[tokio::test]
async fn merge_entry_channels_key_only_candidates_do_not_override_version_metadata() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let marker = test_delete_marker_meta_entry("key", time);
for entries in [
vec![test_meta_entry("key"), marker.clone()],
vec![marker.clone(), test_meta_entry("key")],
] {
let mut merged = merge_test_object_entries(entries)
.await
.expect("merge a name with resolved metadata");
assert!(merged.is_latest_delete_marker());
assert_eq!(merged.file_info_versions("bucket").expect("decode marker").versions.len(), 1);
}
}
#[tokio::test]
async fn merge_entry_channels_accepts_equivalent_migrated_coding_and_data_dirs() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let first = test_object_meta_entry_with_erasure_versions("key", &[(time, "same-etag", 4, 2)]);
let second = rewrite_test_version(
test_object_meta_entry_with_erasure_versions("key", &[(time, "same-etag", 6, 2)]),
|version| version.object.as_mut().expect("object version").data_dir = Some(Uuid::from_u128(42)),
);
let forward = merge_test_object_entries(vec![first.clone(), second.clone()])
.await
.expect("valid migration copies");
let reverse = merge_test_object_entries(vec![second, first])
.await
.expect("reversed migration copies");
assert_eq!(forward.metadata, reverse.metadata, "representation must not depend on channel order");
let versions = forward.file_info_versions("bucket").expect("merged metadata should decode");
assert_eq!(versions.versions.len(), 1);
assert_eq!(versions.versions[0].metadata.get("etag").map(String::as_str), Some("same-etag"));
}
#[tokio::test]
async fn merge_entry_channels_defers_free_version_while_same_remote_source_is_live() {
let live = test_transitioned_meta_entry("key", "remote/shared", false);
let free = test_transitioned_meta_entry("key", "remote/shared", true);
for inputs in [vec![live.clone(), free.clone()], vec![free.clone(), live.clone()]] {
let merged = merge_test_object_entries(inputs)
.await
.expect("same remote source and cleanup owner should merge");
let versions = merged
.file_info_versions_with_free_versions("bucket")
.expect("merged transition history should decode");
assert_eq!(versions.versions.len(), 1);
assert!(versions.free_versions.is_empty(), "a live remote reference must defer cleanup discovery");
}
let unrelated = test_transitioned_meta_entry("key", "remote/other", false);
let merged = merge_test_object_entries(vec![free, unrelated])
.await
.expect("unrelated remote references should merge");
let versions = merged
.file_info_versions_with_free_versions("bucket")
.expect("merged transition history should decode");
assert_eq!(versions.versions.len(), 1);
assert_eq!(versions.free_versions.len(), 1, "an unrelated source must not suppress cleanup");
}
#[tokio::test]
async fn merge_entry_channels_rejects_conflicting_version_identity_and_metadata() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let original = test_object_meta_entry_with_erasure_versions("key", &[(time, "original", 4, 2)]);
for key in [
"etag",
"x-amz-tagging",
"x-amz-object-lock-mode",
"x-amz-object-lock-retain-until-date",
] {
let changed = rewrite_test_version(original.clone(), |version| {
version
.object
.as_mut()
.expect("object version")
.meta_user
.insert(key.to_string(), "changed".to_string());
});
for pair in [[original.clone(), changed.clone()], [changed, original.clone()]] {
let err = merge_test_object_entries(pair.to_vec())
.await
.expect_err("conflicting copies must fail");
assert_eq!(err, StorageError::FileCorrupt, "conflict in {key} must not become arbitrary metadata");
}
}
let marker = rewrite_test_version(test_delete_marker_meta_entry("key", time), |version| {
version.delete_marker.as_mut().expect("delete marker").version_id = Some(Uuid::from_u128(1));
});
assert_eq!(
merge_test_object_entries(vec![original, marker])
.await
.expect_err("UUID type conflict"),
StorageError::FileCorrupt
);
}
#[tokio::test]
async fn merge_entry_channels_reconciles_null_overwrite_without_losing_uuid_history() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let history = test_object_meta_entry_with_erasure_versions("key", &[(time, "history", 4, 2)]);
let old_null = rewrite_test_version(history.clone(), |version| {
version.object.as_mut().expect("null object").version_id = None;
});
let marker = rewrite_test_version(test_delete_marker_meta_entry("key", time + time::Duration::seconds(1)), |version| {
version.delete_marker.as_mut().expect("null marker").version_id = Some(Uuid::nil());
});
for inputs in [
vec![old_null.clone(), marker.clone(), history.clone()],
vec![history, marker, old_null],
] {
let entry = merge_test_object_entries(inputs)
.await
.expect("new null slot should replace old null slot");
let versions = entry.file_info_versions("bucket").expect("null versions should decode");
assert_eq!(versions.versions.len(), 2);
assert!(versions.versions[0].deleted && versions.versions[0].is_latest);
assert!(versions.versions[0].version_id.is_none_or(|id| id.is_nil()));
assert_eq!(versions.versions[1].version_id, Some(Uuid::from_u128(1)));
}
}
#[tokio::test]
async fn merge_entry_channels_rejects_corrupt_version_headers_and_empty_stacks() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let original = test_object_meta_entry_with_erasure_versions("key", &[(time, "etag", 4, 2)]);
for empty in [false, true] {
let mut corrupt = original.clone();
let meta = corrupt.cached.as_mut().expect("fixture metadata");
if empty {
meta.versions.clear();
} else {
meta.versions[0].header.version_id = Some(Uuid::from_u128(99));
}
corrupt.metadata = meta.marshal_msg().expect("encode corrupt fixture");
assert_eq!(
merge_test_object_entries(vec![original.clone(), corrupt])
.await
.expect_err("corrupt candidate must fail"),
StorageError::FileCorrupt
);
}
let mut malformed = original.clone();
malformed.cached = None;
malformed.metadata = vec![0xff];
assert_eq!(
merge_test_object_entries(vec![original, malformed])
.await
.expect_err("malformed metadata must fail"),
StorageError::FileCorrupt
);
}
#[tokio::test]
async fn merge_entry_channels_does_not_combine_subquorum_markers_across_erasure_sets() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let old = test_object_meta_entry_with_erasure_versions("key", &[(time, "history", 4, 2)]);
let marked = test_object_with_delete_marker_meta_entry("key", time, time + time::Duration::seconds(1));
let mut inputs = Vec::new();
for marker_copies in [1, 2] {
let resolver = list_metadata_resolution_params("bucket".to_string(), 3, 3, true, 0);
let copies = (0..3)
.map(|index| Some(if index < marker_copies { marked.clone() } else { old.clone() }))
.collect();
let entry = resolve_listing_entries(MetaCacheEntries(copies), resolver, false)
.expect("each set independently retains its quorum-backed history");
inputs.push(entry);
}
let merged = merge_test_object_entries(inputs).await.expect("merge resolved histories");
let versions = merged.file_info_versions("bucket").expect("decode merged history");
assert_eq!(
versions.versions.len(),
1,
"three marker copies across two EC domains do not form a quorum"
);
assert!(!versions.versions[0].deleted);
}
#[tokio::test]
async fn listing_merge_preserves_error_after_partial_output_without_cancelling_request() {
let time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let (first_tx, first_rx) = mpsc::channel(2);
let (second_tx, second_rx) = mpsc::channel(1);
first_tx.send(test_meta_entry("a/")).await.expect("queue preceding prefix");
first_tx
.send(test_object_meta_entry_with_erasure_versions("b", &[(time, "one", 4, 2)]))
.await
.expect("queue first copy");
second_tx
.send(test_object_meta_entry_with_erasure_versions("b", &[(time, "two", 4, 2)]))
.await
.expect("queue conflicting copy");
drop(first_tx);
drop(second_tx);
let request = CancellationToken::new();
let workers = request.child_token();
let (sender, mut receiver) = mpsc::channel(1);
let task = super::spawn_listing_merge(workers.clone(), vec![first_rx, second_rx], sender);
assert_eq!(receiver.recv().await.expect("preceding result should arrive").name, "a/");
assert!(receiver.recv().await.is_none());
assert_eq!(
task.await
.expect("merge task must not panic")
.expect_err("conflict must propagate"),
StorageError::FileCorrupt
);
assert!(workers.is_cancelled(), "failed merge must stop the disk producers");
assert!(!request.is_cancelled(), "the API must still observe the merge error");
}
#[tokio::test]
async fn merge_entry_channels_handles_single_channel() {
let (tx, rx) = mpsc::channel(4);
+573 -5
View File
@@ -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
@@ -4791,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;
@@ -4917,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) => {
@@ -4954,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);
}
@@ -4976,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;
}
@@ -5001,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);
}
}
@@ -8291,6 +8319,546 @@ mod tests {
}
}
async fn multipool_version_test_store(bucket: &str) -> (Vec<tempfile::TempDir>, Arc<ECStore>) {
let ctx = Arc::new(crate::runtime::instance::InstanceContext::new());
let (mut dirs, first_set) = make_local_set_disks_with_ctx(4, 2, Arc::clone(&ctx)).await;
let (second_dirs, second_set) = make_local_set_disks_with_ctx(4, 2, Arc::clone(&ctx)).await;
dirs.extend(second_dirs);
let store = Arc::new(new_prepared_reader_test_store_with_ctx(&[first_set, second_set], ctx).await);
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
store
.handle_make_bucket(
bucket,
&MakeBucketOptions {
versioning_enabled: true,
..Default::default()
},
)
.await
.expect("create the versioned bucket in both pools");
(dirs, store)
}
#[tokio::test]
async fn multipool_delete_marker_stays_with_existing_versions() {
let bucket = "multipool-marker-routing";
let object = "history.bin";
let (_dirs, store) = multipool_version_test_store(bucket).await;
let mut expected_versions = Vec::new();
for value in 1..=3_u8 {
let written = store.pools[1]
.put_object(
bucket,
object,
&mut PutObjReader::from_vec(vec![value; 4097]),
&ObjectOptions {
versioned: true,
user_defined: HashMap::from([
(rustfs_utils::http::AMZ_OBJECT_TAGGING.to_string(), format!("generation={value}")),
("x-amz-meta-generation".to_string(), value.to_string()),
]),
..Default::default()
},
)
.await
.expect("write a historical version deterministically to pool 1");
expected_versions.push(written.version_id.expect("versioned PUT must acknowledge a UUID"));
}
let marker = store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("delete the current version");
assert!(marker.delete_marker);
let marker_id = marker.version_id.expect("DELETE must acknowledge a marker UUID");
assert!(!expected_versions.contains(&marker_id));
let local = store.pools[1]
.clone()
.inner_list_object_versions(bucket, object, None, None, None, 10)
.await
.expect("list the object-owning pool");
assert_eq!(local.objects.len(), 4, "the marker must be committed beside the three existing versions");
assert_eq!(local.objects[0].version_id, Some(marker_id));
assert!(local.objects[0].delete_marker && local.objects[0].is_latest);
let versions = store
.clone()
.inner_list_object_versions(bucket, object, None, None, None, 10)
.await
.expect("list all pools");
assert_eq!(versions.objects.len(), 4);
assert_eq!(versions.objects.iter().filter(|version| version.is_latest).count(), 1);
for (index, version_id) in expected_versions.iter().copied().enumerate() {
assert!(
versions
.objects
.iter()
.any(|version| version.version_id == Some(version_id) && !version.is_latest)
);
let mut reader = store
.handle_get_object_reader(
bucket,
object,
None,
HeaderMap::new(),
&ObjectOptions {
version_id: Some(version_id.to_string()),
..Default::default()
},
)
.await
.expect("historical version should remain readable");
let mut payload = Vec::new();
reader
.stream
.read_to_end(&mut payload)
.await
.expect("read all historical bytes");
let value = u8::try_from(index + 1).expect("fixture generation fits u8");
assert_eq!(payload, vec![value; 4097]);
let listed = versions
.objects
.iter()
.find(|version| version.version_id == Some(version_id))
.expect("listed historical version");
assert_eq!(listed.user_tags.as_str(), format!("generation={value}"));
assert_eq!(listed.user_defined.get("x-amz-meta-generation"), Some(&value.to_string()));
}
let repeated = store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("repeat simple DELETE");
assert!(repeated.delete_marker);
assert_ne!(
repeated.version_id,
Some(marker_id),
"each enabled-versioning DELETE creates a new marker"
);
store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
version_id: repeated.version_id.map(|id| id.to_string()),
..Default::default()
},
)
.await
.expect("remove the newest marker by identity");
store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
version_id: Some(marker_id.to_string()),
..Default::default()
},
)
.await
.expect("remove the original marker by identity");
let current = store
.get_object_info(bucket, object, &ObjectOptions::default())
.await
.expect("previous version becomes current");
assert_eq!(current.version_id, expected_versions.last().copied());
assert!(!current.delete_marker);
}
#[tokio::test]
async fn multipool_existing_split_history_lists_and_paginates_without_metadata_writes() {
for marker_pool in 0..2 {
let bucket = format!("multipool-split-history-{marker_pool}");
let object = "history.bin";
let (dirs, store) = multipool_version_test_store(&bucket).await;
let mut expected = Vec::new();
for value in 1..=3_u8 {
let version = store.pools[1 - marker_pool]
.put_object(
&bucket,
object,
&mut PutObjReader::from_vec(vec![value; 4097]),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed history through the owning pool's normal write path");
expected.push(version.version_id);
}
let marker = store.pools[marker_pool]
.delete_object(
&bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("reproduce the previously committed split marker with a normal pool DELETE");
assert!(marker.delete_marker);
expected.push(marker.version_id);
expected.reverse();
let mut before = Vec::new();
for dir in &dirs {
before.push(
tokio::fs::read(dir.path().join(&bucket).join(object).join("xl.meta"))
.await
.expect("snapshot persisted version metadata"),
);
}
for max_keys in [1, 2, 4, 10] {
let mut key_marker = None;
let mut version_marker = None;
let mut listed = Vec::new();
let mut completed = false;
for _ in 0..6 {
let page = store
.clone()
.inner_list_object_versions(&bucket, object, key_marker.clone(), version_marker.clone(), None, max_keys)
.await
.expect("read a complete merged version page");
listed.extend(page.objects);
if !page.is_truncated {
completed = true;
break;
}
assert_ne!(
(&page.next_marker, &page.next_version_idmarker),
(&key_marker, &version_marker),
"version cursor must advance"
);
key_marker = page.next_marker;
version_marker = page.next_version_idmarker;
}
assert!(completed, "pagination must terminate");
assert_eq!(listed.iter().map(|version| version.version_id).collect::<Vec<_>>(), expected);
assert!(listed[0].delete_marker && listed[0].is_latest);
assert!(listed[1..].iter().all(|version| !version.delete_marker && !version.is_latest));
}
let visible = store
.clone()
.list_objects_generic(&bucket, "", None, None, 10, false)
.await
.expect("list current objects");
assert!(visible.objects.is_empty(), "the global current marker hides the object");
for (dir, before) in dirs.iter().zip(before) {
assert_eq!(
tokio::fs::read(dir.path().join(&bucket).join(object).join("xl.meta"))
.await
.expect("read unchanged metadata"),
before
);
}
}
}
#[tokio::test]
async fn multipool_conflicting_version_metadata_fails_the_listing_request() {
let bucket = "multipool-version-conflict";
let (_dirs, store) = multipool_version_test_store(bucket).await;
store.pools[0]
.put_object(
bucket,
"a.bin",
&mut PutObjReader::from_vec(b"preceding result".to_vec()),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed an entry before the conflict");
let version_id = Uuid::new_v4();
let mod_time = OffsetDateTime::now_utc();
let mut etags = Vec::new();
for (pool_idx, value) in [(0, 1), (1, 2)] {
let written = store.pools[pool_idx]
.put_object(
bucket,
"z.bin",
&mut PutObjReader::from_vec(vec![value; 4097]),
&ObjectOptions {
versioned: true,
version_id: Some(version_id.to_string()),
mod_time: Some(mod_time),
..Default::default()
},
)
.await
.expect("persist independent conflicting copies");
assert_eq!(written.version_id, Some(version_id));
etags.push(written.etag);
}
assert_ne!(etags[0], etags[1], "fixture must contain a semantic conflict");
let err = store
.clone()
.inner_list_object_versions(bucket, "", None, None, None, 10)
.await
.expect_err("partial output must not hide the merge error");
assert_eq!(err, StorageError::FileCorrupt);
let cancellation = tokio_util::sync::CancellationToken::new();
let (sender, mut receiver) = tokio::sync::mpsc::channel(1);
let walk = store
.clone()
.walk(cancellation.clone(), bucket, "", sender, WalkOptions::default());
let drain = async { while receiver.recv().await.is_some() {} };
let (walk_result, ()) = tokio::time::timeout(Duration::from_secs(5), async { tokio::join!(walk, drain) })
.await
.expect("bounded walk output must drain and terminate on a merge error");
assert_eq!(
walk_result.expect_err("walk must report the same metadata conflict"),
StorageError::FileCorrupt
);
assert!(!cancellation.is_cancelled(), "worker failure must not cancel the caller's request");
}
#[tokio::test]
async fn multipool_marker_rejects_unwritable_owner_without_falling_back() {
let bucket = "multipool-unwritable-owner";
let object = "history.bin";
let (_dirs, store) = multipool_version_test_store(bucket).await;
store.pools[1]
.put_object(
bucket,
object,
&mut PutObjReader::from_vec(vec![1; 4097]),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed the nonzero owner");
*store.pool_meta.write().await = PoolMeta {
pools: vec![prepared_pool_test_status(0, false), prepared_pool_test_status(1, true)],
..Default::default()
};
let err = store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect_err("suspended owner must reject marker creation");
assert_eq!(err, StorageError::SlowDown);
*store.pool_meta.write().await = PoolMeta::default();
let mut rebalancing = crate::services::rebalance::RebalanceStats {
participating: true,
..Default::default()
};
rebalancing.info.status = crate::services::rebalance::RebalStatus::Started;
*store.rebalance_meta.write().await = Some(crate::services::rebalance::RebalanceMeta {
pool_stats: vec![crate::services::rebalance::RebalanceStats::default(), rebalancing],
..Default::default()
});
let error = store
.delete_object(
bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await;
*store.rebalance_meta.write().await = None;
assert_eq!(
error.expect_err("rebalance must not hide the owner during marker lookup"),
StorageError::SlowDown
);
// Isolate object quorum from the bucket metadata preflight: disabling
// a whole pool can otherwise fail before object ownership is looked up.
let quorum_bucket = RUSTFS_META_BUCKET;
store.pools[1]
.put_object(
quorum_bucket,
object,
&mut PutObjReader::from_vec(vec![1; 4097]),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed the object-quorum fixture");
for pool in &store.pools {
pool.put_object(
quorum_bucket,
"split.bin",
&mut PutObjReader::from_vec(vec![2; 4097]),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed another key with history in both pools");
}
let owner = &store.pools[1].disk_set[0];
let healthy = owner.disks.read().await.clone();
for disk in owner.disks.write().await.iter_mut().skip(1) {
*disk = None;
}
let error = store
.delete_object(
quorum_bucket,
object,
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await;
let split_error = store
.delete_object(
quorum_bucket,
"split.bin",
ObjectOptions {
versioned: true,
..Default::default()
},
)
.await;
*owner.disks.write().await = healthy;
assert_eq!(
error.expect_err("subquorum owner must not become a new-pool marker"),
StorageError::ErasureWriteQuorum
);
assert_eq!(
split_error.expect_err("a readable older pool does not prove the global current version"),
StorageError::ErasureWriteQuorum
);
let split = store.pools[0]
.clone()
.inner_list_object_versions(quorum_bucket, "split.bin", None, None, None, 10)
.await
.expect("inspect the readable older pool");
assert_eq!(split.objects.len(), 1);
assert!(!split.objects[0].delete_marker);
let other = store.pools[0]
.clone()
.inner_list_object_versions(quorum_bucket, object, None, None, None, 10)
.await
.expect("inspect the other pool");
assert!(other.objects.is_empty());
let history = store.pools[1]
.clone()
.inner_list_object_versions(quorum_bucket, object, None, None, None, 10)
.await
.expect("inspect the restored owner");
assert_eq!(history.objects.len(), 1);
assert!(!history.objects[0].delete_marker);
}
#[tokio::test]
async fn multipool_suspended_and_batch_markers_keep_the_null_slot_on_the_owner() {
let bucket = "multipool-suspended-markers";
let object = "history.bin";
let (_dirs, store) = multipool_version_test_store(bucket).await;
let original = store.pools[1]
.put_object(
bucket,
object,
&mut PutObjReader::from_vec(vec![1; 4097]),
&ObjectOptions {
versioned: true,
..Default::default()
},
)
.await
.expect("seed the UUID history in pool 1");
store
.update_bucket_metadata_config(
bucket,
crate::bucket::metadata::BUCKET_VERSIONING_CONFIG,
b"<VersioningConfiguration><Status>Suspended</Status></VersioningConfiguration>".to_vec(),
)
.await
.expect("persist suspended bucket versioning");
let null = store
.put_object(
bucket,
object,
&mut PutObjReader::from_vec(vec![2; 4097]),
&ObjectOptions {
version_suspended: true,
..Default::default()
},
)
.await
.expect("write the null version beside the history");
assert!(null.version_id.is_none_or(|id| id.is_nil()));
for _ in 0..2 {
let marker = store
.delete_object(
bucket,
object,
ObjectOptions {
version_suspended: true,
..Default::default()
},
)
.await
.expect("replace the null slot with a marker");
assert!(marker.delete_marker);
assert!(marker.version_id.is_none_or(|id| id.is_nil()));
}
let (deleted, errors) = store
.delete_objects(
bucket,
vec![ObjectToDelete {
object_name: object.to_string(),
..Default::default()
}],
ObjectOptions::default(),
)
.await;
assert!(errors.iter().all(Option::is_none), "batch DELETE should succeed: {errors:?}");
assert_eq!(deleted.len(), 1);
assert!(deleted[0].delete_marker);
let versions = store.pools[1]
.clone()
.inner_list_object_versions(bucket, object, None, None, None, 10)
.await
.expect("inspect the owner after batch DELETE");
assert_eq!(versions.objects.len(), 2, "only one null marker and the UUID history remain");
assert!(versions.objects[0].delete_marker && versions.objects[0].is_latest);
assert_eq!(versions.objects[1].version_id, original.version_id);
let other = store.pools[0]
.clone()
.inner_list_object_versions(bucket, object, None, None, None, 10)
.await
.expect("inspect the unused pool");
assert!(other.objects.is_empty());
}
async fn assert_prepared_reader_blocks_writer(store: &ECStore, bucket: &str, object: &str) {
assert_pool_writer_is_blocked(store, 0, bucket, object).await;
}
+50 -1
View File
@@ -611,7 +611,18 @@ impl ECStore {
object: &str,
opts: &ObjectOptions,
) -> Result<(PoolObjInfo, Vec<PoolErr>)> {
self.internal_get_pool_info_existing_with_opts(bucket, object, opts).await
self.internal_get_pool_info_existing_with_opts(bucket, object, opts, false)
.await
}
pub(super) async fn get_pool_info_for_delete_marker(
&self,
bucket: &str,
object: &str,
opts: &ObjectOptions,
) -> Result<(PoolObjInfo, Vec<PoolErr>)> {
self.internal_get_pool_info_existing_with_opts(bucket, object, opts, true)
.await
}
async fn internal_get_pool_info_existing_with_opts(
@@ -619,6 +630,7 @@ impl ECStore {
bucket: &str,
object: &str,
opts: &ObjectOptions,
require_all_pool_reads: bool,
) -> Result<(PoolObjInfo, Vec<PoolErr>)> {
let mut futures = Vec::new();
for pool in self.pools.iter() {
@@ -647,6 +659,15 @@ impl ECStore {
});
}
Err(e) => {
// A readable older pool cannot prove ownership of the
// current version while another pool is unreadable. Check
// both raw and object-scoped quorum errors before sorting.
if require_all_pool_reads && !is_err_object_not_found(&e) && !is_err_version_not_found(&e) {
return Err(match e {
Error::ErasureReadQuorum | Error::InsufficientReadQuorum(_, _) => Error::ErasureWriteQuorum,
err => err,
});
}
ress.push(PoolObjInfo {
index,
err: Some(e),
@@ -656,6 +677,34 @@ impl ECStore {
}
}
if require_all_pool_reads {
let suspended_pools = {
let pool_meta = self.pool_meta.read().await;
(0..self.pools.len())
.map(|idx| pool_meta.is_suspended(idx))
.collect::<Vec<_>>()
};
let candidates = ress
.iter()
.map(|pinfo| LatestObjectInfoCandidate {
info: pinfo.err.is_none().then(|| pinfo.object_info.clone()),
idx: pinfo.index,
err: pinfo.err.clone(),
})
.collect();
let (object_info, index) =
resolve_latest_object_info_candidates_with_pool_state(candidates, &suspended_pools, bucket, object, opts)?;
let pools_with_object = self.pools_with_object(&ress, opts).await;
return Ok((
PoolObjInfo {
index,
object_info,
err: None,
},
pools_with_object,
));
}
ress.sort_by(|a, b| {
let at = a.object_info.mod_time.unwrap_or(OffsetDateTime::UNIX_EPOCH);
let bt = b.object_info.mod_time.unwrap_or(OffsetDateTime::UNIX_EPOCH);
+39 -21
View File
@@ -520,32 +520,50 @@ impl RootHealRecovery {
let _guard = self.mutation.lock().await;
let disks = self.disks().await?;
let existing = Self::find(&disks, &request.id).await?;
let (disk, expected) = match existing {
Some((disk, bytes)) => (disk, Some(bytes)),
None => {
let disk = disks
.first()
.cloned()
.ok_or_else(|| Error::Other("No local disk available for root heal shutdown recovery".to_string()))?;
(disk, None)
}
};
if request.options.no_lock {
return Err(Error::Other("Administrator root heal cannot skip namespace locking".to_string()));
}
let bytes = serde_json::to_vec(&RootHealIntent::from_request(request))
.map_err(|error| Error::Other(format!("Serialize root heal recovery record: {error}")))?;
match EcstoreDiskAPI::compare_and_update_file(
disk.as_ref(),
RUSTFS_META_BUCKET,
&intent_path(&request.id)?,
expected,
Some(bytes.into()),
)
.await?
{
EcstoreConditionalFileUpdate::Updated => Ok(()),
_ => Err(Error::Other(format!("Root heal recovery record changed for {}", request.id))),
let path = intent_path(&request.id)?;
if let Some((disk, expected)) = existing {
return match EcstoreDiskAPI::compare_and_update_file(
disk.as_ref(),
RUSTFS_META_BUCKET,
&path,
Some(expected),
Some(bytes.into()),
)
.await?
{
EcstoreConditionalFileUpdate::Updated => Ok(()),
_ => Err(Error::Other(format!("Root heal recovery record changed for {}", request.id))),
};
}
if disks.is_empty() {
return Err(Error::Other("No local disk available for root heal shutdown recovery".to_string()));
}
let mut last_not_committed = None;
for disk in &disks {
match EcstoreDiskAPI::compare_and_update_file(
disk.as_ref(),
RUSTFS_META_BUCKET,
&path,
None,
Some(bytes.clone().into()),
)
.await
{
Ok(EcstoreConditionalFileUpdate::Updated) => return Ok(()),
Ok(_) => return Err(Error::Other(format!("Root heal recovery record changed for {}", request.id))),
Err(error) if error.is_conditional_file_not_committed() => last_not_committed = Some(error),
Err(error) => return Err(Error::Disk(error)),
}
}
match last_not_committed {
Some(error) => Err(Error::Disk(error)),
None => Err(Error::Other("No local disk accepted the root heal recovery record".to_string())),
}
}
@@ -17,6 +17,44 @@ use super::*;
use crate::heal::RUSTFS_META_BUCKET;
use std::collections::HashSet;
#[cfg(unix)]
struct RestoreDirectoryMode {
path: std::path::PathBuf,
mode: u32,
}
#[cfg(unix)]
impl RestoreDirectoryMode {
fn read_only(path: std::path::PathBuf) -> Self {
use std::os::unix::fs::PermissionsExt as _;
let mode = std::fs::metadata(&path)
.expect("metadata directory mode")
.permissions()
.mode();
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o555)).expect("make metadata directory read-only");
Self { path, mode }
}
}
#[cfg(unix)]
impl Drop for RestoreDirectoryMode {
fn drop(&mut self) {
use std::os::unix::fs::PermissionsExt as _;
let _ = std::fs::set_permissions(&self.path, std::fs::Permissions::from_mode(self.mode));
}
}
#[cfg(unix)]
fn ordered_recovery_disks(first: DiskStore, second: DiskStore) -> (DiskStore, DiskStore) {
if first.endpoint().to_string() <= second.endpoint().to_string() {
(first, second)
} else {
(second, first)
}
}
async fn recovery_disk() -> (TempDir, DiskStore) {
let temp = TempDir::new().expect("temporary root recovery disk");
let endpoint = Endpoint::try_from(temp.path().to_string_lossy().as_ref()).expect("disk endpoint");
@@ -78,6 +116,94 @@ fn completed_admin_status(heal_type: &HealType, completed_at: SystemTime) -> Com
}
}
#[cfg(unix)]
#[tokio::test]
async fn root_recovery_new_intent_skips_prepublication_read_only_owner() {
let (first_temp, first_disk) = recovery_disk().await;
let (second_temp, second_disk) = recovery_disk().await;
let first_endpoint = first_disk.endpoint().to_string();
let (read_only_disk, writable_disk) = ordered_recovery_disks(first_disk, second_disk);
let read_only_root = if read_only_disk.endpoint().to_string() == first_endpoint {
first_temp.path()
} else {
second_temp.path()
};
let _restore = RestoreDirectoryMode::read_only(read_only_root.join(RUSTFS_META_BUCKET));
let manager = recovery_manager(vec![read_only_disk.clone(), writable_disk.clone()]);
let mut request = admin_request(HealType::Object {
bucket: "bucket".to_string(),
object: "object".to_string(),
version_id: None,
});
let receipt = manager
.submit_heal_request_with_receipt(request.clone())
.await
.expect("a writable local disk should own the admin heal intent");
assert_eq!(receipt.result, HealAdmissionResult::Accepted);
let path = format!("root-heal-{}.json", request.id);
assert!(matches!(
read_only_disk.read_all(RUSTFS_META_BUCKET, &path).await,
Err(DiskError::FileNotFound)
));
assert!(writable_disk.read_all(RUSTFS_META_BUCKET, &path).await.is_ok());
request.retry_attempts = 1;
manager
.root_recovery
.persist(&request)
.await
.expect("an existing fallback owner should remain updateable");
let pending = manager.root_recovery.pending().await.expect("read the single durable owner");
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].id, request.id);
assert_eq!(pending[0].retry_attempts, 1);
}
#[cfg(unix)]
#[tokio::test]
async fn root_recovery_existing_owner_never_migrates_after_write_rejection() {
let (first_temp, first_disk) = recovery_disk().await;
let (second_temp, second_disk) = recovery_disk().await;
let first_endpoint = first_disk.endpoint().to_string();
let (owner_disk, alternate_disk) = ordered_recovery_disks(first_disk, second_disk);
let owner_root = if owner_disk.endpoint().to_string() == first_endpoint {
first_temp.path()
} else {
second_temp.path()
};
let manager = recovery_manager(vec![owner_disk.clone(), alternate_disk.clone()]);
let mut request = root_request();
manager
.root_recovery
.persist(&request)
.await
.expect("create the canonical owner");
let path = format!("root-heal-{}.json", request.id);
let committed = owner_disk
.read_all(RUSTFS_META_BUCKET, &path)
.await
.expect("canonical owner bytes");
let _restore = RestoreDirectoryMode::read_only(owner_root.join(RUSTFS_META_BUCKET));
request.retry_attempts = 1;
assert!(
manager.root_recovery.persist(&request).await.is_err(),
"an existing owner write rejection must fail closed"
);
assert_eq!(
owner_disk
.read_all(RUSTFS_META_BUCKET, &path)
.await
.expect("original owner remains"),
committed
);
assert!(matches!(
alternate_disk.read_all(RUSTFS_META_BUCKET, &path).await,
Err(DiskError::FileNotFound)
));
}
async fn active_root(manager: &HealManager, request: HealRequest) -> Arc<HealTask> {
let task = Arc::new(HealTask::from_request(request, manager.storage.clone()));
*task.status.write().await = HealTaskStatus::Running;