Compare commits

..

7 Commits

Author SHA1 Message Date
overtrue 980f3abbd3 chore: merge main PR evidence guidance 2026-09-05 19:00:05 +08:00
overtrue e36650827b chore: integrate shared quick checks for E2E validation 2026-09-05 18:59:07 +08:00
overtrue 09c8e10d5e feat(test): verify the E2E server build and source identity 2026-09-05 18:58:01 +08:00
overtrue 0ff03c596c chore: merge main after ECStore compile repair 2026-09-05 18:46:20 +08:00
overtrue a74919db8e fix(ci): reject dependencies on required quick checks 2026-09-05 18:17:13 +08:00
overtrue e77c6f0ca5 fix(ci): install actionlint from its verified release 2026-09-05 17:42:14 +08:00
overtrue 3149411943 fix(ci): share quick checks and lint workflows 2026-09-05 17:37:20 +08:00
23 changed files with 1108 additions and 804 deletions
+1
View File
@@ -38,6 +38,7 @@ script-tests: ## Run shell script tests
./scripts/test_python_bin.sh
./scripts/check_embedded_secrets.sh --self-test
$(RUSTFS_PYTHON_BIN) ./scripts/check_test_wiring.py --self-test
$(RUSTFS_PYTHON_BIN) ./scripts/test_e2e_binary.py
$(RUSTFS_PYTHON_BIN) ./scripts/check_security_coverage.py --self-test
$(RUSTFS_PYTHON_BIN) ./scripts/check_scheduled_validation_freshness.py --self-test
$(RUSTFS_PYTHON_BIN) ./scripts/test_security_workflow.py
+116
View File
@@ -0,0 +1,116 @@
# 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/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.
name: Quick Checks
description: Run the shared compile-free RustFS quality checks.
runs:
using: composite
steps:
- name: Install quality tools
uses: taiki-e/install-action@bffeee26d4db9be238a4ea78d8826604ebcb594d # v2
with:
tool: |
ripgrep@15.2.0
shellcheck@0.11.0
- name: Install actionlint
shell: bash
run: |
actionlint_dir="$(mktemp -d "${RUNNER_TEMP}/actionlint.XXXXXX")"
curl --fail --location --silent --show-error \
--output "$actionlint_dir/actionlint.tar.gz" \
https://github.com/rhysd/actionlint/releases/download/v1.7.12/actionlint_1.7.12_linux_amd64.tar.gz
echo "8aca8db96f1b94770f1b0d72b6dddcb1ebb8123cb3712530b08cc387b349a3d8 $actionlint_dir/actionlint.tar.gz" | sha256sum --check --status
tar -xzf "$actionlint_dir/actionlint.tar.gz" -C "$actionlint_dir" actionlint
rm "$actionlint_dir/actionlint.tar.gz"
echo "$actionlint_dir" >> "$GITHUB_PATH"
- name: Install Rust toolchain
uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 # stable
with:
components: rustfmt
- name: Check workflow syntax and shell scripts
shell: bash
run: shellcheck --version && actionlint
- name: Check code formatting
shell: bash
run: cargo fmt --all --check
- name: Check unsafe code allowances
shell: bash
run: ./scripts/check_unsafe_code_allowances.sh
- name: Check layered dependencies
shell: bash
run: ./scripts/check_layer_dependencies.sh
- name: Check architecture migration rules
shell: bash
run: ./scripts/check_architecture_migration_rules.sh
- name: Check logging guardrails
shell: bash
run: ./scripts/check_logging_guardrails.sh
- name: Check error other(format!) ratchet
shell: bash
run: ./scripts/check_error_other_format_ratchet.sh
- name: Check tokio io-uring feature guard
shell: bash
run: ./scripts/check_no_tokio_io_uring.sh
- name: Check extension schema boundaries
shell: bash
run: ./scripts/check_extension_schema_boundaries.sh
- name: Check body-cache whitelist guard
shell: bash
run: ./scripts/check_body_cache_whitelist.sh
- name: Check s3s footprint ratchet
shell: bash
run: ./scripts/check_s3s_footprint.sh
- name: Check cryptographic capability wording
shell: bash
run: ./scripts/check_fips_wording.sh
- name: Check no embedded secret material
shell: bash
run: ./scripts/check_embedded_secrets.sh
- name: Check test wiring
shell: bash
run: |
python3 ./scripts/check_test_wiring.py --self-test
python3 ./scripts/test_e2e_binary.py
python3 ./scripts/check_scheduled_validation_freshness.py --self-test
python3 ./scripts/test_security_workflow.py
python3 ./scripts/check_test_wiring.py
- name: Check no planning docs committed
shell: bash
run: ./scripts/check_no_planning_docs.sh
- name: Check CI paths stay in sync
shell: bash
run: ./scripts/check_ci_paths_sync.sh
- name: Check io_uring lane --lib precondition
shell: bash
run: ./scripts/check_uring_lane_lib_only.sh
+6 -89
View File
@@ -12,24 +12,10 @@
# See the License for the specific language governing permissions and
# limitations under the License.
# Companion to ci.yml for required status checks.
#
# ci.yml skips docs-only pull requests via paths-ignore, but the branch ruleset
# requires a check named "Test and Lint" — without this workflow a docs-only PR
# would wait on it forever. This workflow triggers on exactly the paths ci.yml
# ignores and reports success under the same job name. Mixed PRs trigger both
# workflows and the real check still gates: a required check with any failing
# run blocks the merge.
# https://docs.github.com/en/repositories/configuring-branches-and-merges-in-your-repository/defining-the-mergeability-of-pull-requests/troubleshooting-required-status-checks#handling-skipped-but-required-checks
#
# "Quick Checks" is mirrored here ahead of the ruleset change that will make it
# required too (rustfs/backlog#1599). Until that change lands this job is
# inert; mirroring it first is what lets the ruleset change happen without
# stranding docs-only PRs on a check nobody reports.
#
# Keep the paths list below in sync with the pull_request paths-ignore list
# in ci.yml, and keep the quick-checks steps below byte-identical to the
# quick-checks job in ci.yml.
# Reports the existing required checks for paths excluded by ci.yml.
# Mixed PRs can trigger both workflows; their Quick Checks jobs use one shared
# action to keep validation coverage aligned. Keep this paths list in sync with
# ci.yml's pull_request.paths-ignore via scripts/check_ci_paths_sync.sh.
name: Continuous Integration (docs only)
@@ -59,19 +45,6 @@ permissions:
contents: read
jobs:
# Deliberately NOT a bare `echo`. Once "Quick Checks" becomes a required
# check, ci.yml gates every expensive job behind it, so a mixed PR reports
# two check runs with this name: the real one (45-51s) and this companion.
# GitHub has no written contract for how it picks between same-named
# required check runs ("latest wins" vs "any failure blocks"), so instead of
# relying on ordering we make both runs execute the same commands against
# the same merge ref — their conclusions are then necessarily identical and
# the choice does not matter. Keep these steps byte-identical to the
# quick-checks job in ci.yml (a guard script that asserts this, and the paths
# sync below, is tracked in rustfs/backlog#1603).
#
# For a genuinely docs-only PR this adds no strictness (no code changed, so
# fmt and the guards always pass) and costs ~50s of ubuntu-latest.
quick-checks:
name: Quick Checks
runs-on: ubuntu-latest
@@ -82,64 +55,8 @@ jobs:
with:
persist-credentials: false
- name: Install ripgrep
uses: taiki-e/install-action@bffeee26d4db9be238a4ea78d8826604ebcb594d # v2
with:
tool: ripgrep@15.2.0
- name: Install Rust toolchain
uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 # stable
with:
components: rustfmt
- name: Check code formatting
run: cargo fmt --all --check
- name: Check unsafe code allowances
run: ./scripts/check_unsafe_code_allowances.sh
- name: Check layered dependencies
run: ./scripts/check_layer_dependencies.sh
- name: Check architecture migration rules
run: ./scripts/check_architecture_migration_rules.sh
- name: Check logging guardrails
run: ./scripts/check_logging_guardrails.sh
- name: Check tokio io-uring feature guard
run: ./scripts/check_no_tokio_io_uring.sh
- name: Check extension schema boundaries
run: ./scripts/check_extension_schema_boundaries.sh
- name: Check body-cache whitelist guard
run: ./scripts/check_body_cache_whitelist.sh
- name: Check s3s footprint ratchet
run: ./scripts/check_s3s_footprint.sh
- name: Check cryptographic capability wording
run: ./scripts/check_fips_wording.sh
- name: Check no embedded secret material
run: ./scripts/check_embedded_secrets.sh
- name: Check test wiring
run: |
python3 ./scripts/check_test_wiring.py --self-test
python3 ./scripts/check_scheduled_validation_freshness.py --self-test
python3 ./scripts/test_security_workflow.py
python3 ./scripts/check_test_wiring.py
- name: Check no planning docs committed
run: ./scripts/check_no_planning_docs.sh
- name: Check CI paths stay in sync
run: ./scripts/check_ci_paths_sync.sh
- name: Check io_uring lane --lib precondition
run: ./scripts/check_uring_lane_lib_only.sh
- name: Run shared quick checks
uses: ./.github/actions/quick-checks
test-and-lint:
name: Test and Lint
+17 -77
View File
@@ -100,12 +100,7 @@ jobs:
- name: Typos check with custom config file
uses: crate-ci/typos@37bb98842b0d8c4ffebdb75301a13db0267cef89 # master
# Fast, compile-free checks that fail early so contributors get feedback in
# ~1 minute instead of waiting for the full test job.
#
# These steps are mirrored byte-for-byte in ci-docs-only.yml so that a mixed
# PR, which reports two check runs named "Quick Checks", cannot get one red
# and one green. Edit both jobs together.
# Fail early with compile-free checks shared with docs-only CI.
quick-checks:
name: Quick Checks
if: github.event_name != 'pull_request' || github.event.action != 'closed'
@@ -117,67 +112,8 @@ jobs:
with:
persist-credentials: false
- name: Install ripgrep
uses: taiki-e/install-action@bffeee26d4db9be238a4ea78d8826604ebcb594d # v2
with:
tool: ripgrep@15.2.0
- name: Install Rust toolchain
uses: dtolnay/rust-toolchain@29eef336d9b2848a0b548edc03f92a220660cdb8 # stable
with:
components: rustfmt
- name: Check code formatting
run: cargo fmt --all --check
- name: Check unsafe code allowances
run: ./scripts/check_unsafe_code_allowances.sh
- name: Check layered dependencies
run: ./scripts/check_layer_dependencies.sh
- name: Check architecture migration rules
run: ./scripts/check_architecture_migration_rules.sh
- name: Check logging guardrails
run: ./scripts/check_logging_guardrails.sh
- name: Check error other(format!) ratchet
run: ./scripts/check_error_other_format_ratchet.sh
- name: Check tokio io-uring feature guard
run: ./scripts/check_no_tokio_io_uring.sh
- name: Check extension schema boundaries
run: ./scripts/check_extension_schema_boundaries.sh
- name: Check body-cache whitelist guard
run: ./scripts/check_body_cache_whitelist.sh
- name: Check s3s footprint ratchet
run: ./scripts/check_s3s_footprint.sh
- name: Check cryptographic capability wording
run: ./scripts/check_fips_wording.sh
- name: Check no embedded secret material
run: ./scripts/check_embedded_secrets.sh
- name: Check test wiring
run: |
python3 ./scripts/check_test_wiring.py --self-test
python3 ./scripts/check_scheduled_validation_freshness.py --self-test
python3 ./scripts/test_security_workflow.py
python3 ./scripts/check_test_wiring.py
- name: Check no planning docs committed
run: ./scripts/check_no_planning_docs.sh
- name: Check CI paths stay in sync
run: ./scripts/check_ci_paths_sync.sh
- name: Check io_uring lane --lib precondition
run: ./scripts/check_uring_lane_lib_only.sh
- name: Run shared quick checks
uses: ./.github/actions/quick-checks
test-and-lint:
name: Test and Lint
@@ -646,13 +582,15 @@ jobs:
install-build-packaging-tools: 'false'
- name: Build debug binary
run: cargo build -p rustfs --bins --features e2e-test-hooks
run: python3 scripts/e2e_binary.py build --bins --features e2e-test-hooks
- name: Upload debug binary
uses: actions/upload-artifact@b7c566a772e6b6bfb58ed0dc250532a479d7789f # v6
with:
name: rustfs-debug-binary
path: target/debug/rustfs
path: |
target/debug/rustfs
target/debug/rustfs.e2e.json
if-no-files-found: error
retention-days: 1
@@ -684,13 +622,15 @@ jobs:
install-build-packaging-tools: 'false'
- name: Build debug binary with rio-v2
run: cargo build -p rustfs --bins --features rio-v2,e2e-test-hooks
run: python3 scripts/e2e_binary.py build --bins --features rio-v2,e2e-test-hooks
- name: Upload debug binary
uses: actions/upload-artifact@b7c566a772e6b6bfb58ed0dc250532a479d7789f # v6
with:
name: rustfs-debug-binary-rio-v2
path: target/debug/rustfs
path: |
target/debug/rustfs
target/debug/rustfs.e2e.json
if-no-files-found: error
retention-days: 1
@@ -839,7 +779,7 @@ jobs:
NEXTEST_ARCHIVE: ${{ runner.temp }}/rustfs-e2e-smoke.tar.zst
RUSTFS_E2E_LOG_DIR: ${{ runner.temp }}/rustfs-e2e-smoke-logs
run: |
cargo nextest run --profile e2e-smoke --archive-file "${NEXTEST_ARCHIVE}" \
python3 scripts/e2e_binary.py run --features e2e-test-hooks -- cargo nextest run --profile e2e-smoke --archive-file "${NEXTEST_ARCHIVE}" \
--status-level all --final-status-level all --failure-output final
- name: Upload e2e smoke diagnostics
@@ -875,7 +815,7 @@ jobs:
RUSTFS_TEST_PORT="$(python3 -c 'import socket; s=socket.socket(); s.bind(("127.0.0.1", 0)); print(s.getsockname()[1]); s.close()')"
RUSTFS_TEST_PORT="${RUSTFS_TEST_PORT}" \
RUSTFS_TEST_LOG="${RUN_ROOT}/rustfs.log" \
./scripts/e2e-run.sh ./target/debug/rustfs "${RUN_ROOT}/data"
python3 scripts/e2e_binary.py run --features e2e-test-hooks -- ./scripts/e2e-run.sh ./target/debug/rustfs "${RUN_ROOT}/data"
- name: Upload test logs
if: failure()
@@ -977,7 +917,7 @@ jobs:
# extend that filter, never add ad-hoc e2e jobs here. Reuses the downloaded
# debug binary; each test spawns its own rustfs server on a random port.
- name: Run e2e full suite
run: cargo nextest run --profile e2e-full -p e2e_test
run: python3 scripts/e2e_binary.py run --features e2e-test-hooks -- cargo nextest run --profile e2e-full -p e2e_test
- name: Upload junit
if: always()
@@ -1038,7 +978,7 @@ jobs:
- name: Run end-to-end tests
run: |
s3s-e2e --version
./scripts/e2e-run.sh ./target/debug/rustfs /tmp/rustfs
python3 scripts/e2e_binary.py run --features rio-v2,e2e-test-hooks -- ./scripts/e2e-run.sh ./target/debug/rustfs /tmp/rustfs
- name: Upload test logs
if: failure()
@@ -1081,7 +1021,7 @@ jobs:
S3_PORT="${S3_PORT}" \
DATA_ROOT="${RUN_ROOT}" \
S3TESTS_CONF=artifacts/s3tests-single/s3tests.conf \
./scripts/s3-tests/run.sh
python3 scripts/e2e_binary.py run --features e2e-test-hooks -- ./scripts/s3-tests/run.sh
- name: Upload s3 test artifacts
if: always()
@@ -1163,7 +1103,7 @@ jobs:
S3_PORT="${S3_PORT}" \
DATA_ROOT="${RUN_ROOT}" \
S3TESTS_CONF=artifacts/s3tests-single/s3tests.conf \
./scripts/s3-tests/run.sh
python3 scripts/e2e_binary.py run --features e2e-test-hooks -- ./scripts/s3-tests/run.sh
- name: Upload s3 test artifacts
if: always()
+9 -11
View File
@@ -89,14 +89,10 @@ jobs:
- name: Verify awscurl
run: test -x "$AWSCURL_PATH"
# Build the rustfs binary once up front. The e2e tests spawn it as a
# child process (crates/e2e_test/src/common.rs) and will build it on
# demand otherwise, but a single explicit build avoids several parallel
# nextest test processes racing to build it at once.
# Build once and carry its source/binary identity into the test invocation.
- name: Build rustfs binary
run: |
cargo build -p rustfs --bins
: > target/debug/rustfs.features
python3 scripts/e2e_binary.py build --bins
- name: Verify replication e2e membership
env:
@@ -108,7 +104,7 @@ jobs:
- name: Run replication e2e nightly suite
env:
RUSTFS_E2E_LOG_DIR: ${{ runner.temp }}/rustfs-e2e-repl-nightly-logs
run: cargo nextest run --profile e2e-repl-nightly -p e2e_test
run: python3 scripts/e2e_binary.py run -- cargo nextest run --profile e2e-repl-nightly -p e2e_test
- name: Upload nextest junit report
if: always()
@@ -144,8 +140,7 @@ jobs:
- name: Build rustfs binary
run: |
cargo build -p rustfs --bins --features e2e-test-hooks
: > target/debug/rustfs.features
python3 scripts/e2e_binary.py build --bins --features e2e-test-hooks
- name: Verify cluster fault e2e membership
env:
@@ -157,7 +152,7 @@ jobs:
- name: Run cluster fault e2e nightly suite
env:
RUSTFS_E2E_LOG_DIR: ${{ runner.temp }}/rustfs-e2e-nightly-logs
run: cargo nextest run --profile e2e-nightly -p e2e_test
run: python3 scripts/e2e_binary.py run --features e2e-test-hooks -- cargo nextest run --profile e2e-nightly -p e2e_test
- name: Upload cluster fault diagnostics
if: always()
@@ -198,6 +193,9 @@ jobs:
sudo apt-get install -y -qq iproute2
ss -tn state CLOSE-WAIT >/dev/null
- name: Build protocol server
run: python3 scripts/e2e_binary.py build --features "$RUSTFS_BUILD_FEATURES"
# The suite owns fixed protocol ports and serializes its internal cases.
- name: Verify protocol e2e membership
env:
@@ -210,7 +208,7 @@ jobs:
env:
RUSTFS_E2E_LOG_DIR: ${{ runner.temp }}/rustfs-protocol-e2e-logs
run: >-
cargo nextest run -j 1 --profile e2e-protocols -p e2e_test --no-capture
python3 scripts/e2e_binary.py run --features "$RUSTFS_BUILD_FEATURES" -- cargo nextest run -j 1 --profile e2e-protocols -p e2e_test --no-capture
- name: Upload protocol diagnostics
if: always()
+2 -3
View File
@@ -98,12 +98,11 @@ jobs:
- name: Build current RustFS binary
run: |
cargo build --locked -p rustfs --bin rustfs
: > target/debug/rustfs.features
python3 scripts/e2e_binary.py build
- name: Run upgrade compatibility test
run: |
cargo test --locked -p e2e_test \
python3 scripts/e2e_binary.py run -- cargo test --locked -p e2e_test \
"upgrade_compatibility_test::${{ matrix.test }}" \
-- --ignored --exact --nocapture
@@ -132,7 +132,7 @@ jobs:
s3api create-bucket --bucket "${RUSTFS_ODM_INTEROP_BUCKET}"
- name: Build the RustFS binary under test
run: cargo build --locked -p rustfs --bins
run: python3 scripts/e2e_binary.py build --bins
# The lane selects tests by module, so a rename would quietly shrink it.
# The committed digest in .config/e2e-odm-interop-selection.txt fails
@@ -143,7 +143,7 @@ jobs:
python3 ./scripts/check_test_wiring.py --check-profile e2e-odm-interop "${NEXTEST_LISTING}"
- name: Run the interop cases against MinIO
run: cargo nextest run --profile e2e-odm-interop -p e2e_test --no-tests=fail
run: python3 scripts/e2e_binary.py run -- cargo nextest run --profile e2e-odm-interop -p e2e_test --no-tests=fail
- name: Build the MinIO interop report
if: always()
@@ -251,7 +251,7 @@ jobs:
- name: Build the RustFS binary under test
if: steps.credentials.outputs.present == 'true'
run: cargo build --locked -p rustfs --bins
run: python3 scripts/e2e_binary.py build --bins
# A filterset that matches nothing is valid, so the count is asserted
# rather than inferred from a green run.
@@ -272,7 +272,7 @@ jobs:
- name: Run the three-case minimum
if: steps.credentials.outputs.present == 'true'
run: |
cargo nextest run --profile e2e-odm-interop -p e2e_test \
python3 scripts/e2e_binary.py run -- cargo nextest run --profile e2e-odm-interop -p e2e_test \
-E "${CLOUD_CASE_FILTER}" --no-tests=fail
- name: Build the ${{ matrix.provider }} interop report
+33 -43
View File
@@ -1,7 +1,7 @@
# e2e_test
End-to-end test suite for RustFS. Each test spawns a **real `rustfs` binary**
(built on demand from the workspace) and drives it over the network with the
(built and identified before the test invocation) and drives it over the network with the
AWS SDK (`aws-sdk-s3`), raw HTTP (`reqwest` / `awscurl`), or a protocol client
(FTPS / WebDAV / SFTP). This is the black-box integration layer: exhaustive
end-to-end behavior lives here, unit behavior stays in the source crates
@@ -31,32 +31,28 @@ Registered in [`src/lib.rs`](src/lib.rs). Grouped by concern:
## How to run
All commands assume repo root. `cargo test` triggers an on-demand build of the
`rustfs` binary from [`src/common.rs`](src/common.rs) (`rustfs_binary_path`) on
first use — the first invocation is slow, later ones reuse the binary.
All commands assume repo root and Python 3.9 or newer on Linux or macOS. Build the server once through the provenance entry point, then run the test command through the same script:
```bash
# Whole crate (default = ignored tests skipped)
cargo nextest run -p e2e_test
python3 scripts/e2e_binary.py build --features e2e-test-hooks
# Whole crate (ignored tests remain skipped)
python3 scripts/e2e_binary.py run --features e2e-test-hooks -- cargo nextest run -p e2e_test
# One module
cargo nextest run -p e2e_test -E 'test(list_objects_v2_pagination_test)'
# PR smoke subset (see "CI smoke subset" below)
cargo nextest run --profile e2e-smoke -p e2e_test
# ILM serial lane — ignored lifecycle tests, single-threaded (mirrors CI)
cargo nextest run -j1 --run-ignored ignored-only -p rustfs-scanner -p rustfs \
-E 'binary(lifecycle_integration_test) or (package(rustfs) and test(lifecycle_transition_api_test))'
python3 scripts/e2e_binary.py run --features e2e-test-hooks -- cargo nextest run -p e2e_test -E 'test(list_objects_v2_pagination_test)'
# PR smoke subset
python3 scripts/e2e_binary.py run --features e2e-test-hooks -- cargo nextest run --profile e2e-smoke -p e2e_test
```
The protocols suite has its own contract (fixed bind ports 90229301,
single-worker execution, feature-gated scheduling) documented in
[`src/protocols/README.md`](src/protocols/README.md). `RUSTFS_BUILD_FEATURES`
selects which features the spawned binary is built with; leave it unset to run
every protocol entry. Use the exact profile command under
[Troubleshooting](#troubleshooting) for CI-equivalent execution.
`build` records the source contents, HEAD, resolved Cargo features, profile, toolchain, and binary SHA-256 beside the executable in `rustfs.e2e.json`. `run` validates that identity before and after the command, preserves command failures, and removes its temporary run receipt on completion. The Rust harness checks that receipt before starting each server; it never compiles a server inside a test process. Source or binary changes during a run invalidate the result, even when the test command succeeds. Use an isolated worktree and keep it unchanged until the command finishes.
The additional `--features` arguments must match between `build` and `run`; Cargo defaults remain enabled. The wrapper supplies `RUSTFS_BUILD_FEATURES` from Cargo's resolved feature list, including features enabled by `full`. Protocol helpers require a subset of that list. `CARGO_TARGET_DIR` and `--profile release` are supported. An in-workspace target directory must be Git-ignored; tracked files are always included in the source identity. `build --bins` preserves CI lanes that compile all RustFS binary targets. For a downloaded artifact, copy both the executable and its sidecar, then use `run`; do not generate a new identity for an arbitrary prebuilt binary. `CARGO_BIN_EXE_rustfs` cannot override the verified executable.
Each build/run holds an exclusive `rustfs.e2e.lock` marker beside the binary; concurrent wrappers fail immediately. Use a private target directory and do not run ordinary Cargo builds against it while tests are active: Cargo does not honor this marker. Interrupted runs fail and terminate their command group. After an uncatchable kill, inspect the PID recorded in a leftover marker and remove it only after confirming its owner has stopped. Embedded file symlinks are hashed through their target; embedded directory symlinks are rejected because their contents cannot be enumerated safely by this entry point.
The protocols suite has its own fixed-port and single-worker contract in [`src/protocols/README.md`](src/protocols/README.md). Use its command under [Troubleshooting](#troubleshooting).
### `#[ignore]` semantics
@@ -122,7 +118,7 @@ via `create_s3_client(idx)` / `create_all_clients()`. See
| `wait_for_server_ready` | Poll readiness before issuing requests |
| `create_s3_client` / `create_test_bucket` / `delete_test_bucket` | aws-sdk-s3 client + bucket lifecycle |
| `find_available_port` | Random free port (isolation primitive) |
| `rustfs_binary_path` / `_with_features` | Locate/build the binary; honors `RUSTFS_BUILD_FEATURES` |
| `rustfs_binary_path` / `_with_features` | Verify this run's binary receipt and required feature subset |
| `requested_rustfs_build_features` / `rustfs_build_feature_enabled` | Feature-gate a test to what the binary was built with |
| `execute_awscurl` / `awscurl_post` / `_get` / `_put` / `_delete` / `awscurl_post_sts_form_urlencoded` | Admin/STS API calls via `awscurl`; missing binaries are test failures |
| `replication_fast_env` | Env vars that shrink replication timers (from repl-4); pass to `start_rustfs_server_with_env` |
@@ -185,32 +181,26 @@ the wiring source of truth. Committed test-ID digests under
**Reproduce a CI failure locally** — run the exact profile/lane:
```bash
# Smoke (e2e-tests job) — includes the 20 fast replication tests
cargo nextest run --profile e2e-smoke -p e2e_test
# Full single-node merge/main lane
cargo nextest run --profile e2e-full -p e2e_test
# Cluster fault nightly lane
cargo nextest run --profile e2e-nightly -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
RUSTFS_BUILD_FEATURES=ftps,webdav,sftp \
cargo nextest run -j 1 --profile e2e-protocols -p e2e_test --no-capture
# ILM serial lane
# Smoke, full, and cluster lanes share a server with fault-test hooks.
python3 scripts/e2e_binary.py build --features e2e-test-hooks
python3 scripts/e2e_binary.py run --features e2e-test-hooks -- cargo nextest run --profile e2e-smoke -p e2e_test
python3 scripts/e2e_binary.py run --features e2e-test-hooks -- cargo nextest run --profile e2e-full -p e2e_test
python3 scripts/e2e_binary.py run --features e2e-test-hooks -- cargo nextest run --profile e2e-nightly -p e2e_test
# Replication nightly uses the default server; awscurl is required for STS.
python3 scripts/e2e_binary.py build
python3 scripts/e2e_binary.py run -- cargo nextest run --profile e2e-repl-nightly -p e2e_test
# Protocol nightly owns fixed ports.
python3 scripts/e2e_binary.py build --features ftps,webdav,sftp
python3 scripts/e2e_binary.py run --features ftps,webdav,sftp -- cargo nextest run -j 1 --profile e2e-protocols -p e2e_test --no-capture
# The ILM serial lane does not use this server harness.
cargo nextest run -j1 --run-ignored ignored-only -p rustfs-scanner -p rustfs \
-E 'binary(lifecycle_integration_test) or (package(rustfs) and test(lifecycle_transition_api_test))'
# s3s-e2e black box
./scripts/e2e-run.sh ./target/debug/rustfs /tmp/rustfs-e2e-data
```
**Stale binary.** Tests build the `rustfs` binary once and reuse it. To avoid
rebuilding while iterating on tests, `common.rs` reuses an existing binary when
running *inside* the e2e test process even if sources changed
(`can_reuse_inside_e2e`, [`src/common.rs`](src/common.rs) line 98). Downside: if
you changed **server** code, force a rebuild with
`cargo build -p rustfs` (or `touch` a source file outside the reuse window)
before re-running, or CI's freshly built artifact will diverge from your local
one.
**Stale or unverified binary.** Re-run the matching `build` command after changing source or features, then invoke tests through `run`. A missing receipt, copied old executable, or mismatched build identity is a prerequisite failure. Bare Cargo invocations that start a server deliberately fail; unit tests that do not start a server can still run directly.
**Port already in use / orphan processes.** A hard-killed run can leak a
`rustfs` child holding its port. Find and kill it:
+123 -149
View File
@@ -31,7 +31,6 @@ use rustfs_signer::constants::UNSIGNED_PAYLOAD;
use rustfs_signer::sign_v4;
use s3s::Body;
use serde_json;
use std::ffi::OsStr;
use std::fs as stdfs;
use std::io::ErrorKind;
use std::net::SocketAddr;
@@ -44,7 +43,6 @@ use tokio::net::TcpStream;
use tokio::time::sleep;
use tracing::{error, info, warn};
use uuid::Uuid;
use walkdir::WalkDir;
// Common constants for all E2E tests
pub const DEFAULT_ACCESS_KEY: &str = "rustfsadmin";
@@ -365,59 +363,75 @@ fn resolve_rustfs_binary_path(workspace: &Path, configured_target_dir: Option<&P
path
}
/// Resolve the RustFS binary relative to the workspace, optionally requesting build features.
/// Resolve the server verified by `scripts/e2e_binary.py run` for this test invocation.
/// Requested features are a required subset of the server's resolved Cargo features.
pub fn rustfs_binary_path_with_features(requested_features: Option<&str>) -> PathBuf {
if let Some(path) = std::env::var_os("CARGO_BIN_EXE_rustfs") {
return PathBuf::from(path);
}
let requested_features = requested_features.and_then(normalize_rustfs_build_features);
let workspace = workspace_root();
let configured_target_dir = std::env::var_os("CARGO_TARGET_DIR").map(PathBuf::from);
let binary_path = resolve_rustfs_binary_path(&workspace, configured_target_dir.as_deref());
let binary_path = std::env::var_os("CARGO_BIN_EXE_rustfs")
.map(PathBuf::from)
.unwrap_or_else(|| resolve_rustfs_binary_path(&workspace, configured_target_dir.as_deref()));
let receipt_path = std::env::var_os("RUSTFS_E2E_BINARY_RECEIPT").map(PathBuf::from);
receipt_path
.ok_or_else(|| std::io::Error::new(ErrorKind::NotFound, "missing E2E run receipt"))
.and_then(|receipt| verify_e2e_binary_receipt(&receipt, &workspace, &binary_path, requested_features))
.unwrap_or_else(|error| {
panic!(
"E2E server prerequisite failed: {error}. Build with `python3 scripts/e2e_binary.py build --features <features>` and run tests with `python3 scripts/e2e_binary.py run --features <features> -- cargo nextest run ...`"
)
})
}
let features_match = binary_features_match(&binary_path, requested_features.as_deref());
let source_is_newer = workspace_sources_newer_than_binary(&binary_path);
let can_reuse_inside_e2e = running_inside_e2e_test_binary() && requested_features.is_none() && features_match;
if binary_path.is_file() && features_match && (!source_is_newer || can_reuse_inside_e2e) {
if source_is_newer {
warn!(
"RustFS binary at {:?} appears older than workspace sources; reusing it inside cargo test to avoid nested builds",
binary_path
);
}
info!("Using existing RustFS binary at {:?}", binary_path);
return binary_path;
#[derive(serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct E2eBinaryReceipt {
schema: u32,
workspace: PathBuf,
binary: PathBuf,
size: u64,
modified_ns: u128,
features: Vec<String>,
}
fn verify_e2e_binary_receipt(
receipt_path: &Path,
workspace: &Path,
binary_path: &Path,
requested_features: Option<&str>,
) -> std::io::Result<PathBuf> {
let receipt: E2eBinaryReceipt = serde_json::from_slice(&stdfs::read(receipt_path)?)?;
let binary = binary_path.canonicalize()?;
let metadata = binary.metadata()?;
let modified_ns = metadata
.modified()?
.duration_since(std::time::UNIX_EPOCH)
.map_err(std::io::Error::other)?
.as_nanos();
// The runner hashes source and binary before/after the entire suite. Each
// nextest process checks only this invocation's path, features, and file stat.
if receipt.schema != 1
|| receipt.workspace != workspace.canonicalize()?
|| receipt.binary != binary
|| !metadata.is_file()
|| receipt.size != metadata.len()
|| receipt.modified_ns != modified_ns
{
return Err(std::io::Error::new(
ErrorKind::InvalidData,
"E2E server differs from this run's verified binary",
));
}
info!("Building RustFS binary to ensure it's up to date...");
build_rustfs_binary(requested_features.as_deref(), &binary_path);
info!("Using RustFS binary at {:?}", binary_path);
binary_path
}
fn workspace_sources_newer_than_binary(binary_path: &PathBuf) -> bool {
let Ok(binary_meta) = std::fs::metadata(binary_path) else {
return true;
};
let Ok(binary_modified) = binary_meta.modified() else {
return true;
};
let workspace = workspace_root();
let watch_roots = [
workspace.join("Cargo.toml"),
workspace.join("Cargo.lock"),
workspace.join("rustfs"),
workspace.join("crates"),
];
watch_roots.iter().any(|path| path_is_newer_than(binary_modified, path))
}
fn running_inside_e2e_test_binary() -> bool {
std::env::var("CARGO_PKG_NAME").is_ok_and(|value| value == "e2e_test")
if let Some(requested) = requested_features.and_then(normalize_rustfs_build_features)
&& requested
.split(',')
.any(|feature| !receipt.features.iter().any(|actual| actual == feature))
{
return Err(std::io::Error::new(
ErrorKind::InvalidInput,
"E2E server is missing a requested build feature",
));
}
Ok(binary)
}
pub fn requested_rustfs_build_features() -> Option<String> {
@@ -447,96 +461,6 @@ pub fn rustfs_build_feature_enabled(requested_features: Option<&str>, required_f
.any(|feature| feature.eq_ignore_ascii_case(RUSTFS_FULL_FEATURE) || feature.eq_ignore_ascii_case(required_feature))
}
fn rustfs_binary_features_stamp_path(binary_path: &Path) -> PathBuf {
binary_path.with_extension("features")
}
fn binary_features_match(binary_path: &Path, requested_features: Option<&str>) -> bool {
let stamp_path = rustfs_binary_features_stamp_path(binary_path);
let recorded = stdfs::read_to_string(stamp_path)
.ok()
.and_then(|value| normalize_rustfs_build_features(&value));
let requested = requested_features.and_then(normalize_rustfs_build_features);
match requested.as_deref() {
Some(features) => recorded.as_deref() == Some(features),
None => recorded.is_none(),
}
}
fn path_is_newer_than(binary_modified: std::time::SystemTime, path: &Path) -> bool {
if path.is_file() {
return std::fs::metadata(path)
.and_then(|meta| meta.modified())
.map(|modified| modified > binary_modified)
.unwrap_or(false);
}
if !path.is_dir() {
return false;
}
WalkDir::new(path)
.into_iter()
.filter_entry(|entry| {
let name = entry.file_name();
name != OsStr::new("target") && name != OsStr::new(".git")
})
.filter_map(Result::ok)
.filter(|entry| entry.file_type().is_file())
.any(|entry| {
std::fs::metadata(entry.path())
.and_then(|meta| meta.modified())
.map(|modified| modified > binary_modified)
.unwrap_or(false)
})
}
/// Build the RustFS binary using cargo
fn build_rustfs_binary(requested_features: Option<&str>, binary_path: &Path) {
let workspace = workspace_root();
info!("Building RustFS binary from workspace: {:?}", workspace);
let _profile = if cfg!(debug_assertions) {
info!("Building in debug mode");
"dev"
} else {
info!("Building in release mode");
"release"
};
let mut cmd = Command::new("cargo");
cmd.current_dir(&workspace).args(["build", "--bin", "rustfs"]);
if let Some(features) = requested_features {
cmd.arg("--features").arg(features);
info!("Building with features: {}", features);
}
if !cfg!(debug_assertions) {
cmd.arg("--release");
}
info!(
"Executing: cargo build --bin rustfs {}",
if cfg!(debug_assertions) { "" } else { "--release" }
);
let output = cmd.output().expect("Failed to execute cargo build command");
if !output.status.success() {
let stderr = String::from_utf8_lossy(&output.stderr);
panic!("Failed to build RustFS binary. Error: {stderr}");
}
let stamp_path = rustfs_binary_features_stamp_path(binary_path);
if let Err(err) = stdfs::write(&stamp_path, requested_features.unwrap_or_default()) {
warn!("Failed to write RustFS feature stamp {:?}: {}", stamp_path, err);
}
info!("✅ RustFS binary built successfully");
}
fn awscurl_binary_path() -> PathBuf {
std::env::var_os("AWSCURL_PATH")
.map(PathBuf::from)
@@ -2073,16 +1997,66 @@ mod tests {
}
#[test]
fn binary_feature_stamp_matching_uses_normalized_features() {
let binary_path = std::env::temp_dir().join(format!("rustfs-feature-stamp-test-{}", Uuid::new_v4()));
let stamp_path = rustfs_binary_features_stamp_path(&binary_path);
fn explicit_binary_without_run_receipt_is_rejected() {
const CHILD_ENV: &str = "RUSTFS_E2E_RECEIPT_TEST_CHILD";
if std::env::var_os(CHILD_ENV).is_some() {
rustfs_binary_path_with_features(None);
return;
}
let executable = std::env::current_exe().expect("locate isolated test process");
let output = Command::new(&executable)
.args([
"--exact",
"common::tests::explicit_binary_without_run_receipt_is_rejected",
"--nocapture",
])
.env(CHILD_ENV, "1")
.env("CARGO_BIN_EXE_rustfs", &executable)
.env_remove("RUSTFS_E2E_BINARY_RECEIPT")
.output()
.expect("run the missing-receipt scenario with isolated environment variables");
assert!(!output.status.success(), "an explicit binary must not bypass run verification");
assert!(String::from_utf8_lossy(&output.stderr).contains("missing E2E run receipt"));
}
stdfs::write(&stamp_path, " SFTP, ftps ").expect("write feature stamp");
assert!(binary_features_match(&binary_path, Some("sftp,ftps")));
assert!(binary_features_match(&binary_path, Some(" SFTP, FTPS ")));
assert!(!binary_features_match(&binary_path, Some("sftp")));
stdfs::remove_file(stamp_path).ok();
#[test]
fn e2e_run_receipt_rejects_replaced_binary_and_missing_features() {
let directory = std::env::temp_dir().join(format!("rustfs-e2e-receipt-test-{}", Uuid::new_v4()));
stdfs::create_dir(&directory).expect("create receipt fixture");
let binary = directory.join("rustfs");
let receipt = directory.join("receipt.json");
stdfs::write(&binary, "server").expect("write fixture binary");
let metadata = binary.metadata().expect("stat fixture binary");
let record = serde_json::json!({
"schema": 1,
"workspace": directory.canonicalize().expect("canonical workspace"),
"binary": binary.canonicalize().expect("canonical binary"),
"size": metadata.len(),
"modified_ns": metadata.modified().expect("modified time").duration_since(std::time::UNIX_EPOCH).expect("positive timestamp").as_nanos(),
"features": ["default", "full", "ftps", "webdav", "sftp"]
});
stdfs::write(&receipt, serde_json::to_vec(&record).expect("serialize receipt")).expect("write receipt");
verify_e2e_binary_receipt(&receipt, &directory, &binary, Some("sftp,webdav")).expect("resolved feature subset");
verify_e2e_binary_receipt(&receipt, &directory, &binary, Some("full")).expect("full was actually requested");
assert_eq!(
verify_e2e_binary_receipt(&receipt, &directory, &binary, Some("rio-v2"))
.expect_err("full does not enable rio-v2")
.kind(),
ErrorKind::InvalidInput
);
let other = directory.join("old-server");
stdfs::write(&other, "server").expect("write alternate binary");
assert!(verify_e2e_binary_receipt(&receipt, &directory, &other, None).is_err());
stdfs::write(&binary, "different server").expect("replace fixture binary");
assert!(verify_e2e_binary_receipt(&receipt, &directory, &binary, None).is_err());
stdfs::remove_file(&receipt).expect("remove expired receipt");
assert_eq!(
verify_e2e_binary_receipt(&receipt, &directory, &binary, None)
.expect_err("expired receipt")
.kind(),
ErrorKind::NotFound
);
stdfs::remove_dir_all(directory).expect("remove receipt fixture");
}
/// Build a cluster environment struct in-memory (no ports, no processes) so
+3 -5
View File
@@ -17,15 +17,13 @@ Use the canonical CI-equivalent protocol command in the parent
For targeted debugging of the core suite only:
```bash
RUSTFS_BUILD_FEATURES=ftps,webdav,sftp cargo test --package e2e_test test_protocol_core_suite -- --test-threads=1 --nocapture
python3 scripts/e2e_binary.py build --features ftps,webdav,sftp
python3 scripts/e2e_binary.py run --features ftps,webdav,sftp -- cargo test --package e2e_test test_protocol_core_suite -- --test-threads=1 --nocapture
```
This targeted command does not cover the full `e2e-protocols` profile.
`RUSTFS_BUILD_FEATURES` controls which features the test rustfs binary is
built with. When this variable is set, the protocol test runner schedules
only entries whose protocol is present in the requested feature list. Leave
it unset to run every protocol entry.
`e2e_binary.py` supplies `RUSTFS_BUILD_FEATURES` from the verified server's resolved Cargo features. The protocol runner schedules only entries present in that feature list; helpers check that their required features are available without rebuilding the server.
`--test-threads=1` is required because every entry spawns a rustfs server
on fixed bind ports.
-23
View File
@@ -774,29 +774,6 @@ impl DataUsageCache {
(visited == expected_entries).then_some(entry)
}
pub(crate) fn has_complete_root_inventory(&self, bucket_keys: &HashSet<String>) -> bool {
let Some(root) = self.find(DATA_USAGE_ROOT) else {
return false;
};
// Set roots only connect bucket entries. Scalar data at the root, an
// extra bucket, or an orphan must not disappear during bucket folding.
root.children.len() == bucket_keys.len()
&& bucket_keys.iter().all(|key| root.children.contains(key))
&& root.size == 0
&& root.objects == 0
&& root.versions == 0
&& root.delete_markers == 0
&& root.failed_objects == 0
&& !root.compacted
&& root.obj_sizes.is_empty()
&& root.obj_versions.is_empty()
&& root.replication_stats.is_none()
&& root.all_tier_stats.is_none()
&& root.unknown_tier_stats.is_none()
&& root.tier_accounting_proof.is_none()
&& self.checked_flatten_complete(DATA_USAGE_ROOT).is_some()
}
fn checked_flatten_inner(&self, path: &str) -> Option<(DataUsageEntry, usize)> {
let root_key = hash_path(path).key();
let (root_key, root) = self.cache.get_key_value(&root_key)?;
-30
View File
@@ -4758,28 +4758,6 @@ async fn usage_bootstrap_does_not_overwrite_concurrent_replacement() {
#[serial]
async fn scanner_usage_state_reset_publishes_fenced_bootstrap_marker() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
let quota_ledger_path = "config/quota-ledger/reserved-bucket.json";
let quota_ledger = serde_json::to_vec(&serde_json::json!({
"version": 1,
"bucket_incarnation": "00000000-0000-0000-0000-000000000001",
"quota_revision_unix_nanos": 1,
"accounted_usage": 100,
"reservations": {
"00000000-0000-0000-0000-000000000002": {
"object": "pending-object",
"old_size": 0,
"new_size": 64,
"created_at": 1,
"pool_index": 0,
"set_index": 0,
"commit_started": true
}
}
}))
.expect("quota ledger fixture should encode");
save_config(store.clone(), quota_ledger_path, quota_ledger.clone())
.await
.expect("independent quota reservations should persist");
let cycle = CurrentCycle {
current: 41,
next: 42,
@@ -4838,14 +4816,6 @@ async fn scanner_usage_state_reset_publishes_fenced_bootstrap_marker() {
assert!(!data_usage_info_has_persisted_baseline_identity(&usage));
assert_eq!(usage.scanner_epoch, Some(9));
assert_eq!(
read_config(store.clone(), quota_ledger_path)
.await
.expect("quota ledger must remain readable after scanner reset"),
quota_ledger,
"scanner reset must preserve incarnation and outstanding reserved bytes exactly"
);
for path in [
usage_backup_path.as_str(),
LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str(),
-24
View File
@@ -312,30 +312,6 @@ fn scanner_bucket_plan_digest(buckets: &[BucketInfo], activity_digest: [u8; 32])
DataUsageScanPlanDigest(hasher.finalize().into())
}
fn scanner_bucket_inventory_is_complete(
all_buckets: &[BucketInfo],
buckets_by_source: &HashMap<DataUsageCacheSource, Vec<BucketInfo>>,
) -> bool {
let inventory = all_buckets
.iter()
.map(|bucket| (bucket.name.as_str(), bucket.created))
.collect::<HashMap<_, _>>();
if inventory.len() != all_buckets.len() || inventory.keys().any(|name| name.is_empty() || *name == DATA_USAGE_ROOT) {
return false;
}
let mut covered = HashSet::with_capacity(inventory.len());
for buckets in buckets_by_source.values() {
let mut set_names = HashSet::with_capacity(buckets.len());
for bucket in buckets {
if !set_names.insert(bucket.name.as_str()) || inventory.get(bucket.name.as_str()) != Some(&bucket.created) {
return false;
}
covered.insert(bucket.name.as_str());
}
}
covered.len() == inventory.len()
}
fn scanner_bucket_cache_digest(
scan_plan_digest: DataUsageScanPlanDigest,
dirty_generation: Option<u64>,
+25 -84
View File
@@ -213,85 +213,10 @@ pub(super) fn cache_snapshot_is_current(
)
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) struct ScannerSnapshotIdentity {
pub(super) cycle: u64,
pub(super) leader_epoch: u64,
pub(super) plan_digest: DataUsageScanPlanDigest,
pub(super) tier_registry_generation: Option<u64>,
}
pub(super) struct ScannerSnapshotScope<'a> {
pub(super) sources: &'a HashSet<DataUsageCacheSource>,
pub(super) buckets: &'a [String],
pub(super) identity: ScannerSnapshotIdentity,
}
#[derive(Debug, PartialEq, Eq, thiserror::Error)]
pub(super) enum ScannerSnapshotValidationError {
#[error("scanner snapshot does not cover the expected complete sets")]
IncompleteSets,
#[error("scanner snapshot does not match the requested generation")]
GenerationMismatch,
#[error("scanner snapshot bucket inventory is invalid")]
InvalidInventory,
#[error("scanner snapshot root is incomplete or corrupt")]
InvalidRoot,
}
struct ValidatedScannerSnapshot<'a> {
results: &'a [DataUsageCache],
last_update: SystemTime,
}
impl<'a> ValidatedScannerSnapshot<'a> {
fn validate(
results: &'a [DataUsageCache],
scope: &ScannerSnapshotScope<'_>,
) -> std::result::Result<Self, ScannerSnapshotValidationError> {
if !scanner_results_form_complete_snapshot(results, scope.sources) {
return Err(ScannerSnapshotValidationError::IncompleteSets);
}
let bucket_keys = scope
.buckets
.iter()
.map(|bucket| crate::hash_path(bucket).key())
.collect::<HashSet<_>>();
if bucket_keys.len() != scope.buckets.len()
|| scope
.buckets
.iter()
.any(|bucket| bucket.is_empty() || bucket == DATA_USAGE_ROOT)
{
return Err(ScannerSnapshotValidationError::InvalidInventory);
}
for result in results {
if result.info.next_cycle != scope.identity.cycle
|| result.info.leader_epoch != scope.identity.leader_epoch
|| result.info.scan_plan_digest != Some(scope.identity.plan_digest)
|| result.info.tier_registry_generation != scope.identity.tier_registry_generation
{
return Err(ScannerSnapshotValidationError::GenerationMismatch);
}
if result.info.name != DATA_USAGE_ROOT
|| result.info.cache_key_format != DATA_USAGE_CACHE_KEY_FORMAT
|| !result.has_complete_root_inventory(&bucket_keys)
{
return Err(ScannerSnapshotValidationError::InvalidRoot);
}
}
let last_update = results
.iter()
.filter_map(|result| result.info.last_update)
.max()
.ok_or(ScannerSnapshotValidationError::IncompleteSets)?;
Ok(Self { results, last_update })
}
}
pub(super) fn completed_data_usage_info(
results: &[DataUsageCache],
scope: &ScannerSnapshotScope<'_>,
expected_sources: &HashSet<DataUsageCacheSource>,
all_buckets: &[String],
tier_registry_names: &[String],
bucket_plan_complete: bool,
budget_elapsed: bool,
@@ -304,10 +229,26 @@ pub(super) fn completed_data_usage_info(
if !should_publish_completed_snapshot(completed_set_count, results.len(), budget_elapsed, cancelled) {
return None;
}
let validated = ValidatedScannerSnapshot::validate(results, scope).ok()?;
let results = validated.results;
let all_buckets = scope.buckets;
let registry_generation = scope.identity.tier_registry_generation;
if !scanner_results_form_complete_snapshot(results, expected_sources) {
return None;
}
// A generation is comparable across nodes because it is derived from the
// frozen registry names. Cycle and leader fencing remain separate cache
// metadata. Legacy peers omit the generation; an all-legacy result remains
// readable, but mixing legacy and new (or two new generations) would make
// the per-tier accounting ambiguous.
let registry_generation = results.first()?.info.tier_registry_generation;
if results.iter().any(|result| match registry_generation {
Some(generation) => result.info.tier_registry_generation != Some(generation),
None => result.info.tier_registry_generation.is_some(),
}) {
return None;
}
if results.iter().any(|result| result.root().is_none()) {
return None;
}
let mut total = DataUsageEntry::default();
let mut bucket_entries = HashMap::with_capacity(all_buckets.len());
@@ -332,7 +273,7 @@ pub(super) fn completed_data_usage_info(
return None;
}
let merged_last_update = validated.last_update;
let merged_last_update = results.iter().filter_map(|result| result.info.last_update).max()?;
let buckets_usage = bucket_entries
.iter()
.map(|(bucket, entry)| Some((bucket.clone(), checked_bucket_usage_info(entry)?)))
@@ -359,8 +300,8 @@ pub(super) fn completed_data_usage_info(
usage_snapshot_set_states.sort_by_key(|state| (state.pool_index, state.set_index));
let data_usage_info = DataUsageInfo {
last_update: Some(merged_last_update),
scanner_cycle: Some(scope.identity.cycle),
scanner_epoch: Some(scope.identity.leader_epoch),
scanner_cycle: Some(results.first()?.info.next_cycle),
scanner_epoch: Some(results.first()?.info.leader_epoch),
objects_total_count: u64::try_from(total.objects).ok()?,
versions_total_count: u64::try_from(total.versions).ok()?,
delete_markers_total_count: u64::try_from(total.delete_markers).ok()?,
+11 -8
View File
@@ -39,12 +39,6 @@ pub(super) fn prepare_scoped_set_scan(
else {
return None;
};
// The existing cache does not bind each bucket to a durable incarnation.
// Listing creation times can come from volume metadata, so even Some(time)
// cannot prove that an unselected same-name bucket is the cached bucket.
if all_buckets.iter().any(|bucket| !selected_buckets.contains(&bucket.name)) {
return None;
}
if selected_buckets.is_empty()
|| !old_cache.info.snapshot_complete
|| old_cache.info.last_update.is_none()
@@ -55,7 +49,7 @@ pub(super) fn prepare_scoped_set_scan(
|| old_cache.info.source != Some(generation.source)
|| old_cache.info.scan_plan_digest != Some(baseline_scan_plan_digest)
|| old_cache.info.cache_key_format != DATA_USAGE_CACHE_KEY_FORMAT
|| !old_cache.has_complete_root_inventory(&old_cache.find(DATA_USAGE_ROOT)?.children)
|| old_cache.checked_flatten_complete_scope(DATA_USAGE_ROOT).is_none()
{
return None;
}
@@ -80,12 +74,21 @@ pub(super) fn prepare_scoped_set_scan(
cache: HashMap::new(),
};
cache.replace(DATA_USAGE_ROOT, "", DataUsageEntry::default());
let root_hash = crate::hash_path(DATA_USAGE_ROOT);
let mut current_bucket_names = HashSet::with_capacity(all_buckets.len());
for bucket in all_buckets {
if !current_bucket_names.insert(bucket.name.as_str()) {
return None;
}
cache.replace(&bucket.name, DATA_USAGE_ROOT, DataUsageEntry::default());
if selected_buckets.contains(&bucket.name) {
cache.replace(&bucket.name, DATA_USAGE_ROOT, DataUsageEntry::default());
continue;
}
let bucket_hash = crate::hash_path(&bucket.name);
old_cache.find(&bucket.name)?;
cache.copy_with_children(old_cache, &bucket_hash, &Some(root_hash.clone()));
cache.find(&bucket.name)?;
}
Some(PreparedScopedSetScan {
+2 -11
View File
@@ -260,7 +260,6 @@ where
}
}
bucket_plan_complete &= buckets_by_source.keys().copied().collect::<HashSet<_>>() == *expected_sources;
bucket_plan_complete &= scanner_bucket_inventory_is_complete(&all_buckets, &buckets_by_source);
let scan_plan_digest =
scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_structural_digest(&activity_before));
let dirty_usage_snapshot = Arc::new(snapshot_dirty_usage_buckets(&all_buckets, dirty_generation_before_bucket_list));
@@ -537,16 +536,8 @@ where
let all_bucket_names = all_buckets.iter().map(|bucket| bucket.name.clone()).collect::<Vec<_>>();
let completed_usage = completed_data_usage_info(
&results,
&ScannerSnapshotScope {
sources: &expected_sources,
buckets: &all_bucket_names,
identity: ScannerSnapshotIdentity {
cycle: want_cycle,
leader_epoch,
plan_digest: scan_plan_digest,
tier_registry_generation: Some(tier_registry_generation),
},
},
&expected_sources,
&all_bucket_names,
&tier_registry.names,
bucket_plan_complete,
budget_elapsed,
@@ -18,32 +18,6 @@ use rustfs_data_usage::{ReplicationAllStats, ReplicationTargetUsage, TierAccount
const TEST_PLAN_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([7; 32]);
#[test]
fn scanner_bucket_inventory_requires_exact_unique_set_union() {
let first = BucketInfo {
name: "first".to_string(),
..Default::default()
};
let second = BucketInfo {
name: "second".to_string(),
..Default::default()
};
let source = DataUsageCacheSource::new(0, 0);
let mut sets = HashMap::from([(source, vec![first.clone()])]);
assert!(scanner_bucket_inventory_is_complete(std::slice::from_ref(&first), &sets));
assert!(!scanner_bucket_inventory_is_complete(&[first.clone(), second.clone()], &sets));
assert!(!scanner_bucket_inventory_is_complete(&[], &sets));
assert!(!scanner_bucket_inventory_is_complete(&[first.clone(), first.clone()], &sets));
sets.insert(source, vec![first.clone(), first.clone()]);
assert!(!scanner_bucket_inventory_is_complete(std::slice::from_ref(&first), &sets));
sets.insert(source, vec![second]);
assert!(!scanner_bucket_inventory_is_complete(std::slice::from_ref(&first), &sets));
let mut recreated = first.clone();
recreated.created = Some(OffsetDateTime::UNIX_EPOCH);
sets.insert(source, vec![recreated]);
assert!(!scanner_bucket_inventory_is_complete(&[first], &sets));
}
#[test]
fn should_publish_completed_snapshot_requires_full_clean_cycle() {
assert!(should_publish_completed_snapshot(3, 3, false, false));
@@ -134,133 +108,7 @@ fn completed_data_usage_info_for_test(
cancelled: bool,
) -> Option<(DataUsageInfo, SystemTime)> {
let expected_sources = results.iter().filter_map(|result| result.info.source).collect::<HashSet<_>>();
completed_usage_for_scope(results, &expected_sources, all_buckets, &[], true, budget_elapsed, cancelled)
}
fn completed_usage_for_scope(
results: &[DataUsageCache],
expected_sources: &HashSet<DataUsageCacheSource>,
all_buckets: &[String],
tier_registry_names: &[String],
bucket_plan_complete: bool,
budget_elapsed: bool,
cancelled: bool,
) -> Option<(DataUsageInfo, SystemTime)> {
let first = results.first()?;
completed_data_usage_info(
results,
&ScannerSnapshotScope {
sources: expected_sources,
buckets: all_buckets,
identity: ScannerSnapshotIdentity {
cycle: first.info.next_cycle,
leader_epoch: first.info.leader_epoch,
plan_digest: TEST_PLAN_DIGEST,
tier_registry_generation: first.info.tier_registry_generation,
},
},
tier_registry_names,
bucket_plan_complete,
budget_elapsed,
cancelled,
)
}
#[test]
fn completed_data_usage_info_rejects_duplicate_bucket_inventory() {
let set = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0));
let buckets = vec!["bucket".to_string(), "bucket".to_string()];
assert!(completed_data_usage_info_for_test(&[set], &buckets, false, false).is_none());
}
#[test]
fn completed_data_usage_info_rejects_extra_or_detached_bucket_data() {
let buckets = vec!["bucket".to_string()];
let mut set = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0));
set.replace(
"unlisted",
DATA_USAGE_ROOT,
DataUsageEntry {
objects: 1,
..Default::default()
},
);
assert!(completed_data_usage_info_for_test(&[set.clone()], &buckets, false, false).is_none());
set.cache
.get_mut(DATA_USAGE_ROOT)
.expect("set root")
.children
.remove(&hash_path("unlisted").key());
assert!(
completed_data_usage_info_for_test(&[set], &buckets, false, false).is_none(),
"orphaned data must not disappear from authoritative accounting"
);
}
#[test]
fn completed_data_usage_info_rejects_disconnected_expected_bucket() {
let buckets = vec!["bucket".to_string()];
let mut set = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0));
set.cache.get_mut(DATA_USAGE_ROOT).expect("set root").children.clear();
assert!(completed_data_usage_info_for_test(&[set], &buckets, false, false).is_none());
}
#[test]
fn completed_data_usage_info_rejects_root_scalar_data_and_unknown_key_format() {
let buckets = vec!["bucket".to_string()];
let set = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0));
let mut scalar_root = set.clone();
scalar_root.cache.get_mut(DATA_USAGE_ROOT).expect("set root").size = 10;
assert!(completed_data_usage_info_for_test(&[scalar_root], &buckets, false, false).is_none());
let mut future_format = set;
future_format.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT + 1;
assert!(completed_data_usage_info_for_test(&[future_format], &buckets, false, false).is_none());
}
#[test]
fn completed_data_usage_info_binds_all_results_to_requested_identity() {
let buckets = vec!["bucket".to_string()];
let source = DataUsageCacheSource::new(0, 0);
let sources = HashSet::from([source]);
let set = completed_root_cache("bucket", 2, 10, source);
let identity = ScannerSnapshotIdentity {
cycle: 0,
leader_epoch: 0,
plan_digest: TEST_PLAN_DIGEST,
tier_registry_generation: None,
};
let results = [set];
for expected in [
ScannerSnapshotIdentity { cycle: 1, ..identity },
ScannerSnapshotIdentity {
leader_epoch: 1,
..identity
},
ScannerSnapshotIdentity {
plan_digest: DataUsageScanPlanDigest([9; 32]),
..identity
},
ScannerSnapshotIdentity {
tier_registry_generation: Some(1),
..identity
},
] {
let scope = ScannerSnapshotScope {
sources: &sources,
buckets: &buckets,
identity: expected,
};
assert!(completed_data_usage_info(&results, &scope, &[], true, false, false).is_none());
}
let scope = ScannerSnapshotScope {
sources: &sources,
buckets: &buckets,
identity,
};
let (usage, _) = completed_data_usage_info(&results, &scope, &[], true, false, false)
.expect("the requested complete scope remains publishable");
assert_eq!(usage.objects_total_count, 2);
assert!(usage.is_complete_bucket_usage_snapshot());
completed_data_usage_info(results, &expected_sources, all_buckets, &[], true, budget_elapsed, cancelled)
}
fn lkg_root_cache(bucket: &str, objects: usize, source: DataUsageCacheSource) -> DataUsageCache {
@@ -288,7 +136,7 @@ fn partial_usage_is_observational_not_authoritative_for_quota() {
let expected = HashSet::from([current_source, stalled_source]);
assert!(
completed_usage_for_scope(&[current.clone(), stalled.clone()], &expected, &all_buckets, &[], true, false, false)
completed_data_usage_info(&[current.clone(), stalled.clone()], &expected, &all_buckets, &[], true, false, false)
.is_none()
);
let (observed, _) = observational_data_usage_info(&[current, stalled], &expected, &all_buckets, &[], TEST_PLAN_DIGEST, 8, 3)
@@ -661,7 +509,7 @@ fn completed_data_usage_info_accepts_unknown_only_with_current_registry_generati
let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0)]);
assert!(
completed_usage_for_scope(&[set], &expected_sources, &all_buckets, &["WARM".to_string()], true, false, false,).is_some()
completed_data_usage_info(&[set], &expected_sources, &all_buckets, &["WARM".to_string()], true, false, false,).is_some()
);
}
@@ -699,7 +547,7 @@ fn completed_data_usage_info_rejects_non_registry_tier_in_current_generation() {
let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0)]);
assert!(
completed_usage_for_scope(&[set], &expected_sources, &all_buckets, &["WARM".to_string()], true, false, false,).is_none()
completed_data_usage_info(&[set], &expected_sources, &all_buckets, &["WARM".to_string()], true, false, false,).is_none()
);
}
@@ -850,7 +698,6 @@ fn completed_data_usage_info_publishes_confirmed_empty_namespace() {
source: Some(DataUsageCacheSource::new(0, 0)),
snapshot_complete: true,
scan_plan_digest: Some(TEST_PLAN_DIGEST),
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
..Default::default()
},
..Default::default()
@@ -1008,7 +855,7 @@ fn completed_data_usage_info_requires_exact_topology_sources() {
let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0), DataUsageCacheSource::new(1, 0)]);
assert!(
completed_usage_for_scope(&[first_set, unexpected_set], &expected_sources, &all_buckets, &[], true, false, false)
completed_data_usage_info(&[first_set, unexpected_set], &expected_sources, &all_buckets, &[], true, false, false)
.is_none()
);
}
@@ -1019,7 +866,7 @@ fn completed_data_usage_info_rejects_incomplete_bucket_plan() {
let set = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0));
let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0)]);
assert!(completed_usage_for_scope(&[set], &expected_sources, &all_buckets, &[], false, false, false).is_none());
assert!(completed_data_usage_info(&[set], &expected_sources, &all_buckets, &[], false, false, false).is_none());
}
#[test]
+8 -63
View File
@@ -797,13 +797,6 @@ fn complete_set_usage_cache(buckets: &[(&str, usize)], scan_plan_digest: DataUsa
cache
}
fn bucket_info_with_created_time(name: &str) -> BucketInfo {
BucketInfo {
created: Some(time::OffsetDateTime::UNIX_EPOCH),
..bucket_info(name)
}
}
fn complete_usage_baseline(
source: DataUsageCacheSource,
scan_plan_digest: DataUsageScanPlanDigest,
@@ -994,7 +987,7 @@ fn verified_remote_dirty_usage_buckets_rejects_incomplete_or_stale_peer_state()
}
#[test]
fn scoped_set_scan_rebuilds_selected_buckets_and_drops_deleted_buckets() {
fn scoped_set_scan_preserves_unselected_usage_and_drops_deleted_buckets() {
let baseline_digest = DataUsageScanPlanDigest([1; 32]);
let current_digest = DataUsageScanPlanDigest([2; 32]);
let mut old_cache = complete_set_usage_cache(&[("stable", 10), ("dirty", 20), ("deleted", 30)], baseline_digest);
@@ -1007,11 +1000,8 @@ fn scoped_set_scan_rebuilds_selected_buckets_and_drops_deleted_buckets() {
..Default::default()
},
);
let all_buckets = vec![
bucket_info_with_created_time("stable"),
bucket_info_with_created_time("dirty"),
];
let selected_buckets = Arc::new(HashSet::from(["stable".to_string(), "dirty".to_string(), "deleted".to_string()]));
let all_buckets = vec![bucket_info("stable"), bucket_info("dirty")];
let selected_buckets = Arc::new(HashSet::from(["dirty".to_string(), "deleted".to_string()]));
let prepared = prepare_scoped_set_scan(
&old_cache,
@@ -1031,16 +1021,12 @@ fn scoped_set_scan_rebuilds_selected_buckets_and_drops_deleted_buckets() {
)
.expect("complete matching set cache should support a scoped scan");
assert_eq!(
prepared.buckets.iter().map(|bucket| bucket.name.as_str()).collect::<Vec<_>>(),
["stable", "dirty"]
);
assert_eq!(prepared.buckets.iter().map(|bucket| bucket.name.as_str()).collect::<Vec<_>>(), ["dirty"]);
let stable = prepared
.cache
.checked_flatten("stable")
.expect("selected bucket placeholder should exist");
assert_eq!((stable.size, stable.objects), (0, 0));
assert!(prepared.cache.find("stable/prefix").is_none());
.expect("unselected bucket subtree should be retained");
assert_eq!((stable.size, stable.objects), (15, 2));
assert_eq!(prepared.cache.find("dirty").map(|entry| (entry.size, entry.objects)), Some((0, 0)));
assert!(prepared.cache.find("deleted").is_none());
assert_eq!(prepared.cache.info.scan_plan_digest, Some(current_digest));
@@ -1051,41 +1037,11 @@ fn scoped_set_scan_rebuilds_selected_buckets_and_drops_deleted_buckets() {
assert_eq!(prepared.cache.info.lkg_scan_plan_digest, Some(baseline_digest));
}
#[test]
fn scoped_set_scan_rejects_unbound_bucket_incarnations() {
let baseline_digest = DataUsageScanPlanDigest([1; 32]);
let old_cache = complete_set_usage_cache(&[("stable", 10), ("dirty", 20)], baseline_digest);
let scope = ScannerBucketScanScope {
selected_buckets: Some(Arc::new(HashSet::from(["dirty".to_string()]))),
baseline_scan_plan_digest: Some(baseline_digest),
};
let generation = ScannerSetCacheGeneration {
want_cycle: 8,
leader_epoch: 11,
tier_registry_generation: 13,
source: DataUsageCacheSource::new(1, 2),
scan_plan_digest: DataUsageScanPlanDigest([2; 32]),
};
for created in [
None,
Some(OffsetDateTime::UNIX_EPOCH),
Some(OffsetDateTime::UNIX_EPOCH + time::Duration::days(1)),
] {
let mut stable = bucket_info("stable");
stable.created = created;
let buckets = vec![stable, bucket_info_with_created_time("dirty")];
assert!(
prepare_scoped_set_scan(&old_cache, &buckets, &buckets, &scope, generation).is_none(),
"missing identity, volume timestamps and same-name recreation must all rebuild"
);
}
}
#[test]
fn scoped_set_scan_falls_back_when_an_unselected_bucket_has_no_baseline() {
let baseline_digest = DataUsageScanPlanDigest([3; 32]);
let old_cache = complete_set_usage_cache(&[("stable", 10)], baseline_digest);
let all_buckets = vec![bucket_info_with_created_time("stable"), bucket_info_with_created_time("new")];
let all_buckets = vec![bucket_info("stable"), bucket_info("new")];
assert!(
prepare_scoped_set_scan(
@@ -1111,7 +1067,7 @@ fn scoped_set_scan_falls_back_when_an_unselected_bucket_has_no_baseline() {
#[test]
fn scoped_set_scan_requires_an_exact_complete_baseline() {
let baseline_digest = DataUsageScanPlanDigest([5; 32]);
let all_buckets = vec![bucket_info_with_created_time("dirty")];
let all_buckets = vec![bucket_info("dirty")];
let scope = ScannerBucketScanScope {
selected_buckets: Some(Arc::new(HashSet::from(["dirty".to_string()]))),
baseline_scan_plan_digest: Some(baseline_digest),
@@ -1132,10 +1088,6 @@ fn scoped_set_scan_requires_an_exact_complete_baseline() {
not_durable.info.last_update = None;
assert!(prepare_scoped_set_scan(&not_durable, &all_buckets, &all_buckets, &scope, generation).is_none());
let mut unscoped_usage = complete_set_usage_cache(&[("dirty", 10)], baseline_digest);
unscoped_usage.cache.get_mut(DATA_USAGE_ROOT).expect("set root").objects = 1;
assert!(prepare_scoped_set_scan(&unscoped_usage, &all_buckets, &all_buckets, &scope, generation).is_none());
let mut wrong_digest = complete_set_usage_cache(&[("dirty", 10)], baseline_digest);
wrong_digest.info.scan_plan_digest = Some(DataUsageScanPlanDigest([7; 32]));
assert!(prepare_scoped_set_scan(&wrong_digest, &all_buckets, &all_buckets, &scope, generation).is_none());
@@ -1146,13 +1098,6 @@ fn scoped_set_scan_requires_an_exact_complete_baseline() {
};
let complete = complete_set_usage_cache(&[("dirty", 10)], baseline_digest);
assert!(prepare_scoped_set_scan(&complete, &all_buckets, &all_buckets, &empty_scope, generation).is_none());
assert!(prepare_scoped_set_scan(&complete, &all_buckets, &all_buckets, &scope, generation).is_some());
let unidentified_buckets = vec![bucket_info("dirty")];
assert!(
prepare_scoped_set_scan(&complete, &unidentified_buckets, &unidentified_buckets, &scope, generation).is_some(),
"fully selected buckets are rebuilt without reusing an unproven incarnation"
);
let mut future_cache = complete_set_usage_cache(&[("dirty", 10)], baseline_digest);
future_cache.info.next_cycle = generation.want_cycle.saturating_add(1);
+201 -6
View File
@@ -5,7 +5,9 @@ from __future__ import annotations
import hashlib
import json
import os
import re
import subprocess
import sys
import tempfile
import tomllib
@@ -481,18 +483,20 @@ def yaml_block(lines: list[str], key: str, indent: int) -> list[str] | None:
return lines[start:end]
def workflow_step_block(job_lines: list[str], action: str) -> tuple[int, list[str]] | None:
def workflow_step_block(
job_lines: list[str], value: str, key: str = "uses", indent: int = 6
) -> tuple[int, list[str]] | None:
uses_index = next(
(
index
for index, line in enumerate(job_lines)
if (
line.split("#", 1)[0].strip() == f"- uses: {action}"
and len(line) - len(line.lstrip()) == 6
line.split("#", 1)[0].strip() == f"- {key}: {value}"
and len(line) - len(line.lstrip()) == indent
)
or (
line.split("#", 1)[0].strip() == f"uses: {action}"
and len(line) - len(line.lstrip()) == 8
line.split("#", 1)[0].strip() == f"{key}: {value}"
and len(line) - len(line.lstrip()) == indent + 2
)
),
None,
@@ -520,6 +524,67 @@ def workflow_step_block(job_lines: list[str], action: str) -> tuple[int, list[st
return start, job_lines[start:end]
def yaml_scalar_continues(lines: list[str], index: int, indent: int) -> bool:
following = next(
(line for line in lines[index + 1:] if line.strip() and not line.lstrip().startswith("#")), None
)
return following is not None and len(following) - len(following.lstrip()) > indent
def check_quick_checks(root: Path) -> list[str]:
errors: list[str] = []
bypass_key = r'''(?:if|continue-on-error|needs|"if"|"continue-on-error"|"needs"|'if'|'continue-on-error'|'needs')\s*:'''
for name in ("ci.yml", "ci-docs-only.yml"):
relative = f".github/workflows/{name}"
path = root / relative
job = yaml_block(path.read_text().splitlines(), "quick-checks", 2) if path.is_file() else None
if job is None:
errors.append(f"{relative}: missing Quick Checks job")
continue
conditions = [index for index, line in enumerate(job) if re.match(rf"^ {bypass_key}", line)]
expected = ["if: github.event_name != 'pull_request' || github.event.action != 'closed'"] if name == "ci.yml" else []
if [job[index].strip() for index in conditions] != expected or any(
yaml_scalar_continues(job, index, 4) for index in conditions
):
errors.append(f"{relative}: Quick Checks job must not add dependencies, bypass failures, or change its event condition")
checkout = workflow_step_block(job, "actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0")
action = workflow_step_block(job, "./.github/actions/quick-checks")
if checkout is None or action is None:
errors.append(f"{relative}: Quick Checks requires checkout and the shared quick-checks action")
continue
if checkout[0] >= action[0]:
errors.append(f"{relative}: checkout must run before shared Quick Checks")
if " persist-credentials: false" not in checkout[1]:
errors.append(f"{relative}: Quick Checks checkout must disable persisted credentials")
for step in (checkout, action):
if any(re.match(rf"^\s+(?:- )?{bypass_key}", line) for line in step[1]):
errors.append(f"{relative}: Quick Checks checkout and shared action must run without bypasses")
relative = ".github/actions/quick-checks/action.yml"
path = root / relative
runs = yaml_block(path.read_text().splitlines(), "runs", 0) if path.is_file() else None
if runs is None or " using: composite" not in runs:
errors.append(f"{relative}: missing composite action")
return errors
steps = yaml_block(runs, "steps", 2) or []
for command in ("shellcheck --version && actionlint", "./scripts/check_error_other_format_ratchet.sh"):
step = workflow_step_block(steps, command, key="run", indent=4)
if step is None:
errors.append(f"{relative}: missing direct execution of {command}")
continue
if " shell: bash" not in step[1] or any(
re.match(rf"^\s+(?:- )?{bypass_key}", line) for line in step[1]
):
errors.append(f"{relative}: {command} must use bash without a condition or continue-on-error")
run_index = next(
index for index, line in enumerate(step[1])
if line.split("#", 1)[0].rstrip() in (f" run: {command}", f" - run: {command}")
)
if yaml_scalar_continues(step[1], run_index, 6):
errors.append(f"{relative}: {command} must remain a single-line run scalar")
return errors
def alert_step_errors(
job_lines: list[str],
expected_action_if: str | None,
@@ -820,10 +885,139 @@ def validate(root: Path) -> list[str]:
errors.extend(check_workflow_readiness(root))
errors.extend(check_profile_definitions(root))
errors.extend(check_scheduled_alerts(root))
errors.extend(check_quick_checks(root))
return errors
class SelfTests(unittest.TestCase):
def test_quick_checks_rejects_caller_and_execution_bypasses(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
caller = (
"jobs:\n quick-checks:\n steps:\n"
" - uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0\n"
" with:\n persist-credentials: false\n"
" - uses: ./.github/actions/quick-checks\n"
)
action = (
"runs:\n using: composite\n steps:\n"
" - uses: taiki-e/install-action@pinned\n"
" with:\n tool: actionlint@1.7.12\n"
" - name: Lint workflows\n shell: bash\n run: shellcheck --version && actionlint\n"
" - name: Error format ratchet\n shell: bash\n"
" run: ./scripts/check_error_other_format_ratchet.sh\n"
)
sources = {
".github/workflows/ci.yml": caller.replace(
" steps:", " if: github.event_name != 'pull_request' || github.event.action != 'closed'\n steps:"
),
".github/workflows/ci-docs-only.yml": caller,
".github/actions/quick-checks/action.yml": action,
}
for relative, source in sources.items():
path = root / relative
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(source)
self.assertEqual(check_quick_checks(root), [])
for relative in (".github/workflows/ci.yml", ".github/workflows/ci-docs-only.yml"):
source = sources[relative]
mutations = {
"different action": source.replace("./.github/actions/quick-checks", "./.github/actions/other"),
"conditional call": source + " if: false\n",
"ignored call failure": source + " continue-on-error: true\n",
"conditional checkout": source.replace(" with:", " if: false\n with:"),
"ignored job failure": source.replace(" steps:", " continue-on-error: true\n steps:"),
"changed job condition": (
source.replace("github.event_name != 'pull_request' || github.event.action != 'closed'", "false")
if relative.endswith("/ci.yml") else source.replace(" steps:", " if: false\n steps:")
),
"persisted credentials": source.replace("persist-credentials: false", "persist-credentials: true"),
"late checkout": source.replace(" - uses: ./.github/actions/quick-checks\n", "").replace(
" steps:\n", " steps:\n - uses: ./.github/actions/quick-checks\n"
),
"missing job": source.replace(" quick-checks:", " other-checks:"),
}
for key in ("'if' : false", '"if": false', "'continue-on-error': true", '"continue-on-error" : true'):
mutations[f"quoted call {key}"] = source + f" {key}\n"
mutations[f"quoted checkout {key}"] = source.replace(" with:", f" {key}\n with:")
job_source = source.replace(
" if: github.event_name != 'pull_request' || github.event.action != 'closed'\n", ""
) if "if" in key else source
mutations[f"quoted job {key}"] = job_source.replace(" steps:", f" {key}\n steps:")
for dependency in ("needs: prerequisite", "needs: [prerequisite]", "needs:\n - prerequisite", "'needs' : [prerequisite]", '"needs": [prerequisite]'):
for condition in ("false", "true"):
prerequisite = f"\n prerequisite:\n if: {condition}\n runs-on: ubuntu-latest\n steps:\n - run: exit 1\n"
mutations[f"job dependency {dependency} if {condition}"] = source.replace(" steps:", f" {dependency}\n steps:") + prerequisite
if relative.endswith("/ci.yml"):
for separator in ("", "\n", " # continued condition\n"):
mutations[f"continued job condition {separator!r}"] = source.replace(
" steps:", f"{separator} && false\n steps:"
)
for case, mutated in mutations.items():
with self.subTest(path=relative, case=case):
(root / relative).write_text(mutated)
self.assertTrue(check_quick_checks(root))
(root / relative).write_text(source)
relative = ".github/actions/quick-checks/action.yml"
mutations = {
"not composite": action.replace("using: composite", "using: node24"),
"only installed actionlint": action.replace("run: shellcheck --version && actionlint", "run: echo actionlint"),
"missing shellcheck preflight": action.replace("shellcheck --version && ", ""),
"missing ratchet": action.replace("run: ./scripts/check_error_other_format_ratchet.sh", "run: echo skipped"),
"swallowed lint failure": action.replace("&& actionlint", "&& actionlint || true"),
"swallowed ratchet failure": action.replace("ratchet.sh", "ratchet.sh || true"),
"conditional lint": action.replace("run: shellcheck", "if: false\n run: shellcheck"),
"ignored ratchet failure": action.replace("run: ./scripts/", "continue-on-error: true\n run: ./scripts/"),
"non-failing shell": action.replace("shell: bash", "shell: bash {0}"),
"run text in step name": action.replace(
"name: Lint workflows", "name: |\n run: shellcheck --version && actionlint"
).replace("\n run: shellcheck --version && actionlint\n", "\n run: shellcheck --version && actionlint\n || true\n"),
}
for command in ("shellcheck --version && actionlint", "./scripts/check_error_other_format_ratchet.sh"):
for key in ("'if' : false", '"if": false', "'continue-on-error': true", '"continue-on-error" : true'):
mutations[f"quoted {command} {key}"] = action.replace(f"run: {command}", f"{key}\n run: {command}")
for separator in ("", "\n", " # continued command\n"):
mutations[f"continued {command} {separator!r}"] = action.replace(
f"run: {command}\n", f"run: {command}\n{separator} || true\n"
)
for case, mutated in mutations.items():
with self.subTest(case=case):
(root / relative).write_text(mutated)
self.assertTrue(check_quick_checks(root))
(root / relative).unlink()
self.assertTrue(check_quick_checks(root))
def test_quick_checks_commands_propagate_failure(self) -> None:
runs = yaml_block((ROOT / ".github/actions/quick-checks/action.yml").read_text().splitlines(), "runs", 0)
steps = yaml_block(runs or [], "steps", 2) or []
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
(root / "scripts").mkdir()
commands = ("shellcheck", "actionlint", "./scripts/check_error_other_format_ratchet.sh")
for failing in commands:
with self.subTest(command=failing):
run = "shellcheck --version && actionlint" if failing != commands[-1] else failing
step = workflow_step_block(steps, run, key="run", indent=4)
self.assertIsNotNone(step)
run_index = next(index for index, line in enumerate(step[1]) if line.startswith(" run:"))
self.assertFalse(yaml_scalar_continues(step[1], run_index, 6))
body = step[1][run_index].removeprefix(" run: ")
for command in commands:
shim = root / command
shim.write_text(f"#!/bin/sh\nexit {17 if command == failing else 0}\n")
shim.chmod(0o755)
result = subprocess.run(
["bash", "--noprofile", "--norc", "-e", "-o", "pipefail", "-c", body],
cwd=root, env=dict(os.environ, PATH=f"{root}{os.pathsep}{os.environ['PATH']}"),
capture_output=True, text=True,
)
self.assertEqual(result.returncode, 17, result.stderr)
def test_validate_includes_quick_checks(self) -> None:
error = "Quick Checks wiring regression"
with mock.patch(__name__ + ".check_quick_checks", return_value=[error]):
self.assertIn(error, validate(ROOT))
def test_core_gate_rejects_missing_ignored_filtered_and_corrupt_inputs(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
@@ -1058,6 +1252,7 @@ class SelfTests(unittest.TestCase):
mock.patch(__name__ + ".check_profile_definitions", return_value=[]),
mock.patch(__name__ + ".check_ilm_build_budget", return_value=[]),
mock.patch(__name__ + ".check_scheduled_alerts", return_value=[]),
mock.patch(__name__ + ".check_quick_checks", return_value=[]),
):
self.assertEqual(len(validate(root)), 1)
@@ -1498,7 +1693,7 @@ def main() -> int:
for error in errors:
print(f"ERROR: {error}", file=sys.stderr)
return 1
print("OK: e2e modules, runner selection, fuzz matrices, profiles, and scheduled alerts are wired")
print("OK: e2e modules, runner selection, fuzz matrices, profiles, scheduled alerts, and Quick Checks are wired")
return 0
+253
View File
@@ -0,0 +1,253 @@
#!/usr/bin/env python3
"""Build an identified E2E server and verify it around one test invocation."""
import argparse
from contextlib import contextmanager
import hashlib
import json
import os
from pathlib import Path
import stat
import signal
import subprocess
import sys
import tempfile
ROOT = Path(__file__).resolve().parent.parent
RECEIPT_ENV = "RUSTFS_E2E_BINARY_RECEIPT"
def feature_set(value):
return sorted(set(part.strip() for part in value.split(",") if part.strip()))
def file_hash(path):
digest = hashlib.sha256()
with path.open("rb") as source:
for chunk in iter(lambda: source.read(1024 * 1024), b""):
digest.update(chunk)
return digest.hexdigest()
def source_identity():
head = subprocess.check_output(["git", "rev-parse", "HEAD"], cwd=ROOT, text=True).strip()
tracked = subprocess.check_output(["git", "ls-files", "--cached", "--others", "--exclude-standard", "-z"], cwd=ROOT)
paths = set(tracked.decode("utf-8").rstrip("\0").split("\0")) - {""}
# RustEmbed consumes ignored console assets as well as tracked Rust sources.
static_dir = ROOT / "rustfs/static"
if static_dir.is_symlink():
raise ValueError("The embedded static directory must not be a symlink")
if static_dir.is_dir():
for path in static_dir.rglob("*"):
if path.is_symlink() and path.is_dir():
raise ValueError(f"Unsupported embedded directory symlink: {path}")
if not path.is_dir():
paths.add(str(path.relative_to(ROOT)))
elif static_dir.exists():
paths.add("rustfs/static")
digest = hashlib.sha256()
digest.update(b"static-present\0" if static_dir.is_dir() else b"static-absent\0")
for name in sorted(paths):
path = ROOT / name
digest.update(name.encode("utf-8") + b"\0")
try:
metadata = path.lstat()
except FileNotFoundError:
digest.update(b"deleted\0")
continue
if stat.S_ISLNK(metadata.st_mode):
digest.update(b"symlink\0" + os.fsencode(os.readlink(path)) + b"\0")
if path.is_dir():
target = path.resolve()
if ROOT not in target.parents:
raise ValueError(f"Directory link escapes the source inventory: {name}")
# Directory aliases such as .claude/skills share already-hashed inputs.
for child in target.rglob("*"):
if child.is_dir() and not child.is_symlink():
continue
if child.is_dir() or str(child.relative_to(ROOT)) not in paths:
raise ValueError(f"Directory link contains an unrecorded input: {child}")
digest.update(b"directory\0" + str(target.relative_to(ROOT)).encode("utf-8") + b"\0")
continue
elif not stat.S_ISREG(metadata.st_mode):
raise ValueError(f"Unsupported build input: {name}")
digest.update(str(metadata.st_mode & 0o111).encode() + b"\0")
digest.update(file_hash(path).encode() + b"\0")
return {"head": head, "sha256": digest.hexdigest()}
def sidecar_path(binary):
return binary.with_name(binary.name + ".e2e.json")
def validate_target_directory(target_dir):
if target_dir == ROOT or target_dir in ROOT.parents:
raise ValueError("CARGO_TARGET_DIR must not contain the source workspace")
if ROOT in target_dir.parents:
ignored = subprocess.run(["git", "check-ignore", "--quiet", "--no-index", str(target_dir.relative_to(ROOT))], cwd=ROOT)
if ignored.returncode != 0:
raise ValueError("An in-workspace CARGO_TARGET_DIR must be Git-ignored; use target/ or an external directory")
@contextmanager
def exclusive_binary(binary):
marker = binary.with_name(binary.name + ".e2e.lock")
try:
descriptor = os.open(marker, os.O_CREAT | os.O_EXCL | os.O_WRONLY, 0o600)
except FileExistsError as error:
raise ValueError(f"Another E2E build/run owns {marker}; do not share a target directory between concurrent runs") from error
try:
identity = os.fstat(descriptor)
with os.fdopen(descriptor, "w") as lock:
lock.write(f"pid={os.getpid()}\n")
yield
finally:
current = marker.stat()
if (current.st_dev, current.st_ino) != (identity.st_dev, identity.st_ino):
raise ValueError("The E2E ownership marker changed during the command")
marker.unlink()
def terminate_command(process):
if process.poll() is not None:
return
try:
os.killpg(process.pid, signal.SIGTERM)
except ProcessLookupError:
return
try:
process.wait(timeout=5)
except subprocess.TimeoutExpired:
os.killpg(process.pid, signal.SIGKILL)
process.wait()
def build(binary, target_dir, profile, requested, all_bins):
sidecar = sidecar_path(binary)
sidecar.unlink(missing_ok=True)
before = source_identity()
command = ["cargo", "build", "--locked", "-p", "rustfs", "--target-dir", str(target_dir), "--message-format=json-render-diagnostics"]
command.extend(["--bins"] if all_bins else ["--bin", "rustfs"])
if requested:
command.extend(["--features", ",".join(requested)])
if profile == "release":
command.append("--release")
artifact = None
with subprocess.Popen(command, cwd=ROOT, stdout=subprocess.PIPE, text=True, start_new_session=True) as process:
try:
for line in process.stdout:
message = json.loads(line)
if message.get("reason") == "compiler-message":
print(message["message"].get("rendered", ""), end="", file=sys.stderr)
if message.get("reason") == "compiler-artifact" and message.get("target", {}).get("name") == "rustfs" and "bin" in message.get("target", {}).get("kind", []):
artifact = message
if process.wait() != 0:
raise ValueError("RustFS build failed; no E2E identity was recorded")
except BaseException:
terminate_command(process)
raise
if not artifact or Path(artifact.get("executable", "")).resolve() != binary:
raise ValueError("Cargo did not produce the requested RustFS executable")
if source_identity() != before:
raise ValueError("Build inputs changed during compilation; finish preparing embedded assets and rebuild in an isolated worktree")
record = {
"schema": 1,
"source": before,
"requested_features": requested,
"features": sorted(artifact["features"]),
"profile": profile,
"rustc": subprocess.check_output(["rustc", "-Vv"], text=True),
"binary_sha256": file_hash(binary),
}
sidecar.write_text(json.dumps(record, sort_keys=True) + "\n")
print(f"Built E2E server: {binary}\nIdentity: {sidecar}", file=sys.stderr)
def verify(binary, profile, requested):
record = json.loads(sidecar_path(binary).read_text())
if not isinstance(record, dict) or set(record) != {"schema", "source", "requested_features", "features", "profile", "rustc", "binary_sha256"} or type(record["schema"]) is not int or record["schema"] != 1:
raise ValueError("Missing or unsupported E2E binary identity; run the build command")
if not isinstance(record["rustc"], str) or not record["rustc"].strip():
raise ValueError("Missing E2E build toolchain identity")
if record["requested_features"] != requested or record["profile"] != profile:
raise ValueError("E2E binary build features/profile differ from this test invocation")
if not isinstance(record["features"], list) or not all(isinstance(item, str) for item in record["features"]) or not set(requested) <= set(record["features"]):
raise ValueError("Invalid resolved E2E binary features")
if record["source"] != source_identity():
raise ValueError("E2E binary was built from different inputs; rebuild before testing")
if record["binary_sha256"] != file_hash(binary):
raise ValueError("E2E binary content differs from its build identity")
return record
def run(binary, profile, requested, command):
if not command:
raise ValueError("run requires a test command after --")
override = os.environ.get("CARGO_BIN_EXE_rustfs")
if override and Path(override).resolve() != binary:
raise ValueError("CARGO_BIN_EXE_rustfs selects a different server; use --binary explicitly")
record = verify(binary, profile, requested)
metadata = binary.stat()
with tempfile.TemporaryDirectory(prefix="rustfs-e2e-receipt-") as directory:
receipt = Path(directory) / "receipt.json"
receipt.write_text(json.dumps({
"schema": 1,
"workspace": str(ROOT),
"binary": str(binary),
"size": metadata.st_size,
"modified_ns": metadata.st_mtime_ns,
"features": record["features"],
}))
env = dict(os.environ, CARGO_BIN_EXE_rustfs=str(binary), RUSTFS_BUILD_FEATURES=",".join(record["features"]))
env[RECEIPT_ENV] = str(receipt)
with subprocess.Popen(command, cwd=ROOT, env=env, start_new_session=True) as process:
try:
status = process.wait()
except (KeyboardInterrupt, SystemExit):
terminate_command(process)
raise
try:
if verify(binary, profile, requested) != record:
raise ValueError("E2E build identity changed during testing")
except (OSError, ValueError, subprocess.SubprocessError) as error:
print(f"E2E validation invalidated: {error}", file=sys.stderr)
return status if status else 1
return status
def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("mode", choices=("build", "run"))
parser.add_argument("--features", default="", help="additional Cargo features; defaults remain enabled")
parser.add_argument("--profile", choices=("debug", "release"), default="debug")
parser.add_argument("--binary", type=Path, help="prebuilt server path for run")
parser.add_argument("--bins", action="store_true", help="build all RustFS binary targets, preserving the CI build matrix")
# Parse the child command separately so its options are never interpreted here.
args = sys.argv[1:]
separator = args.index("--") if "--" in args else len(args)
command = args[separator + 1:] if separator < len(args) else []
options = parser.parse_args(args[:separator])
target_dir = Path(os.environ.get("CARGO_TARGET_DIR", ROOT / "target")).resolve()
binary = (options.binary or target_dir / options.profile / ("rustfs.exe" if os.name == "nt" else "rustfs")).resolve()
try:
validate_target_directory(target_dir)
requested = feature_set(options.features)
if options.mode == "build":
binary.parent.mkdir(parents=True, exist_ok=True)
with exclusive_binary(binary):
if options.mode == "build":
if options.binary or command:
raise ValueError("build does not accept --binary or a child command")
build(binary, target_dir, options.profile, requested, options.bins)
return 0
if options.bins:
raise ValueError("--bins is a build option")
return run(binary, options.profile, requested, command)
except (OSError, ValueError, subprocess.SubprocessError) as error:
print(f"E2E prerequisite failed: {error}", file=sys.stderr)
return 1
if __name__ == "__main__":
signal.signal(signal.SIGTERM, lambda signum, frame: sys.exit(128 + signum))
raise SystemExit(main())
+12 -3
View File
@@ -14,7 +14,12 @@ NC='\033[0m' # No Color
# Default values
PROJECT_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)"
TARGET_DIR="$PROJECT_ROOT/target/debug"
CARGO_TARGET_DIR="${CARGO_TARGET_DIR:-$PROJECT_ROOT/target}"
if [[ "$CARGO_TARGET_DIR" != /* ]]; then
CARGO_TARGET_DIR="$PROJECT_ROOT/$CARGO_TARGET_DIR"
fi
export CARGO_TARGET_DIR
TARGET_DIR="$CARGO_TARGET_DIR/debug"
RUSTFS_BINARY="$TARGET_DIR/rustfs"
DATA_DIR="$TARGET_DIR/rustfs_test_data"
RUSTFS_PID=""
@@ -94,7 +99,7 @@ build_rustfs() {
print_info "Building RustFS..."
cd "$PROJECT_ROOT"
if ! cargo build --bin rustfs --features "$RUSTFS_BUILD_FEATURES"; then
if ! python3 scripts/e2e_binary.py build --features "$RUSTFS_BUILD_FEATURES"; then
print_error "Failed to build RustFS"
exit 1
fi
@@ -115,6 +120,10 @@ check_dependencies() {
missing_tools+=("curl")
fi
if ! command -v python3 >/dev/null 2>&1; then
missing_tools+=("python3")
fi
if ! command -v cargo >/dev/null 2>&1; then
missing_tools+=("cargo")
fi
@@ -203,7 +212,7 @@ run_tests() {
print_info "Test command: ${test_cmd[*]}"
if "${test_cmd[@]}"; then
if python3 scripts/e2e_binary.py run --features "$RUSTFS_BUILD_FEATURES" -- "${test_cmd[@]}"; then
print_success "All tests passed!"
return 0
else
+13 -12
View File
@@ -243,9 +243,10 @@ run_quick_e2e_steps() {
return
fi
run_step "e2e-reliability-disk-fault" cargo test --package e2e_test reliability_disk_fault_test -- --nocapture
run_step "e2e-heal-erasure-disk-rebuild" cargo test --package e2e_test heal_erasure_disk_rebuild_test -- --nocapture
run_step "e2e-namespace-lock-quorum" cargo test --package e2e_test namespace_lock_quorum_test -- --nocapture
run_step "build-e2e-server" python3 scripts/e2e_binary.py build
run_step "e2e-reliability-disk-fault" python3 scripts/e2e_binary.py run -- cargo test --package e2e_test reliability_disk_fault_test -- --nocapture
run_step "e2e-heal-erasure-disk-rebuild" python3 scripts/e2e_binary.py run -- cargo test --package e2e_test heal_erasure_disk_rebuild_test -- --nocapture
run_step "e2e-namespace-lock-quorum" python3 scripts/e2e_binary.py run -- cargo test --package e2e_test namespace_lock_quorum_test -- --nocapture
}
run_quick_profile() {
@@ -313,15 +314,15 @@ write_blackbox_matrix() {
{
printf 'profile\tscenario\tgate\tcommand\tfixture_env\tstatus\n'
printf 'quick\tsingle-node disk fault read/write\tblack-box\tcargo test --package e2e_test reliability_disk_fault_test -- --nocapture\tnone\t%s\n' "$e2e_status"
printf 'quick\theal degraded erasure disk rebuild\tblack-box\tcargo test --package e2e_test heal_erasure_disk_rebuild_test -- --nocapture\tnone\t%s\n' "$e2e_status"
printf 'quick\tnamespace lock quorum under EC ops\tblack-box\tcargo test --package e2e_test namespace_lock_quorum_test -- --nocapture\tnone\t%s\n' "$e2e_status"
printf 'quick\tsingle-node disk fault read/write\tblack-box\tpython3 scripts/e2e_binary.py run -- cargo test --package e2e_test reliability_disk_fault_test -- --nocapture\tnone\t%s\n' "$e2e_status"
printf 'quick\theal degraded erasure disk rebuild\tblack-box\tpython3 scripts/e2e_binary.py run -- cargo test --package e2e_test heal_erasure_disk_rebuild_test -- --nocapture\tnone\t%s\n' "$e2e_status"
printf 'quick\tnamespace lock quorum under EC ops\tblack-box\tpython3 scripts/e2e_binary.py run -- cargo test --package e2e_test namespace_lock_quorum_test -- --nocapture\tnone\t%s\n' "$e2e_status"
printf 'full\tlegacy bitrot read fixture restore\tfixture\tcargo test -p rustfs-ecstore --test legacy_bitrot_read_test -- --nocapture\tRUSTFS_LEGACY_TEST_ROOT,RUSTFS_LEGACY_TEST_DISK\t%s\n' "$legacy_status"
printf 'full\tMinIO generated encrypted read and negative restore fixture\tfixture\tcargo test -p rustfs --features rio-v2 storage::minio_generated_read_test --lib -- --ignored --nocapture\tRUSTFS_MINIO_FIXTURE_ROOT,RUSTFS_MINIO_STATIC_KMS_KEY_B64\t%s\n' "$minio_status"
printf 'full\tS3 multipart range versioning delete subset\tblack-box\tenv TESTEXPR=\"multipart or range or versioning or delete\" DEPLOY_MODE=build MAXFAIL=0 ./scripts/s3-tests/run.sh\tnone\t%s\n' "$s3_status"
printf 'destructive\tdistributed cluster concurrency\tblack-box\tcargo test --package e2e_test cluster_concurrency_test -- --nocapture\tnone\t%s\n' "$destructive_status"
printf 'destructive\tstale multipart cleanup cluster\tblack-box\tcargo test --package e2e_test stale_multipart_cleanup_cluster_test -- --nocapture\tnone\t%s\n' "$destructive_status"
printf 'destructive\tdelete marker migration semantics\tblack-box\tcargo test --package e2e_test delete_marker_migration_semantics_test -- --nocapture\tnone\t%s\n' "$destructive_status"
printf 'destructive\tdistributed cluster concurrency\tblack-box\tpython3 scripts/e2e_binary.py run -- cargo test --package e2e_test cluster_concurrency_test -- --nocapture\tnone\t%s\n' "$destructive_status"
printf 'destructive\tstale multipart cleanup cluster\tblack-box\tpython3 scripts/e2e_binary.py run -- cargo test --package e2e_test stale_multipart_cleanup_cluster_test -- --nocapture\tnone\t%s\n' "$destructive_status"
printf 'destructive\tdelete marker migration semantics\tblack-box\tpython3 scripts/e2e_binary.py run -- cargo test --package e2e_test delete_marker_migration_semantics_test -- --nocapture\tnone\t%s\n' "$destructive_status"
} >"$BLACKBOX_MATRIX"
}
@@ -566,9 +567,9 @@ run_destructive_profile() {
return
fi
run_step "e2e-cluster-concurrency" cargo test --package e2e_test cluster_concurrency_test -- --nocapture
run_step "e2e-stale-multipart-cleanup-cluster" cargo test --package e2e_test stale_multipart_cleanup_cluster_test -- --nocapture
run_step "e2e-delete-marker-migration-semantics" cargo test --package e2e_test delete_marker_migration_semantics_test -- --nocapture
run_step "e2e-cluster-concurrency" python3 scripts/e2e_binary.py run -- cargo test --package e2e_test cluster_concurrency_test -- --nocapture
run_step "e2e-stale-multipart-cleanup-cluster" python3 scripts/e2e_binary.py run -- cargo test --package e2e_test stale_multipart_cleanup_cluster_test -- --nocapture
run_step "e2e-delete-marker-migration-semantics" python3 scripts/e2e_binary.py run -- cargo test --package e2e_test delete_marker_migration_semantics_test -- --nocapture
}
run_fuzz_profile() {
+263
View File
@@ -0,0 +1,263 @@
#!/usr/bin/env python3
"""Exercise the E2E build/run boundary without compiling RustFS."""
import json
import os
from pathlib import Path
import shutil
import signal
import subprocess
import sys
import tempfile
import unittest
class BinaryProvenanceTests(unittest.TestCase):
def setUp(self):
self.temp = tempfile.TemporaryDirectory()
self.addCleanup(self.temp.cleanup)
self.root = Path(self.temp.name)
(self.root / "scripts").mkdir()
shutil.copy(Path(__file__).with_name("e2e_binary.py"), self.root / "scripts/e2e_binary.py")
(self.root / "Cargo.toml").write_text("[workspace]\n")
(self.root / "source.rs").write_text("original source\n")
(self.root / ".gitignore").write_text("/target/\n/rustfs/static/\n")
(self.root / ".agents/skills").mkdir(parents=True)
(self.root / ".agents/skills/SKILL.md").write_text("tracked instructions\n")
(self.root / ".claude").mkdir()
(self.root / ".claude/skills").symlink_to("../.agents/skills", target_is_directory=True)
subprocess.run(["git", "init", "-q", str(self.root)], check=True)
for args in (["add", "."], ["-c", "user.name=Test", "-c", "user.email=test@example.com", "commit", "-qm", "fixture"]):
subprocess.run(["git", "-C", str(self.root), *args], check=True)
self.commands = self.root / "target/commands"
self.commands.mkdir(parents=True)
cargo = self.commands / "cargo"
cargo.write_text(f"#!{sys.executable}\n" + '''import json, os, pathlib, sys
if os.environ.get("FAKE_BUILD_FAIL"):
raise SystemExit(23)
args = sys.argv[1:]
target = pathlib.Path(args[args.index("--target-dir") + 1])
binary = target / ("release" if "--release" in args else "debug") / "rustfs"
binary.parent.mkdir(parents=True, exist_ok=True)
binary.write_text("#!/bin/sh\\nexit 0\\n")
binary.chmod(0o755)
features = ["default", "ftps", "webdav"]
if "--features" in args:
features.extend(args[args.index("--features") + 1].split(","))
if "full" in features:
features.extend(["sftp", "swift", "metrics-gpu", "pyroscope"])
print(json.dumps({"reason": "compiler-artifact", "target": {"name": "rustfs", "kind": ["bin"]}, "executable": str(binary), "features": sorted(set(features))}))
if os.environ.get("FAKE_BUILD_MUTATE"):
pathlib.Path("source.rs").write_text("changed during build")
''')
cargo.chmod(0o755)
rustc = self.commands / "rustc"
rustc.write_text("#!/bin/sh\nprintf 'rustc fixture\\nhost: fixture\\n'\n")
rustc.chmod(0o755)
self.env = dict(os.environ, PATH=f"{self.commands}{os.pathsep}{os.environ['PATH']}")
for name in ("CARGO_TARGET_DIR", "CARGO_BIN_EXE_rustfs", "RUSTFS_BUILD_FEATURES", "RUSTFS_E2E_BINARY_RECEIPT"):
self.env.pop(name, None)
self.binary = self.root / "target/debug/rustfs"
self.sidecar = self.binary.with_name("rustfs.e2e.json")
def invoke(self, *args, env=None):
return subprocess.run([sys.executable, str(self.root / "scripts/e2e_binary.py"), *args], cwd=self.root, env=env or self.env, text=True, capture_output=True)
def build(self, features=""):
result = self.invoke("build", "--features", features)
self.assertEqual(result.returncode, 0, result.stderr)
def run_code(self, code="pass", features="", env=None):
return self.invoke("run", "--features", features, "--", sys.executable, "-c", code, env=env)
def test_build_run_and_receipt_cleanup(self):
self.build("full,e2e-test-hooks")
result = self.run_code("import os,pathlib; print(os.environ['RUSTFS_E2E_BINARY_RECEIPT']); assert pathlib.Path(os.environ['CARGO_BIN_EXE_rustfs']).is_file(); assert 'sftp' in os.environ['RUSTFS_BUILD_FEATURES']", "e2e-test-hooks,full")
self.assertEqual(result.returncode, 0, result.stderr)
self.assertFalse(Path(result.stdout.strip()).exists(), "run receipts must not survive their command")
self.assertIn("sftp", json.loads(self.sidecar.read_text())["features"])
def test_source_changes_are_not_hidden_by_timestamps_or_head(self):
self.build()
path = self.root / "source.rs"
old = path.stat()
path.write_text("different bytes\n")
os.utime(path, ns=(old.st_atime_ns, old.st_mtime_ns))
self.assertNotEqual(self.run_code().returncode, 0)
def test_deleted_untracked_and_ignored_embedded_inputs(self):
for mutation in ("delete", "untracked", "static"):
with self.subTest(mutation=mutation):
self.build()
path = self.root / "source.rs"
if mutation == "delete":
path.unlink()
elif mutation == "untracked":
(self.root / "new.rs").write_text("new source")
else:
static = self.root / "rustfs/static"
static.mkdir(parents=True)
(static / "index.html").write_text("embedded content")
self.assertNotEqual(self.run_code().returncode, 0)
path.write_text("original source\n")
def test_wrong_binary_features_and_manifest_fail_closed(self):
self.build("sftp")
self.assertNotEqual(self.run_code(features="webdav").returncode, 0)
self.binary.write_text("old server")
self.assertNotEqual(self.run_code(features="sftp").returncode, 0)
self.sidecar.write_text("{}")
self.assertNotEqual(self.run_code(features="sftp").returncode, 0)
self.sidecar.unlink()
self.assertNotEqual(self.run_code(features="sftp").returncode, 0)
def test_build_failure_or_source_race_does_not_leave_a_receipt(self):
for failure in ("FAKE_BUILD_FAIL", "FAKE_BUILD_MUTATE"):
self.build()
result = self.invoke("build", env=dict(self.env, **{failure: "1"}))
self.assertNotEqual(result.returncode, 0)
self.assertFalse(self.sidecar.exists())
def test_child_failure_and_changes_during_run_fail(self):
self.build()
failed = self.run_code("raise SystemExit(37)")
self.assertEqual(failed.returncode, 37, failed.stderr)
for code in ("import pathlib; pathlib.Path('source.rs').write_text('changed while testing')", "import pathlib; pathlib.Path('target/debug/rustfs').write_text('different server')"):
self.build()
self.assertNotEqual(self.run_code(code).returncode, 0)
def test_override_cannot_select_an_unverified_server(self):
self.build()
result = self.run_code(env=dict(self.env, CARGO_BIN_EXE_rustfs="/some/old/server"))
self.assertNotEqual(result.returncode, 0)
def test_artifact_moves_between_clean_checkouts(self):
self.build()
with tempfile.TemporaryDirectory() as destination:
clone = Path(destination) / "clone"
subprocess.run(["git", "clone", "-q", str(self.root), str(clone)], check=True)
(clone / "target/debug").mkdir(parents=True)
shutil.copy2(self.binary, clone / "target/debug/rustfs")
shutil.copy2(self.sidecar, clone / "target/debug/rustfs.e2e.json")
result = subprocess.run([sys.executable, str(clone / "scripts/e2e_binary.py"), "run", "--", sys.executable, "-c", "pass"], cwd=clone, env=self.env, text=True, capture_output=True)
self.assertEqual(result.returncode, 0, result.stderr)
def test_target_directory_and_profile_are_explicit(self):
env = dict(self.env, CARGO_TARGET_DIR="target/custom")
built = self.invoke("build", "--profile", "release", env=env)
self.assertEqual(built.returncode, 0, built.stderr)
run = self.invoke("run", "--profile", "release", "--", sys.executable, "-c", "pass", env=env)
self.assertEqual(run.returncode, 0, run.stderr)
self.assertNotEqual(self.invoke("run", "--", sys.executable, "-c", "pass", env=env).returncode, 0)
def test_target_directory_cannot_hide_source_inputs(self):
for target in (str(self.root), str(self.root / "crates"), str(self.root.parent)):
with self.subTest(target=target):
result = self.invoke("build", env=dict(self.env, CARGO_TARGET_DIR=target))
self.assertNotEqual(result.returncode, 0)
self.assertIn("CARGO_TARGET_DIR", result.stderr)
tracked = self.root / "target/tracked.rs"
tracked.write_text("tracked build input")
subprocess.run(["git", "add", "-f", "target/tracked.rs"], cwd=self.root, check=True)
self.build()
tracked.write_text("changed tracked build input")
self.assertNotEqual(self.run_code().returncode, 0)
def test_unsupported_embedded_directory_links_fail_closed(self):
self.build()
destination = self.root / "target/embedded-assets"
destination.mkdir()
(destination / "index.html").write_text("untracked embedded input")
static = self.root / "rustfs/static"
static.mkdir(parents=True)
(static / "linked-assets").symlink_to(destination, target_is_directory=True)
self.assertNotEqual(self.run_code().returncode, 0)
def test_directory_aliases_cannot_hide_unrecorded_inputs(self):
self.build()
target = self.root / ".agents/skills/SKILL.md"
target.write_text("changed instructions\n")
self.assertNotEqual(self.run_code().returncode, 0)
self.build()
(target.parent / ".gitignore").write_text("hidden.rs\n")
(target.parent / "hidden.rs").write_text("ignored build input\n")
result = self.invoke("build")
self.assertNotEqual(result.returncode, 0)
self.assertIn("unrecorded input", result.stderr)
alias = self.root / ".claude/skills"
alias.unlink()
with tempfile.TemporaryDirectory() as external:
alias.symlink_to(external, target_is_directory=True)
result = self.invoke("build")
self.assertNotEqual(result.returncode, 0)
self.assertIn("escapes the source inventory", result.stderr)
def test_directory_alias_indirection_is_part_of_the_identity(self):
for name in ("first", "second"):
directory = self.root / name
directory.mkdir()
(directory / "input.rs").write_text(name)
selection = self.root / "target/selection"
selection.symlink_to(self.root / "first", target_is_directory=True)
(self.root / "source-alias").symlink_to("target/selection", target_is_directory=True)
self.build()
selection.unlink()
selection.symlink_to(self.root / "second", target_is_directory=True)
self.assertNotEqual(self.run_code().returncode, 0)
def test_existing_embedded_files_and_symlink_targets_are_hashed(self):
static = self.root / "rustfs/static"
static.mkdir(parents=True)
index = static / "index.html"
index.write_text("embedded version one")
external = self.root / "target/embedded-file"
external.write_text("linked version one")
(static / "linked.html").symlink_to(external)
self.build()
index.write_text("embedded version two")
self.assertNotEqual(self.run_code().returncode, 0)
self.build()
external.write_text("linked version two")
self.assertNotEqual(self.run_code().returncode, 0)
def test_each_run_hashes_binary_twice_and_never_calls_cargo(self):
script = self.root / "scripts/e2e_binary.py"
script.write_text(script.read_text().replace("def file_hash(path):\n", "def file_hash(path):\n if path.name == 'rustfs':\n with (ROOT / 'target/hash-count').open('a') as count:\n count.write('hash\\n')\n"))
self.build()
count = self.root / "target/hash-count"
count.write_text("")
result = self.run_code(env=dict(self.env, FAKE_BUILD_FAIL="1"))
self.assertEqual(result.returncode, 0, result.stderr)
self.assertEqual(count.read_text().splitlines(), ["hash", "hash"])
def test_concurrent_build_or_run_is_rejected(self):
self.build()
command = [sys.executable, str(self.root / "scripts/e2e_binary.py"), "run", "--", sys.executable, "-c", "print('ready', flush=True); input()"]
with subprocess.Popen(command, cwd=self.root, env=self.env, stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True) as process:
self.assertEqual(process.stdout.readline().strip(), "ready")
try:
for args in (("build", "--features", "sftp"), ("run", "--", sys.executable, "-c", "pass")):
rejected = self.invoke(*args)
self.assertNotEqual(rejected.returncode, 0)
self.assertIn("Another E2E build/run", rejected.stderr)
finally:
output, error = process.communicate("\n", timeout=10)
self.assertEqual(process.returncode, 0, error + output)
self.assertFalse(self.binary.with_name("rustfs.e2e.lock").exists())
def test_interruption_cleans_receipt_and_releases_ownership(self):
self.build()
for signum in (signal.SIGINT, signal.SIGTERM):
command = [sys.executable, str(self.root / "scripts/e2e_binary.py"), "run", "--", sys.executable, "-c", "import os; print(os.environ['RUSTFS_E2E_BINARY_RECEIPT'], flush=True); input()"]
with subprocess.Popen(command, cwd=self.root, env=self.env, stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True) as process:
receipt = Path(process.stdout.readline().strip())
self.assertTrue(receipt.is_file())
process.send_signal(signum)
process.communicate(timeout=10)
self.assertNotEqual(process.returncode, 0)
self.assertFalse(receipt.exists())
self.assertFalse(self.binary.with_name("rustfs.e2e.lock").exists())
if __name__ == "__main__":
unittest.main()