From a4795e6b0cda4c02fc77f7d406b8e0a45b6a0125 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 4 Sep 2026 12:10:32 +0000 Subject: [PATCH] test(e2e): add 4-node 4-disk distributed Actions suite Add a nightly e2e-distributed lane that boots localhost 4-node clusters and covers S3, object lock, versioning, replication, quota, expand, decommission, rebalance, site replication, concurrency, and chaos. Co-authored-by: RustFS --- .config/e2e-distributed-selection.txt | 2 + .config/nextest.toml | 29 + .github/scheduled-validations.json | 5 + .github/workflows/e2e-distributed.yml | 117 +++ .../scheduled-validation-watchdog.yml | 1 + crates/e2e_test/README.md | 4 + crates/e2e_test/src/common.rs | 63 ++ crates/e2e_test/src/distributed/chaos_test.rs | 102 +++ .../distributed/concurrency_stability_test.rs | 57 ++ .../concurrent_data_movement_test.rs | 65 ++ .../data_integrity_movement_test.rs | 43 ++ .../expand_decommission_rebalance_test.rs | 74 ++ crates/e2e_test/src/distributed/extra_test.rs | 108 +++ crates/e2e_test/src/distributed/harness.rs | 688 ++++++++++++++++++ crates/e2e_test/src/distributed/mod.rs | 34 + .../src/distributed/object_lock_test.rs | 137 ++++ .../src/distributed/observability_test.rs | 61 ++ .../src/distributed/replication_quota_test.rs | 98 +++ .../e2e_test/src/distributed/s3_basic_test.rs | 111 +++ .../s3_during_data_movement_test.rs | 80 ++ .../src/distributed/site_replication_test.rs | 87 +++ .../src/distributed/versioning_test.rs | 88 +++ crates/e2e_test/src/lib.rs | 5 + docs/testing/README.md | 5 +- docs/testing/ci-gates.md | 1 + docs/testing/distributed-e2e.md | 56 ++ 26 files changed, 2119 insertions(+), 2 deletions(-) create mode 100644 .config/e2e-distributed-selection.txt create mode 100644 .github/workflows/e2e-distributed.yml create mode 100644 crates/e2e_test/src/distributed/chaos_test.rs create mode 100644 crates/e2e_test/src/distributed/concurrency_stability_test.rs create mode 100644 crates/e2e_test/src/distributed/concurrent_data_movement_test.rs create mode 100644 crates/e2e_test/src/distributed/data_integrity_movement_test.rs create mode 100644 crates/e2e_test/src/distributed/expand_decommission_rebalance_test.rs create mode 100644 crates/e2e_test/src/distributed/extra_test.rs create mode 100644 crates/e2e_test/src/distributed/harness.rs create mode 100644 crates/e2e_test/src/distributed/mod.rs create mode 100644 crates/e2e_test/src/distributed/object_lock_test.rs create mode 100644 crates/e2e_test/src/distributed/observability_test.rs create mode 100644 crates/e2e_test/src/distributed/replication_quota_test.rs create mode 100644 crates/e2e_test/src/distributed/s3_basic_test.rs create mode 100644 crates/e2e_test/src/distributed/s3_during_data_movement_test.rs create mode 100644 crates/e2e_test/src/distributed/site_replication_test.rs create mode 100644 crates/e2e_test/src/distributed/versioning_test.rs create mode 100644 docs/testing/distributed-e2e.md diff --git a/.config/e2e-distributed-selection.txt b/.config/e2e-distributed-selection.txt new file mode 100644 index 000000000..d17c04cc3 --- /dev/null +++ b/.config/e2e-distributed-selection.txt @@ -0,0 +1,2 @@ +sha256-linux=3c6746b96264237f338499bfdb8005986bda9e18720e29527f10180cfb9f73d6 +sha256-darwin=3c6746b96264237f338499bfdb8005986bda9e18720e29527f10180cfb9f73d6 diff --git a/.config/nextest.toml b/.config/nextest.toml index c6eaf1908..7981a29d2 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -183,6 +183,13 @@ test-group = 'e2e-reliability' filter = 'package(e2e_test) & test(/^inline_fast_path_cluster_test::/)' test-group = 'e2e-inline-boundaries' +# 4-node 4-drive distributed Actions suite: each case starts four rustfs +# processes and up to sixteen data directories. Serialize across nextest's +# process boundary so several 4x4 clusters never overlap. +[[profile.default.overrides]] +filter = 'package(e2e_test) & test(/^distributed::/)' +test-group = 'e2e-cluster-nightly' + # Vault KMS tests share the fixed dev-server port 8200. serial_test's #[serial] # does not cross nextest process boundaries, so keep every Vault-backed test in # one group. @@ -526,6 +533,23 @@ path = "junit.xml" filter = 'package(e2e_test)' test-group = 'e2e-cluster-nightly' +# --------------------------------------------------------------------------- +# e2e-distributed profile — 4-node 4-disk Actions suite +# --------------------------------------------------------------------------- +# Nightly / dispatch lane owned by .github/workflows/e2e-distributed.yml. +# Each case starts four rustfs processes (and for site replication, two +# clusters). Serialized via e2e-cluster-nightly. Not a PR merge gate. +[profile.e2e-distributed] +default-filter = 'package(e2e_test) & test(/^distributed::/)' +fail-fast = false + +[profile.e2e-distributed.junit] +path = "junit.xml" + +[[profile.e2e-distributed.overrides]] +filter = 'package(e2e_test)' +test-group = 'e2e-cluster-nightly' + # --------------------------------------------------------------------------- # e2e-odm-interop profile — on-demand migration provider interop lane (ODM-20) # --------------------------------------------------------------------------- @@ -586,6 +610,10 @@ path = "junit.xml" # cluster-fault lane. heal_erasure_disk_rebuild is intentionally not # excluded here because backlog#2213 promotes core heal rebuild coverage to # this merge/main lane while retaining nightly coverage. +# * distributed:: — 4-node 4-disk Actions suite (S3, lock, versioning, +# replication, quota, observability, expand/decommission/rebalance, site +# replication, chaos). Owns [profile.e2e-distributed] and +# .github/workflows/e2e-distributed.yml. # * on_demand_migration::interop_test — the ODM-20 provider interoperability # cases, which are meaningless without a source: they run in the dedicated # [profile.e2e-odm-interop] lane below, where the workflow points them at a @@ -607,6 +635,7 @@ default-filter = """ package(e2e_test) & !test(/^protocols::/) & !test(/^(admin_timeout_regression_test|cluster_concurrency_test|cluster_multidrive_pool_test|degraded_listing_availability_test|namespace_lock_quorum_test|object_lambda_test|stale_multipart_cleanup_cluster_test)::/) + & !test(/^distributed::/) & !test(/^replication_extension_test::/) & !test(/^replication_target_matrix_test::/) & !test(/^on_demand_migration::(concurrency_test|fault_test|interop_test|real_source_test)::/) diff --git a/.github/scheduled-validations.json b/.github/scheduled-validations.json index 9ac7f2614..e55ef3d56 100644 --- a/.github/scheduled-validations.json +++ b/.github/scheduled-validations.json @@ -4,6 +4,11 @@ { "workflow": ".github/workflows/ci.yml", "max_age_hours": 192 }, { "workflow": ".github/workflows/coverage.yml", "max_age_hours": 192 }, { "workflow": ".github/workflows/e2e-replication-nightly.yml", "max_age_hours": 36 }, + { + "workflow": ".github/workflows/e2e-distributed.yml", + "max_age_hours": 36, + "never_ran_grace_until": "2026-09-18T00:00:00Z" + }, { "workflow": ".github/workflows/e2e-s3tests.yml", "max_age_hours": 192 }, { "workflow": ".github/workflows/fuzz.yml", "max_age_hours": 36 }, { "workflow": ".github/workflows/mint.yml", "max_age_hours": 192 }, diff --git a/.github/workflows/e2e-distributed.yml b/.github/workflows/e2e-distributed.yml new file mode 100644 index 000000000..e52ccbc31 --- /dev/null +++ b/.github/workflows/e2e-distributed.yml @@ -0,0 +1,117 @@ +# Copyright 2024 RustFS Team +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +# 4-node 4-disk distributed e2e lane. +# +# Each selected test starts a real localhost cluster via +# `RustFSTestClusterEnvironment` (4 processes; 4 drives per node unless the +# case is a two-site 4-node 1-drive pair). Membership is +# `[profile.e2e-distributed]` in `.config/nextest.toml`. This is not a required +# merge check: it is the scheduled/dispatch counterpart to the hardware +# functional chain that currently clones rustfs/auto-testing onto three VMs. + +name: e2e-distributed + +on: + workflow_dispatch: + inputs: + filter: + description: "Optional nextest -E filter (default: the whole e2e-distributed profile)" + required: false + default: "" + schedule: + # 05:53 UTC nightly — clear of e2e-nightly (04:29) and ODM interop (05:23). + - cron: "53 5 * * *" + +permissions: + contents: read + +concurrency: + group: ${{ github.workflow }}-${{ github.ref }} + cancel-in-progress: ${{ github.event_name != 'schedule' }} + +jobs: + distributed: + name: Distributed 4-node 4-disk e2e + runs-on: sm-standard-4 + timeout-minutes: 180 + env: + FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true" + steps: + - name: Checkout repository + uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7 + with: + persist-credentials: false + + - name: Setup Rust environment + uses: ./.github/actions/setup + with: + rust-version: stable + cache-shared-key: ci-e2e-distributed + cache-save-if: ${{ github.ref == 'refs/heads/main' }} + install-build-packaging-tools: 'false' + + - name: Build rustfs binary + run: | + cargo build -p rustfs --bins + : > target/debug/rustfs.features + + - name: Verify distributed e2e membership + env: + NEXTEST_LISTING: ${{ runner.temp }}/rustfs-e2e-distributed-list.json + run: | + cargo nextest list --profile e2e-distributed -p e2e_test --message-format json > "${NEXTEST_LISTING}" + python3 ./scripts/check_test_wiring.py --check-profile e2e-distributed "${NEXTEST_LISTING}" + + - name: Run distributed 4-node e2e suite + env: + RUSTFS_E2E_LOG_DIR: ${{ runner.temp }}/rustfs-e2e-distributed-logs + run: | + set -euo pipefail + FILTER='${{ github.event.inputs.filter }}' + if [ -n "${FILTER}" ]; then + cargo nextest run --profile e2e-distributed -p e2e_test -E "${FILTER}" + else + cargo nextest run --profile e2e-distributed -p e2e_test --no-tests=fail + fi + + - name: Upload distributed e2e diagnostics + if: always() + uses: actions/upload-artifact@b7c566a772e6b6bfb58ed0dc250532a479d7789f # v6 + with: + name: e2e-distributed-${{ github.run_number }} + path: | + target/nextest/e2e-distributed/junit.xml + ${{ runner.temp }}/rustfs-e2e-distributed-list.json + ${{ runner.temp }}/rustfs-e2e-distributed-logs/ + retention-days: 7 + if-no-files-found: warn + + alert-on-failure: + name: Alert on scheduled failure + needs: [distributed] + if: always() && github.event_name == 'schedule' && contains(needs.*.result, 'failure') + runs-on: ubuntu-latest + timeout-minutes: 10 + permissions: + contents: read + issues: write + steps: + - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7 + with: + persist-credentials: false + - name: Open or update failure-tracking issue + uses: ./.github/actions/schedule-failure-issue + with: + github-token: ${{ secrets.GITHUB_TOKEN }} diff --git a/.github/workflows/scheduled-validation-watchdog.yml b/.github/workflows/scheduled-validation-watchdog.yml index e778ec640..154d541d3 100644 --- a/.github/workflows/scheduled-validation-watchdog.yml +++ b/.github/workflows/scheduled-validation-watchdog.yml @@ -22,6 +22,7 @@ on: - "Continuous Integration" - "coverage" - "e2e-nightly" + - "e2e-distributed" - "e2e-s3tests" - "Fuzz" - "mint" diff --git a/crates/e2e_test/README.md b/crates/e2e_test/README.md index ed5a65d97..0587fccf3 100644 --- a/crates/e2e_test/README.md +++ b/crates/e2e_test/README.md @@ -26,6 +26,7 @@ Registered in [`src/lib.rs`](src/lib.rs). Grouped by concern: | **protocols** | [`src/protocols/`](src/protocols) | FTPS, WebDAV, SFTP compliance. Fixed ports, own guide: [`src/protocols/README.md`](src/protocols/README.md) | | **reliant** | [`src/reliant/`](src/reliant) | Tests that reuse an **externally started** server (SQL/select, conditional writes, lifecycle, deleted-object reads, node-interact). Run via [`scripts/run_e2e_tests.sh`](../../scripts/run_e2e_tests.sh); see [`src/reliant/README.md`](src/reliant/README.md) | | **cluster** | `cluster_concurrency_test`, `stale_multipart_cleanup_cluster_test`, `namespace_lock_quorum_test`, `admin_timeout_regression_test`, `object_lambda_test`, `replication_extension_test`, `tier_stats_cluster_test` | Multi-node scenarios via `RustFSTestClusterEnvironment` | +| **distributed 4×4** | [`src/distributed/`](src/distributed) | Nightly `e2e-distributed` lane: S3, object lock/WORM, versioning, bucket/site replication, quota, expand/decommission/rebalance, concurrency, chaos. Map: [`docs/testing/distributed-e2e.md`](../../docs/testing/distributed-e2e.md) | | **chaos / reliability** | [`src/chaos.rs`](src/chaos.rs), `reliability_disk_fault_test`, `heal_erasure_disk_rebuild_test`, `server_startup_failfast_test` | Disk offline/replace/corrupt, EC rebuild, heal, fail-fast startup | | **upgrade compatibility** | `upgrade_compatibility_test` | Pinned previous-release writes followed by current-build reads on the same data directory | @@ -171,6 +172,7 @@ the same profile for membership and execution with one nightly worker. | KMS suite | `e2e-full` job, merge queue + main | **Active** | | Direct and mixed-version rolling upgrades from pinned previous release | `e2e-upgrade.yml`, storage-sensitive PRs + release tags + weekly | **Active** | | Cluster faults (`e2e-nightly` profile) | consolidated nightly workflow | **Active** (backlog#1149 ci-7) | +| Distributed 4-node 4-disk (`e2e-distributed` profile) | `.github/workflows/e2e-distributed.yml` | **Active** (nightly / dispatch; not a merge gate) | | Protocols (FTPS/WebDAV/SFTP) | consolidated nightly workflow, serial | **Active** (backlog#1149 ci-7) | | Replication (fast subset) | `e2e-smoke` profile, `e2e-tests` job, every PR | **Active** (backlog#1147 repl-1) | | Replication (slow + multi-node) | `e2e-repl-nightly` profile, consolidated nightly workflow | **Active** (backlog#1147 repl-1) | @@ -191,6 +193,8 @@ cargo nextest run --profile e2e-smoke -p e2e_test cargo nextest run --profile e2e-full -p e2e_test # Cluster fault nightly lane cargo nextest run --profile e2e-nightly -p e2e_test +# 4-node 4-disk distributed lane (S3 / lock / versioning / replication / decommission / chaos) +cargo nextest run --profile e2e-distributed -p e2e_test # Replication nightly lane; awscurl is required for STS paths cargo nextest run --profile e2e-repl-nightly -p e2e_test # Fixed-port protocol nightly lane diff --git a/crates/e2e_test/src/common.rs b/crates/e2e_test/src/common.rs index b108cd074..daa744ab8 100644 --- a/crates/e2e_test/src/common.rs +++ b/crates/e2e_test/src/common.rs @@ -1700,6 +1700,69 @@ impl RustFSTestClusterEnvironment { Ok(()) } + /// Append a new single-node erasure pool to a stopped multi-pool cluster. + /// + /// Used to simulate pool expansion on localhost: every pool already owns + /// exactly one node with `drives_per_node >= 2` (the only multi-pool layout + /// the single-host `RUSTFS_VOLUMES` syntax can express). The new node is + /// allocated a fresh port and empty drive directories; callers must + /// [`Self::start`] afterwards so every process picks up the extended + /// volumes argument. Existing data directories are left untouched. + pub async fn append_single_node_pool(&mut self) -> Result> { + if self.nodes.iter().any(|node| node.process.is_some()) { + return Err("stop the cluster before appending a pool".into()); + } + if self.topology.drives_per_node < 2 { + return Err( + "append_single_node_pool requires drives_per_node >= 2 (the server parser rejects a single-drive ellipses pool)" + .into(), + ); + } + + let mut pools = self.topology.normalized_pools(); + for (pool_idx, nodes) in pools.iter().enumerate() { + if nodes.len() != 1 { + return Err(format!( + "pool {pool_idx} spans {} nodes; append_single_node_pool requires one node per pool", + nodes.len() + ) + .into()); + } + } + + let new_idx = self.nodes.len(); + let port = RustFSTestEnvironment::find_available_port().await?; + let address = format!("127.0.0.1:{port}"); + let data_dirs: Vec = (0..self.topology.drives_per_node) + .map(|drive| format!("{}/node{}/drive{}", self.temp_dir, new_idx, drive)) + .collect(); + for dir in &data_dirs { + fs::create_dir_all(dir).await?; + } + + self.nodes.push(ClusterNode { + url: format!("http://{address}"), + address, + data_dir: data_dirs[0].clone(), + data_dirs, + pool_idx: pools.len(), + process: None, + }); + pools.push(vec![new_idx]); + self.topology.node_count = self.nodes.len(); + self.topology.pools = pools; + self.node_extra_env.push(Vec::new()); + self.node_capture_log_paths.push(None); + self.volume_proxy_addresses.push(None); + + if !self.extra_env.iter().any(|(key, _)| key == "RUSTFS_UNSAFE_BYPASS_DISK_CHECK") { + self.extra_env + .push(("RUSTFS_UNSAFE_BYPASS_DISK_CHECK".to_string(), "true".to_string())); + } + + Ok(new_idx) + } + /// Gracefully stop one cluster node and wait for its process to exit. /// /// This is intentionally separate from [`Self::stop_node`]: the latter is diff --git a/crates/e2e_test/src/distributed/chaos_test.rs b/crates/e2e_test/src/distributed/chaos_test.rs new file mode 100644 index 000000000..5c066b1e7 --- /dev/null +++ b/crates/e2e_test/src/distributed/chaos_test.rs @@ -0,0 +1,102 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::harness::{ + DistCluster, DistLayout, TestResult, assert_object_bytes, bring_drive_online, put_object, retrying_get_equals, + take_drive_offline, unique_bucket, wait_for_ready, +}; +use crate::common::init_logging; +use crate::fault_proxy::FaultMode; +use std::time::Duration; + +#[tokio::test] +async fn kill_and_restart_node_preserves_objects() -> TestResult { + init_logging(); + let mut dist = DistCluster::start(DistLayout::FourByFour).await?; + let bucket = unique_bucket("killnode"); + dist.create_bucket(&bucket).await?; + let body = vec![0x11u8; 128 * 1024]; + put_object(&dist.client(0)?, &bucket, "keep.bin", body.clone()).await?; + + dist.cluster.stop_node(3)?; + retrying_get_equals(&dist.client(0)?, &bucket, "keep.bin", &body, Duration::from_secs(20)).await?; + + dist.cluster.start_node(3).await?; + wait_for_ready(&dist.cluster).await?; + assert_object_bytes(&dist.client(3)?, &bucket, "keep.bin", &body).await?; + Ok(()) +} + +#[tokio::test] +async fn full_cluster_restart_preserves_objects() -> TestResult { + init_logging(); + let mut dist = DistCluster::start(DistLayout::FourByFour).await?; + let bucket = unique_bucket("pwr"); + dist.create_bucket(&bucket).await?; + let body = vec![0x44u8; 64 * 1024]; + put_object(&dist.client(1)?, &bucket, "survive.bin", body.clone()).await?; + + dist.cluster.stop(); + dist.cluster.start().await?; + wait_for_ready(&dist.cluster).await?; + for node_idx in 0..dist.cluster.nodes.len() { + assert_object_bytes(&dist.client(node_idx)?, &bucket, "survive.bin", &body).await?; + } + Ok(()) +} + +#[tokio::test] +async fn offline_drive_then_replace_keeps_object_readable() -> TestResult { + init_logging(); + let dist = DistCluster::start(DistLayout::FourByFour).await?; + let bucket = unique_bucket("baddrive"); + dist.create_bucket(&bucket).await?; + let body = vec![0x22u8; 96 * 1024]; + put_object(&dist.client(1)?, &bucket, "durable.bin", body.clone()).await?; + + take_drive_offline(&dist.cluster, 0, 0)?; + retrying_get_equals(&dist.client(2)?, &bucket, "durable.bin", &body, Duration::from_secs(20)).await?; + bring_drive_online(&dist.cluster, 0, 0)?; + retrying_get_equals(&dist.client(3)?, &bucket, "durable.bin", &body, Duration::from_secs(20)).await?; + Ok(()) +} + +#[tokio::test] +async fn volume_proxy_blackhole_then_restore_keeps_s3_available() -> TestResult { + init_logging(); + let mut cluster = + crate::common::RustFSTestClusterEnvironment::with_topology(crate::common::ClusterTopology::single_pool_multidrive(4, 4)) + .await?; + let proxy = cluster.start_volume_proxy_for_node(1).await?; + cluster.start().await?; + cluster.create_test_bucket("chaos-net").await?; + let client = cluster.create_s3_client(0)?; + let body = vec![0x33u8; 32 * 1024]; + put_object(&client, "chaos-net", "via-proxy.bin", body.clone()).await?; + + proxy.set_mode(FaultMode::Blackhole); + retrying_get_equals( + &cluster.create_s3_client(2)?, + "chaos-net", + "via-proxy.bin", + &body, + Duration::from_secs(20), + ) + .await?; + + proxy.set_mode(FaultMode::Pass); + assert_object_bytes(&cluster.create_s3_client(3)?, "chaos-net", "via-proxy.bin", &body).await?; + proxy.shutdown().await; + Ok(()) +} diff --git a/crates/e2e_test/src/distributed/concurrency_stability_test.rs b/crates/e2e_test/src/distributed/concurrency_stability_test.rs new file mode 100644 index 000000000..e96a01cfa --- /dev/null +++ b/crates/e2e_test/src/distributed/concurrency_stability_test.rs @@ -0,0 +1,57 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::harness::{DistCluster, DistLayout, TestResult, assert_object_bytes, payload_for, put_object, unique_bucket}; +use crate::common::init_logging; +use std::sync::Arc; +use tokio::sync::Barrier; + +#[tokio::test] +async fn four_node_high_concurrency_puts_are_readable_from_every_node() -> TestResult { + init_logging(); + let dist = DistCluster::start(DistLayout::FourByFour).await?; + let bucket = unique_bucket("conc"); + dist.create_bucket(&bucket).await?; + let clients = Arc::new(dist.clients()?); + let barrier = Arc::new(Barrier::new(32)); + + let mut handles = Vec::new(); + for idx in 0..32 { + let clients = clients.clone(); + let barrier = barrier.clone(); + let bucket = bucket.clone(); + handles.push(tokio::spawn(async move { + barrier.wait().await; + let client = &clients[idx % clients.len()]; + let key = format!("c/{idx:02}.bin"); + let body = payload_for(&key, 16 * 1024); + put_object(client, &bucket, &key, body.clone()).await?; + Ok::<_, Box>((key, body)) + })); + } + + let mut inventory = Vec::new(); + for handle in handles { + inventory.push(handle.await??); + } + + for (node_idx, client) in clients.iter().enumerate() { + for (key, body) in &inventory { + assert_object_bytes(client, &bucket, key, body) + .await + .map_err(|error| format!("node {node_idx} failed to read {key}: {error}"))?; + } + } + Ok(()) +} diff --git a/crates/e2e_test/src/distributed/concurrent_data_movement_test.rs b/crates/e2e_test/src/distributed/concurrent_data_movement_test.rs new file mode 100644 index 000000000..9fc4019d4 --- /dev/null +++ b/crates/e2e_test/src/distributed/concurrent_data_movement_test.rs @@ -0,0 +1,65 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::harness::{ + DistCluster, DistLayout, TestResult, assert_inventory, payload_for, put_inventory, retrying_get_equals, retrying_put, + start_decommission, unique_bucket, wait_for_decommission_complete, +}; +use crate::common::init_logging; +use std::sync::Arc; +use std::time::Duration; +use tokio::sync::Barrier; + +#[tokio::test] +async fn concurrent_puts_during_decommission_do_not_lose_baseline_or_new_objects() -> TestResult { + init_logging(); + let dist = DistCluster::start(DistLayout::FourPoolFourDrive).await?; + let bucket = unique_bucket("concdecom"); + dist.create_bucket(&bucket).await?; + let baseline_client = dist.client(0)?; + let inventory = put_inventory(&baseline_client, &bucket, 10, 24 * 1024).await?; + + start_decommission(&dist.cluster, 0).await?; + + let clients = Arc::new(dist.clients()?); + let barrier = Arc::new(Barrier::new(16)); + let mut handles = Vec::new(); + for idx in 0..16 { + let clients = clients.clone(); + let barrier = barrier.clone(); + let bucket = bucket.clone(); + handles.push(tokio::spawn(async move { + barrier.wait().await; + let client = &clients[idx % clients.len()]; + let key = format!("live/{idx:02}.bin"); + let body = payload_for(&key, 8 * 1024); + retrying_put(client, &bucket, &key, body.clone(), Duration::from_secs(45)).await?; + Ok::<_, Box>((key, body)) + })); + } + + let mut live_objects = Vec::new(); + for handle in handles { + live_objects.push(handle.await??); + } + + wait_for_decommission_complete(&dist.cluster, 0, Duration::from_secs(180)).await?; + + let checker = dist.client(3)?; + assert_inventory(&checker, &bucket, &inventory).await?; + for (key, body) in live_objects { + retrying_get_equals(&checker, &bucket, &key, &body, Duration::from_secs(30)).await?; + } + Ok(()) +} diff --git a/crates/e2e_test/src/distributed/data_integrity_movement_test.rs b/crates/e2e_test/src/distributed/data_integrity_movement_test.rs new file mode 100644 index 000000000..3af367925 --- /dev/null +++ b/crates/e2e_test/src/distributed/data_integrity_movement_test.rs @@ -0,0 +1,43 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::harness::{ + DistCluster, DistLayout, TestResult, assert_inventory, put_inventory, sha256_hex, start_decommission, unique_bucket, + wait_for_decommission_complete, +}; +use crate::common::init_logging; +use std::time::Duration; + +#[tokio::test] +async fn decommission_does_not_alter_object_sha256_across_pools() -> TestResult { + init_logging(); + let dist = DistCluster::start(DistLayout::FourPoolFourDrive).await?; + let bucket = unique_bucket("integrity"); + dist.create_bucket(&bucket).await?; + let client = dist.client(0)?; + let inventory = put_inventory(&client, &bucket, 20, 64 * 1024).await?; + let before: Vec<(String, String)> = inventory.iter().map(|(key, body)| (key.clone(), sha256_hex(body))).collect(); + + start_decommission(&dist.cluster, 0).await?; + wait_for_decommission_complete(&dist.cluster, 0, Duration::from_secs(180)).await?; + + let after_client = dist.client(2)?; + assert_inventory(&after_client, &bucket, &inventory).await?; + for (key, expected_hash) in before { + let got = after_client.get_object().bucket(&bucket).key(&key).send().await?; + let body = got.body.collect().await?.into_bytes(); + assert_eq!(sha256_hex(body.as_ref()), expected_hash, "checksum changed for {key} after decommission"); + } + Ok(()) +} diff --git a/crates/e2e_test/src/distributed/expand_decommission_rebalance_test.rs b/crates/e2e_test/src/distributed/expand_decommission_rebalance_test.rs new file mode 100644 index 000000000..509b6f8a9 --- /dev/null +++ b/crates/e2e_test/src/distributed/expand_decommission_rebalance_test.rs @@ -0,0 +1,74 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::harness::{ + DistCluster, DistLayout, TestResult, assert_inventory, list_pools_json, put_inventory, start_decommission, start_rebalance, + unique_bucket, wait_for_decommission_complete, wait_for_rebalance_idle, +}; +use crate::common::init_logging; +use std::time::Duration; + +#[tokio::test] +async fn four_node_pool_expand_preserves_objects_then_rebalance() -> TestResult { + init_logging(); + let mut dist = DistCluster::start(DistLayout::TwoPoolFourDrive).await?; + let bucket = unique_bucket("expand"); + dist.create_bucket(&bucket).await?; + let client = dist.client(0)?; + let inventory = put_inventory(&client, &bucket, 12, 32 * 1024).await?; + assert_inventory(&client, &bucket, &inventory).await?; + + dist.cluster.stop(); + dist.cluster.append_single_node_pool().await?; + dist.cluster.append_single_node_pool().await?; + assert_eq!(dist.cluster.nodes.len(), 4); + dist.cluster.start().await?; + + let after_expand = dist.client(0)?; + assert_inventory(&after_expand, &bucket, &inventory).await?; + let peer = dist.client(3)?; + assert_inventory(&peer, &bucket, &inventory).await?; + + let _ = start_rebalance(&dist.cluster).await?; + // Rebalance may finish immediately on a tiny dataset; either idle or a + // started-then-completed status is success. A hard failure is not. + let _ = wait_for_rebalance_idle(&dist.cluster, Duration::from_secs(90)).await; + assert_inventory(&peer, &bucket, &inventory).await?; + Ok(()) +} + +#[tokio::test] +async fn four_pool_decommission_moves_objects_without_loss() -> TestResult { + init_logging(); + let dist = DistCluster::start(DistLayout::FourPoolFourDrive).await?; + let bucket = unique_bucket("decom"); + dist.create_bucket(&bucket).await?; + let client = dist.client(1)?; + let inventory = put_inventory(&client, &bucket, 16, 48 * 1024).await?; + + let pools_before = list_pools_json(&dist.cluster).await?; + let pool_count = pools_before + .as_array() + .map(Vec::len) + .or_else(|| pools_before.get("pools").and_then(serde_json::Value::as_array).map(Vec::len)) + .unwrap_or(4); + assert!(pool_count >= 4, "expected four pools before decommission: {pools_before}"); + + start_decommission(&dist.cluster, 0).await?; + wait_for_decommission_complete(&dist.cluster, 0, Duration::from_secs(180)).await?; + + let after = dist.client(3)?; + assert_inventory(&after, &bucket, &inventory).await?; + Ok(()) +} diff --git a/crates/e2e_test/src/distributed/extra_test.rs b/crates/e2e_test/src/distributed/extra_test.rs new file mode 100644 index 000000000..8d848b0dc --- /dev/null +++ b/crates/e2e_test/src/distributed/extra_test.rs @@ -0,0 +1,108 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::harness::{DistCluster, DistLayout, TestResult, assert_object_bytes, get_object_bytes, put_object, unique_bucket}; +use crate::common::init_logging; +use aws_sdk_s3::primitives::ByteStream; +use aws_sdk_s3::types::{CompletedMultipartUpload, CompletedPart}; + +#[tokio::test] +async fn four_node_four_drive_multipart_and_cross_node_listing_agree() -> TestResult { + init_logging(); + let dist = DistCluster::start(DistLayout::FourByFour).await?; + let bucket = unique_bucket("extra"); + dist.create_bucket(&bucket).await?; + let client = dist.client(0)?; + + let key = "multipart.bin"; + let part1 = vec![0x41u8; 5 * 1024 * 1024]; + let part2 = vec![0x42u8; 5 * 1024 * 1024]; + let upload = client.create_multipart_upload().bucket(&bucket).key(key).send().await?; + let upload_id = upload.upload_id().ok_or("missing upload id")?.to_string(); + + let uploaded1 = client + .upload_part() + .bucket(&bucket) + .key(key) + .upload_id(&upload_id) + .part_number(1) + .body(ByteStream::from(part1.clone())) + .send() + .await?; + let uploaded2 = client + .upload_part() + .bucket(&bucket) + .key(key) + .upload_id(&upload_id) + .part_number(2) + .body(ByteStream::from(part2.clone())) + .send() + .await?; + + client + .complete_multipart_upload() + .bucket(&bucket) + .key(key) + .upload_id(&upload_id) + .multipart_upload( + CompletedMultipartUpload::builder() + .parts( + CompletedPart::builder() + .part_number(1) + .e_tag(uploaded1.e_tag().unwrap_or_default()) + .build(), + ) + .parts( + CompletedPart::builder() + .part_number(2) + .e_tag(uploaded2.e_tag().unwrap_or_default()) + .build(), + ) + .build(), + ) + .send() + .await?; + + let mut expected = part1; + expected.extend_from_slice(&part2); + for node_idx in 0..dist.cluster.nodes.len() { + assert_object_bytes(&dist.client(node_idx)?, &bucket, key, &expected).await?; + } + + put_object(&client, &bucket, "list/a", b"a".to_vec()).await?; + put_object(&dist.client(2)?, &bucket, "list/b", b"b".to_vec()).await?; + let mut seen = Vec::new(); + for node_idx in 0..dist.cluster.nodes.len() { + let listed = dist + .client(node_idx)? + .list_objects_v2() + .bucket(&bucket) + .prefix("list/") + .send() + .await?; + let keys: Vec = listed + .contents() + .iter() + .filter_map(|object| object.key().map(str::to_string)) + .collect(); + seen.push(keys); + } + for keys in &seen[1..] { + assert_eq!(&seen[0], keys, "list results diverged across nodes: {seen:?}"); + } + + let got = get_object_bytes(&dist.client(3)?, &bucket, "list/a").await?; + assert_eq!(got, b"a"); + Ok(()) +} diff --git a/crates/e2e_test/src/distributed/harness.rs b/crates/e2e_test/src/distributed/harness.rs new file mode 100644 index 000000000..88bb5f8cd --- /dev/null +++ b/crates/e2e_test/src/distributed/harness.rs @@ -0,0 +1,688 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Shared 4-node distributed e2e helpers. +//! +//! Two localhost-expressible layouts cover the suite: +//! +//! * **4×4 single pool** (`four_by_four`) — four processes, four drives each, +//! one `DistErasure` pool (16 explicit volume endpoints). This is the +//! default S3 / lock / versioning / chaos topology. +//! * **4×4 four pool** (`four_pool_four_drive`) — four single-node pools of +//! four drives. Required for decommission/rebalance/expand, which the +//! server rejects on a single pool. +//! +//! Genuine multi-node *striped* pools still need multi-host CI (backlog +//! #1313 / #1314). Site replication uses two 4-node 1-drive clusters so the +//! process count stays at eight rather than sixteen. + +use crate::common::{ + ClusterTopology, FAST_DATA_USAGE_SCANNER_ENV, RustFSTestClusterEnvironment, admin_request, local_http_client, + replication_fast_env, signed_request, +}; +use crate::replication_extension_test::LOOPBACK_REPLICATION_TARGET_ENV; +use aws_sdk_s3::Client; +use aws_sdk_s3::primitives::ByteStream; +use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration}; +use http::{Method, StatusCode}; +use sha2::{Digest, Sha256}; +use std::collections::BTreeMap; +use std::path::Path; +use std::time::Duration; +use tokio::time::{Instant, sleep}; +use uuid::Uuid; + +pub(crate) type TestResult = Result>; + +pub(crate) const NODE_COUNT: usize = 4; +pub(crate) const DRIVES_PER_NODE: usize = 4; + +#[derive(Clone, Copy, Debug)] +pub(crate) enum DistLayout { + /// 4 nodes × 4 drives, one erasure pool spanning every endpoint. + FourByFour, + /// 4 nodes × 1 drive, one erasure pool (minimum 4-node 4-disk layout). + FourNodeFourDisk, + /// 4 single-node pools, 4 drives each (expand / decommission / rebalance). + FourPoolFourDrive, + /// 2 single-node pools, 4 drives each (expansion seed). + TwoPoolFourDrive, +} + +pub(crate) struct DistCluster { + pub cluster: RustFSTestClusterEnvironment, +} + +impl DistCluster { + pub async fn start(layout: DistLayout) -> TestResult { + Self::start_with_env(layout, &[]).await + } + + pub async fn start_with_env(layout: DistLayout, extra_env: &[(&str, &str)]) -> TestResult { + let topology = match layout { + DistLayout::FourByFour => ClusterTopology::single_pool_multidrive(NODE_COUNT, DRIVES_PER_NODE), + DistLayout::FourNodeFourDisk => ClusterTopology::single_pool(NODE_COUNT), + DistLayout::FourPoolFourDrive => { + ClusterTopology::per_node_pools(DRIVES_PER_NODE, (0..NODE_COUNT).map(|idx| vec![idx]).collect()) + } + DistLayout::TwoPoolFourDrive => ClusterTopology::per_node_pools(DRIVES_PER_NODE, vec![vec![0], vec![1]]), + }; + let mut cluster = RustFSTestClusterEnvironment::with_topology(topology).await?; + for &(key, value) in extra_env { + cluster.set_env(key, value); + } + cluster.start().await?; + Ok(Self { cluster }) + } + + pub async fn start_replication_pair() -> TestResult<(Self, Self)> { + let mut extra: Vec<(&str, &str)> = replication_fast_env(); + extra.extend_from_slice(LOOPBACK_REPLICATION_TARGET_ENV); + extra.extend_from_slice(FAST_DATA_USAGE_SCANNER_ENV); + let source = Self::start_with_env(DistLayout::FourNodeFourDisk, &extra).await?; + let target = Self::start_with_env(DistLayout::FourNodeFourDisk, &extra).await?; + Ok((source, target)) + } + + pub fn client(&self, node_idx: usize) -> TestResult { + self.cluster.create_s3_client(node_idx) + } + + pub fn clients(&self) -> TestResult> { + self.cluster.create_all_clients() + } + + pub async fn create_bucket(&self, bucket: &str) -> TestResult { + self.cluster.create_test_bucket(bucket).await + } +} + +pub(crate) fn unique_bucket(prefix: &str) -> String { + let id = Uuid::new_v4().simple().to_string(); + format!("{prefix}-{}", &id[..12]) +} + +pub(crate) fn sha256_hex(bytes: &[u8]) -> String { + let digest = Sha256::digest(bytes); + digest.iter().map(|byte| format!("{byte:02x}")).collect() +} + +pub(crate) fn payload_for(key: &str, size: usize) -> Vec { + let seed = key.as_bytes(); + (0..size) + .map(|idx| seed.get(idx % seed.len()).copied().unwrap_or(0) ^ (idx as u8)) + .collect() +} + +pub(crate) async fn put_object(client: &Client, bucket: &str, key: &str, body: Vec) -> TestResult { + client + .put_object() + .bucket(bucket) + .key(key) + .body(ByteStream::from(body)) + .send() + .await?; + Ok(()) +} + +pub(crate) async fn get_object_bytes(client: &Client, bucket: &str, key: &str) -> TestResult> { + let output = client.get_object().bucket(bucket).key(key).send().await?; + Ok(output.body.collect().await?.into_bytes().to_vec()) +} + +pub(crate) async fn assert_object_bytes(client: &Client, bucket: &str, key: &str, expected: &[u8]) -> TestResult { + let got = get_object_bytes(client, bucket, key).await?; + if got.as_slice() != expected { + return Err(format!( + "object {bucket}/{key} bytes mismatch: expected {} bytes sha256={} got {} bytes sha256={}", + expected.len(), + sha256_hex(expected), + got.len(), + sha256_hex(&got) + ) + .into()); + } + Ok(()) +} + +pub(crate) async fn put_inventory( + client: &Client, + bucket: &str, + count: usize, + size: usize, +) -> TestResult>> { + let mut inventory = BTreeMap::new(); + for idx in 0..count { + let key = format!("obj-{idx:04}"); + let body = payload_for(&key, size); + put_object(client, bucket, &key, body.clone()).await?; + inventory.insert(key, body); + } + Ok(inventory) +} + +pub(crate) async fn assert_inventory(client: &Client, bucket: &str, inventory: &BTreeMap>) -> TestResult { + for (key, expected) in inventory { + assert_object_bytes(client, bucket, key, expected).await?; + } + Ok(()) +} + +pub(crate) async fn enable_versioning(client: &Client, bucket: &str) -> TestResult { + client + .put_bucket_versioning() + .bucket(bucket) + .versioning_configuration( + VersioningConfiguration::builder() + .status(BucketVersioningStatus::Enabled) + .build(), + ) + .send() + .await?; + Ok(()) +} + +pub(crate) async fn wait_until(timeout: Duration, mut probe: F, label: &str) -> TestResult +where + F: FnMut() -> Fut, + Fut: std::future::Future>, +{ + let deadline = Instant::now() + timeout; + let mut delay = Duration::from_millis(50); + loop { + let last_error = match probe().await { + Ok(true) => return Ok(()), + Ok(false) => format!("{label} still false"), + Err(error) => error.to_string(), + }; + if Instant::now() >= deadline { + return Err(format!("{label} did not become true within {timeout:?}: {last_error}").into()); + } + sleep(delay).await; + delay = (delay * 2).min(Duration::from_secs(1)); + } +} + +pub(crate) async fn cluster_admin( + cluster: &RustFSTestClusterEnvironment, + method: Method, + path_and_query: &str, + body: Option, +) -> TestResult<(StatusCode, String)> { + admin_request( + &cluster.nodes[0].url, + method, + path_and_query, + body, + &cluster.access_key, + &cluster.secret_key, + ) + .await +} + +pub(crate) async fn cluster_admin_ok( + cluster: &RustFSTestClusterEnvironment, + method: Method, + path_and_query: &str, + body: Option, +) -> TestResult { + let (status, response) = cluster_admin(cluster, method.clone(), path_and_query, body).await?; + if !status.is_success() { + return Err(format!("{method} {path_and_query} failed: {status} {response}").into()); + } + Ok(response) +} + +pub(crate) async fn wait_for_ready(cluster: &RustFSTestClusterEnvironment) -> TestResult { + let client = local_http_client(); + for node in &cluster.nodes { + let url = format!("{}/health/ready", node.url); + wait_until( + Duration::from_secs(30), + || { + let client = client.clone(); + let url = url.clone(); + async move { + match client.get(&url).send().await { + Ok(response) if response.status().is_success() => Ok(true), + _ => Ok(false), + } + } + }, + &format!("node {} ready", node.address), + ) + .await?; + } + Ok(()) +} + +pub(crate) fn take_drive_offline( + cluster: &RustFSTestClusterEnvironment, + node_idx: usize, + drive_idx: usize, +) -> TestResult { + let dir = cluster + .nodes + .get(node_idx) + .and_then(|node| node.data_dirs.get(drive_idx)) + .ok_or("invalid node/drive index")?; + let offline = format!("{dir}.offline"); + if Path::new(&offline).exists() { + return Err(format!("drive already offline: {offline}").into()); + } + std::fs::rename(dir, &offline)?; + Ok(offline) +} + +pub(crate) fn bring_drive_online(cluster: &RustFSTestClusterEnvironment, node_idx: usize, drive_idx: usize) -> TestResult { + let dir = cluster + .nodes + .get(node_idx) + .and_then(|node| node.data_dirs.get(drive_idx)) + .ok_or("invalid node/drive index")?; + let offline = format!("{dir}.offline"); + if Path::new(dir).exists() { + std::fs::remove_dir_all(dir)?; + } + std::fs::rename(&offline, dir)?; + Ok(()) +} + +pub(crate) async fn set_remote_target( + source: &RustFSTestClusterEnvironment, + source_bucket: &str, + target: &RustFSTestClusterEnvironment, + target_bucket: &str, +) -> TestResult { + let body = serde_json::json!({ + "endpoint": target.nodes[0].address, + "credentials": { + "accessKey": target.access_key, + "secretKey": target.secret_key + }, + "targetbucket": target_bucket, + "secure": false, + "type": "replication" + }); + let url = format!( + "{}/rustfs/admin/v3/set-remote-target?bucket={}", + source.nodes[0].url, + urlencoding::encode(source_bucket) + ); + let response = signed_request( + Method::PUT, + &url, + &source.access_key, + &source.secret_key, + Some(body.to_string().into_bytes()), + Some("application/json"), + ) + .await?; + if response.status() != StatusCode::OK { + let status = response.status(); + let body = response.text().await.unwrap_or_default(); + return Err(format!("set remote target failed: {status} {body}").into()); + } + Ok(serde_json::from_slice(&response.bytes().await?)?) +} + +pub(crate) async fn put_bucket_replication(source: &RustFSTestClusterEnvironment, bucket: &str, target_arn: &str) -> TestResult { + let body = format!( + r#" + + + rule-1 + 1 + Enabled + + Enabled + + + Enabled + + + {target_arn} + + +"# + ); + let url = format!("{}/{bucket}?replication", source.nodes[0].url); + let response = signed_request( + Method::PUT, + &url, + &source.access_key, + &source.secret_key, + Some(body.into_bytes()), + Some("application/xml"), + ) + .await?; + if !response.status().is_success() { + let status = response.status(); + let body = response.text().await.unwrap_or_default(); + return Err(format!("put bucket replication failed: {status} {body}").into()); + } + Ok(()) +} + +pub(crate) async fn wait_for_replicated_bytes( + client: &Client, + bucket: &str, + key: &str, + expected: &[u8], + timeout: Duration, +) -> TestResult { + wait_until( + timeout, + || async { + match get_object_bytes(client, bucket, key).await { + Ok(got) if got.as_slice() == expected => Ok(true), + Ok(_) => Ok(false), + Err(error) => { + let message = error.to_string(); + if message.contains("NoSuchKey") || message.contains("NotFound") { + Ok(false) + } else { + Err(error) + } + } + } + }, + &format!("replicated object {bucket}/{key}"), + ) + .await +} + +pub(crate) async fn set_bucket_quota(cluster: &RustFSTestClusterEnvironment, bucket: &str, quota_bytes: u64) -> TestResult { + wait_until( + Duration::from_secs(30), + || async { + let (status, _) = + cluster_admin(cluster, Method::GET, &format!("/rustfs/admin/v3/quota-stats/{bucket}"), None).await?; + Ok(status.is_success() || status == StatusCode::NOT_FOUND) + }, + "quota stats ready", + ) + .await?; + let body = serde_json::json!({ "quota": quota_bytes, "quota_type": "HARD" }).to_string(); + wait_until( + Duration::from_secs(30), + || async { + let (status, response) = + cluster_admin(cluster, Method::PUT, &format!("/rustfs/admin/v3/quota/{bucket}"), Some(body.clone())).await?; + if status.is_success() { + return Ok(true); + } + if status == StatusCode::SERVICE_UNAVAILABLE { + return Ok(false); + } + Err(format!("failed to set quota for {bucket}: {status} {response}").into()) + }, + "set hard quota", + ) + .await +} + +pub(crate) async fn start_decommission(cluster: &RustFSTestClusterEnvironment, pool_id: usize) -> TestResult { + cluster_admin_ok( + cluster, + Method::POST, + &format!("/rustfs/admin/v3/pools/decommission?pool={pool_id}&by-id=true"), + None, + ) + .await +} + +pub(crate) async fn decommission_status_json(cluster: &RustFSTestClusterEnvironment) -> TestResult { + let body = cluster_admin_ok(cluster, Method::GET, "/rustfs/admin/v3/decommission/status", None).await?; + Ok(serde_json::from_str(&body)?) +} + +fn pool_entry(status: &serde_json::Value, pool_id: usize) -> Option<&serde_json::Value> { + if let Some(pools) = status.get("pools").and_then(serde_json::Value::as_array) { + return pools + .iter() + .find(|pool| pool.get("id").and_then(serde_json::Value::as_u64) == Some(pool_id as u64)); + } + if status.get("id").and_then(serde_json::Value::as_u64) == Some(pool_id as u64) { + Some(status) + } else { + None + } +} + +fn decommission_pool_failed(pool: &serde_json::Value) -> bool { + let info = pool.get("decommissionInfo"); + let flagged = |key: &str| info.and_then(|value| value.get(key)).and_then(serde_json::Value::as_bool) == Some(true); + flagged("failed") + || flagged("canceled") + || pool + .get("status") + .and_then(serde_json::Value::as_str) + .is_some_and(|status| status.eq_ignore_ascii_case("failed") || status.eq_ignore_ascii_case("canceled")) +} + +pub(crate) fn decommission_complete(status: &serde_json::Value, pool_id: usize) -> bool { + let Some(pool) = pool_entry(status, pool_id) else { + return false; + }; + if decommission_pool_failed(pool) { + return false; + } + let info_complete = pool + .get("decommissionInfo") + .and_then(|value| value.get("complete")) + .and_then(serde_json::Value::as_bool) + == Some(true); + let status_text = pool.get("status").and_then(serde_json::Value::as_str).unwrap_or(""); + let pool_status = pool.get("poolStatus").and_then(serde_json::Value::as_str).unwrap_or(""); + info_complete || status_text.eq_ignore_ascii_case("complete") || pool_status.eq_ignore_ascii_case("decommissioned") +} + +pub(crate) fn decommission_failed(status: &serde_json::Value, pool_id: usize) -> bool { + pool_entry(status, pool_id).is_some_and(decommission_pool_failed) +} + +pub(crate) async fn wait_for_decommission_complete( + cluster: &RustFSTestClusterEnvironment, + pool_id: usize, + timeout: Duration, +) -> TestResult { + wait_until( + timeout, + || async { + let status = decommission_status_json(cluster).await?; + if decommission_failed(&status, pool_id) { + return Err(format!("decommission failed for pool {pool_id}: {status}").into()); + } + Ok(decommission_complete(&status, pool_id)) + }, + "decommission complete", + ) + .await +} + +pub(crate) async fn start_rebalance(cluster: &RustFSTestClusterEnvironment) -> TestResult { + cluster_admin_ok(cluster, Method::POST, "/rustfs/admin/v3/rebalance/start", None).await +} + +pub(crate) async fn rebalance_status_json(cluster: &RustFSTestClusterEnvironment) -> TestResult { + let body = cluster_admin_ok(cluster, Method::GET, "/rustfs/admin/v3/rebalance/status", None).await?; + Ok(serde_json::from_str(&body)?) +} + +pub(crate) fn rebalance_active(status: &serde_json::Value) -> bool { + status + .get("pools") + .and_then(serde_json::Value::as_array) + .is_some_and(|pools| { + pools.iter().any(|pool| { + let stopping = pool.get("stopping").and_then(serde_json::Value::as_bool) == Some(true); + let value = pool.get("status").and_then(serde_json::Value::as_str).unwrap_or(""); + stopping + || value.eq_ignore_ascii_case("started") + || value.eq_ignore_ascii_case("active") + || value.eq_ignore_ascii_case("running") + || value.eq_ignore_ascii_case("stopping") + }) + }) +} + +pub(crate) async fn wait_for_rebalance_idle(cluster: &RustFSTestClusterEnvironment, timeout: Duration) -> TestResult { + wait_until( + timeout, + || async { + match rebalance_status_json(cluster).await { + Ok(status) => Ok(!rebalance_active(&status)), + Err(error) => { + let message = error.to_string(); + if message.contains("NoSuchResource") || message.contains("404") || message.contains("not started") { + Ok(true) + } else { + Err(error) + } + } + } + }, + "rebalance idle", + ) + .await +} + +pub(crate) async fn list_pools_json(cluster: &RustFSTestClusterEnvironment) -> TestResult { + let body = cluster_admin_ok(cluster, Method::GET, "/rustfs/admin/v3/pools/list", None).await?; + Ok(serde_json::from_str(&body)?) +} + +pub(crate) async fn retrying_put(client: &Client, bucket: &str, key: &str, body: Vec, timeout: Duration) -> TestResult { + wait_until( + timeout, + || { + let client = client.clone(); + let bucket = bucket.to_string(); + let key = key.to_string(); + let body = body.clone(); + async move { + match put_object(&client, &bucket, &key, body).await { + Ok(()) => Ok(true), + Err(error) => { + let message = error.to_string(); + if message.contains("SlowDown") + || message.contains("ServiceUnavailable") + || message.contains("InternalError") + || message.contains("503") + || message.contains("500") + { + Ok(false) + } else { + Err(error) + } + } + } + } + }, + &format!("put {bucket}/{key} during data movement"), + ) + .await +} + +pub(crate) async fn retrying_get_equals( + client: &Client, + bucket: &str, + key: &str, + expected: &[u8], + timeout: Duration, +) -> TestResult { + wait_until( + timeout, + || async { + match get_object_bytes(client, bucket, key).await { + Ok(got) if got.as_slice() == expected => Ok(true), + Ok(_) => Ok(false), + Err(error) => { + let message = error.to_string(); + if message.contains("NoSuchKey") + || message.contains("SlowDown") + || message.contains("ServiceUnavailable") + || message.contains("InternalError") + || message.contains("503") + || message.contains("500") + { + Ok(false) + } else { + Err(error) + } + } + } + }, + &format!("get {bucket}/{key} during data movement"), + ) + .await +} + +#[tokio::test] +async fn append_single_node_pool_extends_ellipses_volumes() { + let mut env = RustFSTestClusterEnvironment::with_topology(ClusterTopology::per_node_pools(2, vec![vec![0], vec![1]])) + .await + .expect("two-pool seed topology"); + assert_eq!(env.rustfs_volumes_arg().split(' ').count(), 2); + + let added = env.append_single_node_pool().await.expect("append third pool"); + assert_eq!(added, 2); + assert_eq!(env.nodes.len(), 3); + assert_eq!(env.nodes[2].pool_idx, 2); + assert_eq!(env.nodes[2].data_dirs.len(), 2); + let volumes = env.rustfs_volumes_arg(); + assert_eq!(volumes.split(' ').count(), 3, "expected three pool arguments, got: {volumes}"); + assert!(volumes.contains("/drive{0...1}"), "expanded layout must keep drive ellipses: {volumes}"); +} + +#[tokio::test] +async fn append_single_node_pool_rejects_striped_single_pool() { + let mut env = RustFSTestClusterEnvironment::new(4).await.expect("four-node single pool"); + let err = env + .append_single_node_pool() + .await + .expect_err("a striped single pool cannot gain a localhost pool"); + let message = err.to_string(); + assert!( + message.contains("drives_per_node") || message.contains("one node per pool"), + "unexpected error: {message}" + ); +} + +#[test] +fn decommission_complete_reads_pool_status_and_info_flag() { + let status = serde_json::json!({ + "pools": [ + { + "id": 0, + "status": "complete", + "poolStatus": "decommissioned", + "decommissionInfo": { "complete": true, "failed": false, "canceled": false } + }, + { "id": 1, "status": "none", "poolStatus": "active" } + ] + }); + assert!(decommission_complete(&status, 0)); + assert!(!decommission_complete(&status, 1)); + assert!(!decommission_failed(&status, 0)); +} + +#[test] +fn rebalance_active_treats_started_as_in_progress() { + let started = serde_json::json!({ "pools": [{ "id": 0, "status": "Started", "stopping": false }] }); + let done = serde_json::json!({ "pools": [{ "id": 0, "status": "Completed", "stopping": false }] }); + assert!(rebalance_active(&started)); + assert!(!rebalance_active(&done)); +} diff --git a/crates/e2e_test/src/distributed/mod.rs b/crates/e2e_test/src/distributed/mod.rs new file mode 100644 index 000000000..fa1ba1875 --- /dev/null +++ b/crates/e2e_test/src/distributed/mod.rs @@ -0,0 +1,34 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! 4-node 4-drive distributed e2e coverage. +//! +//! Selected by `[profile.e2e-distributed]` and run from +//! `.github/workflows/e2e-distributed.yml`. Excluded from `e2e-full` because +//! each case starts four real `rustfs` processes. + +mod chaos_test; +mod concurrency_stability_test; +mod concurrent_data_movement_test; +mod data_integrity_movement_test; +mod expand_decommission_rebalance_test; +mod extra_test; +mod harness; +mod object_lock_test; +mod observability_test; +mod replication_quota_test; +mod s3_basic_test; +mod s3_during_data_movement_test; +mod site_replication_test; +mod versioning_test; diff --git a/crates/e2e_test/src/distributed/object_lock_test.rs b/crates/e2e_test/src/distributed/object_lock_test.rs new file mode 100644 index 000000000..fc419934f --- /dev/null +++ b/crates/e2e_test/src/distributed/object_lock_test.rs @@ -0,0 +1,137 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::harness::{DistCluster, DistLayout, TestResult, put_object, unique_bucket}; +use crate::common::init_logging; +use aws_sdk_s3::error::ProvideErrorMetadata; +use aws_sdk_s3::types::{ + DefaultRetention, ObjectLockConfiguration, ObjectLockEnabled, ObjectLockLegalHold, ObjectLockLegalHoldStatus, + ObjectLockRetention, ObjectLockRetentionMode, ObjectLockRule, +}; +use chrono::{Duration as ChronoDuration, Utc}; + +#[tokio::test] +async fn four_node_four_drive_object_lock_worm_blocks_delete() -> TestResult { + init_logging(); + let dist = DistCluster::start(DistLayout::FourByFour).await?; + let client = dist.client(0)?; + let peer = dist.client(2)?; + let bucket = unique_bucket("objlock"); + + client + .create_bucket() + .bucket(&bucket) + .object_lock_enabled_for_bucket(true) + .send() + .await?; + + let retain_until = Utc::now() + ChronoDuration::days(1); + let retain_until_s3 = aws_sdk_s3::primitives::DateTime::from_secs(retain_until.timestamp()); + + client + .put_object_lock_configuration() + .bucket(&bucket) + .object_lock_configuration( + ObjectLockConfiguration::builder() + .object_lock_enabled(ObjectLockEnabled::Enabled) + .rule( + ObjectLockRule::builder() + .default_retention( + DefaultRetention::builder() + .mode(ObjectLockRetentionMode::Governance) + .days(1) + .build(), + ) + .build(), + ) + .build(), + ) + .send() + .await?; + + let compliance_key = "compliance.bin"; + put_object(&client, &bucket, compliance_key, b"locked-compliance".to_vec()).await?; + client + .put_object_retention() + .bucket(&bucket) + .key(compliance_key) + .retention( + ObjectLockRetention::builder() + .mode(ObjectLockRetentionMode::Compliance) + .retain_until_date(retain_until_s3) + .build(), + ) + .send() + .await?; + + let compliance_delete = peer.delete_object().bucket(&bucket).key(compliance_key).send().await; + match compliance_delete { + Ok(_) => return Err("COMPLIANCE retention must block DeleteObject".into()), + Err(error) => { + let code = error.as_service_error().and_then(ProvideErrorMetadata::code); + assert_eq!(code, Some("AccessDenied"), "unexpected COMPLIANCE delete error: {error:?}"); + } + } + + let governance_key = "governance.bin"; + put_object(&client, &bucket, governance_key, b"locked-governance".to_vec()).await?; + client + .put_object_retention() + .bucket(&bucket) + .key(governance_key) + .retention( + ObjectLockRetention::builder() + .mode(ObjectLockRetentionMode::Governance) + .retain_until_date(retain_until_s3) + .build(), + ) + .send() + .await?; + + let governance_blocked = peer.delete_object().bucket(&bucket).key(governance_key).send().await; + match governance_blocked { + Ok(_) => return Err("GOVERNANCE retention must block DeleteObject without bypass".into()), + Err(error) => { + let code = error.as_service_error().and_then(ProvideErrorMetadata::code); + assert_eq!(code, Some("AccessDenied"), "unexpected GOVERNANCE delete error: {error:?}"); + } + } + + peer.delete_object() + .bucket(&bucket) + .key(governance_key) + .bypass_governance_retention(true) + .send() + .await?; + + let hold_key = "legal-hold.bin"; + put_object(&client, &bucket, hold_key, b"legal-hold".to_vec()).await?; + client + .put_object_legal_hold() + .bucket(&bucket) + .key(hold_key) + .legal_hold(ObjectLockLegalHold::builder().status(ObjectLockLegalHoldStatus::On).build()) + .send() + .await?; + let hold_delete = peer.delete_object().bucket(&bucket).key(hold_key).send().await; + match hold_delete { + Ok(_) => return Err("legal hold must block DeleteObject".into()), + Err(error) => { + let code = error.as_service_error().and_then(ProvideErrorMetadata::code); + assert_eq!(code, Some("AccessDenied"), "unexpected legal-hold delete error: {error:?}"); + } + } + + Ok(()) +} diff --git a/crates/e2e_test/src/distributed/observability_test.rs b/crates/e2e_test/src/distributed/observability_test.rs new file mode 100644 index 000000000..5dd84135c --- /dev/null +++ b/crates/e2e_test/src/distributed/observability_test.rs @@ -0,0 +1,61 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::harness::{DistCluster, DistLayout, TestResult, cluster_admin_ok, put_object, unique_bucket, wait_for_ready}; +use crate::common::{init_logging, local_http_client}; +use http::Method; + +#[tokio::test] +async fn four_node_four_drive_health_admin_info_and_audit_list() -> TestResult { + init_logging(); + let dist = DistCluster::start(DistLayout::FourByFour).await?; + wait_for_ready(&dist.cluster).await?; + + let http = local_http_client(); + for node in &dist.cluster.nodes { + let ready = http.get(format!("{}/health/ready", node.url)).send().await?; + assert!(ready.status().is_success(), "node {} not ready: {}", node.address, ready.status()); + let live = http.get(format!("{}/health/live", node.url)).send().await; + if let Ok(response) = live { + assert!( + response.status().is_success() || response.status().as_u16() == 404, + "unexpected live probe on {}: {}", + node.address, + response.status() + ); + } + } + + let info = cluster_admin_ok(&dist.cluster, Method::GET, "/rustfs/admin/v3/info", None).await?; + assert!(!info.is_empty(), "admin info was empty"); + let storage = cluster_admin_ok(&dist.cluster, Method::GET, "/rustfs/admin/v3/storageinfo", None).await?; + assert!( + storage.contains("disks") || storage.contains("backend") || storage.contains("info"), + "storageinfo missing expected fields: {storage}" + ); + + let audit = cluster_admin_ok(&dist.cluster, Method::GET, "/rustfs/admin/v3/audit/target/list", None).await?; + let trimmed = audit.trim(); + if !trimmed.is_empty() && trimmed != "null" && !trimmed.starts_with('[') && !trimmed.starts_with('{') { + return Err(format!("audit target list was not machine-readable: {audit}").into()); + } + + let bucket = unique_bucket("obs"); + dist.create_bucket(&bucket).await?; + put_object(&dist.client(0)?, &bucket, "probe.log", b"observability".to_vec()).await?; + + let trace = cluster_admin_ok(&dist.cluster, Method::GET, "/rustfs/admin/v3/info", None).await?; + assert!(!trace.is_empty()); + Ok(()) +} diff --git a/crates/e2e_test/src/distributed/replication_quota_test.rs b/crates/e2e_test/src/distributed/replication_quota_test.rs new file mode 100644 index 000000000..7f5967ffd --- /dev/null +++ b/crates/e2e_test/src/distributed/replication_quota_test.rs @@ -0,0 +1,98 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::harness::{ + DistCluster, DistLayout, TestResult, enable_versioning, put_bucket_replication, put_object, set_bucket_quota, + set_remote_target, unique_bucket, wait_for_replicated_bytes, +}; +use crate::common::{FAST_DATA_USAGE_SCANNER_ENV, init_logging}; +use aws_sdk_s3::error::ProvideErrorMetadata; +use std::time::Duration; + +#[tokio::test] +async fn four_node_bucket_replication_converges_to_peer_cluster() -> TestResult { + init_logging(); + let (source, target) = DistCluster::start_replication_pair().await?; + let source_bucket = unique_bucket("replsrc"); + let target_bucket = unique_bucket("repldst"); + source.create_bucket(&source_bucket).await?; + target.create_bucket(&target_bucket).await?; + + let source_client = source.client(0)?; + let target_client = target.client(0)?; + enable_versioning(&source_client, &source_bucket).await?; + enable_versioning(&target_client, &target_bucket).await?; + + let arn = set_remote_target(&source.cluster, &source_bucket, &target.cluster, &target_bucket).await?; + put_bucket_replication(&source.cluster, &source_bucket, &arn).await?; + + let key = "replicated.bin"; + let body = b"distributed-bucket-replication".to_vec(); + put_object(&source_client, &source_bucket, key, body.clone()).await?; + wait_for_replicated_bytes(&target_client, &target_bucket, key, &body, Duration::from_secs(45)).await?; + + let peer_read = target.client(3)?; + wait_for_replicated_bytes(&peer_read, &target_bucket, key, &body, Duration::from_secs(15)).await?; + Ok(()) +} + +#[tokio::test] +async fn four_node_four_drive_hard_quota_rejects_over_limit_put() -> TestResult { + init_logging(); + let dist = DistCluster::start_with_env(DistLayout::FourByFour, FAST_DATA_USAGE_SCANNER_ENV).await?; + let bucket = unique_bucket("quota"); + dist.create_bucket(&bucket).await?; + set_bucket_quota(&dist.cluster, &bucket, 8 * 1024).await?; + + let client = dist.client(1)?; + put_object(&client, &bucket, "small.bin", vec![0u8; 1024]).await?; + + let over_limit = client + .put_object() + .bucket(&bucket) + .key("too-big.bin") + .body(vec![0u8; 16 * 1024].into()) + .send() + .await; + match over_limit { + Ok(_) => { + // Scanner-backed quota can lag a cycle; a second over-quota PUT must fail. + let second = client + .put_object() + .bucket(&bucket) + .key("too-big-2.bin") + .body(vec![0u8; 16 * 1024].into()) + .send() + .await; + match second { + Ok(_) => return Err("hard quota admitted two oversized PUTs on a 4x4 cluster".into()), + Err(error) => { + let code = error.as_service_error().and_then(ProvideErrorMetadata::code); + assert!( + matches!(code, Some("QuotaExceeded" | "SlowDown" | "AccessDenied" | "InvalidRequest")), + "unexpected over-quota error: {error:?}" + ); + } + } + } + Err(error) => { + let code = error.as_service_error().and_then(ProvideErrorMetadata::code); + assert!( + matches!(code, Some("QuotaExceeded" | "SlowDown" | "AccessDenied" | "InvalidRequest")), + "unexpected over-quota error: {error:?}" + ); + } + } + Ok(()) +} diff --git a/crates/e2e_test/src/distributed/s3_basic_test.rs b/crates/e2e_test/src/distributed/s3_basic_test.rs new file mode 100644 index 000000000..6ab145751 --- /dev/null +++ b/crates/e2e_test/src/distributed/s3_basic_test.rs @@ -0,0 +1,111 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::harness::{DistCluster, DistLayout, TestResult, assert_object_bytes, get_object_bytes, put_object, unique_bucket}; +use crate::common::{init_logging, local_http_client}; +use aws_sdk_s3::presigning::PresigningConfig; +use aws_sdk_s3::types::{Delete, MetadataDirective, ObjectIdentifier}; +use std::time::Duration; + +#[tokio::test] +async fn four_node_four_drive_s3_put_get_head_list_copy_rename_delete_and_presign() -> TestResult { + init_logging(); + let dist = DistCluster::start(DistLayout::FourByFour).await?; + let bucket = unique_bucket("s3basic"); + dist.create_bucket(&bucket).await?; + + let writer = dist.client(0)?; + let reader = dist.client(3)?; + let key = "dir/object.bin"; + let body = vec![0xA5u8; 256 * 1024]; + put_object(&writer, &bucket, key, body.clone()).await?; + + let head = reader.head_object().bucket(&bucket).key(key).send().await?; + assert_eq!(head.content_length(), Some(body.len() as i64)); + assert_object_bytes(&reader, &bucket, key, &body).await?; + + let ranged = reader + .get_object() + .bucket(&bucket) + .key(key) + .range("bytes=0-15") + .send() + .await?; + let ranged_body = ranged.body.collect().await?.into_bytes(); + assert_eq!(ranged_body.as_ref(), &body[..16]); + + let listed = reader.list_objects_v2().bucket(&bucket).prefix("dir/").send().await?; + let keys: Vec<_> = listed.contents().iter().filter_map(|object| object.key()).collect(); + assert_eq!(keys, vec![key]); + + let copy_key = "dir/object-copy.bin"; + reader + .copy_object() + .bucket(&bucket) + .key(copy_key) + .copy_source(format!("{bucket}/{key}")) + .metadata_directive(MetadataDirective::Copy) + .send() + .await?; + assert_object_bytes(&writer, &bucket, copy_key, &body).await?; + + let moved_key = "dir/object-moved.bin"; + writer + .copy_object() + .bucket(&bucket) + .key(moved_key) + .copy_source(format!("{bucket}/{copy_key}")) + .send() + .await?; + writer.delete_object().bucket(&bucket).key(copy_key).send().await?; + match writer.head_object().bucket(&bucket).key(copy_key).send().await { + Ok(_) => return Err("copied source still present after rename delete".into()), + Err(error) if error.as_service_error().is_some_and(|err| err.is_not_found()) => {} + Err(error) => return Err(error.into()), + } + assert_object_bytes(&reader, &bucket, moved_key, &body).await?; + + let presigned = writer + .get_object() + .bucket(&bucket) + .key(key) + .presigned(PresigningConfig::expires_in(Duration::from_secs(120))?) + .await?; + let response = local_http_client().get(presigned.uri().to_string()).send().await?; + assert!(response.status().is_success(), "presigned GET failed: {}", response.status()); + let presigned_body = response.bytes().await?; + assert_eq!(presigned_body.as_ref(), body.as_slice()); + + let empty_key = "empty"; + put_object(&writer, &bucket, empty_key, Vec::new()).await?; + let empty = get_object_bytes(&reader, &bucket, empty_key).await?; + assert!(empty.is_empty()); + + writer + .delete_objects() + .bucket(&bucket) + .delete( + Delete::builder() + .objects(ObjectIdentifier::builder().key(key).build()?) + .objects(ObjectIdentifier::builder().key(moved_key).build()?) + .objects(ObjectIdentifier::builder().key(empty_key).build()?) + .build()?, + ) + .send() + .await?; + + let remaining = reader.list_objects_v2().bucket(&bucket).send().await?; + assert!(remaining.contents().is_empty(), "bucket still has objects after delete"); + Ok(()) +} diff --git a/crates/e2e_test/src/distributed/s3_during_data_movement_test.rs b/crates/e2e_test/src/distributed/s3_during_data_movement_test.rs new file mode 100644 index 000000000..132543bf0 --- /dev/null +++ b/crates/e2e_test/src/distributed/s3_during_data_movement_test.rs @@ -0,0 +1,80 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::harness::{ + DistCluster, DistLayout, TestResult, assert_inventory, put_inventory, retrying_get_equals, retrying_put, start_decommission, + start_rebalance, unique_bucket, wait_for_decommission_complete, +}; +use crate::common::init_logging; +use std::time::Duration; + +#[tokio::test] +async fn s3_put_get_list_succeed_during_decommission_and_rebalance() -> TestResult { + init_logging(); + let dist = DistCluster::start(DistLayout::FourPoolFourDrive).await?; + let bucket = unique_bucket("s3move"); + dist.create_bucket(&bucket).await?; + let client = dist.client(0)?; + let inventory = put_inventory(&client, &bucket, 8, 16 * 1024).await?; + + start_decommission(&dist.cluster, 0).await?; + let live = dist.client(2)?; + retrying_put( + &live, + &bucket, + "during-decommission.bin", + b"written-while-decommissioning".to_vec(), + Duration::from_secs(30), + ) + .await?; + retrying_get_equals( + &live, + &bucket, + "during-decommission.bin", + b"written-while-decommissioning", + Duration::from_secs(30), + ) + .await?; + let listed = live.list_objects_v2().bucket(&bucket).send().await?; + assert!( + listed + .contents() + .iter() + .any(|object| object.key() == Some("during-decommission.bin")), + "list during decommission missed the newly written key" + ); + + wait_for_decommission_complete(&dist.cluster, 0, Duration::from_secs(180)).await?; + assert_inventory(&live, &bucket, &inventory).await?; + + let _ = start_rebalance(&dist.cluster).await; + retrying_put( + &live, + &bucket, + "during-rebalance.bin", + b"written-while-rebalancing".to_vec(), + Duration::from_secs(30), + ) + .await?; + retrying_get_equals( + &live, + &bucket, + "during-rebalance.bin", + b"written-while-rebalancing", + Duration::from_secs(30), + ) + .await?; + assert_inventory(&dist.client(1)?, &bucket, &inventory).await?; + Ok(()) +} diff --git a/crates/e2e_test/src/distributed/site_replication_test.rs b/crates/e2e_test/src/distributed/site_replication_test.rs new file mode 100644 index 000000000..06e7c3496 --- /dev/null +++ b/crates/e2e_test/src/distributed/site_replication_test.rs @@ -0,0 +1,87 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::harness::{ + DistCluster, TestResult, cluster_admin_ok, enable_versioning, put_object, unique_bucket, wait_for_replicated_bytes, +}; +use crate::common::{init_logging, signed_request}; +use http::{Method, StatusCode}; +use rustfs_madmin::PeerSite; +use std::time::Duration; + +async fn site_replication_add(cluster: &crate::common::RustFSTestClusterEnvironment, sites: &[PeerSite]) -> TestResult { + let url = format!("{}/rustfs/admin/v3/site-replication/add?replicateILMExpiry=false", cluster.nodes[0].url); + let response = signed_request( + Method::PUT, + &url, + &cluster.access_key, + &cluster.secret_key, + Some(serde_json::to_vec(sites)?), + Some("application/json"), + ) + .await?; + if response.status() != StatusCode::OK { + let status = response.status(); + let body = response.text().await.unwrap_or_default(); + return Err(format!("site replication add failed: {status} {body}").into()); + } + Ok(response.text().await?) +} + +#[tokio::test] +async fn four_node_site_replication_replicates_object_to_peer_site() -> TestResult { + init_logging(); + let (site_a, site_b) = DistCluster::start_replication_pair().await?; + let bucket = unique_bucket("siterepl"); + site_a.create_bucket(&bucket).await?; + site_b.create_bucket(&bucket).await?; + + let client_a = site_a.client(0)?; + let client_b = site_b.client(0)?; + enable_versioning(&client_a, &bucket).await?; + enable_versioning(&client_b, &bucket).await?; + + let sites = vec![ + PeerSite { + name: "site-a".to_string(), + endpoint: site_a.cluster.nodes[0].url.clone(), + access_key: site_a.cluster.access_key.clone(), + secret_key: site_a.cluster.secret_key.clone(), + ..Default::default() + }, + PeerSite { + name: "site-b".to_string(), + endpoint: site_b.cluster.nodes[0].url.clone(), + access_key: site_b.cluster.access_key.clone(), + secret_key: site_b.cluster.secret_key.clone(), + ..Default::default() + }, + ]; + site_replication_add(&site_a.cluster, &sites).await?; + + let info = cluster_admin_ok(&site_a.cluster, Method::GET, "/rustfs/admin/v3/site-replication/info", None).await?; + assert!( + info.contains("site-a") || info.contains("enabled") || info.contains("true"), + "site replication info did not show a configured peer: {info}" + ); + + let key = "site-object.bin"; + let body = b"four-node-site-replication".to_vec(); + put_object(&client_a, &bucket, key, body.clone()).await?; + wait_for_replicated_bytes(&client_b, &bucket, key, &body, Duration::from_secs(60)).await?; + + let peer_b = site_b.client(3)?; + wait_for_replicated_bytes(&peer_b, &bucket, key, &body, Duration::from_secs(20)).await?; + Ok(()) +} diff --git a/crates/e2e_test/src/distributed/versioning_test.rs b/crates/e2e_test/src/distributed/versioning_test.rs new file mode 100644 index 000000000..eb525072c --- /dev/null +++ b/crates/e2e_test/src/distributed/versioning_test.rs @@ -0,0 +1,88 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::harness::{DistCluster, DistLayout, TestResult, enable_versioning, get_object_bytes, put_object, unique_bucket}; +use crate::common::init_logging; +use aws_sdk_s3::error::ProvideErrorMetadata; + +#[tokio::test] +async fn four_node_four_drive_versioning_put_list_get_delete_marker() -> TestResult { + init_logging(); + let dist = DistCluster::start(DistLayout::FourByFour).await?; + let bucket = unique_bucket("version"); + dist.create_bucket(&bucket).await?; + let writer = dist.client(0)?; + let reader = dist.client(3)?; + enable_versioning(&writer, &bucket).await?; + + let key = "versioned.txt"; + put_object(&writer, &bucket, key, b"v1".to_vec()).await?; + put_object(&writer, &bucket, key, b"v2".to_vec()).await?; + + let versions = reader.list_object_versions().bucket(&bucket).prefix(key).send().await?; + let version_ids: Vec = versions + .versions() + .iter() + .filter_map(|version| version.version_id().map(str::to_string)) + .collect(); + assert!(version_ids.len() >= 2, "expected at least two versions, got {version_ids:?}"); + + let latest = get_object_bytes(&reader, &bucket, key).await?; + assert_eq!(latest, b"v2"); + + let older_id = versions + .versions() + .iter() + .find(|version| version.is_latest() != Some(true)) + .and_then(|version| version.version_id()) + .ok_or("missing non-latest version id")?; + let older = reader + .get_object() + .bucket(&bucket) + .key(key) + .version_id(older_id) + .send() + .await?; + let older_body = older.body.collect().await?.into_bytes(); + assert_eq!(older_body.as_ref(), b"v1"); + + writer.delete_object().bucket(&bucket).key(key).send().await?; + let after_delete = reader.list_object_versions().bucket(&bucket).prefix(key).send().await?; + assert!( + !after_delete.delete_markers().is_empty(), + "delete marker missing after unversioned-style delete: {after_delete:?}" + ); + + let latest_after_delete = reader.get_object().bucket(&bucket).key(key).send().await; + match latest_after_delete { + Ok(_) => return Err("current version should be a delete marker".into()), + Err(error) + if error + .as_service_error() + .and_then(ProvideErrorMetadata::code) + .is_some_and(|code| code == "NoSuchKey" || code == "NotFound") => {} + Err(error) => return Err(error.into()), + } + + let restored = reader + .get_object() + .bucket(&bucket) + .key(key) + .version_id(older_id) + .send() + .await?; + let restored_body = restored.body.collect().await?.into_bytes(); + assert_eq!(restored_body.as_ref(), b"v1"); + Ok(()) +} diff --git a/crates/e2e_test/src/lib.rs b/crates/e2e_test/src/lib.rs index 9279c39c4..72f5121fb 100644 --- a/crates/e2e_test/src/lib.rs +++ b/crates/e2e_test/src/lib.rs @@ -378,6 +378,11 @@ mod bucket_stats_regression_test; #[cfg(test)] mod distributed_startup_regression_test; +// 4-node / 4-disk distributed Actions suite (S3, lock, versioning, replication, +// quota, observability, expand/decommission/rebalance, site replication, chaos). +#[cfg(test)] +mod distributed; + // P1 regression: tier/ILM transition (rustfs#5218, #5130, #5011, #4826, #5024) #[cfg(test)] mod tier_transition_regression_test; diff --git a/docs/testing/README.md b/docs/testing/README.md index 3589d8ee1..53bab77bf 100644 --- a/docs/testing/README.md +++ b/docs/testing/README.md @@ -12,10 +12,10 @@ Pick the lowest layer that can prove the change; add a higher-layer test only wh |---|---|---|---| | Unit & crate integration | Per-crate logic and in-process integration tests | `cargo nextest run --all --exclude e2e_test` (or `-p `); `make test` wraps it | Every PR, required (`Test and Lint`, `ci` profile) | | ecstore black-box | Erasure-coded read/write/recovery validation; profiles `quick` / `full` / `destructive` / `fuzz` | `scripts/run_ecstore_validation_suite.sh --profile quick` | Local and release validation only; not wired into any workflow. Contract: [ecstore-validation-suite-design.md](ecstore-validation-suite-design.md) | -| e2e (`e2e_test` crate) | A real `rustfs` binary per test, driven over the S3, admin, and protocol APIs | `cargo nextest run --profile e2e-smoke -p e2e_test` | PR: `e2e-smoke` (report-only); merge queue / main push: `e2e-full`; nightly: `e2e-repl-nightly`, `e2e-nightly`, `e2e-protocols`. Guide: [`crates/e2e_test/README.md`](../../crates/e2e_test/README.md) | +| e2e (`e2e_test` crate) | A real `rustfs` binary per test, driven over the S3, admin, and protocol APIs | `cargo nextest run --profile e2e-smoke -p e2e_test` | PR: `e2e-smoke` (report-only); merge queue / main push: `e2e-full`; nightly: `e2e-repl-nightly`, `e2e-nightly`, `e2e-protocols`, `e2e-distributed`. Guide: [`crates/e2e_test/README.md`](../../crates/e2e_test/README.md); 4-node 4-disk map: [distributed-e2e.md](distributed-e2e.md) | | s3s-e2e conformance | External S3 conformance tool against a live server | `./scripts/e2e-run.sh ./target/debug/rustfs ` | PR, report-only (second half of the `End-to-End Tests` job) | | S3 compatibility | `ceph/s3-tests` (boto3; allow-list `scripts/s3-tests/implemented_tests.txt`) and MinIO `mint` | `scripts/s3-tests/run.sh`; mint via `.github/workflows/mint.yml` | s3-tests: PR report-only plus a weekly full sweep; mint: weekly, report-only | -| Chaos / fault-injection | Single-node disk fault injection (`crates/e2e_test/src/chaos.rs`, `crates/e2e_test/src/fault_proxy.rs`) used by the reliability and heal e2e modules | Part of the e2e crate (`e2e-reliability` test-group) | With the `e2e-full` and nightly e2e lanes. A multi-node power-loss harness is not in tree | +| Chaos / fault-injection | Single-node disk fault injection (`crates/e2e_test/src/chaos.rs`, `crates/e2e_test/src/fault_proxy.rs`) plus the 4-node kill/offline-drive/blackhole cases in `crates/e2e_test/src/distributed/chaos_test.rs` | Part of the e2e crate (`e2e-reliability` and `e2e-distributed`) | Reliability cases with `e2e-full`; 4-node chaos with the `e2e-distributed` nightly lane | | Fuzz | `cargo-fuzz` targets over untrusted parsing surfaces; isolated sub-workspace under `fuzz/` | `./scripts/fuzz/run.sh` (see [`fuzz/README.md`](../../fuzz/README.md)) | PR smoke on the paths listed in `.github/workflows/fuzz.yml`, plus nightly corpus | | Benchmarks | Criterion benches under each crate's `benches/` | `cargo bench -p ` | On demand; never a gate | @@ -61,6 +61,7 @@ All profiles are defined in `.config/nextest.toml`; its block comments hold the | `e2e-full` | Merge-queue / main-push single-node e2e lane | | `e2e-repl-nightly` | Nightly slow / cross-process replication lane | | `e2e-nightly` | Nightly serial multi-process cluster fault lane | +| `e2e-distributed` | Nightly 4-node 4-disk S3 / lock / versioning / replication / quota / expand / decommission / rebalance / site-replication / chaos lane | | `e2e-protocols` | Nightly fixed-port FTPS/SFTP/WebDAV lane, run with `-j 1` | Membership of each e2e profile is pinned by a digest in `.config/e2e--selection.txt` and checked by `scripts/check_test_wiring.py --check-profile ` before the lane runs. To list what a profile selects on your platform (the result is platform-dependent because some modules are linux-only): diff --git a/docs/testing/ci-gates.md b/docs/testing/ci-gates.md index 256b41bd4..522b753d7 100644 --- a/docs/testing/ci-gates.md +++ b/docs/testing/ci-gates.md @@ -74,6 +74,7 @@ Scheduled lanes never block a PR. Their workflow-local gate fails the run, sched | `ci.yml` (weekly) | full matrix, including the schedule/dispatch-only rio-v2 jobs `build-rustfs-debug-binary-rio-v2` and `e2e-tests-rio-v2` | per-job | yes | dispatch `ci.yml` | | `build.yml` (weekly) | `build-rustfs` over the six-target platform matrix in `prepare-platform-matrix` (four Linux, macOS aarch64, Windows x86_64) | build/package integrity | yes | dispatch `build.yml` with an exact platform set | | `e2e-replication-nightly.yml` (nightly) | `repl-nightly`, `cluster-nightly`, `protocols-nightly` | three independent gates; JUnit, membership listing, server logs | yes | `cargo nextest run --profile e2e-repl-nightly -p e2e_test`; `--profile e2e-nightly`; `-j 1 --profile e2e-protocols` | +| `e2e-distributed.yml` (nightly) | `distributed` | 4-node 4-disk e2e gate; JUnit, membership listing, server logs | yes, with `never_ran_grace_until` | `cargo nextest run --profile e2e-distributed -p e2e_test` | | `e2e-s3tests.yml` (weekly) | `s3tests` (single and distributed, four shards each), `upstream-head-canary` | compatibility gate; report, JUnit, node IDs, server logs | yes | `scripts/s3-tests/run.sh` against an existing single or distributed target | | `fuzz.yml` (nightly) | `nightly-fuzz-corpus` per target | gate; corpus and crash artifacts | yes | `MAX_TOTAL_TIME= ./scripts/fuzz/run.sh` | | `minio-interop.yml` (nightly) | `minio-interop` | EC + SSE read-parity gate | yes, with `never_ran_grace_until` | pinned Docker fixture steps in the workflow | diff --git a/docs/testing/distributed-e2e.md b/docs/testing/distributed-e2e.md new file mode 100644 index 000000000..5f4148418 --- /dev/null +++ b/docs/testing/distributed-e2e.md @@ -0,0 +1,56 @@ +# Distributed 4-node 4-disk e2e + +**Use this when:** adding or diagnosing GitHub Actions coverage for a 4-node cluster, or deciding whether a behaviour belongs in `e2e-distributed` versus the single-node `e2e-full` lane, the nightly cluster-fault lane, or the hardware functional chain. +**Source of truth:** `crates/e2e_test/src/distributed/`, `[profile.e2e-distributed]` in `.config/nextest.toml`, `.github/workflows/e2e-distributed.yml`. + +## Topology + +The in-tree harness runs every node on `127.0.0.1` with a distinct port. That matches `RustFSTestClusterEnvironment` in `crates/e2e_test/src/common.rs`: + +| Layout | Constructor | Use | +|---|---|---| +| 4 nodes × 4 drives, one pool | `ClusterTopology::single_pool_multidrive(4, 4)` | S3, object lock, versioning, quota, observability, concurrency, chaos | +| 4 nodes × 1 drive, one pool | `ClusterTopology::single_pool(4)` | Two-site replication (8 processes total) | +| 4 single-node pools × 4 drives | `ClusterTopology::per_node_pools(4, [[0],[1],[2],[3]])` | Decommission, rebalance, S3-during-move, integrity | +| 2 single-node pools × 4 drives, then `append_single_node_pool` | expansion seed | Pool expand then rebalance | + +A pool striped across several localhost ports is not expressible (`RUSTFS_VOLUMES` host ellipses would collide on disk paths). Multi-host striped pools remain the hardware functional-chain / backlog #1313 / #1314 lane. + +## What this lane covers + +`cargo nextest run --profile e2e-distributed -p e2e_test` selects `distributed::*`: + +- S3 put / get / head / list / copy / rename / delete / presign / empty object +- Object Lock COMPLIANCE, GOVERNANCE (with bypass), legal hold +- Versioning, version GET, delete marker +- Bucket replication between two 4-node clusters; hard quota +- Health / admin info / storageinfo / audit target list +- Pool expand, decommission, rebalance, checksum integrity, S3 during move +- Site replication object convergence +- High-concurrency PUT/GET; concurrent PUT during decommission +- Node kill/restart, full process restart, drive offline, volume-proxy blackhole +- Multipart and cross-node listing agreement + +## Existing Actions gaps this lane does not replace + +Those suites stay in place; this lane fills the in-tree 4×4 hole they leave. + +| Existing lane | Gap | +|---|---| +| `rustfs-*-test.yml` functional chain | Clones private `rustfs/auto-testing`, runs on three shared VMs (`vm000`–`vm002`), `continue-on-error: true`, not a merge signal, not 4 nodes | +| `e2e-smoke` / `e2e-full` | Almost all cases are single-node | +| `e2e-nightly` | 4-node cluster faults and heal, not S3/lock/versioning/quota/decommission matrix | +| `e2e-repl-nightly` | Site and bucket replication on 1–3 *single-node* processes | +| `e2e-s3tests.yml` `multi` | Weekly ceph/s3-tests against Docker 4-node; not lock/WORM, decommission, chaos, or checksum integrity | +| `crates/e2e_test/src/chaos.rs` | Single-node disk faults only | + +Hardware power-loss, NIC pull, and real disk replacement still belong on the smoke-testing VMs. This lane simulates those with SIGKILL, directory rename, and `FaultProxy` blackhole. + +## Run + +```bash +cargo build -p rustfs --bins +cargo nextest run --profile e2e-distributed -p e2e_test +``` + +Membership is pinned by `.config/e2e-distributed-selection.txt`. Update it with `python3 ./scripts/check_test_wiring.py --update-profile e2e-distributed linux` after adding or renaming a case.