Compare commits

...

9 Commits

Author SHA1 Message Date
hector 8552c9f8aa feat(nightly): publish packages as assets of the rolling 'nightly' release (#7592)
Replace the assets-branch scheme with a proper GitHub Release on
rustfs/auto-testing: a single 'nightly' release whose deb/rpm assets
are replaced in place on every build. This is the standard channel —
visible on the repo's Releases page, stable download URLs, no git
history growth (release assets live outside the repository).

- New scripts/release/publish_nightly_assets.sh: resolves-or-creates
  the 'nightly' release via the REST API, deletes same-name assets,
  uploads rustfs-nightly-latest.{deb,rpm}, then PATCHes the release
  body with the build provenance (ref@sha, run link, sizes, SHA256).
  Plain curl + python3, no gh CLI (the build fleet has none — #7586).
- The workflow step shrinks to invoking the script; full flow
  exercised end-to-end against the real release with probe files
  (create / upload / overwrite / download round-trip / body update).
2026-09-09 21:33:58 +08:00
唐小鸭 27d66c159f fix(kms): prevent transient health failures from latching status (#7578)
Keep backend health checks from overwriting the running service lifecycle state so subsequent admin checks and probe-based readiness can recover without a restart.

Add a regression test that fails on the original implementation after backend recovery and verifies the service instance and version remain unchanged.

Validation: 39 focused resilience, lifecycle, concurrency, and service manager tests passed; one existing live AWS test remained ignored. cargo fmt --all --check and git diff --check passed.

Thanks to @stevapple for reporting the issue and providing a detailed diagnosis and reproduction.

Fixes #7554
2026-09-09 11:05:26 +00:00
hector e546ae9c62 fix(nightly): push assets with plain git — the build fleet has no gh CLI (#7572)
The sm-standard-4 runners used by the nightly build have no gh binary
(functional-chain workflows run elsewhere, on the jumpbox). The assets
publish step died on 'gh: command not found' at the credential-helper
setup, and the preceding clone failure had been masked by 2>/dev/null,
misleading the step into the orphan path. Swap clone and remote setup
to plain git with the token embedded in the URL; push semantics are
unchanged.
2026-09-09 16:39:59 +08:00
JaySon 35b5cfcf8d docs: fix invalid docker-buildx.sh usage example in README and README_ZH (#7570)
docs: fix invalid docker-buildx.sh usage in README and README_ZH

Replace the docker-buildx.sh --build-arg RELEASE=latest example, which
the script never accepted as a CLI flag (--build-arg is only used
internally for docker buildx build), with the supported invocations:

- bare ./docker-buildx.sh for the default local build
- ./docker-buildx.sh -p linux/amd64 for a single-platform local build

Also update the surrounding comments to reflect single-platform local
builds and add the multi-arch example comment accordingly.
2026-09-09 16:39:25 +08:00
cxymds d7c2fc7587 fix(ecstore): apply read quorum to bucket validation (#7556)
* fix(ecstore): apply read quorum to bucket validation

* fix(e2e): remove needless borrow in quorum test
2026-09-09 14:50:28 +08:00
Zhengchao An 7db206466b test(ci): reserve runner capacity for bounded rollback probe (#7560) 2026-09-09 14:24:57 +08:00
hector e549252ac6 feat(nightly): make the build branch configurable (NIGHTLY_BRANCH var + dispatch input) (#7557) 2026-09-09 14:24:44 +08:00
Zhengchao An 17807b05fb fix(ecstore): preserve internode error context when cloning (#7548) 2026-09-09 14:24:23 +08:00
cxymds 789f1832a4 fix(rebalance): align activation locks and preserve retryable causes (#7551) 2026-09-09 03:25:05 +00:00
21 changed files with 1398 additions and 124 deletions
+6
View File
@@ -197,6 +197,12 @@ test-group = 'e2e-cluster-nightly'
filter = 'package(e2e_test) & (test(/^kms::kms_vault_test::/) | test(/^kms::kms_rekey_sweep_test::/) | test(/^kms::configured_roundtrip_test::test_configured_vault_kms_admin_and_versioned_cleanup$/))' filter = 'package(e2e_test) & (test(/^kms::kms_vault_test::/) | test(/^kms::kms_rekey_sweep_test::/) | test(/^kms::configured_roundtrip_test::test_configured_vault_kms_admin_and_versioned_cleanup$/))'
test-group = 'e2e-vault' test-group = 'e2e-vault'
# This four-disk, 65-member rollback probe already drives up to 32 concurrent
# durable deletions. Reserve this nextest run's capacity for its progress oracle.
[[profile.default.overrides]]
filter = 'package(rustfs-ecstore) & test(=store::init::tests::dispatch_manifest_rollback_bounded_concurrency_reaches_tail_behind_slow_member)'
threads-required = "num-test-threads"
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# ci profile — the strict CI gate (ci.yml `cargo nextest run --profile ci`) # ci profile — the strict CI gate (ci.yml `cargo nextest run --profile ci`)
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
+122 -5
View File
@@ -19,17 +19,27 @@ on:
- cron: "7 0 * * *" - cron: "7 0 * * *"
timezone: "Asia/Shanghai" timezone: "Asia/Shanghai"
workflow_dispatch: workflow_dispatch:
inputs:
branch:
description: 'Branch/ref to build and publish as the nightly (empty = scheduled source, see NIGHTLY_BUILD_REF)'
required: false
default: ''
permissions: permissions:
contents: read contents: read
# Scheduled builds follow the NIGHTLY_BRANCH repo variable so the channel can
# be pointed at e.g. `release` for the GA cycle and back to `main` afterwards
# without touching this file. Manual runs take the `branch` input, falling
# back to the branch the run was dispatched from.
concurrency: concurrency:
group: nightly-gnu-build-main-${{ github.event_name }} group: nightly-gnu-build-${{ github.event_name }}-${{ github.event_name == 'schedule' && (vars.NIGHTLY_BRANCH || 'main') || (inputs.branch || github.ref_name) }}
cancel-in-progress: ${{ github.event_name == 'workflow_dispatch' }} cancel-in-progress: ${{ github.event_name == 'workflow_dispatch' }}
env: env:
CARGO_TERM_COLOR: always CARGO_TERM_COLOR: always
RUST_BACKTRACE: 1 RUST_BACKTRACE: 1
NIGHTLY_BUILD_REF: ${{ github.event_name == 'schedule' && (vars.NIGHTLY_BRANCH || 'main') || (inputs.branch || github.ref_name) }}
jobs: jobs:
build: build:
@@ -43,6 +53,7 @@ jobs:
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7 uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with: with:
persist-credentials: false persist-credentials: false
ref: ${{ env.NIGHTLY_BUILD_REF }}
- name: Setup Rust environment - name: Setup Rust environment
uses: ./.github/actions/setup uses: ./.github/actions/setup
@@ -152,13 +163,104 @@ jobs:
fakeroot dpkg-deb --build "${PKG_DIR}" fakeroot dpkg-deb --build "${PKG_DIR}"
ls -lh "${DEB_FILE}" ls -lh "${DEB_FILE}"
echo "deb_date=${DEB_DATE}" >> "${GITHUB_OUTPUT}"
echo "deb_file=${DEB_FILE}" >> "${GITHUB_OUTPUT}" echo "deb_file=${DEB_FILE}" >> "${GITHUB_OUTPUT}"
# Same packaging scheme as .github/workflows/package.yml (fpm), but from
# the locally built nightly binary instead of a release artifact, with a
# date-based version that mirrors the DEB.
- name: Build RPM package
id: rpm
shell: bash
env:
DEB_DATE: ${{ steps.deb.outputs.deb_date }}
run: |
set -euo pipefail
if ! command -v fpm >/dev/null 2>&1; then
SUDO=""; [ "$(id -u)" -ne 0 ] && SUDO="sudo -n"
${SUDO} apt-get update -qq && ${SUDO} apt-get install -y -qq ruby ruby-dev build-essential rpm >/dev/null
${SUDO} gem install fpm --no-document >/dev/null
fi
RPM_FILE="rustfs-nightly-${DEB_DATE}.rpm"
RPM_VERSION="0"
RPM_RELEASE="0.nightly.${DEB_DATE//-/.}"
echo "Building RPM: ${RPM_FILE} (version ${RPM_VERSION}-${RPM_RELEASE})"
# fpm wants the config file to exist before packaging.
mkdir -p ./tmp-pkg/etc/default
cat > ./tmp-pkg/etc/default/rustfs << 'ENVEOF'
# RustFS Environment Configuration
# See https://rustfs.com/docs/ for more information
# RUSTFS_VOLUMES=""
# RUSTFS_ROOT_USER=""
# RUSTFS_ROOT_PASSWORD=""
ENVEOF
fpm -s dir -t rpm \
--name rustfs \
--version "$RPM_VERSION" \
--iteration "$RPM_RELEASE" \
--architecture x86_64 \
--package "$RPM_FILE" \
--depends "glibc >= 2.31" \
--maintainer "RustFS Team <support@rustfs.com>" \
--description "High-performance distributed object storage" \
--url "https://rustfs.com" \
--license "Apache-2.0" \
--after-install <(cat << 'POSTINST'
#!/bin/bash
set -e
if ! getent passwd rustfs > /dev/null 2>&1; then
useradd -r -s /bin/false -d /opt/rustfs rustfs
fi
mkdir -p /opt/rustfs /data/rustfs /var/log/rustfs
chown rustfs:rustfs /opt/rustfs /data/rustfs /var/log/rustfs
if [ -d /run/systemd/system ]; then
systemctl daemon-reload
fi
POSTINST
) \
--before-remove <(cat << 'PRERM'
#!/bin/bash
set -e
if [ -d /run/systemd/system ] && systemctl is-active --quiet rustfs; then
systemctl stop rustfs
fi
PRERM
) \
--after-remove <(cat << 'POSTRM'
#!/bin/bash
set -e
if [ -d /run/systemd/system ]; then
systemctl daemon-reload
fi
POSTRM
) \
--config-files /etc/default/rustfs \
"rustfs-nightly-${DEB_DATE}/usr/bin/rustfs=/usr/bin/rustfs" \
./tmp-pkg/etc/default/rustfs=/etc/default/rustfs \
deploy/build/rustfs.service=/lib/systemd/system/rustfs.service \
LICENSE=/usr/share/doc/rustfs/LICENSE \
README.md=/usr/share/doc/rustfs/README.md
[[ -f "$RPM_FILE" ]] || { echo "RPM build failed"; exit 1; }
rpm -qpl "$RPM_FILE" | grep -Fx '/usr/bin/rustfs' >/dev/null
stat --printf='%n %s bytes\n' "$RPM_FILE"
echo "rpm_file=$RPM_FILE" >> "$GITHUB_OUTPUT"
- name: Upload DEB artifact - name: Upload DEB artifact
uses: actions/upload-artifact@b7c566a772e6b6bfb58ed0dc250532a479d7789f # v6 uses: actions/upload-artifact@b7c566a772e6b6bfb58ed0dc250532a479d7789f # v6
with: with:
name: ${{ steps.deb.outputs.deb_file }} name: ${{ steps.deb.outputs.deb_file }}
path: ${{ steps.deb.outputs.deb_file }} path: ${{ steps.deb.outputs.deb_file }}
- name: Upload RPM artifact
uses: actions/upload-artifact@b7c566a772e6b6bfb58ed0dc250532a479d7789f # v6
with:
name: ${{ steps.rpm.outputs.rpm_file }}
path: ${{ steps.rpm.outputs.rpm_file }}
if-no-files-found: error if-no-files-found: error
# Persist the nightly deb on Cloudflare R2 (same channel as package.yml) # Persist the nightly deb on Cloudflare R2 (same channel as package.yml)
@@ -187,11 +289,10 @@ jobs:
export AWS_SECRET_ACCESS_KEY="$R2_SECRET_ACCESS_KEY" export AWS_SECRET_ACCESS_KEY="$R2_SECRET_ACCESS_KEY"
export AWS_DEFAULT_REGION="auto" export AWS_DEFAULT_REGION="auto"
# The candidate manifest must describe the tree that was actually
# built. With a ref override (NIGHTLY_BRANCH / dispatch input) that
# is not necessarily GITHUB_SHA, so always advertise HEAD.
SOURCE_SHA="$(git rev-parse HEAD)" SOURCE_SHA="$(git rev-parse HEAD)"
if [[ "${SOURCE_SHA}" != "${GITHUB_SHA}" ]]; then
echo "Checkout SHA does not match the nightly build run" >&2
exit 1
fi
DEB_SHA256="$(sha256sum "${DEB_FILE}" | cut -d ' ' -f 1)" DEB_SHA256="$(sha256sum "${DEB_FILE}" | cut -d ' ' -f 1)"
CANDIDATE_KEY="artifacts/rustfs/packages/nightly/runs/${GITHUB_RUN_ID}/${GITHUB_RUN_ATTEMPT}/${DEB_SHA256}/rustfs.deb" CANDIDATE_KEY="artifacts/rustfs/packages/nightly/runs/${GITHUB_RUN_ID}/${GITHUB_RUN_ATTEMPT}/${DEB_SHA256}/rustfs.deb"
CANDIDATE_URL="https://dl.rustfs.com/${CANDIDATE_KEY}" CANDIDATE_URL="https://dl.rustfs.com/${CANDIDATE_KEY}"
@@ -247,6 +348,20 @@ jobs:
path: ${{ steps.publish.outputs.candidate_file }} path: ${{ steps.publish.outputs.candidate_file }}
if-no-files-found: error if-no-files-found: error
# Publish the deb/rpm pair to the auto-testing repo's `assets` branch so
# engineers can download and install the nightly directly. The branch is
# a single-commit orphan rewritten on every build, which keeps the repo
# small while the latest files stay reachable at stable raw URLs.
# Publish the deb/rpm pair as assets of the rolling `nightly` release on
# rustfs/auto-testing (see scripts/release/publish_nightly_assets.sh).
- name: Publish packages to auto-testing release assets
env:
ASSETS_TOKEN: ${{ secrets.PF_TESTING_GH_TOKEN }}
DEB_FILE: ${{ steps.deb.outputs.deb_file }}
RPM_FILE: ${{ steps.rpm.outputs.rpm_file }}
DEB_DATE: ${{ steps.deb.outputs.deb_date }}
BUILD_REF: ${{ env.NIGHTLY_BUILD_REF }}
run: bash scripts/release/publish_nightly_assets.sh
# Live-Vault lane for the rustfs-kms suite (rustfs/backlog#1774). # Live-Vault lane for the rustfs-kms suite (rustfs/backlog#1774).
# #
# RUSTFS_KMS_VAULT_TOKEN is the single switch that adds the Vault KV2 and # RUSTFS_KMS_VAULT_TOKEN is the single switch that adds the Vault KV2 and
@@ -284,6 +399,7 @@ jobs:
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7 uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with: with:
persist-credentials: false persist-credentials: false
ref: ${{ env.NIGHTLY_BUILD_REF }}
- name: Setup Rust environment - name: Setup Rust environment
uses: ./.github/actions/setup uses: ./.github/actions/setup
@@ -372,6 +488,7 @@ jobs:
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7 uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with: with:
persist-credentials: false persist-credentials: false
ref: ${{ env.NIGHTLY_BUILD_REF }}
- name: Setup Rust environment - name: Setup Rust environment
uses: ./.github/actions/setup uses: ./.github/actions/setup
+4 -1
View File
@@ -211,7 +211,10 @@ For developers who want to build RustFS Docker images from source with multi-arc
```bash ```bash
# Build multi-architecture images locally # Build multi-architecture images locally
./docker-buildx.sh --build-arg RELEASE=latest ./docker-buildx.sh
# Build a single-platform image locally
./docker-buildx.sh -p linux/amd64
# Build and push to registry # Build and push to registry
./docker-buildx.sh --push ./docker-buildx.sh --push
+4 -1
View File
@@ -150,7 +150,10 @@ docker compose -f docker-compose-simple.yml up -d
```bash ```bash
# 在本地构建多架构镜像 # 在本地构建多架构镜像
./docker-buildx.sh --build-arg RELEASE=latest ./docker-buildx.sh
# 在本地构建单平台镜像
./docker-buildx.sh -p linux/amd64
# 构建并推送到仓库 # 构建并推送到仓库
./docker-buildx.sh --push ./docker-buildx.sh --push
@@ -1047,6 +1047,8 @@ mod tests {
let mut cluster = RustFSTestClusterEnvironment::new(4).await?; let mut cluster = RustFSTestClusterEnvironment::new(4).await?;
cluster.set_env("RUSTFS_UNSAFE_BYPASS_DISK_CHECK", "true"); cluster.set_env("RUSTFS_UNSAFE_BYPASS_DISK_CHECK", "true");
cluster.set_env("RUSTFS_HEAL_ENABLED", "true"); cluster.set_env("RUSTFS_HEAL_ENABLED", "true");
// Capture physical baselines after the PUT rename fanout has drained.
cluster.set_env("RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE", "false");
// Heal control uses the first lexicographically sorted grid host. // Heal control uses the first lexicographically sorted grid host.
// Keep that coordinator distinct from the remote target at index 1. // Keep that coordinator distinct from the remote target at index 1.
cluster.nodes.sort_by(|left, right| left.url.cmp(&right.url)); cluster.nodes.sort_by(|left, right| left.url.cmp(&right.url));
@@ -25,6 +25,190 @@ const KEY: &str = "thumb/79/concurrent-overwrite.jpg";
type TestResult = Result<(), Box<dyn std::error::Error + Send + Sync>>; type TestResult = Result<(), Box<dyn std::error::Error + Send + Sync>>;
async fn assert_quorum_object_body(client: &Client, bucket: &str, key: &str, expected: &[u8]) -> TestResult {
let body = client
.get_object()
.bucket(bucket)
.key(key)
.send()
.await?
.body
.collect()
.await?
.into_bytes();
assert_eq!(body.as_ref(), expected, "quorum read returned incorrect contents for {key}");
Ok(())
}
async fn wait_for_quorum_read_admission(clients: &[Client], bucket: &str) -> TestResult {
// SIGKILL can orphan a granted lease. Wait for shared metadata-lock
// admission before asserting the stable quorum boundary; cold bodies
// remain unread throughout this readiness probe.
let deadline =
tokio::time::Instant::now() + rustfs_lock::fast_lock::DEFAULT_LOCK_TIMEOUT + std::time::Duration::from_secs(15);
loop {
let mut ready = true;
for client in clients {
for key in ["warm-small", "warm-large"] {
match client.head_object().bucket(bucket).key(key).send().await {
Ok(_) => {}
Err(error) if error.raw_response().is_some_and(|response| response.status().as_u16() == 503) => {
ready = false;
break;
}
Err(error) => return Err(error.into()),
}
}
if !ready {
break;
}
}
if ready {
return Ok(());
}
if tokio::time::Instant::now() >= deadline {
return Err(format!("read quorum did not become available after lease convergence for {bucket}").into());
}
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
}
#[tokio::test]
async fn test_degraded_cluster_read_quorum_follows_erasure_layout() -> TestResult {
crate::common::init_logging();
for (node_count, parity) in [(4, 2), (6, 3), (6, 2)] {
let read_quorum = node_count - parity;
let write_quorum = read_quorum + usize::from(read_quorum == parity);
let mut cluster = RustFSTestClusterEnvironment::new(node_count).await?;
cluster.set_env("RUSTFS_STORAGE_CLASS_STANDARD", format!("EC:{parity}"));
// Wait for every seed fanout before removing any physical shard.
cluster.set_env("RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE", "false");
cluster.set_env("RUSTFS_OBS_METRICS_EXPORT_ENABLED", "false");
cluster.set_env("RUST_LOG", "warn,rustfs_lock=debug");
cluster.start().await?;
let clients = cluster
.create_all_clients()?
.into_iter()
.map(|client| {
Client::from_conf(
client
.config()
.to_builder()
.retry_config(aws_sdk_s3::config::retry::RetryConfig::standard().with_max_attempts(1))
.build(),
)
})
.collect::<Vec<_>>();
let bucket = format!("read-quorum-{node_count}-{parity}");
clients[0].create_bucket().bucket(&bucket).send().await?;
let small = b"read quorum is derived from the erasure layout".to_vec();
let large = (0..1_048_576)
.map(|index| u8::try_from(index % 251).expect("bounded payload byte"))
.collect::<Vec<_>>();
for (key, body) in [
("warm-small", &small),
("warm-large", &large),
("cold-small", &small),
("cold-large", &large),
("below-quorum", &large),
] {
clients[node_count - 1]
.put_object()
.bucket(&bucket)
.key(key)
.body(Bytes::copy_from_slice(body).into())
.send()
.await?;
}
for node in &cluster.nodes {
for key in ["warm-small", "warm-large", "cold-small", "cold-large", "below-quorum"] {
let census =
crate::chaos::census_object_version_on_disk(std::path::Path::new(&node.data_dir), &bucket, key, None)?;
assert!(census.is_complete(), "seed shard must be complete before fault injection: {census:?}");
assert_eq!(census.data_blocks, Some(read_quorum));
assert_eq!(census.parity_blocks, Some(parity));
}
}
for client in &clients {
assert_quorum_object_body(client, &bucket, "warm-small", &small).await?;
assert_quorum_object_body(client, &bucket, "warm-large", &large).await?;
}
for offline_node in (read_quorum..node_count).rev() {
cluster.stop_node(offline_node)?;
wait_for_quorum_read_admission(&clients[..offline_node], &bucket).await?;
for client in clients.iter().take(offline_node) {
client.head_bucket().bucket(&bucket).send().await?;
assert_quorum_object_body(client, &bucket, "warm-large", &large).await?;
}
}
// Exercise more than the five-second positive bucket-validation TTL.
// Every sample must succeed; polling must not hide a transient failure.
let validation_deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(6);
loop {
for client in clients.iter().take(read_quorum) {
assert_quorum_object_body(client, &bucket, "warm-small", &small).await?;
assert_quorum_object_body(client, &bucket, "warm-large", &large).await?;
let listing = client.list_objects_v2().bucket(&bucket).send().await?;
for key in ["warm-small", "warm-large", "cold-small", "cold-large", "below-quorum"] {
assert!(listing.contents().iter().any(|entry| entry.key() == Some(key)), "listing omitted {key}");
}
}
if tokio::time::Instant::now() >= validation_deadline {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(250)).await;
}
for client in clients.iter().take(read_quorum) {
assert_quorum_object_body(client, &bucket, "cold-small", &small).await?;
assert_quorum_object_body(client, &bucket, "cold-large", &large).await?;
}
let write = clients[0]
.put_object()
.bucket(&bucket)
.key("quorum-write")
.body(Bytes::copy_from_slice(&small).into())
.send()
.await;
if read_quorum >= write_quorum {
write?;
} else {
let error = write.expect_err("a read quorum must not authorize a write that needs more votes");
assert_eq!(error.as_service_error().and_then(|error| error.meta().code()), Some("ServiceUnavailable"));
}
cluster.stop_node(read_quorum - 1)?;
for client in clients.iter().take(read_quorum - 1) {
match client.get_object().bucket(&bucket).key("below-quorum").send().await {
Ok(response) => assert!(
response.body.collect().await.is_err(),
"fewer than {read_quorum} valid fragments must not reconstruct an uncached object"
),
Err(error) => assert_eq!(
error.as_service_error().and_then(|error| error.meta().code()),
Some("ServiceUnavailable"),
"a quorum loss must not be mistaken for a missing object"
),
}
}
for node in 0..read_quorum - 1 {
cluster.stop_node(node)?;
}
cluster.start().await?;
for client in &clients {
assert_quorum_object_body(client, &bucket, "warm-large", &large).await?;
assert_quorum_object_body(client, &bucket, "below-quorum", &large).await?;
}
}
Ok(())
}
async fn put_object(client: Client, payload: Vec<u8>, writer_id: usize) -> Result<(), String> { async fn put_object(client: Client, payload: Vec<u8>, writer_id: usize) -> Result<(), String> {
client client
.put_object() .put_object()
+57 -8
View File
@@ -3463,6 +3463,11 @@ impl PoolRebalanceActivationFence {
} }
} }
#[cfg(test)]
tokio::task_local! {
pub(crate) static REBALANCE_ACTIVATION_LOCK_ATTEMPT: Arc<tokio::sync::Notify>;
}
pub(crate) async fn acquire_pool_rebalance_activation_locks<S>( pub(crate) async fn acquire_pool_rebalance_activation_locks<S>(
pool: Arc<S>, pool: Arc<S>,
fleet_proof: Option<crate::services::notification_sys::CrossPoolFenceFleetProofToken>, fleet_proof: Option<crate::services::notification_sys::CrossPoolFenceFleetProofToken>,
@@ -3473,17 +3478,21 @@ where
NamespaceLock = rustfs_lock::NamespaceLockWrapper, NamespaceLock = rustfs_lock::NamespaceLockWrapper,
>, >,
{ {
// Activation lock order is always pool.bin -> rebalance.bin. // Match entry admission: rebalance.bin -> pool.bin. An entry retains its
// run read fence while target mutations acquire the pool metadata fence;
// activation must not hold pool.bin while waiting for that entry to drain.
let rebalance_meta_lock = pool.new_ns_lock(RUSTFS_META_BUCKET, REBAL_META_NAME).await?;
#[cfg(test)]
let _ = REBALANCE_ACTIVATION_LOCK_ATTEMPT.try_with(|attempted| attempted.notify_one());
let rebalance_meta_guard = rebalance_meta_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.map_err(activation_rebalance_meta_lock_error)?;
let pool_meta_lock = pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?; let pool_meta_lock = pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?;
let pool_meta_guard = pool_meta_lock let pool_meta_guard = pool_meta_lock
.get_write_lock(get_lock_acquire_timeout()) .get_write_lock(get_lock_acquire_timeout())
.await .await
.map_err(activation_pool_meta_lock_error)?; .map_err(activation_pool_meta_lock_error)?;
let rebalance_meta_lock = pool.new_ns_lock(RUSTFS_META_BUCKET, REBAL_META_NAME).await?;
let rebalance_meta_guard = rebalance_meta_lock
.get_write_lock(get_lock_acquire_timeout())
.await
.map_err(activation_rebalance_meta_lock_error)?;
Ok(PoolRebalanceActivationFence { Ok(PoolRebalanceActivationFence {
pool_meta_guard, pool_meta_guard,
@@ -22094,7 +22103,7 @@ mod pools_tests {
.resources .resources
.lock() .lock()
.expect("activation lock recorder should not be poisoned"), .expect("activation lock recorder should not be poisoned"),
vec![POOL_META_NAME.to_string(), REBAL_META_NAME.to_string()] vec![REBAL_META_NAME.to_string(), POOL_META_NAME.to_string()]
); );
let mut second_acquire = Box::pin(acquire_pool_rebalance_activation_locks(second.clone(), None)); let mut second_acquire = Box::pin(acquire_pool_rebalance_activation_locks(second.clone(), None));
@@ -22110,10 +22119,50 @@ mod pools_tests {
.resources .resources
.lock() .lock()
.expect("activation lock recorder should not be poisoned"), .expect("activation lock recorder should not be poisoned"),
vec![POOL_META_NAME.to_string(), REBAL_META_NAME.to_string()] vec![REBAL_META_NAME.to_string(), POOL_META_NAME.to_string()]
); );
} }
#[tokio::test]
async fn test_activation_cancellation_releases_rebalance_fence_while_pool_fence_is_contended() {
use crate::storage_api_contracts::namespace::NamespaceLocking as _;
let pool = Arc::new(ActivationLockRecorder {
lock_manager: Arc::new(rustfs_lock::GlobalLockManager::new()),
owner: "activation-cancellation",
resources: StdMutex::new(Vec::new()),
});
let pool_lock = pool
.new_ns_lock(crate::disk::RUSTFS_META_BUCKET, POOL_META_NAME)
.await
.expect("pool lock should be created");
let pool_reader = pool_lock
.get_read_lock(std::time::Duration::from_secs(5))
.await
.expect("ordinary mutation should hold the pool read fence");
pool.resources.lock().expect("recorder should not be poisoned").clear();
let mut activation = Box::pin(acquire_pool_rebalance_activation_locks(Arc::clone(&pool), None));
assert!(matches!(futures::poll!(&mut activation), Poll::Pending));
assert_eq!(
*pool.resources.lock().expect("recorder should not be poisoned"),
vec![REBAL_META_NAME.to_string(), POOL_META_NAME.to_string()],
"activation must hold the run fence before waiting for the pool fence",
);
drop(activation);
let rebalance_lock = pool
.new_ns_lock(crate::disk::RUSTFS_META_BUCKET, REBAL_META_NAME)
.await
.expect("run lock should be created");
let run_writer = rebalance_lock
.get_write_lock(std::time::Duration::from_secs(5))
.await
.expect("cancelling activation must release its already-acquired run fence");
assert!(
!pool_reader.is_released(),
"cancelling activation must not release another caller's pool fence"
);
assert!(!run_writer.is_lock_lost());
}
#[test] #[test]
fn decommission_receipt_run_token_changes_with_persisted_start_time() { fn decommission_receipt_run_token_changes_with_persisted_start_time() {
let first = OffsetDateTime::from_unix_timestamp(1_000).expect("first run timestamp should be valid"); let first = OffsetDateTime::from_unix_timestamp(1_000).expect("first run timestamp should be valid");
+107 -1
View File
@@ -627,7 +627,13 @@ impl From<tokio::task::JoinError> for DiskError {
impl Clone for DiskError { impl Clone for DiskError {
fn clone(&self) -> Self { fn clone(&self) -> Self {
match self { match self {
DiskError::Io(io_error) => DiskError::Io(std::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)
// The helper derives a kind from the source; Clone must retain the original outer kind.
.map(|source| std::io::Error::new(io_error.kind(), source))
.unwrap_or_else(|| std::io::Error::new(io_error.kind(), io_error.to_string())),
),
DiskError::MaxVersionsExceeded => DiskError::MaxVersionsExceeded, DiskError::MaxVersionsExceeded => DiskError::MaxVersionsExceeded,
DiskError::Unexpected => DiskError::Unexpected, DiskError::Unexpected => DiskError::Unexpected,
DiskError::CorruptedFormat => DiskError::CorruptedFormat, DiskError::CorruptedFormat => DiskError::CorruptedFormat,
@@ -1265,6 +1271,49 @@ mod tests {
assert!(!bad_request.is_retryable_internode_write_failure()); assert!(!bad_request.is_retryable_internode_write_failure());
} }
#[test]
fn test_internode_http_clone_preserves_retryability_status_and_context() {
use http::StatusCode;
use rustfs_rio::InternodeHttpErrorKind::{ConnectionRefused, ConnectionReset, HttpStatus, Unknown};
for (kind, retryable) in [
(ConnectionRefused, true),
(ConnectionReset, true),
(HttpStatus(StatusCode::TOO_MANY_REQUESTS), true),
(HttpStatus(StatusCode::SERVICE_UNAVAILABLE), true),
(HttpStatus(StatusCode::CONFLICT), true),
(Unknown, false),
(HttpStatus(StatusCode::BAD_REQUEST), false),
(HttpStatus(StatusCode::INTERNAL_SERVER_ERROR), false),
] {
let original = DiskError::from(rustfs_rio::new_test_internode_http_io_error(kind));
assert_eq!(original.internode_http_error_kind(), Some(kind));
assert_eq!(original.is_retryable_internode_write_failure(), retryable);
let cloned = original.clone();
assert_eq!(cloned, original, "clone must preserve the error bucket for {kind:?}");
assert_eq!(
cloned.is_retryable_internode_write_failure(),
retryable,
"clone changed retryability for {kind:?}"
);
assert_eq!(cloned.internode_http_error_kind(), Some(kind));
if let HttpStatus(status) = kind {
assert!(cloned.is_internode_http_status(status.as_u16()));
}
let DiskError::Io(io_error) = &cloned else {
panic!("unmarked internode error must remain Io: {cloned:?}");
};
let source = io_error
.get_ref()
.and_then(|source| source.downcast_ref::<InternodeHttpError>())
.expect("clone must retain the structured internode error");
assert_eq!(source.context().method(), "PUT");
assert_eq!(source.context().target(), "/rustfs/rpc/put_file_stream");
assert_eq!(source.context().operation(), Some(INTERNODE_OPERATION_PUT_FILE_STREAM));
}
}
#[tokio::test] #[tokio::test]
async fn read_stream_conflict_is_not_a_retryable_put_file_failure() { async fn read_stream_conflict_is_not_a_retryable_put_file_failure() {
use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::io::{AsyncReadExt, AsyncWriteExt};
@@ -1309,11 +1358,57 @@ mod tests {
!error.is_retryable_internode_write_failure(), !error.is_retryable_internode_write_failure(),
"read-operation 409 must not trigger put-file retry" "read-operation 409 must not trigger put-file retry"
); );
let cloned = error.clone();
let reduced = crate::disk::error_reduce::reduce_write_quorum_errs(&[Some(error)], &[], 1)
.expect("the read conflict must remain the dominant error");
for preserved in [&cloned, &reduced] {
assert!(
!preserved.is_retryable_internode_write_failure(),
"cloning or reducing a read conflict must not turn it into a PUT retry"
);
assert!(preserved.is_internode_http_status(409));
let DiskError::Io(io_error) = preserved else {
panic!("read conflict must remain Io: {preserved:?}");
};
let source = io_error
.get_ref()
.and_then(|source| source.downcast_ref::<InternodeHttpError>())
.expect("read conflict must retain its request context");
assert_eq!(source.context().method(), "GET");
assert_eq!(source.context().target(), "/rustfs/rpc/read_file_stream");
assert_eq!(
source.context().operation(),
Some(rustfs_io_metrics::internode_metrics::INTERNODE_OPERATION_READ_FILE_STREAM)
);
}
}) })
.await .await
.expect("isolated read-conflict test must finish within its budget"); .expect("isolated read-conflict test must finish within its budget");
} }
#[test]
fn test_internode_http_clone_preserves_outer_io_kind_and_message() {
let source = rustfs_rio::new_test_internode_http_io_error(InternodeHttpErrorKind::ConnectionReset)
.into_inner()
.expect("the internode helper must provide a typed source");
let original_io = io::Error::new(io::ErrorKind::InvalidData, source);
let message = original_io.to_string();
let original = DiskError::from(original_io);
assert_eq!(original.internode_http_error_kind(), Some(InternodeHttpErrorKind::ConnectionReset));
assert!(original.is_retryable_internode_write_failure());
let cloned = original.clone();
let reduced = crate::disk::error_reduce::reduce_write_quorum_errs(&[Some(original)], &[], 1)
.expect("the wrapped internode error must remain the dominant error");
for preserved in [&cloned, &reduced] {
let DiskError::Io(io_error) = preserved else {
panic!("the wrapped error must remain Io: {preserved:?}");
};
assert_eq!(io_error.kind(), io::ErrorKind::InvalidData);
assert_eq!(io_error.to_string(), message);
}
}
#[test] #[test]
fn test_internode_missing_errors_preserve_disk_error_types() { fn test_internode_missing_errors_preserve_disk_error_types() {
let file_missing = DiskError::from(rustfs_rio::new_test_remote_file_not_found_http_io_error()); let file_missing = DiskError::from(rustfs_rio::new_test_remote_file_not_found_http_io_error());
@@ -1325,6 +1420,17 @@ mod tests {
assert_eq!(file_missing, DiskError::FileNotFound); assert_eq!(file_missing, DiskError::FileNotFound);
assert_eq!(volume_missing, DiskError::VolumeNotFound); assert_eq!(volume_missing, DiskError::VolumeNotFound);
assert!(matches!(unmarked_server_error, DiskError::Io(_))); assert!(matches!(unmarked_server_error, DiskError::Io(_)));
for missing in [file_missing, volume_missing] {
assert_eq!(missing.clone(), missing);
assert_eq!(
crate::disk::error_reduce::reduce_write_quorum_errs(
&[Some(missing.clone()), Some(missing.clone()), None],
&[],
2
),
Some(missing)
);
}
} }
#[test] #[test]
+72
View File
@@ -226,6 +226,78 @@ mod tests {
assert_eq!(res, Some(quorum_err)); assert_eq!(res, Some(quorum_err));
} }
#[test]
fn test_write_quorum_reduction_preserves_internode_http_identity() {
use http::StatusCode;
use rustfs_rio::InternodeHttpErrorKind::{ConnectionRefused, HttpStatus, Unknown};
for (kind, retryable) in [
(ConnectionRefused, true),
(HttpStatus(StatusCode::SERVICE_UNAVAILABLE), true),
(HttpStatus(StatusCode::CONFLICT), true),
(Unknown, false),
(HttpStatus(StatusCode::BAD_REQUEST), false),
] {
// Construct both producer errors independently: the reducer owns the first clone.
let first = Error::from(rustfs_rio::new_test_internode_http_io_error(kind));
let second = Error::from(rustfs_rio::new_test_internode_http_io_error(kind));
assert_eq!(first.internode_http_error_kind(), Some(kind));
assert_eq!(second.internode_http_error_kind(), Some(kind));
assert_eq!(first.is_retryable_internode_write_failure(), retryable);
let errors = [Some(first), Some(second), None];
let reduced = reduce_write_quorum_errs(&errors, OBJECT_OP_IGNORED_ERRS, 2)
.expect("two equal producer errors must dominate one successful write");
assert_eq!(Some(&reduced), errors[0].as_ref());
assert_eq!(
reduced.is_retryable_internode_write_failure(),
retryable,
"quorum reduction changed retryability for {kind:?}"
);
assert_eq!(reduced.internode_http_error_kind(), Some(kind));
if let HttpStatus(status) = kind {
assert!(reduced.is_internode_http_status(status.as_u16()));
}
let Error::Io(io_error) = &reduced else {
panic!("the dominant error must remain Io: {reduced:?}");
};
let source = io_error
.get_ref()
.and_then(|source| source.downcast_ref::<rustfs_rio::InternodeHttpError>())
.expect("quorum reduction must retain the structured internode error");
assert_eq!(source.context().method(), "PUT");
assert_eq!(source.context().target(), "/rustfs/rpc/put_file_stream");
assert_eq!(
source.context().operation(),
Some(rustfs_io_metrics::internode_metrics::INTERNODE_OPERATION_PUT_FILE_STREAM)
);
}
}
#[test]
fn test_clone_and_write_quorum_do_not_promote_non_retryable_errors() {
use http::StatusCode;
use rustfs_rio::InternodeHttpErrorKind::{HttpStatus, Unknown};
for original in [
Error::from(rustfs_rio::new_test_internode_http_io_error(Unknown)),
Error::from(rustfs_rio::new_test_internode_http_io_error(HttpStatus(StatusCode::BAD_REQUEST))),
Error::from(rustfs_rio::new_test_internode_http_io_error(HttpStatus(StatusCode::FORBIDDEN))),
Error::from(rustfs_rio::new_test_internode_http_io_error(HttpStatus(StatusCode::NOT_FOUND))),
Error::from(rustfs_rio::new_test_internode_http_io_error(HttpStatus(
StatusCode::INTERNAL_SERVER_ERROR,
))),
err_io("internode connection reset: PUT /rustfs/rpc/put_file_stream"),
] {
assert!(!original.is_retryable_internode_write_failure());
let cloned = original.clone();
let reduced =
reduce_write_quorum_errs(&[Some(original)], &[], 1).expect("a non-retryable error must remain an error");
assert!(!cloned.is_retryable_internode_write_failure());
assert!(!reduced.is_retryable_internode_write_failure());
}
}
#[test] #[test]
fn test_count_errs() { fn test_count_errs() {
let e1 = err_io("a"); let e1 = err_io("a");
@@ -572,7 +572,7 @@ impl ECStore {
where where
S: EcstoreObjectIO + StorageNamespaceLocking<Error = Error, NamespaceLock = rustfs_lock::NamespaceLockWrapper>, S: EcstoreObjectIO + StorageNamespaceLocking<Error = Error, NamespaceLock = rustfs_lock::NamespaceLockWrapper>,
{ {
// Lock order: pool_meta_save_gate -> pool.bin -> rebalance.bin. // Lock order: pool_meta_save_gate -> rebalance.bin -> pool.bin.
let mut pool_meta_guard = self.pool_meta_save_gate.lock().await; let mut pool_meta_guard = self.pool_meta_save_gate.lock().await;
pool_meta_guard.ensure_write_safe("rebalance worker activation")?; pool_meta_guard.ensure_write_safe("rebalance worker activation")?;
// Classify the durable rebalance record while holding both namespace // Classify the durable rebalance record while holding both namespace
@@ -50,6 +50,11 @@ fn ensure_rebalance_entry_active(cancel: &CancellationToken) -> Result<()> {
Ok(()) Ok(())
} }
#[cfg(test)]
tokio::task_local! {
static REBALANCE_ENTRY_RUN_FENCE_BARRIER: (Arc<tokio::sync::Notify>, Arc<tokio::sync::Notify>);
}
#[derive(Debug)] #[derive(Debug)]
struct RebalanceEntryTarget { struct RebalanceEntryTarget {
bucket: String, bucket: String,
@@ -256,9 +261,15 @@ impl ECStore {
.sort_by_key(|v| (v.mod_time.is_none(), std::cmp::Reverse(v.mod_time))); .sort_by_key(|v| (v.mod_time.is_none(), std::cmp::Reverse(v.mod_time)));
// Entry lock order is bucket incarnation -> activation_gate -> rebalance.bin -> movement gate. // Entry lock order is bucket incarnation -> activation_gate -> rebalance.bin -> movement gate.
// Target capacity admission can then acquire pool.bin under the run fence.
// Stop waits for in-flight entries through cleanup, but not for entries admitted later. // Stop waits for in-flight entries through cleanup, but not for entries admitted later.
ensure_rebalance_entry_active(&cancel)?; ensure_rebalance_entry_active(&cancel)?;
let run_guard = self.rebalance_run_guard(rebalance_id.as_ref(), "rebalance entry").await?; let run_guard = self.rebalance_run_guard(rebalance_id.as_ref(), "rebalance entry").await?;
#[cfg(test)]
if let Ok((arrived, release)) = REBALANCE_ENTRY_RUN_FENCE_BARRIER.try_with(Clone::clone) {
arrived.notify_one();
release.notified().await;
}
let lock_lost_signal = run_guard.lock_lost_signal(); let lock_lost_signal = run_guard.lock_lost_signal();
#[cfg(test)] #[cfg(test)]
let _run_signal_test_fence = lock_lost_signal let _run_signal_test_fence = lock_lost_signal
@@ -1237,6 +1248,130 @@ mod tests {
assert_eq!(pool_stats.cleanup_warnings.count, 1, "deferred cleanup must not add a permanent warning"); assert_eq!(pool_stats.cleanup_warnings.count, 1, "deferred cleanup must not add a permanent warning");
} }
#[tokio::test]
#[serial_test::serial]
async fn real_rebalance_entry_progresses_while_peer_activation_waits_for_run_fence() {
const REBALANCE_ID: &str = "rebalance-peer-activation-lock-order";
let (_temp_dirs, store, peer) = crate::services::rebalance::test_two_pool_stores_with_isolated_node_contexts(Some(
active_rebalance_meta(REBALANCE_ID),
))
.await;
assert!(!Arc::ptr_eq(&store.ctx, &peer.ctx), "node-local movement gates must be independent");
{
let mut meta = peer.rebalance_meta.write().await;
let meta = meta.as_mut().expect("peer should know the durable run");
meta.activation_gate = Arc::default();
meta.cancel = None;
}
let bucket = crate::disk::RUSTFS_META_BUCKET;
let object = "rebalance-peer-activation-object";
let version_id = uuid::Uuid::new_v4();
let payload = b"entry must drain before peer activation takes the pool fence".repeat(1024);
let source_set = store.pools[0].get_disks_by_key(object);
let target_set = store.pools[1].get_disks_by_key(object);
let opts = ObjectOptions {
versioned: true,
version_id: Some(version_id.to_string()),
..Default::default()
};
let mut writer = PutObjReader::from_vec(payload.clone());
let source_before = source_set
.put_object(bucket, object, &mut writer, &opts)
.await
.expect("source version should be written");
let entry = metacache_entry_from_source(&source_set, bucket, object).await;
let arrived = Arc::new(tokio::sync::Notify::new());
let release = Arc::new(tokio::sync::Notify::new());
// JoinSet aborts both scoped tasks if an assertion or timeout fails.
let mut tasks = tokio::task::JoinSet::new();
let entry_store = Arc::clone(&store);
tasks.spawn(
REBALANCE_ENTRY_RUN_FENCE_BARRIER.scope((Arc::clone(&arrived), Arc::clone(&release)), async move {
entry_store
.rebalance_entry(
RebalanceEntryTarget {
bucket: bucket.to_string(),
pool_index: 0,
},
entry,
source_set,
Arc::new(RebalanceBucketConfigs::default()),
Arc::from(REBALANCE_ID),
CancellationToken::new(),
)
.await
}),
);
tokio::time::timeout(StdDuration::from_secs(30), arrived.notified())
.await
.expect("real entry must acquire its persisted run read fence");
let attempted = Arc::new(tokio::sync::Notify::new());
let peer_pool = Arc::clone(&peer.pools[0]);
let (activation_done, activation_result) = tokio::sync::oneshot::channel();
tasks.spawn(
crate::core::pools::REBALANCE_ACTIVATION_LOCK_ATTEMPT.scope(Arc::clone(&attempted), async move {
let result = peer.fence_rebalance_worker_activation(peer_pool, REBALANCE_ID).await;
let result = result.map(|fence| match fence {
super::super::control::RebalanceWorkerActivationFence::Ready(fence) => {
fence.ensure_held().expect("peer activation must retain both fences");
}
super::super::control::RebalanceWorkerActivationFence::NotStartedTerminal => {
panic!("the paused entry's run must still require activation");
}
});
activation_done.send(result).expect("activation receiver should remain alive");
Ok(RebalanceEntryOutcome::Completed)
}),
);
tokio::time::timeout(StdDuration::from_secs(30), attempted.notified())
.await
.expect("peer activation must attempt the persisted rebalance write fence");
release.notify_one();
tokio::time::timeout(StdDuration::from_secs(30), async {
while let Some(result) = tasks.join_next().await {
assert!(matches!(
result
.expect("scoped task must not panic")
.expect("entry must not fail or defer"),
RebalanceEntryOutcome::Completed
));
}
})
.await
.expect("entry and peer activation must both make progress");
activation_result
.await
.expect("peer activation result should be sent")
.expect("peer activation must not time out behind the entry it blocks");
let mut reader = target_set
.get_object_reader(bucket, object, None, HeaderMap::new(), &opts)
.await
.expect("the exact target version must be readable");
let mut actual = Vec::new();
reader
.stream
.read_to_end(&mut actual)
.await
.expect("target body should drain completely");
assert_eq!(actual, payload);
assert_eq!(reader.object_info.version_id, source_before.version_id);
assert_eq!(reader.object_info.etag, source_before.etag);
assert_eq!(reader.object_info.mod_time, source_before.mod_time);
let source_error = store.pools[0]
.get_object_info(bucket, object, &opts)
.await
.expect_err("completed entry must clean up the source version");
assert!(crate::error::is_err_object_not_found(&source_error) || crate::error::is_err_version_not_found(&source_error));
let meta = store.rebalance_meta.read().await;
let stats = &meta.as_ref().expect("local run must remain installed").pool_stats[0];
assert_eq!(stats.num_objects, 1);
assert_eq!(stats.num_versions, 1);
assert_eq!(stats.cleanup_warnings.count, 0);
}
#[tokio::test] #[tokio::test]
#[serial_test::serial] #[serial_test::serial]
async fn real_rebalance_run_fence_loss_before_target_commit_preserves_target_and_source() { async fn real_rebalance_run_fence_loss_before_target_commit_preserves_target_and_source() {
@@ -1907,6 +1907,124 @@ fn test_is_transient_rebalance_error_accepts_wrapped_disk_timeout() {
assert!(is_transient_rebalance_error(&Error::Io(std::io::Error::other(DiskError::Timeout)))); assert!(is_transient_rebalance_error(&Error::Io(std::io::Error::other(DiskError::Timeout))));
} }
#[test]
fn test_rebalance_stage_wrapped_transient_errors_remain_retryable() {
let cases = [
Error::Lock(rustfs_lock::LockError::timeout(".rustfs.sys/pool.bin@latest", Duration::from_secs(5))),
Error::Lock(rustfs_lock::LockError::network(
"peer unavailable",
std::io::Error::from(std::io::ErrorKind::ConnectionReset),
)),
Error::SlowDown,
Error::ErasureReadQuorum,
Error::ErasureWriteQuorum,
Error::Io(std::io::Error::other(DiskError::Timeout)),
Error::Io(std::io::Error::from(std::io::ErrorKind::TimedOut)),
];
for mut error in cases {
for depth in 0..=3 {
assert!(is_transient_rebalance_error(&error), "transient source lost at depth {depth}: {error:?}");
assert!(
should_defer_rebalance_entry_failure(&error),
"exhausted transient entries must be deferred"
);
assert!(should_retry_rebalance_listing(&error, 0, 3));
assert!(
!should_retry_rebalance_listing(&error, 2, 3),
"wrapping must not bypass the attempt limit"
);
error = data_movement::data_movement_stage_error_for_test(
"rebalance_object",
"put_object",
"bucket",
"baseline/00042.bin",
error,
);
}
}
}
#[test]
fn test_rebalance_stage_wrapped_terminal_errors_remain_terminal() {
let cases = [
Error::FileAccessDenied,
Error::FileCorrupt,
Error::OperationCanceled,
Error::DataMovementOverwriteErr("bucket".to_string(), "object".to_string(), "version".to_string()),
Error::Lock(rustfs_lock::LockError::already_locked("bucket/object", "owner")),
Error::other("permission denied"),
];
for mut error in cases {
for depth in 0..=3 {
assert!(
!is_transient_rebalance_error(&error),
"terminal source must survive depth {depth}: {error:?}"
);
assert!(!should_defer_rebalance_entry_failure(&error));
// Object names are untrusted context, not evidence of a transient failure.
error = data_movement::data_movement_stage_error_for_test(
"rebalance_object",
"put_object",
"bucket",
"remote lock rpc timed out",
error,
);
}
}
}
#[tokio::test]
async fn test_rebalance_stage_wrapped_lock_timeout_retries_real_migration_loop() {
for succeeds_on_retry in [true, false] {
let backend = MigrationBackendSpy::new(None, None);
let attempts = AtomicUsize::new(0);
let waits = AtomicUsize::new(0);
let mut transfer = |_, _, _| {
let attempt = attempts.fetch_add(1, Ordering::SeqCst);
async move {
if succeeds_on_retry && attempt > 0 {
return Ok(());
}
Err(data_movement::data_movement_stage_error_for_test(
"rebalance_object",
"put_object",
"bucket",
"baseline/00042.bin",
Error::Lock(rustfs_lock::LockError::timeout(".rustfs.sys/pool.bin@latest", Duration::from_secs(5))),
))
}
};
let version = version_normal();
let result = migrate_entry_version_with_retry_wait(
&backend,
"bucket".to_string(),
0,
&version,
None,
3,
false,
&mut transfer,
|_: String, _: String, _: ObjectOptions| async { Ok::<_, Error>(ObjectInfo::default()) },
|_| {
waits.fetch_add(1, Ordering::SeqCst);
std::future::ready(())
},
)
.await;
assert_eq!(result.moved, succeeds_on_retry);
assert_eq!(result.failed, !succeeds_on_retry);
assert_eq!(attempts.load(Ordering::SeqCst), if succeeds_on_retry { 2 } else { 3 });
assert_eq!(backend.get_calls(), attempts.load(Ordering::SeqCst));
assert_eq!(waits.load(Ordering::SeqCst), attempts.load(Ordering::SeqCst) - 1);
if !succeeds_on_retry {
assert_eq!(result.stage, Some("write_target"));
assert!(should_defer_rebalance_entry_failure(
result.error.as_ref().expect("exhaustion must retain its source error")
));
}
}
}
#[test] #[test]
fn test_is_transient_rebalance_error_accepts_io_timeout_message() { fn test_is_transient_rebalance_error_accepts_io_timeout_message() {
assert!(is_transient_rebalance_error(&Error::Io(std::io::Error::other("timeout")))); assert!(is_transient_rebalance_error(&Error::Io(std::io::Error::other("timeout"))));
@@ -244,6 +244,7 @@ pub(super) fn resolve_rebalance_bucket_result(
} }
pub(super) fn is_transient_rebalance_error(err: &Error) -> bool { pub(super) fn is_transient_rebalance_error(err: &Error) -> bool {
let err = rebalance_error_source(err);
match err { match err {
Error::SlowDown Error::SlowDown
| Error::ErasureReadQuorum | Error::ErasureReadQuorum
@@ -256,6 +257,15 @@ pub(super) fn is_transient_rebalance_error(err: &Error) -> bool {
} }
} }
fn rebalance_error_source(mut err: &Error) -> &Error {
// Stage context contains object names, so classify the preserved source,
// not timeout-like text supplied by an object name. Iterate nested stages.
while let Some(source) = crate::data_movement::data_movement_stage_source(err) {
err = source;
}
err
}
fn is_rebalance_transient_lock_error(err: &rustfs_lock::LockError) -> bool { fn is_rebalance_transient_lock_error(err: &rustfs_lock::LockError) -> bool {
match err { match err {
rustfs_lock::LockError::Timeout { .. } | rustfs_lock::LockError::Network { .. } => true, rustfs_lock::LockError::Timeout { .. } | rustfs_lock::LockError::Network { .. } => true,
@@ -309,6 +319,7 @@ pub(super) fn rebalance_listing_retry_delay(attempt: usize) -> Duration {
} }
fn is_rebalance_lock_or_rpc_timeout(err: &Error) -> bool { fn is_rebalance_lock_or_rpc_timeout(err: &Error) -> bool {
let err = rebalance_error_source(err);
match err { match err {
Error::Lock(rustfs_lock::LockError::Timeout { .. }) | Error::Lock(rustfs_lock::LockError::Network { .. }) => true, Error::Lock(rustfs_lock::LockError::Timeout { .. }) | Error::Lock(rustfs_lock::LockError::Network { .. }) => true,
Error::Io(io_err) => is_rebalance_lock_or_rpc_timeout_message(&io_err.to_string()), Error::Io(io_err) => is_rebalance_lock_or_rpc_timeout_message(&io_err.to_string()),
@@ -585,3 +596,48 @@ impl SetDisks {
Ok(()) Ok(())
} }
} }
#[cfg(test)]
mod error_source_tests {
use super::*;
#[test]
fn stage_wrapped_errors_select_the_source_backoff_policy() {
let cases = [
(
Error::Lock(rustfs_lock::LockError::timeout(".rustfs.sys/pool.bin@latest", Duration::from_secs(5))),
true,
),
(
Error::Lock(rustfs_lock::LockError::network(
"peer unavailable",
std::io::Error::from(std::io::ErrorKind::ConnectionReset),
)),
true,
),
(Error::other("remote lock rpc timed out"), true),
(Error::SlowDown, false),
(Error::Io(std::io::Error::other(DiskError::Timeout)), false),
(Error::FileAccessDenied, false),
];
for (mut error, lock_backoff) in cases {
for depth in 0..=3 {
assert_eq!(
is_rebalance_lock_or_rpc_timeout(&error),
lock_backoff,
"wrong backoff at depth {depth}: {error:?}"
);
if !lock_backoff {
assert_eq!(rebalance_migration_retry_delay(1, &error), REBALANCE_MIGRATION_RETRY_BASE_DELAY * 2);
}
error = crate::data_movement::data_movement_stage_error_for_test(
"rebalance_object",
"put_object",
"bucket",
"remote lock rpc timed out",
error,
);
}
}
}
}
+42 -7
View File
@@ -326,14 +326,15 @@ impl SetDisks {
let parity_blocks = Self::common_parity(&parities, default_parity_count as i32); let parity_blocks = Self::common_parity(&parities, default_parity_count as i32);
if parity_blocks < 0 { if parity_blocks < 0 {
// No parity value reached read quorum. Distinguish two cases: // A consistent layout can require more replies than the initial
// enough disks answered with valid-looking metadata that simply // half-set probe. Reaching that probe alone is not corruption;
// cannot be reconciled (corrupt/foreign entries — retrying cannot // only invalid or conflicting healthy replies establish that.
// help, and heal should see Corrupt, rustfs#5801) versus too few
// healthy answers (a genuine quorum condition where retry may
// succeed once disks recover).
let healthy_replies = errs.iter().filter(|err| err.is_none()).count(); let healthy_replies = errs.iter().filter(|err| err.is_none()).count();
if healthy_replies >= expected_rquorum { let consistent_parity = parities
.iter()
.find(|&&parity| parity >= 0)
.filter(|&&parity| parities.iter().filter(|&&candidate| candidate == parity).count() == healthy_replies);
if healthy_replies >= expected_rquorum && consistent_parity.is_none() {
error!( error!(
"object_quorum_from_meta: irreconcilable parity across {healthy_replies} healthy replies (corrupt metadata), errs={errs:?}" "object_quorum_from_meta: irreconcilable parity across {healthy_replies} healthy replies (corrupt metadata), errs={errs:?}"
); );
@@ -1652,6 +1653,40 @@ mod tests {
assert_eq!(err, DiskError::FileCorrupt); assert_eq!(err, DiskError::FileCorrupt);
} }
#[test]
fn consistent_parity_below_its_data_shard_quorum_is_not_corruption() {
for (drive_count, parity) in [(6, 2), (8, 2), (12, 4)] {
let data = drive_count - parity;
let mut metas = (1..=drive_count)
.map(|index| {
let mut info = FileInfo::new("bucket/object", data, parity);
info.size = 1024;
info.erasure.index = index;
info
})
.collect::<Vec<_>>();
let mut errs = vec![Some(DiskError::DiskNotFound); drive_count];
errs[..data].fill(None);
assert_eq!(
SetDisks::object_quorum_from_meta(&metas, &errs, parity).expect("exact data quorum should resolve"),
(data as i32, data as i32)
);
errs[data - 1] = Some(DiskError::DiskNotFound);
assert_eq!(
SetDisks::object_quorum_from_meta(&metas, &errs, parity).expect_err("one fewer shard cannot resolve"),
DiskError::ErasureReadQuorum,
"layout {drive_count}/{parity} has consistent metadata but insufficient shards"
);
metas[0].erasure.parity_blocks = usize::MAX;
assert_eq!(
SetDisks::object_quorum_from_meta(&metas, &errs, parity).expect_err("corrupt healthy replies must be rejected"),
DiskError::FileCorrupt
);
}
}
/// Too few healthy replies remains a genuine quorum condition where a /// Too few healthy replies remains a genuine quorum condition where a
/// retry may succeed once disks recover. /// retry may succeed once disks recover.
#[test] #[test]
+1
View File
@@ -865,6 +865,7 @@ pub(crate) use core::io_primitives::{ENV_RUSTFS_PUT_RENAME_EARLY_ACK_ENABLE, ren
mod ctx; mod ctx;
mod metadata; mod metadata;
mod ops; mod ops;
pub(crate) use ops::bucket::BucketInfoQuorum;
#[cfg(test)] #[cfg(test)]
pub(crate) use ops::hermetic_set_disks_isolated; pub(crate) use ops::hermetic_set_disks_isolated;
+65 -52
View File
@@ -21,12 +21,72 @@
use super::super::{ use super::super::{
BUCKET_OP_IGNORED_ERRS, BucketInfo, BucketOperations, BucketOptions, DeleteBucketOptions, DiskError, Error, HashMap, BUCKET_OP_IGNORED_ERRS, BucketInfo, BucketOperations, BucketOptions, DeleteBucketOptions, DiskError, Error, HashMap,
MakeBucketOptions, Result, SetDisks, is_reserved_or_invalid_bucket, join_all, reduce_write_quorum_errs, MakeBucketOptions, Result, SetDisks, is_reserved_or_invalid_bucket, join_all, reduce_read_quorum_errs,
reduce_write_quorum_errs,
}; };
use crate::api::bucket::metadata_sys; use crate::api::bucket::metadata_sys;
use crate::disk::DiskAPI; use crate::disk::DiskAPI;
#[derive(Clone, Copy)]
pub(crate) enum BucketInfoQuorum {
Read,
Write,
}
impl SetDisks { impl SetDisks {
pub(crate) async fn stat_bucket_with_quorum(&self, bucket: &str, quorum: BucketInfoQuorum) -> Result<BucketInfo> {
let disks = self.disk_inventory().await;
let disk_count = disks.len();
let mut futures = Vec::with_capacity(disk_count);
for disk in disks {
let bucket = bucket.to_string();
futures.push(async move {
match disk {
Some(disk) => disk.stat_volume(&bucket).await,
None => Err(DiskError::DiskNotFound),
}
});
}
let results = join_all(futures).await;
let mut infos = Vec::with_capacity(results.len());
let mut errs = Vec::with_capacity(results.len());
for result in results {
match result {
Ok(info) => {
infos.push(Some(info));
errs.push(None);
}
Err(err) => {
infos.push(None);
errs.push(Some(err));
}
}
}
let error = match quorum {
// Bucket mutations use a majority regardless of object storage
// class. A namespace read must intersect that majority; object
// readers still enforce the persisted layout's data-shard quorum.
BucketInfoQuorum::Read => reduce_read_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, disk_count.div_ceil(2).max(1)),
BucketInfoQuorum::Write => reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, disk_count / 2 + 1),
};
if let Some(err) = error {
return Err(err.into());
}
infos
.into_iter()
.flatten()
.next()
.map(|info| BucketInfo {
name: info.name,
created: info.created,
..Default::default()
})
.ok_or(Error::VolumeNotFound)
}
pub(crate) async fn list_bucket_for_scanner(&self, _opts: &BucketOptions) -> Result<(Vec<BucketInfo>, bool)> { pub(crate) async fn list_bucket_for_scanner(&self, _opts: &BucketOptions) -> Result<(Vec<BucketInfo>, bool)> {
let disks = self.disk_inventory().await; let disks = self.disk_inventory().await;
let write_quorum = (disks.len() / 2) + 1; let write_quorum = (disks.len() / 2) + 1;
@@ -131,59 +191,12 @@ impl BucketOperations for SetDisks {
#[tracing::instrument(skip(self))] #[tracing::instrument(skip(self))]
async fn get_bucket_info(&self, bucket: &str, _opts: &BucketOptions) -> Result<BucketInfo> { async fn get_bucket_info(&self, bucket: &str, _opts: &BucketOptions) -> Result<BucketInfo> {
let disks = self.disk_inventory().await; let mut info = self.stat_bucket_with_quorum(bucket, BucketInfoQuorum::Write).await?;
let write_quorum = (disks.len() / 2) + 1;
let mut futures = Vec::with_capacity(disks.len());
for disk in disks {
let bucket = bucket.to_string();
futures.push(async move {
match disk {
Some(disk) => disk.stat_volume(&bucket).await,
None => Err(DiskError::DiskNotFound),
}
});
}
let results = join_all(futures).await;
let mut infos = Vec::with_capacity(results.len());
let mut errs = Vec::with_capacity(results.len());
for result in results {
match result {
Ok(info) => {
infos.push(Some(info));
errs.push(None);
}
Err(err) => {
infos.push(None);
errs.push(Some(err));
}
}
}
if let Some(err) = reduce_write_quorum_errs(&errs, BUCKET_OP_IGNORED_ERRS, write_quorum) {
return Err(err.into());
}
let mut versioning = false;
let mut object_locking = false;
if let Ok(sys) = metadata_sys::get(bucket).await { if let Ok(sys) = metadata_sys::get(bucket).await {
versioning = sys.versioning(); info.versioning = sys.versioning();
object_locking = sys.object_locking(); info.object_locking = sys.object_locking();
} }
Ok(info)
infos
.into_iter()
.flatten()
.next()
.map(|info| BucketInfo {
name: info.name,
created: info.created,
versioning,
object_locking,
..Default::default()
})
.ok_or(Error::VolumeNotFound)
} }
#[tracing::instrument(skip(self))] #[tracing::instrument(skip(self))]
+286 -15
View File
@@ -19,7 +19,7 @@ use crate::bucket::{
}; };
use crate::error::is_err_bucket_not_found; use crate::error::is_err_bucket_not_found;
use crate::runtime::sources as runtime_sources; use crate::runtime::sources as runtime_sources;
use crate::set_disk::get_lock_acquire_timeout; use crate::set_disk::{BucketInfoQuorum, get_lock_acquire_timeout};
use crate::storage_api_contracts::bucket::{BUCKET_LIFECYCLE_LOCK_OBJECT, SRBucketDeleteOp}; use crate::storage_api_contracts::bucket::{BUCKET_LIFECYCLE_LOCK_OBJECT, SRBucketDeleteOp};
use crate::storage_api_contracts::namespace::NamespaceLocking as _; use crate::storage_api_contracts::namespace::NamespaceLocking as _;
use futures::stream::{self, StreamExt}; use futures::stream::{self, StreamExt};
@@ -772,17 +772,30 @@ impl ECStore {
#[instrument(skip(self))] #[instrument(skip(self))]
pub(crate) async fn get_bucket_info_from_sets(&self, bucket: &str, opts: &BucketOptions) -> Result<BucketInfo> { pub(crate) async fn get_bucket_info_from_sets(&self, bucket: &str, opts: &BucketOptions) -> Result<BucketInfo> {
self.get_bucket_info_from_sets_with_quorum(bucket, opts, BucketInfoQuorum::Write)
.await
}
async fn get_bucket_info_from_sets_with_quorum(
&self,
bucket: &str,
opts: &BucketOptions,
quorum: BucketInfoQuorum,
) -> Result<BucketInfo> {
// One host may participate in several pools after expansion. Resolve the // One host may participate in several pools after expansion. Resolve the
// namespace against each erasure set so disks from different pools can // namespace against each erasure set so disks from different pools can
// never be combined into one bucket quorum. // never be combined into one bucket quorum.
// Bucket validation is request-path IO. Keep the previous peer fanout's // Bucket validation is request-path IO. Keep the previous peer fanout's
// latency shape by probing every set concurrently; scanner listings use // latency shape by probing every set concurrently; scanner listings use
// a separate bounded path below because they run continuously. // a separate bounded path below because they run continuously.
let mut scoped_results = let mut scoped_results = futures::future::join_all(self.bucket_sets().map(|(pool_index, set_index, set)| async move {
futures::future::join_all(self.bucket_sets().map(|(pool_index, set_index, set)| async move { let result = match quorum {
(pool_index, set_index, set.get_bucket_info(bucket, opts).await) BucketInfoQuorum::Read => set.stat_bucket_with_quorum(bucket, quorum).await,
})) BucketInfoQuorum::Write => set.get_bucket_info(bucket, opts).await,
.await; };
(pool_index, set_index, result)
}))
.await;
scoped_results.sort_unstable_by_key(|(pool_index, set_index, _)| (*pool_index, *set_index)); scoped_results.sort_unstable_by_key(|(pool_index, set_index, _)| (*pool_index, *set_index));
let mut first_info = None; let mut first_info = None;
@@ -806,7 +819,11 @@ impl ECStore {
#[instrument(skip(self))] #[instrument(skip(self))]
pub(super) async fn handle_get_bucket_info(&self, bucket: &str, opts: &BucketOptions) -> Result<BucketInfo> { pub(super) async fn handle_get_bucket_info(&self, bucket: &str, opts: &BucketOptions) -> Result<BucketInfo> {
let mut info = self.get_bucket_info_from_sets(bucket, opts).await?; let mut info = match self.get_bucket_info_from_sets(bucket, opts).await {
Ok(info) => info,
Err(Error::ErasureWriteQuorum) => return self.get_bucket_info_at_read_quorum(bucket, opts).await,
Err(err) => return Err(err),
};
if let Ok(sys) = metadata_sys::get_in(&self.ctx, bucket).await { if let Ok(sys) = metadata_sys::get_in(&self.ctx, bucket).await {
if should_override_created_from_metadata(sys.created) { if should_override_created_from_metadata(sys.created) {
@@ -819,6 +836,35 @@ impl ECStore {
Ok(info) Ok(info)
} }
async fn get_bucket_info_at_read_quorum(&self, bucket: &str, opts: &BucketOptions) -> Result<BucketInfo> {
// Lock order: bucket lifecycle -> internal metadata object read locks.
// Keep create/delete from changing the namespace while a read quorum
// confirms both physical presence and persisted bucket metadata.
let guard = self.acquire_bucket_lifecycle_read_lock(bucket).await?;
await_bucket_namespace_operation(Some(&guard), bucket, "bucket read quorum validation", async {
let mut info = self
.get_bucket_info_from_sets_with_quorum(bucket, opts, BucketInfoQuorum::Read)
.await?;
let (metadata, persisted) = metadata_sys::get_config_from_disk_with_presence_in(&self.ctx, bucket).await?;
if !persisted {
// A minority of directories left by failed creation is not an
// authoritative bucket. Never turn fabricated defaults into
// permission to serve degraded reads.
return Err(Error::ErasureReadQuorum);
}
if metadata.name != bucket {
return Err(Error::FileCorrupt);
}
if should_override_created_from_metadata(metadata.created) {
info.created = Some(metadata.created);
}
info.versioning = metadata.versioning();
info.object_locking = metadata.object_locking();
Ok(info)
})
.await
}
#[instrument(skip(self))] #[instrument(skip(self))]
pub(super) async fn handle_list_bucket(&self, opts: &BucketOptions) -> Result<Vec<BucketInfo>> { pub(super) async fn handle_list_bucket(&self, opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
// TODO(backlog): support cached bucket listing via opts.cached // TODO(backlog): support cached bucket listing via opts.cached
@@ -1049,7 +1095,7 @@ mod tests {
run_physical_bucket_deletion, scan_metadata_less_residue, scan_metadata_less_residue_with_budget, run_physical_bucket_deletion, scan_metadata_less_residue, scan_metadata_less_residue_with_budget,
should_override_created_from_metadata, validate_table_bucket_delete_allowed, should_override_created_from_metadata, validate_table_bucket_delete_allowed,
}; };
use crate::bucket::metadata::table_bucket_catalog_metadata_prefix; use crate::bucket::metadata::{BucketMetadata, table_bucket_catalog_metadata_prefix};
use crate::bucket::metadata_sys; use crate::bucket::metadata_sys;
use crate::cluster::rpc::peer_s3_client::install_delete_bucket_empty_scan_barrier; use crate::cluster::rpc::peer_s3_client::install_delete_bucket_empty_scan_barrier;
use crate::disk::{BUCKET_META_PREFIX, DiskAPI, RUSTFS_META_BUCKET, STORAGE_FORMAT_FILE}; use crate::disk::{BUCKET_META_PREFIX, DiskAPI, RUSTFS_META_BUCKET, STORAGE_FORMAT_FILE};
@@ -1076,6 +1122,7 @@ mod tests {
use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::atomic::{AtomicBool, Ordering};
use std::time::{Duration, SystemTime}; use std::time::{Duration, SystemTime};
use time::OffsetDateTime; use time::OffsetDateTime;
use tokio::io::AsyncReadExt;
use tokio::sync::{Notify, OnceCell}; use tokio::sync::{Notify, OnceCell};
use tokio_util::sync::CancellationToken; use tokio_util::sync::CancellationToken;
use uuid::Uuid; use uuid::Uuid;
@@ -1359,11 +1406,18 @@ mod tests {
} }
async fn setup_multi_pool_bucket_test_env() -> (tempfile::TempDir, Arc<ECStore>) { async fn setup_multi_pool_bucket_test_env() -> (tempfile::TempDir, Arc<ECStore>) {
setup_bucket_quorum_test_env(&[4, 4], None).await
}
async fn setup_bucket_quorum_test_env(
drives_per_pool: &[usize],
standard_parity: Option<usize>,
) -> (tempfile::TempDir, Arc<ECStore>) {
let temp_dir = tempfile::tempdir().expect("multi-pool bucket test directory should be created"); let temp_dir = tempfile::tempdir().expect("multi-pool bucket test directory should be created");
let mut pools = Vec::new(); let mut pools = Vec::new();
for pool_index in 0..2 { for (pool_index, &drive_count) in drives_per_pool.iter().enumerate() {
let mut endpoints = Vec::new(); let mut endpoints = Vec::new();
for disk_index in 0..4 { for disk_index in 0..drive_count {
let disk_path = temp_dir.path().join(format!("pool{pool_index}-disk{disk_index}")); let disk_path = temp_dir.path().join(format!("pool{pool_index}-disk{disk_index}"));
tokio::fs::create_dir_all(&disk_path) tokio::fs::create_dir_all(&disk_path)
.await .await
@@ -1378,7 +1432,7 @@ mod tests {
pools.push(PoolEndpoints { pools.push(PoolEndpoints {
legacy: false, legacy: false,
set_count: 1, set_count: 1,
drives_per_set: 4, drives_per_set: drive_count,
endpoints: Endpoints::from(endpoints), endpoints: Endpoints::from(endpoints),
cmd_line: format!("bucket-test-pool-{pool_index}"), cmd_line: format!("bucket-test-pool-{pool_index}"),
platform: format!("OS: {} | Arch: {}", std::env::consts::OS, std::env::consts::ARCH), platform: format!("OS: {} | Arch: {}", std::env::consts::OS, std::env::consts::ARCH),
@@ -1399,9 +1453,12 @@ mod tests {
) )
.await .await
.expect("multi-pool ECStore should initialize"); .expect("multi-pool ECStore should initialize");
let storage_class = let mut storage_class_kvs = rustfs_config::server_config::KVS::new();
crate::config::storageclass::lookup_config_for_pools_without_env(&rustfs_config::server_config::KVS::new(), &[4, 4]) if let Some(parity) = standard_parity {
.expect("multi-pool storage class should match both four-disk pools"); storage_class_kvs.insert(crate::config::storageclass::CLASS_STANDARD.to_string(), format!("EC:{parity}"));
}
let storage_class = crate::config::storageclass::lookup_config_for_pools_without_env(&storage_class_kvs, drives_per_pool)
.expect("storage class should match every test erasure set");
for pool in &ecstore.pools { for pool in &ecstore.pools {
for set in &pool.disk_set { for set in &pool.disk_set {
set.set_test_storage_class_config(storage_class.clone()); set.set_test_storage_class_config(storage_class.clone());
@@ -2067,6 +2124,218 @@ mod tests {
} }
} }
#[tokio::test]
#[serial]
async fn bucket_info_read_quorum_tracks_erasure_layout() {
for (drive_count, parity) in [(2, 1), (3, 1), (4, 2), (5, 2), (6, 3), (8, 4), (6, 2), (12, 6)] {
let (_temp_dir, store) = setup_bucket_quorum_test_env(&[drive_count], Some(parity)).await;
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("read-quorum-{drive_count}-{parity}");
let object = "uncached-object";
let body = b"erasure read quorum must follow the persisted layout".repeat(32_768);
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("healthy namespace should accept bucket creation");
store
.put_object(&bucket, object, &mut PutObjReader::from_vec(body.clone()), &ObjectOptions::default())
.await
.expect("healthy erasure set should accept the seed object");
let set = &store.pools[0].disk_set[0];
let lock = set
.new_ns_lock(&bucket, object)
.await
.expect("seed namespace lock should resolve");
drop(
lock.get_write_lock(Duration::from_secs(30))
.await
.expect("seed physical fanout must finish before taking disks offline"),
);
if (drive_count, parity) == (6, 3) {
let mut kvs = rustfs_config::server_config::KVS::new();
kvs.insert(crate::config::storageclass::CLASS_STANDARD.to_string(), "EC:2".to_string());
set.set_test_storage_class_config(
crate::config::storageclass::lookup_config_for_pools_without_env(&kvs, &[drive_count])
.expect("a later storage-class change must not raise old objects' read quorum"),
);
}
let offline_indexes = (0..parity).collect::<Vec<_>>();
let offline = take_set_disks_offline(&store, set, &offline_indexes).await;
let info = store
.get_bucket_info(&bucket, &BucketOptions::default())
.await
.expect("bucket validation must admit the object's exact read quorum");
assert_eq!(info.name, bucket);
let mut reader = store
.get_object_reader(&bucket, object, None, Default::default(), &ObjectOptions::default())
.await
.expect("the persisted layout should remain readable at its exact data-shard quorum");
let mut restored = Vec::new();
reader
.stream
.read_to_end(&mut restored)
.await
.expect("quorum read should reconstruct the body");
assert_eq!(restored, body, "layout {drive_count}/{parity} must retain exact object contents");
drop(reader);
if drive_count - parity == drive_count / 2 {
let error = store
.get_bucket_info_from_sets(&bucket, &BucketOptions::default())
.await
.expect_err("bucket mutations must retain their majority namespace check");
assert_eq!(error, StorageError::ErasureWriteQuorum);
}
let below_quorum = take_set_disks_offline(&store, set, &[parity]).await;
let read = store
.get_object_reader(&bucket, object, None, Default::default(), &ObjectOptions::default())
.await;
match read {
Ok(mut reader) => assert!(
reader.stream.read_to_end(&mut Vec::new()).await.is_err(),
"layout {drive_count}/{parity} must reject fewer than its data-shard quorum"
),
Err(error) => assert!(
matches!(error, StorageError::ErasureReadQuorum | StorageError::InsufficientReadQuorum(_, _)),
"a missing shard must report read quorum loss, got {error}"
),
}
restore_set_disks(&store, set, below_quorum).await;
restore_set_disks(&store, set, offline).await;
}
}
#[tokio::test]
#[serial]
async fn bucket_info_read_quorum_is_scoped_to_each_erasure_set() {
let (_temp_dir, store) = setup_bucket_quorum_test_env(&[4, 6], None).await;
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = "read-quorum-mixed-pools";
store
.make_bucket(bucket, &MakeBucketOptions::default())
.await
.expect("healthy pools should accept bucket creation");
let first_set = &store.pools[0].disk_set[0];
let second_set = &store.pools[1].disk_set[0];
let first_offline = take_set_disks_offline(&store, first_set, &[0, 1]).await;
let second_offline = take_set_disks_offline(&store, second_set, &[0, 1, 2]).await;
store
.get_bucket_info(bucket, &BucketOptions::default())
.await
.expect("each set independently satisfies its namespace read quorum");
for (set, extra_disk) in [(first_set, 2), (second_set, 3)] {
let extra_offline = take_set_disks_offline(&store, set, &[extra_disk]).await;
assert_eq!(
store
.get_bucket_info(bucket, &BucketOptions::default())
.await
.expect_err("another pool must not subsidize a set below its read quorum"),
StorageError::ErasureReadQuorum
);
restore_set_disks(&store, set, extra_offline).await;
}
restore_set_disks(&store, first_set, first_offline).await;
restore_set_disks(&store, second_set, second_offline).await;
}
#[tokio::test]
#[serial]
async fn bucket_info_read_quorum_requires_authoritative_metadata() {
for state in ["missing", "corrupt", "foreign", "incarnation"] {
let (_temp_dir, store) = setup_bucket_quorum_test_env(&[4], None).await;
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = format!("read-quorum-{state}-metadata");
let mut metadata = if state == "missing" {
store
.make_bucket_on_sets(&bucket, &MakeBucketOptions::default())
.await
.expect("simulate directories left before bucket metadata is published");
BucketMetadata::new(&bucket)
} else {
store
.make_bucket(&bucket, &MakeBucketOptions::default())
.await
.expect("healthy bucket should publish metadata");
metadata_sys::get_in(&store.ctx, &bucket)
.await
.expect("seed metadata should be cached")
.as_ref()
.clone()
};
let path = metadata.save_file_path();
match state {
"corrupt" => crate::config::com::save_config(store.clone(), &path, b"corrupt".to_vec())
.await
.expect("persist corrupt metadata while the cached copy remains valid"),
"foreign" => {
metadata.name = "different-bucket".to_string();
let mut encoded = vec![1, 0, 1, 0];
encoded.extend(metadata.marshal_msg().expect("foreign metadata should encode"));
crate::config::com::save_config(store.clone(), &path, encoded)
.await
.expect("persist metadata for a different bucket at the requested path");
}
"incarnation" => crate::bucket::metadata::save_bucket_incarnation(store.clone(), &bucket, Uuid::new_v4())
.await
.expect("persist a different bucket generation"),
_ => {}
}
let set = &store.pools[0].disk_set[0];
let offline = take_set_disks_offline(&store, set, &[0, 1]).await;
let error = store
.get_bucket_info(&bucket, &BucketOptions::default())
.await
.expect_err("read admission must not trust residual directories or cached metadata");
match state {
"missing" => assert_eq!(error, StorageError::ErasureReadQuorum),
"foreign" => assert_eq!(error, StorageError::FileCorrupt),
"incarnation" => assert!(error.to_string().contains("sidecar does not match bucket metadata")),
"corrupt" => assert!(error.to_string().contains("format invalid"), "unexpected corruption error: {error}"),
_ => unreachable!(),
}
restore_set_disks(&store, set, offline).await;
}
}
#[tokio::test]
#[serial]
async fn bucket_info_read_quorum_accepts_persisted_legacy_metadata() {
let (_temp_dir, store) = setup_bucket_quorum_test_env(&[4], None).await;
metadata_sys::init_bucket_metadata_sys(store.clone(), Vec::new()).await;
let bucket = "interop";
store
.make_bucket_on_sets(bucket, &MakeBucketOptions::default())
.await
.expect("legacy bucket directories should exist");
let hex = include_str!("../../tests/fixtures/minio/bucket_metadata.blob.hex")
.split_whitespace()
.collect::<String>();
let body = (0..hex.len())
.step_by(2)
.map(|index| u8::from_str_radix(&hex[index..index + 2], 16).expect("pinned MinIO metadata fixture"))
.collect();
crate::config::com::save_config(store.clone(), &BucketMetadata::new(bucket).save_file_path(), body)
.await
.expect("legacy metadata should be persisted without an incarnation sidecar");
let set = &store.pools[0].disk_set[0];
let offline = take_set_disks_offline(&store, set, &[0, 1]).await;
let info = store
.get_bucket_info(bucket, &BucketOptions::default())
.await
.expect("persisted MinIO metadata should authorize reads at the namespace read quorum");
assert_eq!(info.name, bucket);
assert!(info.versioning);
assert!(info.object_locking);
restore_set_disks(&store, set, offline).await;
}
#[tokio::test] #[tokio::test]
#[serial] #[serial]
async fn bucket_namespace_reads_report_missing_when_every_set_is_absent() { async fn bucket_namespace_reads_report_missing_when_every_set_is_absent() {
@@ -2090,6 +2359,7 @@ mod tests {
#[serial] #[serial]
async fn bucket_namespace_reads_fail_closed_when_any_set_loses_quorum() { async fn bucket_namespace_reads_fail_closed_when_any_set_loses_quorum() {
let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await; let (_temp_dir, ecstore) = setup_multi_pool_bucket_test_env().await;
metadata_sys::init_bucket_metadata_sys(ecstore.clone(), Vec::new()).await;
let bucket = format!("degraded-expansion-{}", Uuid::new_v4().simple()); let bucket = format!("degraded-expansion-{}", Uuid::new_v4().simple());
ecstore.pools[0].disk_set[0] ecstore.pools[0].disk_set[0]
.make_bucket(&bucket, &MakeBucketOptions::default()) .make_bucket(&bucket, &MakeBucketOptions::default())
@@ -2097,6 +2367,7 @@ mod tests {
.expect("bucket should be created in the original pool only"); .expect("bucket should be created in the original pool only");
ecstore.pools[1].disk_set[0].disks.write().await[0] = None; ecstore.pools[1].disk_set[0].disks.write().await[0] = None;
ecstore.pools[1].disk_set[0].disks.write().await[1] = None; ecstore.pools[1].disk_set[0].disks.write().await[1] = None;
ecstore.pools[1].disk_set[0].disks.write().await[2] = None;
let list_err = ecstore let list_err = ecstore
.list_bucket(&BucketOptions::default()) .list_bucket(&BucketOptions::default())
@@ -2108,7 +2379,7 @@ mod tests {
.get_bucket_info(&bucket, &BucketOptions::default()) .get_bucket_info(&bucket, &BucketOptions::default())
.await .await
.expect_err("bucket validation must fail when an expansion pool is unavailable"); .expect_err("bucket validation must fail when an expansion pool is unavailable");
assert_eq!(info_err, StorageError::ErasureWriteQuorum); assert_eq!(info_err, StorageError::ErasureReadQuorum);
} }
#[tokio::test] #[tokio::test]
+5 -28
View File
@@ -577,13 +577,16 @@ impl KmsServiceManager {
Some(service_version.probe_worker.as_ref()?.status()) Some(service_version.probe_worker.as_ref()?.status())
} }
/// Health check for the KMS service /// Check backend health without changing the service lifecycle state.
///
/// A transient backend failure leaves the published service available for
/// subsequent checks and operations. Readiness uses the background probe
/// to evaluate backend availability independently of lifecycle state.
pub async fn health_check(&self) -> Result<bool> { pub async fn health_check(&self) -> Result<bool> {
let checked_state = self.state.load_full(); let checked_state = self.state.load_full();
match checked_state.current_service.as_ref() { match checked_state.current_service.as_ref() {
Some(service_version) => { Some(service_version) => {
let manager = service_version.manager.clone(); let manager = service_version.manager.clone();
let checked_version = service_version.version;
// Perform health check on the backend // Perform health check on the backend
match manager.health_check().await { match manager.health_check().await {
Ok(healthy) => { Ok(healthy) => {
@@ -594,8 +597,6 @@ impl KmsServiceManager {
} }
Err(e) => { Err(e) => {
error!("KMS health check error: {}", e); error!("KMS health check error: {}", e);
let _guard = self.lifecycle_mutex.lock().await;
self.mark_health_error_if_current(checked_version, &e);
Err(e) Err(e)
} }
} }
@@ -739,17 +740,6 @@ impl KmsServiceManager {
task: std::sync::Mutex::new(Some(task)), task: std::sync::Mutex::new(Some(task)),
})) }))
} }
fn mark_health_error_if_current(&self, checked_version: u64, error: &KmsError) {
let current = self.state.load_full();
if current.current_service.as_ref().map(|version| version.version) == Some(checked_version) {
self.state.store(Arc::new(RuntimeState {
config: current.config.clone(),
status: KmsServiceStatus::Error(format!("Health check failed: {error}")),
current_service: current.current_service.clone(),
}));
}
}
} }
impl Default for KmsServiceManager { impl Default for KmsServiceManager {
@@ -1004,19 +994,6 @@ mod tests {
assert!(manager.get_service_version().await.expect("restarted version") > first_version); assert!(manager.get_service_version().await.expect("restarted version") > first_version);
} }
#[tokio::test]
async fn stale_health_failure_cannot_poison_new_service_status() {
let manager = KmsServiceManager::new();
manager.configure(static_config("key-a", 0x11)).await.expect("configure");
manager.start().await.expect("start");
let old_version = manager.get_service_version().await.expect("old version");
manager.restart().await.expect("restart");
manager.mark_health_error_if_current(old_version, &KmsError::backend_error("stale failure"));
assert_eq!(manager.get_status().await, KmsServiceStatus::Running);
}
#[tokio::test] #[tokio::test]
async fn forbidden_local_master_key_change_preserves_running_config_and_service() { async fn forbidden_local_master_key_change_preserves_running_config_and_service() {
use crate::types::{CreateKeyRequest, KeyUsage}; use crate::types::{CreateKeyRequest, KeyUsage};
+38
View File
@@ -75,6 +75,44 @@ fn unreachable_vault_config() -> KmsConfig {
} }
} }
#[tokio::test]
async fn transient_health_failure_does_not_latch_the_service_status() {
let kms = TestKms::local().await;
let manager = kms.manager();
let service = manager.get_encryption_service().await.expect("running service");
let version = manager.get_service_version().await.expect("running version");
assert!(manager.health_check().await.expect("initial backend health"));
// Move only this test's keys out of reach, then restore the same backend.
let key_dir = kms.key_dir().expect("local key directory");
let outage = tempfile::TempDir::new().expect("temporary outage directory");
let hidden_keys = outage.path().join("keys");
tokio::fs::rename(&key_dir, &hidden_keys)
.await
.expect("make backend unavailable");
let failure = manager.health_check().await;
let outage_status = manager.get_status().await;
tokio::fs::rename(&hidden_keys, &key_dir).await.expect("restore backend");
assert!(failure.is_err(), "the outage must surface as a health-check error");
assert!(manager.health_check().await.expect("backend recovers without restart"));
assert!(Arc::ptr_eq(
&service,
&manager.get_encryption_service().await.expect("service survives the outage")
));
assert_eq!(manager.get_service_version().await, Some(version));
assert_eq!(
manager.get_status().await,
KmsServiceStatus::Running,
"a recovered backend must not leave service-status and readiness latched in Error"
);
assert_eq!(
outage_status,
KmsServiceStatus::Running,
"backend health does not change the running service's lifecycle state"
);
}
#[tokio::test] #[tokio::test]
async fn starting_against_an_unreachable_backend_fails_without_publishing_a_service() { async fn starting_against_an_unreachable_backend_fails_without_publishing_a_service() {
let manager = KmsServiceManager::new(); let manager = KmsServiceManager::new();
+87
View File
@@ -0,0 +1,87 @@
#!/usr/bin/env bash
# Publish the nightly DEB/RPM as assets of the rolling `nightly` release on
# rustfs/auto-testing, replacing the previous build's files in place.
#
# Required environment:
# ASSETS_TOKEN token with contents:write on rustfs/auto-testing
# DEB_FILE path to the built .deb
# RPM_FILE path to the built .rpm
# DEB_DATE build date (YYYY-MM-DD)
# BUILD_REF branch/ref the nightly was built from
# Optional environment:
# GITHUB_SHA / GITHUB_RUN_ID / GITHUB_REPOSITORY / GITHUB_SERVER_URL
#
# Plain curl + python3 by design: the nightly build fleet has no gh CLI.
set -euo pipefail
: "${ASSETS_TOKEN:?ASSETS_TOKEN is required}"
: "${DEB_FILE:?DEB_FILE is required}"
: "${RPM_FILE:?RPM_FILE is required}"
: "${DEB_DATE:?DEB_DATE is required}"
: "${BUILD_REF:?BUILD_REF is required}"
for f in "${DEB_FILE}" "${RPM_FILE}"; do
[ -f "$f" ] || { echo "missing package: $f" >&2; exit 1; }
done
API="https://api.github.com/repos/rustfs/auto-testing"
UPLOADS="https://uploads.github.com/repos/rustfs/auto-testing/releases"
AUTH="Authorization: token ${ASSETS_TOKEN}"
SOURCE_SHA="$(git rev-parse HEAD 2>/dev/null || echo "${GITHUB_SHA:-unknown}")"
release_id="$(curl -fsS --retry 3 -H "${AUTH}" "${API}/releases/tags/nightly" \
| python3 -c 'import json,sys; print(json.load(sys.stdin).get("id", ""))' 2>/dev/null || true)"
if [ -z "${release_id}" ]; then
echo "creating the rolling nightly release"
release_id="$(curl -fsS --retry 3 -X POST -H "${AUTH}" -H "Content-Type: application/json" \
-d '{"tag_name":"nightly","name":"Nightly builds","body":"Rolling nightly builds. Assets are replaced on every build; the release body documents the provenance of the current files."}' \
"${API}/releases" | python3 -c 'import json,sys; print(json.load(sys.stdin)["id"])')"
fi
[ -n "${release_id}" ] || { echo "could not resolve the nightly release id" >&2; exit 1; }
upload_asset() {
local name="$1" file="$2" asset_id
asset_id="$(curl -fsS --retry 3 -H "${AUTH}" "${API}/releases/tags/nightly" \
| ASSET_NAME="${name}" python3 -c '
import json, sys, os
d = json.load(sys.stdin)
name = os.environ["ASSET_NAME"]
print(next((a["id"] for a in d.get("assets", []) if a["name"] == name), ""))')"
if [ -n "${asset_id}" ]; then
curl -fsS --retry 3 -X DELETE -H "${AUTH}" "${API}/releases/assets/${asset_id}" >/dev/null
fi
curl -fsS --retry 3 --max-time 900 -X POST \
-H "${AUTH}" -H "Content-Type: application/octet-stream" \
--data-binary "@${file}" \
"${UPLOADS}/${release_id}/assets?name=${name}" >/dev/null
echo "uploaded ${name}"
}
upload_asset "rustfs-nightly-latest.deb" "${DEB_FILE}"
upload_asset "rustfs-nightly-latest.rpm" "${RPM_FILE}"
DEB_SHA="$(sha256sum "${DEB_FILE}" | cut -d ' ' -f 1)"
RPM_SHA="$(sha256sum "${RPM_FILE}" | cut -d ' ' -f 1)"
RUN_URL="${GITHUB_SERVER_URL:-https://github.com}/${GITHUB_REPOSITORY:-/rustfs/rustfs}/actions/runs/${GITHUB_RUN_ID:-0}"
export BUILD_REF SOURCE_SHA DEB_DATE DEB_FILE RPM_FILE DEB_SHA RPM_SHA RUN_URL
python3 - << 'PY' > /tmp/release-body.json
import json, os
e = os.environ
deb_mb = os.path.getsize(e["DEB_FILE"]) // 1048576
rpm_mb = os.path.getsize(e["RPM_FILE"]) // 1048576
body = (
f"Nightly build from `{e['BUILD_REF']}@{e['SOURCE_SHA'][:12]}`, built {e['DEB_DATE']}.\n\n"
f"[Build run]({e['RUN_URL']}). The `latest` assets are replaced in place on every nightly.\n\n"
f"| Asset | Size | SHA256 |\n|---|---|---|\n"
f"| rustfs-nightly-latest.deb (={e['DEB_FILE']}) | {deb_mb} MB | `{e['DEB_SHA']}` |\n"
f"| rustfs-nightly-latest.rpm (={e['RPM_FILE']}) | {rpm_mb} MB | `{e['RPM_SHA']}` |\n"
)
print(json.dumps({"body": body}))
PY
curl -fsS --retry 3 -X PATCH -H "${AUTH}" -H "Content-Type: application/json" \
--data-binary @/tmp/release-body.json "${API}/releases/${release_id}" >/dev/null
echo "published ${DEB_FILE} and ${RPM_FILE} to the rustfs/auto-testing 'nightly' release"
+6 -5
View File
@@ -163,12 +163,13 @@ SH
self.assertFalse(self.store.exists()) self.assertFalse(self.store.exists())
self.assertEqual(list(self.root.glob("nightly-awscli.*")), []) self.assertEqual(list(self.root.glob("nightly-awscli.*")), [])
def test_checkout_sha_mismatch_fails_before_upload(self): def test_manifest_advertises_checked_out_head_even_when_github_sha_differs(self):
# With a ref override (NIGHTLY_BRANCH variable / dispatch `branch`
# input) the checked-out HEAD intentionally differs from GITHUB_SHA;
# the candidate manifest must record the tree that was built.
result = self.run_publish(GITHUB_SHA="f" * 40) result = self.run_publish(GITHUB_SHA="f" * 40)
self.assertNotEqual(result.returncode, 0) self.assertEqual(result.returncode, 0, result.stderr)
self.assertIn("Checkout SHA", result.stderr) self.assertEqual(self.manifest()["source_sha"], self.sha)
self.assertFalse(self.output.exists())
self.assertFalse((self.root / "aws.log").exists())
def test_same_date_builds_and_reruns_keep_distinct_candidates(self): def test_same_date_builds_and_reruns_keep_distinct_candidates(self):
urls = [] urls = []