diff --git a/.config/e2e-distributed-selection.txt b/.config/e2e-distributed-selection.txt index 68936d5e5..c034893ad 100644 --- a/.config/e2e-distributed-selection.txt +++ b/.config/e2e-distributed-selection.txt @@ -1,2 +1,2 @@ -sha256-linux=21ab66fd172018682ce825a53f26aa8f286d57f3f687d8616d1850e56474dcd8 -sha256-darwin=21ab66fd172018682ce825a53f26aa8f286d57f3f687d8616d1850e56474dcd8 +sha256-linux=6fc377fa1f9f06e065f077c7bd53d04efbccf4453f95dde67afaa47778185377 +sha256-darwin=6fc377fa1f9f06e065f077c7bd53d04efbccf4453f95dde67afaa47778185377 diff --git a/.config/make/tests.mak b/.config/make/tests.mak index 2788c0f2a..72f39e663 100644 --- a/.config/make/tests.mak +++ b/.config/make/tests.mak @@ -52,6 +52,7 @@ script-tests: ## Run shell script tests ./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/test_migration_gate_evidence.py $(RUSTFS_PYTHON_BIN) ./scripts/ci_gate.py --self-test $(RUSTFS_PYTHON_BIN) ./scripts/check_security_coverage.py --self-test $(RUSTFS_PYTHON_BIN) ./scripts/check_scheduled_validation_freshness.py --self-test @@ -61,6 +62,7 @@ script-tests: ## Run shell script tests $(RUSTFS_PYTHON_BIN) ./scripts/test_functional_chain_health.py $(RUSTFS_PYTHON_BIN) ./scripts/test_ci_timing_report.py $(RUSTFS_PYTHON_BIN) ./scripts/s3-tests/test_report_compat.py + $(RUSTFS_PYTHON_BIN) ./scripts/s3-tests/test_runner_tools.py bash -n ./scripts/validate_object_data_cache_cold_stampede.sh $(RUSTFS_PYTHON_BIN) ./scripts/check_object_data_cache_follower_samples.py --self-test ./scripts/validate_object_data_cache_cold_stampede.sh --self-test diff --git a/.config/nextest.toml b/.config/nextest.toml index 6a7d83bbc..70efc267d 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -149,6 +149,13 @@ test-group = 'embedded-test-ports' filter = 'package(rustfs) & (binary(connect_inventory) | binary(connect_perf_drive) | (binary(connect_perf_object) & test(=real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put)))' threads-required = "num-test-threads" +# These concurrent state writers keep the production 5s lock-acquisition limit +# while their peer completes durable IO. Reserve capacity from unrelated test +# processes while preserving each test's internal two-writer race. +[[profile.default.overrides]] +filter = 'package(rustfs) & test(/^admin::handlers::site_replication::tests::test_(peer_edit_generations_are_unique_across_nodes|recreated_state_object_allocates_over_the_previous_lifetimes_mark|state_object_lock_serializes_writers_from_separate_nodes|retry_event_persist_must_not_wipe_concurrent_locked_rmw)$/)' +threads-required = "num-test-threads" + # Reserve capacity for these short heartbeat TLS/fsync probes (rustfs#8129). # Their existing deadlines, assertions, and retry policy remain in force. [[profile.default.overrides]] @@ -359,6 +366,12 @@ test-group = 'embedded-test-ports' filter = 'package(rustfs) & (binary(connect_inventory) | binary(connect_perf_drive) | (binary(connect_perf_object) & test(=real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put)))' threads-required = "num-test-threads" +# Keep the same state-writer capacity reservation in CI without changing the +# production lock deadline, internal concurrency, assertions, or retry policy. +[[profile.ci.overrides]] +filter = 'package(rustfs) & test(/^admin::handlers::site_replication::tests::test_(peer_edit_generations_are_unique_across_nodes|recreated_state_object_allocates_over_the_previous_lifetimes_mark|state_object_lock_serializes_writers_from_separate_nodes|retry_event_persist_must_not_wipe_concurrent_locked_rmw)$/)' +threads-required = "num-test-threads" + # Reserve capacity for these short heartbeat TLS/fsync probes (rustfs#8129). # Their existing deadlines, assertions, and retry policy remain in force. [[profile.ci.overrides]] diff --git a/.github/actions/quick-checks/action.yml b/.github/actions/quick-checks/action.yml index ca761f71f..e6650ed47 100644 --- a/.github/actions/quick-checks/action.yml +++ b/.github/actions/quick-checks/action.yml @@ -100,12 +100,7 @@ runs: - name: Check test wiring shell: bash - 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/test_nightly_candidate.py - python3 ./scripts/check_test_wiring.py + run: python3 ./scripts/check_test_wiring.py - name: Check no planning docs committed shell: bash diff --git a/.github/actions/setup/action.yml b/.github/actions/setup/action.yml index 712afb3d0..eb8f8b377 100644 --- a/.github/actions/setup/action.yml +++ b/.github/actions/setup/action.yml @@ -107,6 +107,8 @@ runs: - name: Install cargo-nextest if: inputs.install-test-tools == 'true' uses: taiki-e/install-action@96c7780c1d8a2b8723e12031def873a434d39d8d # nextest + with: + tool: cargo-nextest@0.9.138 - name: Setup Rust cache uses: Swatinem/rust-cache@6323deb102c322ba6fcbdcafc7e3dddab59af2b6 # v2 diff --git a/.github/workflows/architecture-migration-rules.yml b/.github/workflows/architecture-migration-rules.yml deleted file mode 100644 index e71459a4b..000000000 --- a/.github/workflows/architecture-migration-rules.yml +++ /dev/null @@ -1,61 +0,0 @@ -# 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: Architecture Migration Rules - -on: - pull_request: - types: [ opened, synchronize, reopened, closed ] - branches: [ main ] - paths: - - "ARCHITECTURE.md" - - "docs/architecture/**" - - "scripts/check_architecture_migration_rules.sh" - - ".github/workflows/architecture-migration-rules.yml" - workflow_dispatch: - -permissions: - contents: read - -concurrency: - group: ${{ github.workflow }}-${{ github.event.pull_request.number || github.ref }} - cancel-in-progress: true - -jobs: - cancel-closed-pr-runs: - name: Cancel Closed PR Runs - if: github.event_name == 'pull_request' && github.event.action == 'closed' - runs-on: sm-standard-2 - timeout-minutes: 10 - steps: - - name: Explain cancellation run - run: echo "PR closed; this run only cancels older runs in the same concurrency group." - - architecture-migration-rules: - name: Architecture Migration Rules - if: github.event_name != 'pull_request' || github.event.action != 'closed' - runs-on: sm-standard-2 - timeout-minutes: 10 - steps: - - uses: actions/checkout@f548e57e544e1ff5a4c46bf1e1b8685f8e4a348a # v7 - with: - persist-credentials: false - - - name: Install ripgrep - uses: taiki-e/install-action@7623a79cdfecb99d681017af368ca353d9f49bb5 # v2 - with: - tool: ripgrep@15.2.0 - - - name: Check architecture migration rules - run: ./scripts/check_architecture_migration_rules.sh diff --git a/.github/workflows/cache-warm.yml b/.github/workflows/cache-warm.yml index 0ef180908..699ba2990 100644 --- a/.github/workflows/cache-warm.yml +++ b/.github/workflows/cache-warm.yml @@ -15,16 +15,16 @@ # Sole writer of the Rust dependency caches that ci.yml restores. # # Why this is a separate workflow rather than steps inside ci.yml: ci.yml's -# concurrency group cancels in-progress runs on main pushes, and merges land far -# faster than its 70-minute pipeline. Measured over 15 consecutive main pushes: +# concurrency group originally cancelled in-progress main runs, while merges +# landed faster than its 70-minute pipeline. Over 15 consecutive main pushes: # 12 cancelled, 2 failed, 0 succeeded. A cancelled run never reaches # Swatinem/rust-cache's post step (cache-on-failure does not cover cancellation), # so the writer lanes were saving nothing and every PR paid a cold restore — # 11.8-20.9 minutes of "Setup Rust environment" against 0.7-3.4 warm. # -# Splitting cache writing out of the test pipeline lets ci.yml keep cancelling -# superseded runs (which is correct — nobody needs test results for a commit -# that is already three merges behind) while the caches still get written. +# Keep cache writing separate from validation: only these jobs publish the +# complete feature closure for each key. Main validation now also finishes its +# running baseline, while superseded PR attempts can still be cancelled. # # The group below deliberately does NOT cancel in progress; see the comment on # it for how that bounds concurrency and why it is scoped by event. @@ -256,7 +256,7 @@ jobs: # belonged to it. warm-ci-uring: name: Warm ci-uring - runs-on: sm-standard-2 + runs-on: ubuntu-latest timeout-minutes: 60 steps: - name: Checkout repository diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 759a3c510..90a32f534 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -47,12 +47,12 @@ permissions: contents: read # Concurrency groups are scoped per event so different triggers never cancel -# each other: PR pushes cancel the previous run of that PR, main pushes keep -# latest-wins semantics among themselves, and scheduled runs always complete -# (a shared group used to let every merge kill the weekly scheduled run). +# each other. PR pushes cancel superseded attempts; main and scheduled runs +# finish so a busy merge stream cannot starve the complete baseline. For main, +# GitHub retains one running run and replaces the pending run with the latest. concurrency: group: ${{ github.workflow }}-${{ github.event_name }}-${{ github.event.pull_request.number || github.ref }} - cancel-in-progress: ${{ github.event_name != 'schedule' }} + cancel-in-progress: ${{ github.event_name == 'pull_request' }} env: CARGO_TERM_COLOR: always @@ -128,7 +128,7 @@ jobs: test-and-lint: name: Workspace Test and Lint - if: needs.classify-changes.outputs.mode == 'full' && (github.event_name != 'pull_request' || github.event.action != 'closed') + if: contains(fromJSON('["full", "e2e"]'), needs.classify-changes.outputs.mode) && (github.event_name != 'pull_request' || github.event.action != 'closed') needs: [ quick-checks, classify-changes ] runs-on: sm-standard-4 timeout-minutes: 90 @@ -146,11 +146,9 @@ jobs: uses: ./.github/actions/setup with: rust-version: stable - # Every lane in this workflow reads its cache and none writes it. - # cache-warm.yml is the sole writer for all four keys: this workflow - # cancels superseded runs on main, and a cancelled run never reaches - # rust-cache's post step, so writing from here saved nothing (12 of 15 - # consecutive main-push runs were cancelled). See rustfs/backlog#1600. + # Every lane reads its cache; cache-warm.yml remains the sole writer + # with the matching feature closure. Tests never compete to replace + # shared caches, including when a PR attempt is cancelled. cache-shared-key: ci-dev cache-save-if: 'false' install-build-packaging-tools: 'false' @@ -281,21 +279,10 @@ jobs: - name: Check log-analyzer rule anchors run: ./scripts/check_log_analyzer_rules.sh - # Explicit gate for migration-critical suites. These tests already ran in - # the full nextest pass above; a single filtered nextest invocation keeps - # the named gate without rebuilding or re-running them one package at a time. - # - # The gate selects tests by name substring (data_movement / rebalance / - # decommission / source_cleanup / delete_marker), so renames can silently - # thin it. The script owns the filter expression and first verifies the - # selected-test count against the committed floor in - # .config/migration-gate-floor.txt before running the gate; renames or - # removals must update that file consciously (see the script header). - # Kept on the default profile (no --profile ci): a second --profile ci run - # would clobber target/nextest/ci/junit.xml, and none of these tests are - # quarantined so they gain nothing from the ci profile's retry overrides. - - name: Run rebalance/decommission migration proofs - run: ./scripts/check_migration_gate_count.sh + # Preserve the migration floor and require successful execution evidence + # from the workspace run, without compiling or executing its subset again. + - name: Verify rebalance/decommission migration proofs + run: ./scripts/check_migration_gate_count.sh evidence artifacts/test-and-lint/core-test-listing.json target/nextest/ci/junit.xml # This gate builds into a fresh target directory to isolate the E2E root. # Give its cold build a separate budget from workspace tests and migration proofs. @@ -525,13 +512,9 @@ jobs: runs-on: sm-standard-4 timeout-minutes: 90 strategy: - # On a PR, one failing protocol leg is enough to know the PR is not ready, - # so stop the sibling leg instead of paying another ~40 minutes for it. - # Everywhere else (main pushes, the merge queue, the weekly schedule) keep - # the full signal: there we want to know whether swift AND sftp are broken, - # not just whichever failed first. This is the only part of the early-stop - # work that also covers fork PRs, since it needs no token. - fail-fast: ${{ github.event_name == 'pull_request' }} + # Preserve both independent results when one protocol fails, so diagnosis + # and a focused fix do not require rebuilding an interrupted sibling. + fail-fast: false matrix: features: - name: swift @@ -562,6 +545,8 @@ jobs: cargo clippy -p rustfs -p rustfs-protocols --all-targets ${{ matrix.features.flags }} -- -D warnings - name: Run tests with ${{ matrix.features.name }} + id: protocol-tests + shell: bash env: # Keep feature-test linking under the same bounded concurrency as the # main nextest lane; Clippy is metadata-only and needs no such limit. @@ -570,11 +555,26 @@ jobs: # --profile ci so the quarantine list (and its junit flaky markers) # covers this leg too; the default profile is the local no-retry # profile and silently ignored quarantined flakes here (rustfs#6703). - cargo nextest run --profile ci -p rustfs -p rustfs-protocols ${{ matrix.features.flags }} + mkdir -p artifacts/protocol-tests + rm -f target/nextest/ci/junit.xml + cargo nextest run --profile ci -p rustfs -p rustfs-protocols ${{ matrix.features.flags }} \ + 2>&1 | tee artifacts/protocol-tests/nextest.log + + - name: Upload protocol test reports and diagnostics + if: >- + always() && contains(fromJSON('["success", "failure", "cancelled"]'), steps.protocol-tests.outcome) + uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1 + with: + name: junit-test-and-lint-${{ matrix.features.name }}-${{ github.run_number }}-${{ github.run_attempt }} + path: | + target/nextest/ci/junit.xml + artifacts/protocol-tests + retention-days: 3 + if-no-files-found: warn build-rustfs-debug-binary: name: Build RustFS Debug Binary - if: needs.classify-changes.outputs.mode == 'full' && (github.event_name != 'pull_request' || github.event.action != 'closed') + if: contains(fromJSON('["full", "e2e"]'), needs.classify-changes.outputs.mode) && (github.event_name != 'pull_request' || github.event.action != 'closed') needs: [ quick-checks, classify-changes ] runs-on: sm-standard-4 timeout-minutes: 30 diff --git a/.github/workflows/e2e-distributed.yml b/.github/workflows/e2e-distributed.yml index 889d84733..b890c5d43 100644 --- a/.github/workflows/e2e-distributed.yml +++ b/.github/workflows/e2e-distributed.yml @@ -38,6 +38,7 @@ on: - "Cargo.lock" - "Cargo.toml" - ".config/nextest.toml" + - ".config/e2e-distributed-selection.txt" - ".github/workflows/e2e-distributed.yml" - "crates/audit/**" - "crates/common/**" diff --git a/.github/workflows/e2e-upgrade.yml b/.github/workflows/e2e-upgrade.yml index 59779003e..21125fe13 100644 --- a/.github/workflows/e2e-upgrade.yml +++ b/.github/workflows/e2e-upgrade.yml @@ -51,8 +51,53 @@ env: UPGRADE_SOURCE_SHA256: 3ee8df71e8edcfada533be452c4135868f697bc515460ae97b027313eade7a3d jobs: + build: + name: Build upgrade candidate + runs-on: sm-standard-2 + timeout-minutes: 60 + env: + FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true" + steps: + - name: Checkout repository + uses: actions/checkout@f548e57e544e1ff5a4c46bf1e1b8685f8e4a348a # v7 + with: + persist-credentials: false + + - name: Setup Rust environment + uses: ./.github/actions/setup + with: + cache-shared-key: e2e-upgrade-server + cache-save-if: ${{ github.ref == 'refs/heads/main' }} + install-build-packaging-tools: "false" + install-test-tools: "false" + + - name: Download pinned previous release + run: | + mkdir -p target/debug + archive="target/debug/$UPGRADE_SOURCE_ASSET" + curl --fail --location --retry 3 --output "$archive" \ + "https://github.com/${GITHUB_REPOSITORY}/releases/download/${UPGRADE_SOURCE_VERSION}/${UPGRADE_SOURCE_ASSET}" + echo "$UPGRADE_SOURCE_SHA256 $archive" | sha256sum --check --strict + + - name: Build current RustFS binary + run: python3 scripts/e2e_binary.py build + + - name: Upload verified upgrade candidate + uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1 + with: + name: upgrade-candidate-${{ github.run_id }} + path: | + target/debug/rustfs + target/debug/rustfs.e2e.json + target/debug/${{ env.UPGRADE_SOURCE_ASSET }} + if-no-files-found: error + retention-days: 3 + compression-level: 0 + overwrite: true + upgrade: name: ${{ matrix.name }} + needs: build strategy: fail-fast: false matrix: @@ -63,31 +108,24 @@ jobs: # the CI required-check names. UPGRADE_SOURCE_VERSION above is the # single source of truth for which release they actually run against. - name: Direct upgrade from the previous release - cache_key: e2e-direct-upgrade test: direct_upgrade_from_rc2_preserves_object_contracts artifact: direct-upgrade - name: Mixed-version rolling upgrade from the previous release - cache_key: e2e-mixed-version-upgrade test: rolling_upgrade_from_rc2_preserves_mixed_version_contracts artifact: mixed-version-upgrade - name: Bucket configuration survives the upgrade - cache_key: e2e-bucket-config-upgrade test: direct_upgrade_from_previous_release_preserves_bucket_configuration artifact: bucket-config-upgrade - name: Rollback reads current bucket metadata - cache_key: e2e-bucket-config-rollback test: rollback_to_previous_release_reads_current_bucket_metadata artifact: bucket-config-rollback - name: ODM configuration recovery after rc.5 rollback - cache_key: e2e-odm-config-rollback test: rc5_rollback_requires_restoring_odm_configuration artifact: odm-config-rollback - name: Multipart layouts survive the rc.5 upgrade - cache_key: e2e-multipart-layout-upgrade test: direct_upgrade_from_rc5_preserves_multipart_layouts artifact: multipart-layout-upgrade - name: rc.5 multipart replication baseline - cache_key: e2e-multipart-layout-baseline test: rc5_baseline_replicates_multipart_layouts artifact: multipart-layout-baseline runs-on: sm-standard-2 @@ -103,19 +141,28 @@ jobs: - name: Setup Rust environment uses: ./.github/actions/setup with: - cache-shared-key: ${{ matrix.cache_key }} - cache-save-if: ${{ github.ref == 'refs/heads/main' }} + cache-shared-key: e2e-upgrade-tests + # One main-branch consumer warms the shared test-harness dependencies. + cache-save-if: ${{ github.ref == 'refs/heads/main' && matrix.artifact == 'direct-upgrade' }} install-build-packaging-tools: "false" + install-test-tools: "false" - - name: Download pinned previous release + - name: Download verified upgrade candidate + uses: actions/download-artifact@3e5f45b2cfb9172054b4087a40e8e0b5a5461e7c # v8.0.1 + with: + name: upgrade-candidate-${{ github.run_id }} + path: target/debug + + - name: Restore candidate executable bit + run: chmod +x target/debug/rustfs + + - name: Verify and unpack pinned previous release env: SOURCE_DIR: ${{ runner.temp }}/rustfs-upgrade-source run: | set -euo pipefail mkdir -p "$SOURCE_DIR" - archive="$SOURCE_DIR/$UPGRADE_SOURCE_ASSET" - curl --fail --location --retry 3 --output "$archive" \ - "https://github.com/${GITHUB_REPOSITORY}/releases/download/${UPGRADE_SOURCE_VERSION}/${UPGRADE_SOURCE_ASSET}" + archive="target/debug/$UPGRADE_SOURCE_ASSET" echo "$UPGRADE_SOURCE_SHA256 $archive" | sha256sum --check --strict unzip -q "$archive" -d "$SOURCE_DIR" chmod +x "$SOURCE_DIR/rustfs" @@ -123,14 +170,12 @@ jobs: echo "RUSTFS_UPGRADE_SOURCE_BINARY=$SOURCE_DIR/rustfs" >> "$GITHUB_ENV" echo "RUSTFS_E2E_LOG_DIR=$RUNNER_TEMP/rustfs-upgrade-logs" >> "$GITHUB_ENV" - - name: Build current RustFS binary - run: | - python3 scripts/e2e_binary.py build - - name: Run upgrade compatibility test env: RUSTFS_SCANNER_HEAL_G09_EVIDENCE_DIR: ${{ runner.temp }}/rustfs-upgrade-g09-evidence/${{ matrix.artifact }} + NO_PROXY: 127.0.0.1,localhost,::1 run: | + export no_proxy="$NO_PROXY" python3 scripts/e2e_binary.py run -- cargo test --locked -p e2e_test \ "upgrade_compatibility_test::${{ matrix.test }}" \ -- --ignored --exact --nocapture diff --git a/crates/e2e_test/src/distributed/harness.rs b/crates/e2e_test/src/distributed/harness.rs index bf53789c9..364034477 100644 --- a/crates/e2e_test/src/distributed/harness.rs +++ b/crates/e2e_test/src/distributed/harness.rs @@ -34,6 +34,9 @@ use crate::common::{ }; use crate::replication_extension_test::LOOPBACK_REPLICATION_TARGET_ENV; use aws_sdk_s3::Client; +use aws_sdk_s3::config::retry::RetryConfig; +use aws_sdk_s3::error::{ProvideErrorMetadata, SdkError}; +use aws_sdk_s3::operation::{get_object::GetObjectError, put_object::PutObjectError}; use aws_sdk_s3::primitives::ByteStream; use aws_sdk_s3::types::{BucketVersioningStatus, VersioningConfiguration}; use http::{Method, StatusCode}; @@ -41,7 +44,7 @@ use sha2::{Digest, Sha256}; use std::collections::BTreeMap; use std::path::{Path, PathBuf}; use std::time::Duration; -use tokio::time::{Instant, sleep}; +use tokio::time::{Instant, sleep, timeout_at}; use uuid::Uuid; pub(crate) type TestResult = Result>; @@ -400,9 +403,9 @@ pub(crate) async fn put_inventory( Ok(inventory) } -/// Retry only transport-level service availability failures while a data -/// movement operation changes the pool map. Generic InternalError responses -/// remain fatal because accepting them would hide server defects. +/// Retry only explicit S3 availability responses while a data movement +/// operation changes the pool map. Transport failures and InternalError +/// responses remain fatal because callers already require a ready endpoint. pub(crate) async fn put_inventory_retrying( client: &Client, bucket: &str, @@ -449,19 +452,48 @@ where let deadline = Instant::now() + timeout; let mut delay = Duration::from_millis(50); loop { - let last_error = match probe().await { - Ok(true) => return Ok(()), - Ok(false) => format!("{label} still false"), - Err(error) => error.to_string(), - }; if Instant::now() >= deadline { - return Err(format!("{label} did not become true within {timeout:?}: {last_error}").into()); + return Err(format!("{label} did not become true within {timeout:?}").into()); } - sleep(delay).await; + // A pending request must consume the same budget as unsuccessful probes. + match timeout_at(deadline, probe()).await { + Err(_) => return Err(format!("{label} did not become true within {timeout:?}").into()), + Ok(Err(error)) => return Err(error), + Ok(Ok(true)) => return Ok(()), + Ok(Ok(false)) => {} + } + tokio::time::sleep_until((Instant::now() + delay).min(deadline)).await; delay = (delay * 2).min(Duration::from_secs(1)); } } +// Boxed SDK errors retain the operation's service code; Display only says +// "service error" and cannot distinguish convergence from a server defect. +pub(crate) fn s3_probe_error_code<'a>(error: &'a (dyn std::error::Error + Send + Sync + 'static)) -> Option<&'a str> { + error + .downcast_ref::>() + .and_then(|error| error.as_service_error()) + .and_then(ProvideErrorMetadata::code) + .or_else(|| { + error + .downcast_ref::>() + .and_then(|error| error.as_service_error()) + .and_then(ProvideErrorMetadata::code) + }) +} + +pub(crate) fn s3_probe_client(client: &Client) -> Client { + // The probe owns the retry policy. SDK retries must not hide a response + // that the probe would classify as fatal. + Client::from_conf( + client + .config() + .to_builder() + .retry_config(RetryConfig::standard().with_max_attempts(1)) + .build(), + ) +} + pub(crate) async fn cluster_admin( cluster: &RustFSTestClusterEnvironment, method: Method, @@ -598,15 +630,20 @@ pub(crate) async fn wait_for_replicated_bytes( expected: &[u8], timeout: Duration, ) -> TestResult { + let client = s3_probe_client(client); wait_until( timeout, || async { - match get_object_bytes(client, bucket, key).await { + match get_object_bytes(&client, bucket, key).await { Ok(got) if got.as_slice() == expected => Ok(true), - Ok(_) => Ok(false), + Ok(got) => Err(format!( + "replicated object {bucket}/{key} bytes mismatch: expected sha256={} got sha256={}", + sha256_hex(expected), + sha256_hex(&got) + ) + .into()), Err(error) => { - let message = error.to_string(); - if message.contains("NoSuchKey") || message.contains("NotFound") { + if matches!(s3_probe_error_code(error.as_ref()), Some("NoSuchKey" | "NotFound")) { Ok(false) } else { Err(error) @@ -988,6 +1025,7 @@ pub(crate) async fn list_pools_json(cluster: &RustFSTestClusterEnvironment) -> T } pub(crate) async fn retrying_put(client: &Client, bucket: &str, key: &str, body: Vec, timeout: Duration) -> TestResult { + let client = s3_probe_client(client); wait_until( timeout, || { @@ -999,8 +1037,7 @@ pub(crate) async fn retrying_put(client: &Client, bucket: &str, key: &str, body: match put_object(&client, &bucket, &key, body).await { Ok(()) => Ok(true), Err(error) => { - let message = error.to_string(); - if message.contains("SlowDown") || message.contains("ServiceUnavailable") || message.contains("503") { + if matches!(s3_probe_error_code(error.as_ref()), Some("SlowDown" | "ServiceUnavailable")) { Ok(false) } else { Err(error) @@ -1021,19 +1058,20 @@ pub(crate) async fn retrying_get_equals( expected: &[u8], timeout: Duration, ) -> TestResult { + let client = s3_probe_client(client); wait_until( timeout, || async { - match get_object_bytes(client, bucket, key).await { + match get_object_bytes(&client, bucket, key).await { Ok(got) if got.as_slice() == expected => Ok(true), - Ok(_) => Ok(false), + Ok(got) => Err(format!( + "object {bucket}/{key} bytes mismatch: expected sha256={} got sha256={}", + sha256_hex(expected), + sha256_hex(&got) + ) + .into()), Err(error) => { - let message = error.to_string(); - if message.contains("NoSuchKey") - || message.contains("SlowDown") - || message.contains("ServiceUnavailable") - || message.contains("503") - { + if matches!(s3_probe_error_code(error.as_ref()), Some("SlowDown" | "ServiceUnavailable")) { Ok(false) } else { Err(error) @@ -1046,6 +1084,162 @@ pub(crate) async fn retrying_get_equals( .await } +#[cfg(test)] +mod retry_tests { + use super::*; + use crate::fake_s3_target::{FAKE_ACCESS_KEY, FAKE_SECRET_KEY, FakeS3Target, FaultAction, Operation, SeedMetadata}; + use std::cell::Cell; + + fn client(target: &FakeS3Target) -> Client { + Client::from_conf(build_test_s3_config( + target.endpoint(), + FAKE_ACCESS_KEY, + FAKE_SECRET_KEY, + None, + "distributed-retry-test", + )) + } + + #[tokio::test] + async fn wait_until_preserves_first_error() { + let attempts = Cell::new(0); + let error = wait_until( + Duration::from_secs(5), + || { + attempts.set(attempts.get() + 1); + std::future::ready(if attempts.get() == 1 { + Err(std::io::Error::new(std::io::ErrorKind::PermissionDenied, "permanent probe failure").into()) + } else { + Ok(true) + }) + }, + "permanent failure", + ) + .await + .expect_err("a later success must not hide a permanent failure"); + assert_eq!(attempts.get(), 1); + assert_eq!( + error + .downcast_ref::() + .expect("preserve the original error") + .kind(), + std::io::ErrorKind::PermissionDenied + ); + } + + #[tokio::test] + async fn wait_until_retries_explicit_pending_state() -> TestResult { + let attempts = Cell::new(0); + wait_until( + Duration::from_secs(5), + || { + attempts.set(attempts.get() + 1); + std::future::ready(Ok(attempts.get() == 2)) + }, + "eventual readiness", + ) + .await?; + assert_eq!(attempts.get(), 2); + Ok(()) + } + + #[tokio::test] + async fn wait_until_bounds_a_pending_probe() { + let error = tokio::time::timeout( + Duration::from_secs(5), + wait_until(Duration::from_millis(10), std::future::pending, "hung probe"), + ) + .await + .expect("the probe deadline must finish before the test watchdog") + .expect_err("a permanently pending probe must time out"); + assert!(error.to_string().contains("hung probe did not become true")); + } + + #[tokio::test] + async fn wait_until_does_not_start_a_probe_after_its_deadline() { + let attempts = Cell::new(0); + wait_until( + Duration::ZERO, + || { + attempts.set(attempts.get() + 1); + std::future::ready(Ok(true)) + }, + "expired budget", + ) + .await + .expect_err("an expired budget must not admit another probe"); + assert_eq!(attempts.get(), 0); + } + + #[tokio::test] + async fn retrying_put_does_not_hide_internal_error() -> TestResult { + let target = FakeS3Target::start().await?; + target.create_bucket("retry-probe"); + target.inject(Operation::PutObject, FaultAction::ResponseStatus(500), 1); + let error = retrying_put(&client(&target), "retry-probe", "key", b"body".to_vec(), Duration::from_secs(5)) + .await + .expect_err("InternalError must remain fatal even if the next PUT would succeed"); + assert_eq!(s3_probe_error_code(error.as_ref()), Some("InternalError")); + assert_eq!(target.requests().len(), 1); + Ok(()) + } + + #[tokio::test] + async fn retrying_put_allows_transient_availability_errors() -> TestResult { + let target = FakeS3Target::start().await?; + target.create_bucket("retry-probe"); + for status in [429, 503] { + target.inject(Operation::PutObject, FaultAction::ResponseStatus(status), 1); + retrying_put(&client(&target), "retry-probe", "key", b"body".to_vec(), Duration::from_secs(5)).await?; + assert_eq!(target.take_requests().len(), 2); + } + Ok(()) + } + + #[tokio::test] + async fn object_polls_reject_wrong_object_bytes() -> TestResult { + let target = FakeS3Target::start().await?; + target.create_bucket("retry-probe"); + target.put_seed_object("retry-probe", "key", "wrong bytes", &SeedMetadata::default()); + let error = retrying_get_equals(&client(&target), "retry-probe", "key", b"expected bytes", Duration::from_secs(5)) + .await + .expect_err("a successful response with wrong bytes must fail immediately"); + assert!(error.to_string().contains("bytes mismatch")); + assert_eq!(target.take_requests().len(), 1); + let error = wait_for_replicated_bytes(&client(&target), "retry-probe", "key", b"expected bytes", Duration::from_secs(5)) + .await + .expect_err("replication may lag, but it must not return corrupt bytes"); + assert!(error.to_string().contains("bytes mismatch")); + assert_eq!(target.requests().len(), 1); + Ok(()) + } + + #[tokio::test] + async fn object_polls_distinguish_committed_and_replicated_keys() -> TestResult { + let target = FakeS3Target::start().await?; + target.create_bucket("retry-probe"); + target.put_seed_object("retry-probe", "key", "expected bytes", &SeedMetadata::default()); + let client = client(&target); + for (status, expected_code) in [(404, "NoSuchKey"), (500, "InternalError"), (403, "AccessDenied")] { + target.inject(Operation::GetObject, FaultAction::ResponseStatus(status), 1); + let error = retrying_get_equals(&client, "retry-probe", "key", b"expected bytes", Duration::from_secs(5)) + .await + .expect_err("an acknowledged object must not disappear or fail before a later successful read"); + assert_eq!(s3_probe_error_code(error.as_ref()), Some(expected_code)); + assert_eq!(target.take_requests().len(), 1); + } + for status in [429, 503] { + target.inject(Operation::GetObject, FaultAction::ResponseStatus(status), 1); + retrying_get_equals(&client, "retry-probe", "key", b"expected bytes", Duration::from_secs(5)).await?; + assert_eq!(target.take_requests().len(), 2); + } + target.inject(Operation::GetObject, FaultAction::ResponseStatus(404), 1); + wait_for_replicated_bytes(&client, "retry-probe", "key", b"expected bytes", Duration::from_secs(5)).await?; + assert_eq!(target.requests().len(), 2); + Ok(()) + } +} + #[tokio::test] async fn append_single_node_pool_extends_ellipses_volumes() { let mut env = RustFSTestClusterEnvironment::with_topology(ClusterTopology::per_node_pools(2, vec![vec![0], vec![1]])) diff --git a/crates/e2e_test/src/distributed/upgrade_test.rs b/crates/e2e_test/src/distributed/upgrade_test.rs index 33d3c1c7a..d29c6f61f 100644 --- a/crates/e2e_test/src/distributed/upgrade_test.rs +++ b/crates/e2e_test/src/distributed/upgrade_test.rs @@ -26,7 +26,7 @@ use super::harness::{ DistCluster, DistLayout, TestResult, assert_object_bytes, cluster_admin_ok, enable_versioning, get_object_bytes, put_object, - unique_bucket, wait_until, + s3_probe_client, s3_probe_error_code, sha256_hex, unique_bucket, wait_until, }; use crate::common::{ AdminTransport, admin_add_canned_policy_via, admin_attach_user_policy_via, admin_create_user_via, init_logging, @@ -111,6 +111,7 @@ async fn create_iam_user(dist: &DistCluster, user: &str, secret: &str, policy_na } async fn wait_for_put(client: &Client, bucket: &str, key: &str, body: Vec, label: &str) -> TestResult { + let client = s3_probe_client(client); wait_until( CREDENTIAL_TIMEOUT, || { @@ -119,8 +120,18 @@ async fn wait_for_put(client: &Client, bucket: &str, key: &str, body: Vec, l let key = key.to_string(); let body = body.clone(); async move { - put_object(&client, &bucket, &key, body).await?; - Ok(true) + match put_object(&client, &bucket, &key, body).await { + Ok(()) => Ok(true), + Err(error) + if matches!( + s3_probe_error_code(error.as_ref()), + Some("AccessDenied" | "InvalidAccessKeyId" | "SlowDown" | "ServiceUnavailable") + ) => + { + Ok(false) + } + Err(error) => Err(error), + } } }, label, @@ -129,6 +140,7 @@ async fn wait_for_put(client: &Client, bucket: &str, key: &str, body: Vec, l } async fn wait_for_bytes(client: &Client, bucket: &str, key: &str, expected: &[u8], label: &str) -> TestResult { + let client = s3_probe_client(client); wait_until( CREDENTIAL_TIMEOUT, || { @@ -137,8 +149,24 @@ async fn wait_for_bytes(client: &Client, bucket: &str, key: &str, expected: &[u8 let key = key.to_string(); let expected = expected.to_vec(); async move { - let got = get_object_bytes(&client, &bucket, &key).await?; - Ok(got == expected) + match get_object_bytes(&client, &bucket, &key).await { + Ok(got) if got == expected => Ok(true), + Ok(got) => Err(format!( + "{label}: object {bucket}/{key} bytes mismatch: expected sha256={} got sha256={}", + sha256_hex(&expected), + sha256_hex(&got) + ) + .into()), + Err(error) + if matches!( + s3_probe_error_code(error.as_ref()), + Some("AccessDenied" | "InvalidAccessKeyId" | "SlowDown" | "ServiceUnavailable") + ) => + { + Ok(false) + } + Err(error) => Err(error), + } } }, label, @@ -146,6 +174,57 @@ async fn wait_for_bytes(client: &Client, bucket: &str, key: &str, expected: &[u8 .await } +#[cfg(test)] +mod retry_tests { + use super::*; + use crate::common::build_test_s3_config; + use crate::fake_s3_target::{FAKE_ACCESS_KEY, FAKE_SECRET_KEY, FakeS3Target, FaultAction, Operation}; + + #[tokio::test] + async fn iam_polls_only_retry_auth_propagation_and_availability() -> TestResult { + let target = FakeS3Target::start().await?; + target.create_bucket("upgrade-retry-probe"); + let client = Client::from_conf(build_test_s3_config( + target.endpoint(), + FAKE_ACCESS_KEY, + FAKE_SECRET_KEY, + None, + "upgrade-retry-test", + )); + + for status in [403, 503] { + target.inject(Operation::PutObject, FaultAction::ResponseStatus(status), 1); + wait_for_put(&client, "upgrade-retry-probe", "key", b"body".to_vec(), "IAM PUT").await?; + assert_eq!(target.take_requests().len(), 2); + + target.inject(Operation::GetObject, FaultAction::ResponseStatus(status), 1); + wait_for_bytes(&client, "upgrade-retry-probe", "key", b"body", "IAM GET").await?; + assert_eq!(target.take_requests().len(), 2); + } + + target.inject(Operation::PutObject, FaultAction::ResponseStatus(500), 1); + let error = wait_for_put(&client, "upgrade-retry-probe", "key", b"body".to_vec(), "IAM PUT") + .await + .expect_err("IAM convergence must not hide InternalError on PUT"); + assert_eq!(s3_probe_error_code(error.as_ref()), Some("InternalError")); + assert_eq!(target.take_requests().len(), 1); + + target.inject(Operation::GetObject, FaultAction::ResponseStatus(500), 1); + let error = wait_for_bytes(&client, "upgrade-retry-probe", "key", b"body", "IAM GET") + .await + .expect_err("IAM convergence must not hide InternalError on GET"); + assert_eq!(s3_probe_error_code(error.as_ref()), Some("InternalError")); + assert_eq!(target.take_requests().len(), 1); + + let error = wait_for_bytes(&client, "upgrade-retry-probe", "key", b"wrong bytes", "IAM GET") + .await + .expect_err("IAM convergence must not hide a payload mismatch"); + assert!(error.to_string().contains("bytes mismatch")); + assert_eq!(target.requests().len(), 1); + Ok(()) + } +} + async fn seed_history_and_iam(dist: &DistCluster) -> TestResult { let history_bucket = unique_bucket("upg-hist"); let versioned_bucket = unique_bucket("upg-ver"); diff --git a/crates/heal/src/heal/mrf_queue.rs b/crates/heal/src/heal/mrf_queue.rs index d706a99c1..94a98154c 100644 --- a/crates/heal/src/heal/mrf_queue.rs +++ b/crates/heal/src/heal/mrf_queue.rs @@ -930,17 +930,38 @@ pub fn spawn_mrf_consumer(manager: Arc) { tracing::info!(target: "rustfs::heal::mrf", "MRF intent consumer started"); } -/// Replay the durable journal into a fresh pending queue and submit whatever -/// it armed. Returns the number of intact intents replayed. Duplicates are -/// merged by the manager's dedup key; the journal is retained whenever replay -/// cannot fully hand off a successor in-memory snapshot (torn tails truncate -/// via the per-record CRC). Public for integration tests; the live consumer -/// invokes this through [`replay_into`] at startup. +/// Replay the journal once for integration tests, publishing a durable +/// successor before dispatching partial writes. Returns the number of intact +/// intents replayed. Ordinary-only admission retries remain in the startup journal; +/// this helper does not run the consumer's retry loop or proof-driven cleanup. +/// The live consumer invokes [`replay_into`] directly at startup. pub async fn replay_journal_once(manager: &Arc) -> usize { let config = MrfConsumerConfig::default(); let mut queue = MrfQueue::new(config.queue_capacity, config.journal_max_bytes); let mut backoff_until: Option = None; - replay_into(manager, &mut queue, &mut backoff_until).await.replayed + let replay = replay_into(manager, &mut queue, &mut backoff_until).await; + if replay.partial_writes.is_empty() { + return replay.replayed; + } + let mut runtime = MrfRuntime { + partial_writes: PartialWrites::default(), + queue, + config, + checkpoint_owner: Uuid::new_v4(), + next_checkpoint_sequence: replay.next_checkpoint_sequence, + new_since_flush: 0, + dirty: false, + journal_on_disk: replay.journal_on_disk, + retain_replay_journal: replay.retain_journal_for_replay, + durable_replay_anchors: replay.durable_replay_anchors, + replay_cleanup: replay.cleanup, + runtime_checkpoint: None, + backoff_until, + }; + runtime.adopt_replayed_partial_writes(replay.partial_writes); + runtime.flush().await; + runtime.dispatch(manager).await; + replay.replayed } struct ReplayOutcome { @@ -1551,15 +1572,15 @@ mod tests { let mut backoff_until = None; let replay = replay_into(&manager, &mut queue, &mut backoff_until).await; assert_eq!(replay.replayed, 1, "W13 committed checkpoint must replay one record"); - assert_eq!(queue.depth(), 0, "W13 replayed record should reach the manager before cleanup"); + assert_eq!(queue.depth(), 0, "W13 durable replay must leave the ordinary queue"); + assert_eq!(replay.partial_writes.len(), 1, "W13 replay must retain its executable durable record"); assert_eq!(replay.durable_replay_anchors.len(), 1, "W13 replay must create a proof anchor"); assert_eq!( manager.operations_snapshot().await.queued_by_source.mrf, - 1, - "W13 replayed work must be visible as MRF manager work" + 0, + "W13 durable replay must wait for its successor checkpoint before dispatch" ); - let anchor = replay.durable_replay_anchors[0].clone(); let mut runtime = MrfRuntime { partial_writes: PartialWrites::default(), queue, @@ -1575,6 +1596,45 @@ mod tests { runtime_checkpoint: None, backoff_until, }; + runtime.adopt_replayed_partial_writes(replay.partial_writes); + assert_eq!(runtime.partial_writes.depth(), 1, "durable replay must have one executable owner"); + assert!( + runtime.durable_replay_anchors.is_empty(), + "adoption must remove the duplicate startup anchor" + ); + runtime.dispatch(&manager).await; + assert_eq!( + manager.operations_snapshot().await.queued_by_source.mrf, + 0, + "unpersisted replay must not enter the manager" + ); + assert!(runtime.flush().await, "publish the durable replay successor before dispatch"); + let successor = snapshot::inspect_local_committed_snapshot(runtime.config.journal_max_bytes) + .await + .expect("inspect durable replay successor") + .expect("durable replay must publish a committed successor"); + assert_eq!((successor.owner(), successor.sequence()), (runtime.checkpoint_owner, 12)); + assert_eq!(runtime.runtime_checkpoint, Some((runtime.checkpoint_owner, 12))); + assert!( + env.disk_paths.iter().all(|path| { + [".heal-mrf-commit.0.bin", ".heal-mrf-commit.1.bin"] + .iter() + .all(|manifest| path.join(".rustfs.sys").join(manifest).exists()) + }), + "both startup and successor checkpoints must remain before proof" + ); + runtime.dispatch(&manager).await; + assert_eq!( + manager.operations_snapshot().await.queued_by_source.mrf, + 1, + "the checkpointed replay must be visible as one MRF manager request" + ); + let anchor = runtime + .partial_writes + .anchors() + .next() + .expect("dispatched replay has a proof anchor") + .clone(); let retained_before_proof = runtime.retained_replay_journal(); assert!(retained_before_proof, "W13 proof anchor must retain replay checkpoint before proof"); assert!( @@ -2469,12 +2529,13 @@ mod tests { let mut backoff_until = None; let replay = replay_into(&manager, &mut queue, &mut backoff_until).await; assert_eq!(replay.replayed, 1, "the committed replay checkpoint must decode one record"); - assert_eq!(queue.depth(), 0, "the replayed record must be admitted before cleanup is considered"); + assert_eq!(queue.depth(), 0, "durable replay must leave the ordinary queue"); + assert_eq!(replay.partial_writes.len(), 1, "replay must retain its executable durable record"); assert!(backoff_until.is_none(), "the accepted replay must not arm admission backoff"); assert_eq!( manager.operations_snapshot().await.queued_by_source.mrf, - 1, - "the replayed record must be visible as an MRF manager request" + 0, + "durable replay must wait for its successor checkpoint before dispatch" ); assert!( replay.journal_on_disk, @@ -2498,7 +2559,6 @@ mod tests { "cleanup must remember the committed checkpoint generation read at startup" ); - let anchor = replay.durable_replay_anchors[0].clone(); let mut runtime = MrfRuntime { partial_writes: PartialWrites::default(), queue, @@ -2514,6 +2574,45 @@ mod tests { runtime_checkpoint: None, backoff_until, }; + runtime.adopt_replayed_partial_writes(replay.partial_writes); + assert_eq!(runtime.partial_writes.depth(), 1, "durable replay must have one executable owner"); + assert!( + runtime.durable_replay_anchors.is_empty(), + "adoption must remove the duplicate startup anchor" + ); + runtime.dispatch(&manager).await; + assert_eq!( + manager.operations_snapshot().await.queued_by_source.mrf, + 0, + "unpersisted replay must not enter the manager" + ); + assert!(runtime.flush().await, "publish the durable replay successor before dispatch"); + let successor = snapshot::inspect_local_committed_snapshot(runtime.config.journal_max_bytes) + .await + .expect("inspect durable replay successor") + .expect("durable replay must publish a committed successor"); + assert_eq!((successor.owner(), successor.sequence()), (runtime.checkpoint_owner, 12)); + assert_eq!(runtime.runtime_checkpoint, Some((runtime.checkpoint_owner, 12))); + assert!( + env.disk_paths.iter().all(|path| { + [".heal-mrf-commit.0.bin", ".heal-mrf-commit.1.bin"] + .iter() + .all(|manifest| path.join(".rustfs.sys").join(manifest).exists()) + }), + "both startup and successor checkpoints must remain before proof" + ); + runtime.dispatch(&manager).await; + assert_eq!( + manager.operations_snapshot().await.queued_by_source.mrf, + 1, + "the checkpointed replay must be visible as one MRF manager request" + ); + let anchor = runtime + .partial_writes + .anchors() + .next() + .expect("dispatched replay has a proof anchor") + .clone(); assert!(runtime.retained_replay_journal(), "proof-bearing replay anchors must block idle cleanup"); assert!( snapshot::inspect_local_committed_snapshot(runtime.config.journal_max_bytes) @@ -2541,7 +2640,7 @@ mod tests { ); assert!( runtime.delete_idle_recovery_anchors().await, - "idle cleanup must delete the proof-discharged committed replay checkpoint" + "idle cleanup must delete both proof-discharged checkpoint generations" ); runtime.journal_on_disk = false; assert!( @@ -2549,8 +2648,10 @@ mod tests { .await .expect("inspect committed checkpoints after proof cleanup") .is_none(), - "the committed replay checkpoint must be gone after proof-driven cleanup" + "both committed checkpoint generations must be gone after proof-driven cleanup" ); + assert_eq!(read_journal(MRF_SCOPED_JOURNAL_PATH).await, None); + assert_eq!(read_journal(MRF_JOURNAL_PATH).await, None); let restart_manager = Arc::new(HealManager::new( storage, diff --git a/crates/heal/tests/mrf_pipeline_test.rs b/crates/heal/tests/mrf_pipeline_test.rs index 6f0e8081c..7cbf3d331 100644 --- a/crates/heal/tests/mrf_pipeline_test.rs +++ b/crates/heal/tests/mrf_pipeline_test.rs @@ -60,6 +60,16 @@ async fn heal_env() -> (Vec, Arc) { heal_env_at(None).await } +async fn heal_env_with_bucket(bucket: &str) -> (Vec, Arc) { + let env = rustfs_test_utils::TestECStoreEnv::builder() + .prefix("rustfs_heal_mrf_test") + .build() + .await; + env.make_bucket(bucket, false).await; + let storage: Arc = Arc::new(ECStoreHealStorage::new(env.ecstore.clone())); + (env.disk_paths, storage) +} + async fn heal_env_at(base_dir: Option<&Path>) -> (Vec, Arc) { let mut builder = rustfs_test_utils::TestECStoreEnv::builder().prefix("rustfs_heal_mrf_test"); if let Some(base_dir) = base_dir { @@ -356,7 +366,7 @@ async fn decode_failure_intent_maps_to_urgent_mrf_heal_request() { #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[serial] async fn journal_replay_arms_intents_and_retains_unproven_partial_write_anchor() { - let (disk_paths, storage) = heal_env().await; + let (disk_paths, storage) = heal_env_with_bucket("replay-bucket").await; // The journal reader resolves disks through the process-local disk map; // register the environment's disks the same way server startup does. @@ -387,8 +397,14 @@ async fn journal_replay_arms_intents_and_retains_unproven_partial_write_anchor() assert!( disk_paths .iter() - .all(|path| !Path::new(path).join(META_BUCKET).join(SCOPED_JOURNAL_REL).exists()), - "missing authoritative journal remains absent" + .all(|path| Path::new(path).join(META_BUCKET).join(SCOPED_JOURNAL_REL).exists()), + "partial-write dispatch must first publish its authoritative successor" + ); + let successor = journal_record(3, "replay-bucket", "partial-object", None, 1); + assert!( + journal_matches_on_all_disks(&disk_paths, SCOPED_JOURNAL_REL, &successor) + && committed_checkpoint_matches_on_all_disks(&disk_paths, 1, &successor), + "the committed and authoritative successor must preserve the partial-write identity" ); let snapshot = manager.operations_snapshot().await; @@ -402,7 +418,7 @@ async fn journal_replay_arms_intents_and_retains_unproven_partial_write_anchor() #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[serial] async fn committed_snapshot_replay_takes_precedence_over_stale_legacy_mirror() { - let (disk_paths, storage) = heal_env().await; + let (disk_paths, storage) = heal_env_with_bucket("committed-bucket").await; register_local_disks(&disk_paths, "mrf-committed-replay-test").await; let committed = scoped_journal_record(3, "committed-bucket", "committed-object", Some([9u8; 16]), 0, 0, 0); @@ -534,11 +550,13 @@ async fn authoritative_journal_is_not_merged_with_legacy_mirror() { #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[serial] async fn authoritative_journal_replay_preserves_kind_and_scope_identity() { - let (disk_paths, storage) = heal_env().await; + let (disk_paths, storage) = heal_env_with_bucket("identity-bucket").await; register_local_disks(&disk_paths, "mrf-authoritative-identity-test").await; - let mut authoritative = scoped_journal_record(3, "identity-bucket", "same-object", None, 0, 3, 7); - authoritative.extend(scoped_journal_record(3, "identity-bucket", "same-object", None, 0, 3, 8)); + let first_partial = scoped_journal_record(3, "identity-bucket", "same-object", None, 0, 3, 7); + let second_partial = scoped_journal_record(3, "identity-bucket", "same-object", None, 0, 3, 8); + let mut authoritative = first_partial.clone(); + authoritative.extend_from_slice(&second_partial); authoritative.extend(journal_record(2, "identity-bucket", "same-object", None, 0)); authoritative.extend(journal_record(1, "identity-bucket", "same-object", Some([4u8; 16]), 0)); let stale_legacy = journal_record(3, "identity-bucket", "stale-legacy-object", None, 0); @@ -570,11 +588,14 @@ async fn authoritative_journal_replay_preserves_kind_and_scope_identity() { "decode-failure repair must not merge with object repair responsibility" ); assert!( - disk_paths.iter().all(|path| { - Path::new(path).join(META_BUCKET).join(JOURNAL_REL).exists() - && Path::new(path).join(META_BUCKET).join(SCOPED_JOURNAL_REL).exists() - }), - "partial-write responsibilities keep both replay anchors until proof" + journal_contains_on_all_disks(&disk_paths, SCOPED_JOURNAL_REL, &first_partial) + && journal_contains_on_all_disks(&disk_paths, SCOPED_JOURNAL_REL, &second_partial) + && committed_payload_contains_on_all_disks(&disk_paths, &[&first_partial, &second_partial]), + "both scoped partial-write identities must survive in the authoritative and committed successor" + ); + assert!( + journal_matches_on_all_disks(&disk_paths, JOURNAL_REL, &[]), + "the legacy mirror must not misrepresent scoped-only partial-write responsibilities" ); } diff --git a/crates/scanner/src/scanner_folder/tests/mrf_ownership.rs b/crates/scanner/src/scanner_folder/tests/mrf_ownership.rs index ccc77963c..117541ac7 100644 --- a/crates/scanner/src/scanner_folder/tests/mrf_ownership.rs +++ b/crates/scanner/src/scanner_folder/tests/mrf_ownership.rs @@ -21,7 +21,7 @@ use rustfs_common::mrf_channel::{ use rustfs_heal::heal::{ manager::{HealConfig, HealManager}, mrf_queue::spawn_mrf_consumer, - storage::{HealListItem, HealObjectInfo, HealStorageAPI}, + storage::{HealListItem, HealObjectInfo, HealStorageAPI, HealStorageObjectResult}, }; use rustfs_heal_contracts::heal_channel::HealOpts; @@ -239,12 +239,26 @@ impl HealStorageAPI for NoticeStorage { async fn mrf_bucket_incarnation_id(&self, _: &str) -> rustfs_heal::Result> { Ok(Some(self.bucket_incarnation_id)) } + async fn bucket_incarnation_id(&self, _: &str) -> rustfs_heal::Result> { + Ok(Some(self.bucket_incarnation_id)) + } async fn list_buckets(&self) -> rustfs_heal::Result> { Ok(Vec::new()) } async fn object_exists(&self, _: &str, _: &str) -> rustfs_heal::Result { Ok(true) } + async fn heal_object_at_incarnation( + &self, + bucket: &str, + object: &str, + version_id: Option<&str>, + expected: Uuid, + opts: &HealOpts, + ) -> rustfs_heal::Result { + self.validate_bucket_incarnation(bucket, Some(expected)).await?; + self.heal_object_with_receipt(bucket, object, version_id, opts).await + } async fn heal_object( &self, _: &str, @@ -396,8 +410,8 @@ async fn mrf_ownership_manager_completion_preserves_scanner_pending() { if *object == "unknown" { assert_eq!( manager.get_statistics().await.total_objects_healed, - 1, - "legacy healed count is not repair proof" + 0, + "unproved durable repairs must not increment the healed count" ); } assert_eq!( diff --git a/docs/testing/ci-gates.md b/docs/testing/ci-gates.md index 171faf10c..55a5f352e 100644 --- a/docs/testing/ci-gates.md +++ b/docs/testing/ci-gates.md @@ -17,6 +17,8 @@ The `main` ruleset (`6436880`) requires exactly these contexts, with `strict_req Every PR enters `ci.yml`. The `classify-changes` job uses the base revision of `scripts/ci_gate.py` to select a conservative documentation-only path: root Markdown/licenses, `AGENTS.md`, Markdown under `docs/` or `.agents/skills/`, and documentation images. Unknown paths, unavailable Git history, an empty diff, or a missing base policy select the full matrix. Renames include their deleted source path. Documentation-only PRs still run Quick Checks and Typos; the aggregate requires the expensive jobs to be skipped exactly as selected. +PRs changing only Rust sources under `crates/e2e_test/src/` and known E2E selection digests (optionally with those documentation paths) use the `e2e` scope. They keep workspace validation, the server build, E2E smoke, and both S3 lanes; independent workflows retain their path selection. Connect boundary, serial ILM, optional protocol/rio-v2, and real io_uring jobs are skipped because their sources and test inputs are unchanged. Manifests, fixtures, shared nextest configuration, scripts, production sources, and unknown paths still select `full`. A dependency on `e2e_test` from another workspace member or reachable local dependency disables this shortcut. All non-PR events retain their full selection. + `required-checks` runs even after failed or skipped dependencies. `scripts/ci_gate.py verify` rejects missing jobs, unexpected jobs, failure, cancellation, and unexpected skips; optional lanes are required only on their declared events. `Workspace Test and Lint` is the ordinary Rust job, while `Test and Lint` uniquely names the aggregate. New validation jobs must update both its direct dependencies and the script contract. Test this wiring and its failure cases with `python3 scripts/ci_gate.py --self-test`. Verify the live rule before changing merge policy: @@ -28,6 +30,12 @@ gh api repos/rustfs/rustfs/rulesets/6436880 \ The aggregate requires the validation lanes already selected by `ci.yml`; this closes the gap where a failing critical lane left the required workspace check green. Independent workflows remain report-only unless separately required. Before adding a new expensive lane or moving existing PR coverage to a schedule, collect representative execution and regression evidence, establish ownership and a working scheduled replacement, and update this reference with the resulting policy. +Superseded PR attempts are cancelled. A running main validation finishes, with only the newest pending main run retained, so frequent merges cannot continuously cancel the full baseline. Having a `merge_group` trigger does not itself require use of the merge queue. + +Protocol matrix jobs finish independently when a sibling fails. Each executed protocol test step preserves its log and available JUnit under a separate, attempt-specific artifact; a skipped test step cannot upload a cached report. + +The four site-replication state-writer concurrency proofs reserve the available nextest slots in both local and CI profiles. This isolates their durable IO from unrelated test processes while retaining each proof's internal two-writer race, production five-second lock-acquisition limit, assertions, and zero retries. + ## Pull request and merge matrix "Via aggregate" means a wrong result fails the required `Test and Lint` check. "Report-only" means visible and actionable but outside both the required list and aggregate. Budgets are each job's `timeout-minutes` in the named workflow and are not copied here. @@ -35,18 +43,18 @@ The aggregate requires the validation lanes already selected by `ci.yml`; this c | Event | Check name | Workflow / job | Merge status | Reproduce | |---|---|---|---|---| | PR, non-doc change | `Quick Checks` | `ci.yml` `quick-checks` | Required | `make pre-commit` | -| PR, non-doc change | `Workspace Test and Lint` | `ci.yml` `test-and-lint` | Via aggregate | `cargo clippy --all-targets -- -D warnings`; `cargo nextest run --profile ci --all --exclude e2e_test`; `cargo test --all --doc`; `scripts/check_migration_gate_count.sh` | +| PR, non-doc change | `Workspace Test and Lint` | `ci.yml` `test-and-lint` | Via aggregate | `cargo clippy --all-targets -- -D warnings`; `cargo nextest run --profile ci --all --exclude e2e_test`; `cargo test --all --doc`; migration evidence check described below | | PR, non-doc change | `Typos` | `ci.yml` `typos` | Via aggregate | `typos` | -| PR, non-doc change | `ILM Integration (serial)` | `ci.yml` `test-ilm-integration-serial` | Via aggregate | exact command in the job | -| PR, non-doc change | `Test and Lint (rio-v2)`, `Test and Lint (swift)`, `Test and Lint (sftp)` | `ci.yml` `test-and-lint-rio-v2`, `test-and-lint-protocols` | Via aggregate | `cargo nextest run` with the job's `--features` | -| PR, non-doc change | `Connect Short Credential Boundary` | `ci.yml` `connect-short-credential-boundary` | Via aggregate | `cargo test -p rustfs --test connect_registration --features connect-e2e-short-credentials`; `cargo check -p rustfs --release --features connect-e2e-short-credentials` must fail | +| PR, full selection | `ILM Integration (serial)` | `ci.yml` `test-ilm-integration-serial` | Via aggregate | exact command in the job | +| PR, full selection | `Test and Lint (rio-v2)`, `Test and Lint (swift)`, `Test and Lint (sftp)` | `ci.yml` `test-and-lint-rio-v2`, `test-and-lint-protocols` | Via aggregate | `cargo nextest run` with the job's `--features` | +| PR, full selection | `Connect Short Credential Boundary` | `ci.yml` `connect-short-credential-boundary` | Via aggregate | `cargo test -p rustfs --test connect_registration --features connect-e2e-short-credentials`; `cargo check -p rustfs --release --features connect-e2e-short-credentials` must fail | +| PR, full selection | `Offline Enrollment E2E Root Boundary` | `ci.yml` `offline-enrollment-e2e` | Via aggregate | `scripts/check_offline_enrollment_e2e.sh` | | PR, non-doc change | `Build RustFS Debug Binary` | `ci.yml` `build-rustfs-debug-binary` | Via aggregate; prerequisite for black-box jobs | `python3 scripts/e2e_binary.py build --bins --features e2e-test-hooks` (binary plus its `rustfs.e2e.json` sidecar) | -| PR, non-doc change | `io_uring Integration (real)` | `ci.yml` `uring-integration` | Via aggregate | `cargo test -p rustfs-ecstore --lib uring_ -- --test-threads=1 --nocapture` | +| PR, full selection | `io_uring Integration (real)` | `ci.yml` `uring-integration` | Via aggregate | `cargo test -p rustfs-ecstore --lib uring_ -- --test-threads=1 --nocapture` | | PR, non-doc change | `End-to-End Tests` | `ci.yml` `e2e-tests` | Via aggregate | `python3 scripts/e2e_binary.py run --features e2e-test-hooks -- cargo nextest run --profile e2e-smoke -p e2e_test`, then `./scripts/e2e-run.sh ./target/debug/rustfs ` under the same wrapper; membership guards `scripts/check_test_wiring.py --check-profile e2e-smoke ` and `scripts/check_security_smoke_count.sh check ` | | PR, non-doc change | `S3 Implemented Tests` | `ci.yml` `s3-implemented-tests` | Via aggregate | build `rustfs`, then `scripts/s3-tests/run.sh` with the job's `DEPLOY_MODE` / `TEST_MODE` / `MAXFAIL` env | | PR, non-doc change | `S3 Lifecycle Behavior Tests` | `ci.yml` `s3-lifecycle-behavior-tests` | Via aggregate | `scripts/s3-tests/run.sh` with the job's accelerated-scanner env | | PR touching `paths` in `audit.yml` | `Cargo Deny`, `Workflow Pin Report`, `Dependency Review` | `audit.yml` `cargo-deny`, `workflow-pin-report`, `dependency-review` | Report-only | `cargo deny check`; `scripts/security/check_workflow_pins.sh` | -| PR touching `paths` in `architecture-migration-rules.yml` | `Architecture Migration Rules` | `architecture-migration-rules.yml` `architecture-migration-rules` | Report-only | `scripts/check_architecture_migration_rules.sh` | | PR touching `paths` in `nix.yml` | `Nix Build & Check` | `nix.yml` `nix-validation` | Report-only | `nix flake check` | | PR touching `paths` in `fuzz.yml` | `Build Fuzz Harness`, `Smoke / ` | `fuzz.yml` `fuzz-build`, `pr-fuzz-smoke` | Report-only | `MAX_TOTAL_TIME=60 ./scripts/fuzz/run.sh` | | PR touching `paths` in `windows-filesystem.yml` | `Rename Safety` | `windows-filesystem.yml` `rename-safety` | Report-only | the `cargo test -p rustfs-ecstore --lib ` commands in the job, on Windows | @@ -59,6 +67,8 @@ The aggregate requires the validation lanes already selected by `ci.yml`; this c e2e filters live in `.config/nextest.toml`; extend a profile instead of adding a second selector. Before a profile runs, `scripts/check_test_wiring.py` compares its listing to the committed digest in `.config/e2e--selection.txt`, so a silent test drop fails closed. +Architecture migration rules run once in the required Quick Checks job. Migration proofs reuse the successful workspace run: `scripts/check_migration_gate_count.sh evidence ` checks the unchanged name selection and committed floor, and requires exactly one successful execution without retries for each proof. Missing, filtered, skipped, failed, or duplicate evidence fails the gate. Local `check` and `run` modes remain available. Upgrade compatibility builds one current server and shares its binary and provenance sidecar across the seven cases; each case independently verifies the pinned previous release checksum and runs its existing test. + Scanner usage and heal rebuild coverage are intentionally split by risk and cost. `data_usage_test` runs in the PR `e2e-smoke` lane so changes that affect authoritative scanner usage publication, quota-visible usage, or admin usage diff --git a/rustfs/src/connect/diagnostics/trace_runtime.rs b/rustfs/src/connect/diagnostics/trace_runtime.rs index f830f03c6..9417bb76e 100644 --- a/rustfs/src/connect/diagnostics/trace_runtime.rs +++ b/rustfs/src/connect/diagnostics/trace_runtime.rs @@ -2410,13 +2410,11 @@ mod tests { rustfs_credentials::set_global_rpc_secret("top-rpc-local-test-secret".to_owned()).unwrap(); let runtime = spawn_local_trace_capture_runtime(std::path::Path::new(&state), &stop).unwrap(); let emit = async { - tokio::time::timeout(Duration::from_secs(10), async { - while telemetry_trace_subscriber_count() == 0 { - tokio::task::yield_now().await; - } - }) - .await - .expect("top.rpc service subscriber"); + // Executable hashing precedes subscription. The parent's bounded + // request and stdin cancellation govern this readiness wait. + while telemetry_trace_subscriber_count() == 0 { + tokio::task::yield_now().await; + } let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let addr = listener.local_addr().unwrap(); let server = tokio::spawn(async move { @@ -2552,13 +2550,11 @@ mod tests { tokio::runtime::Runtime::new().unwrap().block_on(async { let runtime = spawn_local_trace_capture_runtime(std::path::Path::new(&state), &stop).unwrap(); let emit = async { - tokio::time::timeout(Duration::from_secs(10), async { - while telemetry_trace_subscriber_count() == 0 { - tokio::task::yield_now().await; - } - }) - .await - .expect("top.api service subscriber"); + // Executable hashing precedes subscription. The parent's bounded + // request and stdin cancellation govern this readiness wait. + while telemetry_trace_subscriber_count() == 0 { + tokio::task::yield_now().await; + } for (operation, status) in [ (S3Operation::GetObject, 200), (S3Operation::GetObject, 503), diff --git a/rustfs/tests/connect_top_net.rs b/rustfs/tests/connect_top_net.rs index 42ca6b77d..294cd3642 100644 --- a/rustfs/tests/connect_top_net.rs +++ b/rustfs/tests/connect_top_net.rs @@ -345,9 +345,14 @@ fn local_top_export_is_private_no_clobber_cancel_safe_and_rejects_forged_artifac } #[test] -fn production_cli_exports_top_net_and_fails_closed_for_unavailable_unsupported_and_invalid_runs() { +fn production_cli_exports_top_net_and_fails_closed_for_unavailable_and_invalid_runs() { let directory = tempfile::tempdir().expect("CLI directory"); let state = directory.path().join("state"); + let offline_identity = rustfs::connect::OfflineKeyStore::new(&state) + .load_or_create() + .expect("offline identity"); + let offline_key_id = + hex_simd::encode_to_string(Sha256::digest(offline_identity.public_key_der()), hex_simd::AsciiCase::Lower); let identity = rustfs::connect::IdentityStore::new(state.join("identity")) .load_or_create() .expect("enrolled identity"); @@ -412,26 +417,32 @@ fn production_cli_exports_top_net_and_fails_closed_for_unavailable_unsupported_a .verify(&signed, &signature) .expect("valid ES256 signature"); - let locks_output = directory.path().join("locks.zip"); - let locks = top_command("locks", &state, &locks_output, "019e3ae0-0000-7000-8000-000000000031", 1, true) - .output() - .expect("run top.locks outside the server process"); - assert!(!locks.status.success()); - let stdout = String::from_utf8(locks.stdout).expect("UTF-8 stdout"); - assert!(stdout.contains(r#""outcome":"FAILED""#)); - assert!(stdout.contains(r#""reasonCode":"SOURCE_UNAVAILABLE""#)); - assert!(!locks_output.exists()); - - for (index, tool) in ["api", "rpc"].into_iter().enumerate() { + // Service-backed captures require an explicit offline identity pin before + // connecting to the server or emitting an export. + for (index, tool) in ["locks", "api", "rpc", "disk"].into_iter().enumerate() { let output = directory.path().join(format!("{tool}.zip")); let artifact_uid = format!("019e3ae0-0000-7000-8000-00000000003{}", index + 1); - let run = top_command(tool, &state, &output, &artifact_uid, 1, true) - .output() - .expect("run unsupported top command"); - assert!(!run.status.success()); + let mut command = top_command(tool, &state, &output, &artifact_uid, 1, true); + let run = command.output().expect("run top command without an offline identity pin"); + assert!(!run.status.success(), "{tool} must reject a missing offline identity pin"); let stdout = String::from_utf8(run.stdout).expect("UTF-8 stdout"); - assert!(stdout.contains("\"outcome\":\"UNSUPPORTED\"")); - assert!(stdout.contains("\"reasonCode\":\"UNSUPPORTED_TOOL\"")); + let stderr = String::from_utf8(run.stderr).expect("UTF-8 stderr"); + assert!(stdout.is_empty(), "{tool} stdout: {stdout}"); + assert!( + stderr.contains("--offline-key-id must select an existing offline identity"), + "{tool} stderr: {stderr}" + ); + assert!(!output.exists()); + + let unavailable = command + .args(["--offline-key-id", &offline_key_id]) + .output() + .expect("run service capture without the server"); + assert!(!unavailable.status.success(), "{tool} must reject an unavailable server"); + let stdout = String::from_utf8(unavailable.stdout).expect("UTF-8 stdout"); + let stderr = String::from_utf8(unavailable.stderr).expect("UTF-8 stderr"); + assert!(stdout.is_empty(), "{tool} stdout: {stdout}"); + assert!(stderr.contains("telemetry server runtime is unavailable"), "{tool} stderr: {stderr}"); assert!(!output.exists()); } diff --git a/scripts/check_architecture_migration_rules.sh b/scripts/check_architecture_migration_rules.sh index 7b6a2e3af..8f3983269 100755 --- a/scripts/check_architecture_migration_rules.sh +++ b/scripts/check_architecture_migration_rules.sh @@ -6,7 +6,7 @@ # #1052 — all closed) and now keep the resulting boundaries from rotting # (facade bypasses, compat-shim resurrection, owner-module drift). Closed # migration issues are NOT a reason to retire this script or its pins. -# Runs in ci.yml Quick Checks and .github/workflows/architecture-migration-rules.yml. +# Runs in ci.yml Quick Checks for every pull request. set -euo pipefail diff --git a/scripts/check_migration_gate_count.sh b/scripts/check_migration_gate_count.sh index 7ee471ed1..c3cc3dd43 100755 --- a/scripts/check_migration_gate_count.sh +++ b/scripts/check_migration_gate_count.sh @@ -4,8 +4,8 @@ # ci.yml's migration gate selects tests BY NAME SUBSTRING. A rename that drops # a test out of the filter silently thins the gate — potentially to zero — # without any CI signal (this is how the layered gate died when #878 closed). -# This script owns the filter expression so the count check and the test run -# cannot drift, and fails fast when the number of selected tests drops below +# The shared filter keeps the count check, evidence check, and test run aligned. +# Fail fast when the number of selected tests drops below # the committed floor in .config/migration-gate-floor.txt. # # NAMING CONVENTION DEPENDENCY: migration-gate tests are matched by these @@ -28,24 +28,34 @@ # scripts/check_migration_gate_count.sh # count check + run the gate # scripts/check_migration_gate_count.sh check # count check only # scripts/check_migration_gate_count.sh run # run the gate only +# scripts/check_migration_gate_count.sh evidence CORE_LISTING JUNIT +# Verify the same selection already passed in the workspace test run. set -euo pipefail cd "$(dirname "$0")/.." -# Single source of truth for the migration-gate target and filter. The +# Single source of truth for the migration-gate target and shared filter. The # test-util feature activates migration-critical tests that otherwise leave # their shared fixtures compiled but unused. ci.yml must invoke this script # instead of inlining either selection. MIGRATION_GATE_TARGET_ARGS=(-p rustfs-ecstore --lib --features test-util) -MIGRATION_GATE_FILTER='test(data_movement) or test(rebalance) or test(decommission) or test(source_cleanup) or test(delete_marker)' +MIGRATION_GATE_FILTER="$(python3 scripts/check_migration_gate_evidence.py --filter)" FLOOR_FILE=".config/migration-gate-floor.txt" mode="${1:-all}" case "$mode" in all | check | run) ;; + evidence) + if [[ "$#" -ne 3 ]]; then + echo "usage: $0 evidence CORE_LISTING JUNIT" >&2 + exit 2 + fi + python3 scripts/check_migration_gate_evidence.py "$2" "$3" "$FLOOR_FILE" + exit + ;; *) - echo "usage: $0 [all|check|run]" >&2 + echo "usage: $0 [all|check|run|evidence CORE_LISTING JUNIT]" >&2 exit 2 ;; esac diff --git a/scripts/check_migration_gate_evidence.py b/scripts/check_migration_gate_evidence.py new file mode 100644 index 000000000..f1f5353f3 --- /dev/null +++ b/scripts/check_migration_gate_evidence.py @@ -0,0 +1,90 @@ +#!/usr/bin/env python3 +"""Verify migration proofs already passed in the workspace nextest run.""" + +import json +from pathlib import Path +import sys +import xml.etree.ElementTree as ET + + +SUITE = "rustfs-ecstore" +NAME_PARTS = ("data_movement", "rebalance", "decommission", "source_cleanup", "delete_marker") + + +def unique_object(pairs): + result = {} + for key, value in pairs: + if key in result: + raise ValueError(f"duplicate JSON key: {key}") + result[key] = value + return result + + +def library_suite(path): + listing = json.loads(Path(path).read_text(), object_pairs_hook=unique_object) + suite = listing["rust-suites"][SUITE] + if any(suite.get(key) != value for key, value in { + "package-name": SUITE, "binary-id": SUITE, "kind": "lib", "status": "listed", + }.items()) or not isinstance(suite.get("testcases"), dict): + raise ValueError(f"{path}: expected the listed {SUITE} library test binary") + return suite["testcases"] + + +def verify(core_path, junit_path, floor_path): + floor_lines = [line.strip() for line in Path(floor_path).read_text().splitlines() + if line.strip() and not line.lstrip().startswith("#")] + if len(floor_lines) != 1 or not floor_lines[0].isdigit() or int(floor_lines[0]) <= 0: + raise ValueError("migration floor must contain one positive integer") + floor = int(floor_lines[0]) + expected = set() + for name, case in library_suite(core_path).items(): + if not any(part in name for part in NAME_PARTS) or case.get("ignored") is True: + continue + if (case.get("kind") != "test" or case.get("ignored") is not False + or case.get("filter-match", {}).get("status") != "matches"): + raise ValueError(f"{name}: migration test is malformed or filtered from the core run") + expected.add(name) + if len(expected) < floor: + raise ValueError(f"migration selection has {len(expected)} tests, below the committed floor {floor}") + + report = ET.parse(junit_path).getroot() + if (report.tag != "testsuites" or int(report.get("tests", "0")) < len(expected) + or report.get("failures") != "0" or report.get("errors") != "0"): + raise ValueError("core JUnit must report a successful nonempty nextest run") + suites = [suite for suite in report.findall("testsuite") if suite.get("name") == SUITE] + if len(suites) != 1: + raise ValueError("core JUnit must contain exactly one migration library suite") + suite_cases = set(suites[0].findall("testcase")) + executions = {} + for case in report.iter("testcase"): + if case.get("classname") == SUITE: + executions.setdefault(case.get("name"), []).append(case) + for name in sorted(expected): + cases = executions.get(name, []) + if len(cases) != 1 or cases[0] not in suite_cases: + raise ValueError(f"{name}: expected exactly one core JUnit execution in the migration library suite") + case = cases[0] + if (case.get("status", "passed") != "passed" + or any(child.tag not in ("system-out", "system-err", "properties") for child in case)): + raise ValueError(f"{name}: migration proof failed, skipped, or required a retry") + return len(expected), floor + + +def main(): + if sys.argv[1:] == ["--filter"]: + print(" or ".join(f"test({part})" for part in NAME_PARTS)) + return 0 + if len(sys.argv) != 4: + print("usage: check_migration_gate_evidence.py CORE_LISTING JUNIT FLOOR | --filter", file=sys.stderr) + return 2 + try: + count, floor = verify(*sys.argv[1:]) + except (OSError, ValueError, TypeError, KeyError, AttributeError, ET.ParseError) as error: + print(f"migration gate evidence failed: {error}", file=sys.stderr) + return 1 + print(f"migration gate evidence OK: {count} proofs passed without retries (floor: {floor})") + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/scripts/ci_gate.py b/scripts/ci_gate.py index 2dc1207f6..3979991f7 100644 --- a/scripts/ci_gate.py +++ b/scripts/ci_gate.py @@ -1,5 +1,5 @@ #!/usr/bin/env python3 -"""Select safe documentation-only CI and verify the complete required job set.""" +"""Select conservative PR scopes and verify the complete required job set.""" from __future__ import annotations import json @@ -9,6 +9,8 @@ import re import subprocess import sys import tempfile +import tomllib +from unittest.mock import patch import unittest ROOT = Path(__file__).resolve().parent.parent @@ -19,6 +21,19 @@ CODE_JOBS = ( "build-rustfs-debug-binary", "uring-integration", "e2e-tests", "s3-implemented-tests", "s3-lifecycle-behavior-tests", ) +# These jobs still exercise the server and its black-box harness when only E2E +# Rust sources change. Production code, manifests and shared test configuration +# always select the full matrix. No production package depends on e2e_test. +E2E_JOBS = ( + "test-and-lint", "build-rustfs-debug-binary", "e2e-tests", + "s3-implemented-tests", "s3-lifecycle-behavior-tests", +) +E2E_SELECTION_FILES = { + f".config/{profile}-selection.txt" for profile in ( + "e2e-smoke", "e2e-full", "e2e-nightly", "e2e-repl-nightly", + "e2e-distributed", "e2e-protocols", "e2e-odm-interop", + ) +} OPTIONAL_JOBS = ("build-rustfs-debug-binary-rio-v2", "e2e-tests-rio-v2", "e2e-full") NON_VALIDATION_JOBS = {"required-checks", "cancel-closed-pr-runs", "alert-on-failure"} @@ -36,6 +51,53 @@ def documentation_path(path: str) -> bool: return path.startswith("docs/") and path.endswith((".png", ".jpg", ".svg")) +def e2e_crate_is_isolated(root: Path) -> bool: + """A new dependency on the harness invalidates the test-only shortcut.""" + target = (root / "crates/e2e_test").resolve() + manifests = {root / "Cargo.toml"} + + def references_harness(value: object, directory: Path) -> bool: + if not isinstance(value, dict): + return False + for key, item in value.items(): + if key in ("dependencies", "dev-dependencies", "build-dependencies") and isinstance(item, dict): + for name, dependency in item.items(): + if name == "e2e_test": + return True + if isinstance(dependency, dict): + if dependency.get("package") == "e2e_test": + return True + if isinstance(dependency.get("path"), str): + dependency_root = (directory / dependency["path"]).resolve() + if dependency_root == target or not dependency_root.is_relative_to(root.resolve()): + return True + # Cargo also includes in-tree path dependencies that + # are not explicitly listed as workspace members. + manifests.add(dependency_root / "Cargo.toml") + if references_harness(item, directory): + return True + return False + + try: + workspace = tomllib.loads((root / "Cargo.toml").read_text()) + for member in workspace["workspace"]["members"]: + members = list(root.glob(member)) + if not members: + return False + manifests.update(directory / "Cargo.toml" for directory in members if directory.resolve() != target) + visited = set() + while manifests: + path = manifests.pop().resolve() + if path in visited: + continue + visited.add(path) + if references_harness(tomllib.loads(path.read_text()), path.parent): + return False + return True + except (OSError, ValueError, KeyError, TypeError): + return False + + def select_mode(event: str, base: str, head: str, root: Path) -> str: if event != "pull_request" or not all(re.fullmatch(r"[0-9a-f]{40}", sha) for sha in (base, head)): return "full" @@ -47,16 +109,24 @@ def select_mode(event: str, base: str, head: str, root: Path) -> str: except (subprocess.CalledProcessError, UnicodeError): return "full" paths = changed.rstrip("\0").split("\0") if changed else [] - return "docs" if paths and all(documentation_path(path) for path in paths) else "full" + if paths and all(documentation_path(path) for path in paths): + return "docs" + if paths and all(documentation_path(path) or path in E2E_SELECTION_FILES or ( + path.startswith("crates/e2e_test/src/") and path.endswith(".rs") + and not any(part in (".", "..") for part in path.split("/")) + and not any(ord(char) < 32 for char in path) + ) for path in paths) and e2e_crate_is_isolated(root): + return "e2e" + return "full" def expected_results(mode: str, event: str, ref: str) -> dict[str, str]: if event not in ("pull_request", "push", "merge_group", "schedule", "workflow_dispatch"): raise ValueError(f"unsupported CI event: {event!r}") - if mode not in ("docs", "full") or (mode == "docs" and event != "pull_request"): + if mode not in ("docs", "e2e", "full") or (mode != "full" and event != "pull_request"): raise ValueError(f"invalid CI selection: {mode!r} for {event!r}") expected = {job: "success" for job in ALWAYS_JOBS} - expected.update({job: "success" if mode == "full" else "skipped" for job in CODE_JOBS}) + expected.update({job: "success" if mode == "full" or (mode == "e2e" and job in E2E_JOBS) else "skipped" for job in CODE_JOBS}) rio = mode == "full" and event in ("schedule", "workflow_dispatch") expected.update({job: "success" if rio else "skipped" for job in OPTIONAL_JOBS[:2]}) full = mode == "full" and (event in ("merge_group", "workflow_dispatch") or (event == "push" and ref == "refs/heads/main")) @@ -158,6 +228,62 @@ def check_workflow(root: Path) -> list[str]: class SelfTests(unittest.TestCase): + def test_e2e_shortcut_rejects_new_direct_renamed_and_target_dependencies(self): + self.assertTrue(e2e_crate_is_isolated(ROOT)) + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + (root / "rustfs").mkdir() + (root / "Cargo.toml").write_text('[workspace]\nmembers = ["rustfs"]\n') + path = root / "rustfs/Cargo.toml" + path.write_text('[package]\nname = "rustfs"\n') + self.assertTrue(e2e_crate_is_isolated(root)) + for dependency in ( + '[dependencies]\ne2e_test = "1"\n', + '[dev-dependencies]\nharness = { package = "e2e_test", version = "1" }\n', + '[target.\'cfg(unix)\'.build-dependencies]\nharness = { path = "../crates/e2e_test" }\n', + ): + path.write_text(dependency) + self.assertFalse(e2e_crate_is_isolated(root), dependency) + path.unlink() + self.assertFalse(e2e_crate_is_isolated(root)) + + def test_e2e_shortcut_traverses_implicit_path_dependencies(self): + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + (root / "rustfs").mkdir() + (root / "helper").mkdir() + (root / "Cargo.toml").write_text('[workspace]\nmembers = ["rustfs"]\n') + (root / "rustfs/Cargo.toml").write_text('[dependencies]\nhelper = { path = "../helper" }\n') + helper = root / "helper/Cargo.toml" + helper.write_text('[package]\nname = "helper"\n') + self.assertTrue(e2e_crate_is_isolated(root)) + helper.write_text('[dependencies]\nharness = { path = "../crates/e2e_test" }\n') + self.assertFalse(e2e_crate_is_isolated(root)) + helper.unlink() + self.assertFalse(e2e_crate_is_isolated(root)) + + def test_e2e_scope_keeps_production_and_shared_configuration_full(self): + cases = ( + (["crates/e2e_test/src/distributed/harness.rs"], "e2e"), + (["README.md", "crates/e2e_test/src/common.rs"], "e2e"), + (["crates/e2e_test/src/distributed/harness.rs", ".config/e2e-distributed-selection.txt"], "e2e"), + ([".config/e2e-full-selection.txt"], "e2e"), + ([".config/unrecognized-selection.txt"], "full"), + (["crates/e2e_test/src/common.rs", "crates/ecstore/src/lib.rs"], "full"), + (["crates/e2e_test/src/common.rs", "crates/e2e_test/Cargo.toml"], "full"), + (["crates/e2e_test/build.rs"], "full"), + (["crates/e2e_test/src/fixture.json"], "full"), + ([".config/nextest.toml"], "full"), + (["Cargo.lock"], "full"), + (["scripts/e2e_binary.py"], "full"), + (["crates/e2e_test/src/../Cargo.toml.rs"], "full"), + (["crates/e2e_test/src/unusual\nname.rs"], "full"), + ([], "full"), + ) + for paths, expected in cases: + with self.subTest(paths=paths), patch("subprocess.check_output", return_value=("\0".join(paths) + "\0").encode()): + self.assertEqual(select_mode("pull_request", "a" * 40, "b" * 40, ROOT), expected) + def test_documentation_paths_do_not_hide_build_or_fixture_changes(self): for path in ("README.md", "AGENTS.md", "crates/utils/AGENTS.md", "docs/testing/README.md", "docs/diagram.svg", ".agents/skills/example/SKILL.md"): self.assertTrue(documentation_path(path), path) @@ -192,15 +318,19 @@ class SelfTests(unittest.TestCase): self.assertEqual({job for job, state in ordinary.items() if state == "skipped"}, set(OPTIONAL_JOBS)) docs = expected_results("docs", "pull_request", "refs/pull/1/merge") self.assertEqual({job for job, state in docs.items() if state == "success"}, set(ALWAYS_JOBS)) + e2e = expected_results("e2e", "pull_request", "refs/pull/1/merge") + self.assertEqual({job for job, state in e2e.items() if state == "success"}, set(ALWAYS_JOBS + E2E_JOBS)) for event in ("schedule", "workflow_dispatch", "merge_group", "push"): result = expected_results("full", event, "refs/heads/main") self.assertEqual(result["e2e-full"], "skipped" if event == "schedule" else "success") self.assertEqual(result["e2e-tests-rio-v2"], "success" if event in ("schedule", "workflow_dispatch") else "skipped") with self.assertRaises(ValueError): expected_results("docs", event, "refs/heads/main") + with self.assertRaises(ValueError): + expected_results("e2e", event, "refs/heads/main") def test_every_wrong_result_missing_job_or_selection_fails_closed(self): - for mode, event in (("full", "pull_request"), ("docs", "pull_request"), ("full", "schedule"), ("full", "workflow_dispatch"), ("full", "merge_group")): + for mode, event in (("full", "pull_request"), ("docs", "pull_request"), ("e2e", "pull_request"), ("full", "schedule"), ("full", "workflow_dispatch"), ("full", "merge_group")): good = {job: {"result": value} for job, value in expected_results(mode, event, "refs/heads/main").items()} good["classify-changes"]["outputs"] = {"mode": mode} self.assertEqual(verify_results(good, event, "refs/heads/main"), []) @@ -324,6 +454,9 @@ class SelfTests(unittest.TestCase): for event, changed, base_sha, available, broken, expected in ( ("pull_request", "README.md", "b" * 40, True, False, "docs"), ("pull_request", "src/server.rs", "b" * 40, True, False, "full"), + ("pull_request", "crates/e2e_test/src/distributed/harness.rs", "b" * 40, True, False, "e2e"), + ("pull_request", "crates/e2e_test/Cargo.toml", "b" * 40, True, False, "full"), + ("pull_request", ".config/nextest.toml", "b" * 40, True, False, "full"), ("pull_request", "README.md", "b" * 40, False, False, "full"), ("merge_group", "README.md", "b" * 40, False, False, "full"), ("pull_request", "README.md", "b" * 40, True, True, None), @@ -333,6 +466,7 @@ class SelfTests(unittest.TestCase): root = Path(directory) (root / "scripts").mkdir() (root / "scripts/ci_gate.py").write_text("raise SystemExit(71)\n") + (root / "Cargo.toml").write_text('[workspace]\nmembers = []\n') (root / "python3").symlink_to(sys.executable) base = root / "base-policy.py" base.write_text("raise SystemExit(29)\n" if broken else Path(__file__).read_text()) diff --git a/scripts/s3-tests/run.sh b/scripts/s3-tests/run.sh index ead73e80b..5204ada00 100755 --- a/scripts/s3-tests/run.sh +++ b/scripts/s3-tests/run.sh @@ -929,7 +929,7 @@ install_python_package() { } if ! command -v awscurl >/dev/null 2>&1; then - install_python_package awscurl || { + install_python_package "awscurl==0.44" || { log_error "Failed to install awscurl" exit 1 } @@ -1021,15 +1021,18 @@ fi cd "${PROJECT_ROOT}/s3-tests" -# Install tox if not available -if ! command -v tox >/dev/null 2>&1; then - install_python_package tox || { - log_error "Failed to install tox" +# Match the weekly compatibility workflow even on runners with an older tox. +TOX_VERSION="$(tox --version 2>/dev/null || true)" +if [[ "${TOX_VERSION%% *}" != "4.60.0" ]]; then + install_python_package "tox==4.60.0" || { + log_error "Failed to install tox 4.60.0" exit 1 } - # Add common Python user bin directories to PATH (same as awscurl) - PYTHON_VERSION=$(python3 -c "import sys; print(f'{sys.version_info.major}.{sys.version_info.minor}')" 2>/dev/null || echo "3.14") - export PATH="$HOME/Library/Python/${PYTHON_VERSION}/bin:$HOME/.local/bin:$PATH" +fi +TOX_VERSION="$(tox --version 2>/dev/null || true)" +if [[ "${TOX_VERSION%% *}" != "4.60.0" ]]; then + log_error "Expected tox 4.60.0, found: ${TOX_VERSION:-unavailable}" + exit 1 fi # Step 9: Run ceph s3-tests @@ -1039,10 +1042,11 @@ mkdir -p "${ARTIFACTS_DIR}" XDIST_ARGS="" if [ "${XDIST}" != "0" ]; then # Add pytest-xdist to requirements.txt so tox installs it inside its virtualenv - grep -qxF "pytest-xdist" requirements.txt || echo "pytest-xdist" >> requirements.txt + grep -qxF "pytest-xdist==3.8.0" requirements.txt || echo "pytest-xdist==3.8.0" >> requirements.txt XDIST_ARGS="-n ${XDIST} --dist=loadgroup" fi -grep -qxF "pytest-timeout" requirements.txt || echo "pytest-timeout" >> requirements.txt +grep -qxF "pytest-timeout==2.4.0" requirements.txt || echo "pytest-timeout==2.4.0" >> requirements.txt +grep -qxF "tox==4.60.0" requirements.txt || echo "tox==4.60.0" >> requirements.txt # Resolve config path (absolute path for tox) if [[ "${S3TESTS_CONF}" = /* ]]; then diff --git a/scripts/s3-tests/test_runner_tools.py b/scripts/s3-tests/test_runner_tools.py new file mode 100644 index 000000000..2a146238f --- /dev/null +++ b/scripts/s3-tests/test_runner_tools.py @@ -0,0 +1,119 @@ +#!/usr/bin/env python3 +"""Exercise S3 harness tool setup without installing packages or starting RustFS.""" + +from __future__ import annotations + +import os +import subprocess +import tempfile +import unittest +from pathlib import Path + + +SOURCE = Path(__file__).with_name("run.sh").read_text() +TOX_SETUP = SOURCE[ + SOURCE.index("# Match the weekly compatibility workflow") : SOURCE.index("# Step 9: Run ceph s3-tests") +] +PLUGIN_SETUP = SOURCE[SOURCE.index('XDIST_ARGS=""') : SOURCE.index("# Resolve config path (absolute path for tox)")] +PATH_SETUP = SOURCE[SOURCE.index("# Ensure user-level Python scripts") : SOURCE.index("# Configuration")] +INSTALLER = SOURCE[SOURCE.index("ensure_python_pip() {") : SOURCE.index('if ! command -v awscurl')] + + +class RunnerToolsTests(unittest.TestCase): + def test_real_pip_installer_finds_new_user_binary_without_uv(self) -> None: + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + tools = root / "tools" + tools.mkdir() + python = tools / "python3" + python.write_text(r'''#!/bin/bash +if [[ "$1" == -c ]]; then + printf '3.12\n' +elif [[ "$*" == '-m pip --version' ]]; then + printf 'pip 25.0\n' +elif [[ "$*" == *'tox==4.60.0'* ]]; then + mkdir -p "$TOOL_TEST_HOME/.local/bin" + printf '#!/bin/sh\nprintf "4.60.0 from user-site\\n"\n' > "$TOOL_TEST_HOME/.local/bin/tox" + chmod +x "$TOOL_TEST_HOME/.local/bin/tox" +else + exit 99 +fi +''') + python.chmod(0o755) + awscurl = tools / "awscurl" + awscurl.write_text("#!/bin/sh\nexit 0\n") + awscurl.chmod(0o755) + # Redirect only the extracted home paths into this test's sandbox. + # Keep the real initialization and installer to exercise PATH ordering. + script = (PATH_SETUP + INSTALLER + TOX_SETUP).replace("$HOME", "$TOOL_TEST_HOME") + script = 'log_error() { printf "%s\\n" "$*" >&2; }\n' + script + script += 'command -v tox\n' + result = subprocess.run( + ["/bin/bash", "-euo", "pipefail", "-c", script], + env={**os.environ, "PATH": f"{tools}:/usr/bin:/bin", "TOOL_TEST_HOME": str(root)}, + capture_output=True, + text=True, + ) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertEqual(result.stdout.strip(), str(root / ".local/bin/tox")) + + def test_tox_version_is_enforced_before_collection(self) -> None: + stubs = r''' +tox() { + [[ "$TOOL_TEST_VERSION" != missing ]] || return 127 + printf '%s from /runner/tox\n' "$TOOL_TEST_VERSION" +} +install_python_package() { + printf 'INSTALL %s\n' "$1" + [[ "$TOOL_TEST_INSTALL" != fail ]] || return 1 + if [[ "$TOOL_TEST_INSTALL" != shadowed ]]; then + TOOL_TEST_VERSION=4.60.0 + fi +} +log_error() { printf '%s\n' "$*" >&2; } +''' + for version, install, expected, installs in ( + ("4.60.0", "ok", 0, 0), + ("4.59.0", "ok", 0, 1), + ("missing", "ok", 0, 1), + ("4.59.0", "fail", 1, 1), + ("4.59.0", "shadowed", 1, 1), + ): + with self.subTest(version=version, install=install): + result = subprocess.run( + ["bash", "-euo", "pipefail", "-c", stubs + TOX_SETUP + "printf 'COLLECT\n'"], + env={**os.environ, "TOOL_TEST_VERSION": version, "TOOL_TEST_INSTALL": install}, + capture_output=True, + text=True, + ) + self.assertEqual(result.returncode, expected, result.stderr) + self.assertEqual(result.stdout.count("INSTALL tox==4.60.0"), installs) + self.assertEqual("COLLECT" in result.stdout, expected == 0) + if install == "shadowed": + self.assertIn("Expected tox 4.60.0", result.stderr) + + def test_plugin_pins_preserve_serial_and_parallel_selection(self) -> None: + for workers in ("0", "2"): + with self.subTest(workers=workers), tempfile.TemporaryDirectory() as directory: + requirements = Path(directory) / "requirements.txt" + requirements.write_text("pytest\ntox\n") + result = subprocess.run( + [ + "bash", "-euo", "pipefail", "-c", + PLUGIN_SETUP + PLUGIN_SETUP + 'printf "%s" "$XDIST_ARGS"', + ], + cwd=directory, + env={**os.environ, "XDIST": workers}, + capture_output=True, + text=True, + ) + self.assertEqual(result.returncode, 0, result.stderr) + dependencies = requirements.read_text().splitlines() + self.assertEqual(dependencies.count("tox==4.60.0"), 1) + self.assertEqual(dependencies.count("pytest-timeout==2.4.0"), 1) + self.assertEqual(dependencies.count("pytest-xdist==3.8.0"), int(workers != "0")) + self.assertEqual(result.stdout, "" if workers == "0" else "-n 2 --dist=loadgroup") + + +if __name__ == "__main__": + unittest.main() diff --git a/scripts/test_migration_gate_evidence.py b/scripts/test_migration_gate_evidence.py new file mode 100644 index 000000000..a8103de56 --- /dev/null +++ b/scripts/test_migration_gate_evidence.py @@ -0,0 +1,170 @@ +#!/usr/bin/env python3 +"""Check that reusing core test evidence cannot silently drop migration proofs.""" + +import copy +import json +import os +from pathlib import Path +import shutil +import subprocess +import sys +import tempfile +import unittest +import xml.etree.ElementTree as ET + +from check_migration_gate_evidence import NAME_PARTS, SUITE, verify + + +class MigrationEvidenceTests(unittest.TestCase): + def setUp(self): + directory = tempfile.TemporaryDirectory() + self.addCleanup(directory.cleanup) + self.root = Path(directory.name) + self.listing = self.root / "core.json" + self.junit = self.root / "junit.xml" + self.floor = self.root / "floor.txt" + self.floor.write_text("# Existing floor\n2\n") + self.names = ("store::rebalance_commits", "object::delete_marker_preserves_version") + self.suite = { + "package-name": SUITE, "binary-id": SUITE, "kind": "lib", "status": "listed", + "testcases": {name: self.runnable_case() for name in self.names}, + } + self.data = {"rust-suites": {SUITE: self.suite}} + self.write_listing() + self.write_junit(self.names) + + @staticmethod + def runnable_case(): + return {"kind": "test", "ignored": False, "filter-match": {"status": "matches"}} + + def write_listing(self): + self.listing.write_text(json.dumps(self.data)) + + def write_junit(self, names): + report = ET.Element("testsuites", tests=str(len(names)), failures="0", errors="0") + suite = ET.SubElement(report, "testsuite", name=SUITE) + for name in names: + ET.SubElement(suite, "testcase", name=name, classname=SUITE) + ET.ElementTree(report).write(self.junit) + + def check(self): + return verify(self.listing, self.junit, self.floor) + + def test_successful_core_evidence_meets_the_existing_floor(self): + self.assertEqual(self.check(), (2, 2)) + + def test_substring_selection_matches_the_canonical_nextest_filter(self): + names = [f"nested::prefix_{part}_suffix" for part in NAME_PARTS] + self.suite["testcases"] = {name: self.runnable_case() for name in names} + for name in ("nested::Rebalance", "nested::rebalancing", "nested::ordinary_test"): + self.suite["testcases"][name] = self.runnable_case() + self.data["rust-suites"]["other-package"] = copy.deepcopy(self.suite) + self.write_listing() + self.write_junit(names) + self.assertEqual(self.check(), (5, 2)) + script = Path(__file__).with_name("check_migration_gate_evidence.py") + result = subprocess.run([sys.executable, str(script), "--filter"], capture_output=True, text=True, check=True) + self.assertEqual(result.stdout.strip(), + "test(data_movement) or test(rebalance) or test(decommission) or test(source_cleanup) or test(delete_marker)") + + def test_ignored_tests_do_not_count_towards_the_floor(self): + self.suite["testcases"]["store::rebalance_ignored"] = dict(self.runnable_case(), ignored=True) + self.write_listing() + self.assertEqual(self.check(), (2, 2)) + self.suite["testcases"][self.names[0]]["ignored"] = True + self.write_listing() + with self.assertRaisesRegex(ValueError, "below the committed floor"): + self.check() + + def test_filtered_migration_test_is_rejected_even_above_the_floor(self): + case = self.runnable_case() + case["filter-match"] = {"status": "mismatch", "reason": "expression"} + self.suite["testcases"]["store::rebalance_filtered"] = case + self.write_listing() + with self.assertRaisesRegex(ValueError, "filtered"): + self.check() + + def test_wrong_library_identity_and_unlisted_suite_are_rejected(self): + for key, value in (("package-name", "impostor"), ("binary-id", "other"), ("kind", "test"), ("status", "skipped")): + with self.subTest(key=key): + bad = copy.deepcopy(self.data) + bad["rust-suites"][SUITE][key] = value + self.listing.write_text(json.dumps(bad)) + with self.assertRaisesRegex(ValueError, "library test binary"): + self.check() + + def test_empty_malformed_and_duplicate_listing_inputs_fail(self): + for value in ("", "[]", "{}", '{"rust-suites":{},"rust-suites":{}}'): + with self.subTest(value=value): + self.listing.write_text(value) + with self.assertRaises((ValueError, KeyError, TypeError)): + self.check() + self.suite["testcases"] = {} + self.write_listing() + with self.assertRaisesRegex(ValueError, "below the committed floor"): + self.check() + + def test_missing_duplicate_and_wrong_junit_test_identity_fail(self): + for fault in ("missing", "duplicate", "wrong-class", "wrong-suite", "duplicate-suite"): + with self.subTest(fault=fault): + self.write_junit(self.names) + report = ET.parse(self.junit) + suite = report.getroot().find("testsuite") + case = suite.find("testcase") + if fault == "missing": + suite.remove(case) + elif fault == "duplicate": + suite.append(copy.deepcopy(case)) + elif fault == "wrong-class": + case.set("classname", "impostor") + elif fault == "wrong-suite": + suite.set("name", "impostor") + else: + report.getroot().append(copy.deepcopy(suite)) + report.write(self.junit) + with self.assertRaises(ValueError): + self.check() + + def test_failed_skipped_or_retried_proofs_are_not_successes(self): + for tag in ("failure", "error", "skipped", "rerunFailure", "rerunError", "flakyFailure", "flakyError"): + with self.subTest(tag=tag): + self.write_junit(self.names) + report = ET.parse(self.junit) + ET.SubElement(report.getroot().find("testsuite/testcase"), tag) + report.write(self.junit) + with self.assertRaisesRegex(ValueError, "failed, skipped, or required a retry"): + self.check() + + def test_nonempty_successful_junit_and_positive_floor_are_required(self): + for content in ("", "", ''): + self.junit.write_text(content) + with self.assertRaises((ValueError, ET.ParseError)): + self.check() + self.write_junit(self.names) + for content in ("", "0", "-1", "2\n3", "invalid"): + self.floor.write_text(content) + with self.assertRaisesRegex(ValueError, "positive integer"): + self.check() + + def test_evidence_shell_mode_never_invokes_cargo(self): + scripts = self.root / "scripts" + scripts.mkdir() + source_dir = Path(__file__).resolve().parent + for name in ("check_migration_gate_count.sh", "check_migration_gate_evidence.py"): + shutil.copy(source_dir / name, scripts / name) + (self.root / ".config").mkdir() + shutil.copy(self.floor, self.root / ".config/migration-gate-floor.txt") + commands = self.root / "commands" + commands.mkdir() + cargo = commands / "cargo" + cargo.write_text("#!/bin/sh\necho unexpected cargo invocation >&2\nexit 99\n") + cargo.chmod(0o755) + env = dict(os.environ, PATH=f"{commands}{os.pathsep}{os.environ['PATH']}") + result = subprocess.run(["bash", str(scripts / "check_migration_gate_count.sh"), "evidence", + str(self.listing), str(self.junit)], env=env, capture_output=True, text=True) + self.assertEqual(result.returncode, 0, result.stderr) + self.assertIn("2 proofs passed without retries", result.stdout) + + +if __name__ == "__main__": + unittest.main()