diff --git a/.config/e2e-smoke-selection.txt b/.config/e2e-smoke-selection.txt index a0b32f596..d3c471f36 100644 --- a/.config/e2e-smoke-selection.txt +++ b/.config/e2e-smoke-selection.txt @@ -1 +1 @@ -sha256=7e233cd838efd2688a9cf11ba23ad60297172d124067e91f59502d10c9633464 +sha256=476cb92aa9e385c9bb509ee5df3e5bd07e8c929730040b60d16791b6c0616d00 diff --git a/.config/nextest.toml b/.config/nextest.toml index 3d786bd11..c7939d244 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -155,11 +155,18 @@ threads-required = "num-test-threads" filter = 'package(rustfs) & test(=connect::diagnostics::trace_runtime::tests::local_runtime_profile_cancellation_waits_for_lease_release)' 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. +# These real IAM/state fixtures keep production lock-acquisition limits while +# peers complete durable IO. Reserve capacity from unrelated test processes; +# the shared fixture's durable_iam_state_ prefix keeps the whole family together. [[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)$/)' +filter = 'package(rustfs) & test(/^admin::handlers::site_replication::tests::durable_iam_state_/)' +threads-required = "num-test-threads" + +# Real target-repair and rebalance-stop fixtures commit durable metadata while +# peers wait on production locks or bounded cancellation checks. Reserve IO +# capacity for each fixture; its internal races and deadlines remain intact. +[[profile.default.overrides]] +filter = 'package(rustfs) & test(/^admin::handlers::(replication::target_repair_tests::|rebalance::rebalance_handler_tests::real_admin_stop_)/)' threads-required = "num-test-threads" # Reserve capacity for these short heartbeat TLS/fsync probes (rustfs#8129). @@ -319,13 +326,6 @@ retries = 2 filter = 'package(rustfs-ecstore) & (test(concurrent_resend_same_part_commits_one_generation) | test(concurrent_config_writes_from_separate_nodes_do_not_lose_writes))' test-group = 'ecstore-serial-flaky' -# QUARANTINE: OPEN rustfs#4690 — walk_dir stall-budget accounting test depends -# on producer/consumer timing windows that stretch past the budget on loaded -# CI runners (regression test for rustfs#4644; failed on a zero-Rust-diff PR). -[[profile.ci.overrides]] -filter = 'package(rustfs-ecstore) & test(walk_dir_does_not_charge_consumer_backpressure_to_the_stall_budget)' -retries = 2 - # Serialize the relocated-pool GET resume regression under the ci profile too # (see the matching default-profile override near the top). No longer a # quarantine: the fixture race (rustfs#6701/rustfs#6703) was fixed by #6707, @@ -384,10 +384,16 @@ threads-required = "num-test-threads" filter = 'package(rustfs) & test(=connect::diagnostics::trace_runtime::tests::local_runtime_profile_cancellation_waits_for_lease_release)' threads-required = "num-test-threads" -# Keep the same state-writer capacity reservation in CI without changing the +# Keep the same IAM/state fixture 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)$/)' +filter = 'package(rustfs) & test(/^admin::handlers::site_replication::tests::durable_iam_state_/)' +threads-required = "num-test-threads" + +# Match the real admin fixture capacity reservation without retries or changes +# to production lock deadlines, cancellation barriers, or assertions. +[[profile.ci.overrides]] +filter = 'package(rustfs) & test(/^admin::handlers::(replication::target_repair_tests::|rebalance::rebalance_handler_tests::real_admin_stop_)/)' threads-required = "num-test-threads" # Reserve capacity for these short heartbeat TLS/fsync probes (rustfs#8129). @@ -536,13 +542,10 @@ fail-fast = false [profile.e2e-smoke.junit] path = "junit.xml" -# The pagination boundary cases can stall when a server/listing regression -# prevents the continuation request from completing. Keep the timeout scoped -# to those known failure modes so legitimate lifecycle/tiering waits retain -# their test-level timing budget. +# These storage-heavy fixtures reserve capacity. Pagination tests bound listing +# requests separately, after their seed objects have been written. [[profile.e2e-smoke.overrides]] filter = 'package(e2e_test) & test(/^list_objects_v2_pagination_test::tests::(test_list_objects_v2_delimiter_small_page_traverses_all|test_list_objects_v2_max_keys_above_limit_returns_token|test_list_objects_v2_maxkeys_above_limit_with_delimiter)$/)' -slow-timeout = { period = "60s", terminate-after = 2, grace-period = "10s" } threads-required = "num-test-threads" [[profile.e2e-smoke.overrides]] diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 90a32f534..f2bbe9b39 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -193,13 +193,10 @@ jobs: # Raised to 3 on main pushes and manual dispatches; PRs keep 2 so the # merge path is untouched while the experiment runs. # - # Dispatch is included because push alone cannot supply the samples: - # this workflow cancels superseded runs on main, and only 4 of the last - # 20 push-triggered Test and Lint jobs reached a terminal state — at - # that rate ten samples would take roughly fifty merges. The - # concurrency group is scoped by event_name, so a dispatched run has - # its own group and is not cancelled by merge traffic, which makes the - # sample collectable on demand rather than by waiting. + # During baseline collection, main pushes cancelled superseded runs: + # only 4 of 20 Test and Lint jobs reached a terminal state. Running + # main validations now finish; dispatch remains an independent group + # for collecting additional samples on demand. # # Baseline over 17 samples at 2: # median nextest/clippy step ratio 1.95, spread 1.85-2.06. The gate-2 @@ -248,11 +245,11 @@ jobs: mkdir -p artifacts/test-and-lint set +e timeout --verbose --signal=TERM --kill-after=30s 15m \ - cargo test --all --doc \ + cargo test --all --exclude e2e_test --doc \ 2>&1 | tee artifacts/test-and-lint/doctest.log status=${PIPESTATUS[0]} { - echo "command=cargo test --all --doc" + echo "command=cargo test --all --exclude e2e_test --doc" echo "exit_status=${status}" echo "finished_at=$(date --utc --iso-8601=seconds)" echo @@ -537,11 +534,15 @@ jobs: cache-save-if: 'false' install-build-packaging-tools: 'false' - - name: Protect Connect test home - run: chmod go-w "$(realpath "$HOME")" + - name: Prepare protocol test state + run: | + rm -f target/nextest/ci/junit.xml + chmod go-w "$(realpath "$HOME")" - name: Run clippy with ${{ matrix.features.name }} run: | + ./scripts/ci/resource_sampler.sh start clippy + trap './scripts/ci/resource_sampler.sh stop' EXIT cargo clippy -p rustfs -p rustfs-protocols --all-targets ${{ matrix.features.flags }} -- -D warnings - name: Run tests with ${{ matrix.features.name }} @@ -551,24 +552,26 @@ jobs: # Keep feature-test linking under the same bounded concurrency as the # main nextest lane; Clippy is metadata-only and needs no such limit. CARGO_BUILD_JOBS: "2" + SAMPLER_INTERVAL_SECS: "15" run: | # --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). mkdir -p artifacts/protocol-tests - rm -f target/nextest/ci/junit.xml + ./scripts/ci/resource_sampler.sh start nextest + trap './scripts/ci/resource_sampler.sh stop' EXIT 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) + if: always() 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 + artifacts/test-and-lint retention-days: 3 if-no-files-found: warn @@ -827,6 +830,10 @@ jobs: python3 ./scripts/check_test_wiring.py --check-profile e2e-smoke "${NEXTEST_LISTING}" ./scripts/check_security_smoke_count.sh check "${NEXTEST_LISTING}" + # Keep E2E documentation coverage with its already-built package graph. + - name: Run E2E documentation tests + run: cargo test -p e2e_test --doc + # PR smoke subset of the in-repo e2e suite (backlog#1149 ci-4). The # profile.e2e-smoke default-filter in .config/nextest.toml is the single # wiring mechanism for e2e tests in CI — extend that filter instead of @@ -836,18 +843,22 @@ jobs: env: NEXTEST_ARCHIVE: ${{ runner.temp }}/rustfs-e2e-smoke.tar.zst RUSTFS_E2E_LOG_DIR: ${{ runner.temp }}/rustfs-e2e-smoke-logs + SAMPLER_INTERVAL_SECS: "15" run: | + ./scripts/ci/resource_sampler.sh start e2e-smoke + trap './scripts/ci/resource_sampler.sh stop' EXIT python3 scripts/e2e_binary.py run --features e2e-test-hooks -- cargo nextest run --profile e2e-smoke --archive-file "${NEXTEST_ARCHIVE}" \ --status-level all --final-status-level all --failure-output final - name: Upload e2e smoke diagnostics - if: failure() + if: always() uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1 with: name: e2e-smoke-diagnostics-${{ github.run_number }} path: | ${{ runner.temp }}/rustfs-e2e-smoke-logs/ ${{ runner.temp }}/rustfs-e2e-smoke-list.json + artifacts/test-and-lint if-no-files-found: warn - name: Upload e2e smoke JUnit report diff --git a/crates/e2e_test/src/list_objects_v2_pagination_test.rs b/crates/e2e_test/src/list_objects_v2_pagination_test.rs index 70b00467b..312200817 100644 --- a/crates/e2e_test/src/list_objects_v2_pagination_test.rs +++ b/crates/e2e_test/src/list_objects_v2_pagination_test.rs @@ -31,8 +31,19 @@ mod tests { use aws_sdk_s3::Client; use aws_sdk_s3::primitives::ByteStream; use std::collections::HashSet; + use std::future::Future; + use std::time::Duration; + use tokio::time::{Instant, error::Elapsed, timeout_at}; use tracing::info; + const PAGINATION_TIMEOUT: Duration = Duration::from_secs(120); + + // Fixture population can require thousands of durable PUTs. Bound the + // listing traversal itself with one deadline shared by every page. + async fn listing_request_before(deadline: Instant, request: impl Future) -> Result { + timeout_at(deadline, request).await + } + /// Helper function to create an S3 client for testing fn create_s3_client(env: &RustFSTestEnvironment) -> Client { env.create_s3_client() @@ -57,6 +68,54 @@ mod tests { } } + #[tokio::test] + async fn test_list_objects_v2_deadline_rejects_unresponsive_peer() { + use tokio::io::AsyncReadExt; + use tokio::net::TcpListener; + use tokio::sync::oneshot; + use tokio::time::timeout; + + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("bind unresponsive listing peer"); + let endpoint = format!("http://{}", listener.local_addr().expect("listing peer address")); + let (received_tx, received_rx) = oneshot::channel(); + let server = tokio::spawn(async move { + let (mut stream, _) = listener.accept().await.expect("accept listing request"); + let mut bytes = [0; 4096]; + assert!(stream.read(&mut bytes).await.expect("read listing request") > 0); + received_tx.send(()).expect("signal received listing request"); + std::future::pending::<()>().await; + drop(stream); + }); + let client = Client::from_conf( + crate::common::build_test_s3_config(&endpoint, "test-access", "test-secret", None, "pagination-deadline") + .to_builder() + .retry_config(aws_sdk_s3::config::retry::RetryConfig::standard().with_max_attempts(1)) + .build(), + ); + let request = client.list_objects_v2().bucket("deadline-test").send(); + tokio::pin!(request); + tokio::select! { + response = &mut request => panic!("unresponsive listing peer unexpectedly returned: {response:?}"), + received = timeout(Duration::from_secs(10), received_rx) => { + received.expect("listing request must reach peer").expect("listing peer must remain alive"); + } + } + + // Expire the same absolute deadline after the real request is observed; + // this tests a stalled response without relying on transport timing. + let result = timeout(Duration::from_secs(1), listing_request_before(Instant::now(), request)) + .await + .expect("expired pagination deadline must stop the in-flight request"); + server.abort(); + let _ = server.await; + assert!( + result.is_err(), + "an unresponsive ListObjectsV2 request must exceed the pagination deadline" + ); + } + /// Test for Issue #2775: continuation forwarding must not /// skip a child directory when the prefix component repeats in the key. #[tokio::test] @@ -514,12 +573,12 @@ mod tests { .expect("Failed to put object"); } - let output = client - .list_objects_v2() - .bucket(bucket) - .max_keys(1001) - .send() + eprintln!("Seeded {object_count} objects in {bucket}; starting ListObjectsV2 pagination"); + let deadline = Instant::now() + PAGINATION_TIMEOUT; + + let output = listing_request_before(deadline, client.list_objects_v2().bucket(bucket).max_keys(1001).send()) .await + .expect("ListObjectsV2 pagination exceeded its deadline after fixture population") .expect("Failed to list objects"); assert_eq!(output.contents().len(), 1000); @@ -535,14 +594,18 @@ mod tests { .expect("NextContinuationToken should be present when capped response is truncated") .to_string(); - let output = client - .list_objects_v2() - .bucket(bucket) - .max_keys(1001) - .continuation_token(next_token) - .send() - .await - .expect("Failed to list objects with continuation token"); + let output = listing_request_before( + deadline, + client + .list_objects_v2() + .bucket(bucket) + .max_keys(1001) + .continuation_token(next_token) + .send(), + ) + .await + .expect("ListObjectsV2 continuation exceeded the shared pagination deadline") + .expect("Failed to list objects with continuation token"); assert_eq!(output.contents().len(), 2); assert!(!output.is_truncated().unwrap_or(false)); @@ -765,6 +828,9 @@ mod tests { } } + eprintln!("Seeded {} objects in {bucket}; starting ListObjectsV2 pagination", all_keys.len()); + let deadline = Instant::now() + PAGINATION_TIMEOUT; + // Paginate with delimiter="/" and max_keys=50 // Visible per page: up to 50 CommonPrefixes let mut listed_keys = Vec::new(); @@ -780,7 +846,10 @@ mod tests { request = request.continuation_token(token); } - let output = request.send().await.expect("Failed to list objects"); + let output = listing_request_before(deadline, request.send()) + .await + .expect("ListObjectsV2 delimiter traversal exceeded the shared pagination deadline") + .expect("Failed to list objects"); last_page_is_truncated = output.is_truncated().unwrap_or(false); for obj in output.contents() { @@ -988,15 +1057,18 @@ mod tests { } } + eprintln!( + "Seeded {} objects in {bucket}; starting ListObjectsV2 pagination", + dir_count * files_per_dir + ); + let deadline = Instant::now() + PAGINATION_TIMEOUT; + // With delimiter: 12 CommonPrefixes visible, all fit within capped 1000 - let output = client - .list_objects_v2() - .bucket(bucket) - .delimiter("/") - .max_keys(2000) - .send() - .await - .expect("Failed to list objects"); + let output = + listing_request_before(deadline, client.list_objects_v2().bucket(bucket).delimiter("/").max_keys(2000).send()) + .await + .expect("ListObjectsV2 delimiter listing exceeded its deadline after fixture population") + .expect("Failed to list objects"); assert_eq!( output.common_prefixes().len(), diff --git a/crates/e2e_test/src/reliant/tiering.rs b/crates/e2e_test/src/reliant/tiering.rs index f013bf2f2..8c2c2a3cd 100644 --- a/crates/e2e_test/src/reliant/tiering.rs +++ b/crates/e2e_test/src/reliant/tiering.rs @@ -1056,7 +1056,17 @@ async fn test_hermetic_transition_restore_failure_expiry_and_retry() -> TestResu cold_client.create_bucket().bucket(TIER_BUCKET).send().await?; let mut hot = RustFSTestEnvironment::new().await?; - start_tier_source(&mut hot, &[("RUSTFS_SCANNER_CYCLE", "1"), ("RUSTFS_ILM_DEBUG_DAY_SECS", "5")]).await?; + // Revisit the restored key every cycle; the default 16-cycle sampling is + // independent of the accelerated lifecycle clock used by this fixture. + start_tier_source( + &mut hot, + &[ + ("RUSTFS_SCANNER_CYCLE", "1"), + ("RUSTFS_DATA_USAGE_UPDATE_DIR_CYCLES", "1"), + ("RUSTFS_ILM_DEBUG_DAY_SECS", "5"), + ], + ) + .await?; let hot_client = hot.create_s3_client(); add_rustfs_tier(&hot, &cold).await?; diff --git a/crates/e2e_test/src/upgrade_compatibility_test.rs b/crates/e2e_test/src/upgrade_compatibility_test.rs index c7b607282..e5c938835 100644 --- a/crates/e2e_test/src/upgrade_compatibility_test.rs +++ b/crates/e2e_test/src/upgrade_compatibility_test.rs @@ -51,6 +51,7 @@ const PLAIN_BUCKET: &str = "upgrade-plain-data"; const VERSIONED_BUCKET: &str = "upgrade-versioned-data"; const MIXED_BUCKET: &str = "upgrade-mixed-version-data"; const MIXED_NODE_COUNT: usize = 4; +const UPGRADE_READINESS_BODY: &[u8] = b"upgrade write readiness"; const MULTIPART_WORKERS: usize = 16; const MULTIPART_UPLOADS_PER_WORKER: usize = 16; // Peers keep a restarted node's drive in Suspect/Returning for roughly @@ -442,6 +443,16 @@ async fn exercise_mixed_cluster( Ok(()) } +fn upgrade_probe_client(client: &Client) -> Client { + Client::from_conf( + client + .config() + .to_builder() + .retry_config(aws_sdk_s3::config::retry::RetryConfig::standard().with_max_attempts(1)) + .build(), + ) +} + async fn wait_for_upgrade_write_readiness(clients: &[Client], phase: &str, budget: Duration) -> TestResult { // ListBuckets can succeed before peers recover a restarted disk. Cluster // health also accepts Returning disks whose write health is still FAULTY. @@ -449,13 +460,7 @@ async fn wait_for_upgrade_write_readiness(clients: &[Client], phase: &str, budge // writes still execute once and retain their original assertions. let deadline = Instant::now() + budget; for (node, client) in clients.iter().enumerate() { - let client = Client::from_conf( - client - .config() - .to_builder() - .retry_config(aws_sdk_s3::config::retry::RetryConfig::standard().with_max_attempts(1)) - .build(), - ); + let client = upgrade_probe_client(client); let key = format!(".upgrade-readiness/{phase}/node-{node}"); let mut last_response = "no response".to_string(); loop { @@ -470,7 +475,7 @@ async fn wait_for_upgrade_write_readiness(clients: &[Client], phase: &str, budge .put_object() .bucket(MIXED_BUCKET) .key(&key) - .body(ByteStream::from_static(b"upgrade write readiness")) + .body(ByteStream::from_static(UPGRADE_READINESS_BODY)) .send(), ) .await; @@ -497,6 +502,45 @@ async fn wait_for_upgrade_write_readiness(clients: &[Client], phase: &str, budge Ok(()) } +async fn prepare_previous_release_baseline(cluster: &mut RustFSTestClusterEnvironment, previous_binary: &Path) -> TestResult { + let clients: Vec<_> = cluster.create_all_clients()?.iter().map(upgrade_probe_client).collect(); + // rc.5 can latch a non-elected node's write fence while its first startup + // waits for pool metadata (#7473). First prove the elected writer can + // persist data, then restart each old process once with three peers still + // readable. This preparation ends before any current binary is started. + wait_for_upgrade_write_readiness(&clients[..1], "previous-seed", LISTING_CONVERGENCE_TIMEOUT).await?; + for node in [1, 2, 3, 0] { + cluster.stop_node(node)?; + cluster.start_node_from_binary(node, previous_binary).await?; + wait_for_upgrade_write_readiness( + std::slice::from_ref(&clients[node]), + &format!("previous-restart-{node}"), + LISTING_CONVERGENCE_TIMEOUT, + ) + .await?; + } + + wait_for_upgrade_write_readiness(&clients, "previous-baseline", LISTING_CONVERGENCE_TIMEOUT).await?; + for (reader, client) in clients.iter().enumerate() { + assert_eq!( + read_object(client, MIXED_BUCKET, ".upgrade-readiness/previous-seed/node-0", None) + .await? + .1, + UPGRADE_READINESS_BODY, + "previous-release node {reader} must retain the seed across its rolling restart" + ); + for writer in 0..clients.len() { + let key = format!(".upgrade-readiness/previous-baseline/node-{writer}"); + assert_eq!( + read_object(client, MIXED_BUCKET, &key, None).await?.1, + UPGRADE_READINESS_BODY, + "previous-release node {reader} must read baseline data from writer {writer}" + ); + } + } + Ok(()) +} + #[cfg(test)] mod upgrade_write_readiness_tests { use super::*; @@ -875,6 +919,7 @@ async fn rolling_upgrade_from_rc2_preserves_mixed_version_contracts() -> TestRes configure_cluster_logs(&mut cluster)?; cluster.start_with_binary(&previous_binary).await?; cluster.create_test_bucket(MIXED_BUCKET).await?; + prepare_previous_release_baseline(&mut cluster, &previous_binary).await?; cluster.stop_node(0)?; cluster.start_node_from_binary(0, ¤t_binary).await?; diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 3cc2aac1a..1c62872f3 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -13428,24 +13428,24 @@ mod test { ); } - /// A writer that stalls on every write, standing in for a slow listing - /// consumer (quorum merge, a lagging peer drive). + /// A writer that advances paused time on every write, standing in for a + /// slow listing consumer without depending on host scheduling. struct SlowWriter { delay: Duration, - sleep: Option>>, + advance: Option + Send>>>, } impl AsyncWrite for SlowWriter { fn poll_write(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &[u8]) -> Poll> { - if self.sleep.is_none() { + if self.advance.is_none() { let delay = self.delay; - self.sleep = Some(Box::pin(tokio::time::sleep(delay))); + self.advance = Some(Box::pin(tokio::time::advance(delay))); } - let sleep = self.sleep.as_mut().expect("sleep was just installed"); - match sleep.as_mut().poll(cx) { + let advance = self.advance.as_mut().expect("clock advance was just installed"); + match advance.as_mut().poll(cx) { Poll::Ready(()) => { - self.sleep = None; + self.advance = None; Poll::Ready(Ok(buf.len())) } Poll::Pending => Poll::Pending, @@ -13480,6 +13480,7 @@ mod test { let endpoint = Endpoint::try_from(dir.path().to_str().expect("temp dir should be utf8")).expect("endpoint should parse"); let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should be created"); + disk.wait_for_startup_cleanup().await; let stall = Duration::from_millis(300); let write_delay = Duration::from_millis(150); @@ -13493,12 +13494,21 @@ mod test { let mut writer = SlowWriter { delay: write_delay, - sleep: None, + advance: None, }; - let started = std::time::Instant::now(); + // A live blocking task inhibits Tokio's automatic clock advance while + // real filesystem reads are pending. Only consumer writes advance time. + let (clock_guard_tx, clock_guard_rx) = std::sync::mpsc::channel::<()>(); + let clock_guard = tokio::task::spawn_blocking(move || clock_guard_rx.recv()); + tokio::time::pause(); + let started = Instant::now(); let result = disk.walk_dir(opts, &mut writer).await; let elapsed = started.elapsed(); + drop(clock_guard_tx); + let _ = clock_guard + .await + .expect("clock guard should exit after its sender is dropped"); assert!(result.is_ok(), "a walk making steady progress must not time out, got {result:?}"); assert!( diff --git a/crates/ecstore/src/ecstore_validation_blackbox.rs b/crates/ecstore/src/ecstore_validation_blackbox.rs index 15d30b81b..3722d86bc 100644 --- a/crates/ecstore/src/ecstore_validation_blackbox.rs +++ b/crates/ecstore/src/ecstore_validation_blackbox.rs @@ -378,8 +378,10 @@ async fn blackbox_heal_requests_preserve_repair_scope() { // admission channel and wait for its receipt before acknowledging the write. // Drive the test receiver alongside the PUT so neither waits on the other. let (_put_dirs, put_store) = crate::bucket::metadata_sys::test_support::isolated_store_over_temp_disks().await; - crate::bucket::metadata_sys::init_bucket_metadata_sys(put_store.clone(), Vec::new()).await; - let put_set = put_store.pools[0].disk_set[0].clone(); + crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&put_store), Vec::new()).await; + let put_set = Arc::clone(&put_store.pools[0].disk_set[0]); + assert_eq!(put_set.set_drive_count, 4); + assert_eq!(put_set.default_parity_count, 2); let put_bucket = "bb-put-partial-convergence"; let put_object = "object.bin"; put_store @@ -389,7 +391,7 @@ async fn blackbox_heal_requests_preserve_repair_scope() { let put_incarnation = put_store .bucket_incarnation_id_from_disk(put_bucket) .await - .expect("PUT bucket should have a persisted incarnation"); + .expect("partial repair must bind the real persisted bucket generation"); let offline_disk = { let mut disks = put_set.disks.write().await; disks[0].take() @@ -423,9 +425,10 @@ async fn blackbox_heal_requests_preserve_repair_scope() { .to_string(); assert_eq!(request.object_version_id.as_deref(), Some(committed_version.as_str())); - assert_eq!(request.expected_bucket_incarnation_id, Some(put_incarnation)); assert_eq!(request.pool_index, Some(0)); assert_eq!(request.set_index, Some(0)); + assert_eq!(request.source, HealRequestSource::Mrf); + assert_eq!(request.expected_bucket_incarnation_id, Some(put_incarnation)); let duplicate_request = tokio::time::timeout(std::time::Duration::from_millis(100), async { loop { diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 720b1e2ab..81ab6cb48 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -3502,7 +3502,10 @@ impl SetDisks { let protect_write = opts.shard_integrity_write_enabled(); let source_bucket_incarnation_id = match opts.expected_bucket_incarnation_id { Some(incarnation_id) => Some(incarnation_id), - None => self.bucket_incarnation_id_from_disk(bucket).await.ok(), + // Internal metadata writes may already hold the pool metadata + // lock; resolving a user bucket identity would re-enter it. + None if !is_meta_bucketname(bucket) => self.bucket_incarnation_id_from_disk(bucket).await.ok(), + None => None, }; if publication_fence.is_none() && opts.data_movement @@ -8077,7 +8080,11 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { } #[tracing::instrument(skip(self))] async fn delete_object_version(&self, bucket: &str, object: &str, fi: &FileInfo, force_del_marker: bool) -> Result<()> { - let source_bucket_incarnation_id = self.bucket_incarnation_id_from_disk(bucket).await.ok(); + let source_bucket_incarnation_id = if is_meta_bucketname(bucket) { + None + } else { + self.bucket_incarnation_id_from_disk(bucket).await.ok() + }; self.delete_object_version_with_purge(bucket, object, fi, force_del_marker, None, source_bucket_incarnation_id) .await } @@ -9277,7 +9284,11 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { #[tracing::instrument(skip(self))] async fn add_partial(&self, bucket: &str, object: &str, version_id: &str) -> Result<()> { - let source_bucket_incarnation_id = self.bucket_incarnation_id_from_disk(bucket).await.ok(); + let source_bucket_incarnation_id = if is_meta_bucketname(bucket) { + None + } else { + self.bucket_incarnation_id_from_disk(bucket).await.ok() + }; self.add_partial_with_source_incarnation(bucket, object, version_id, source_bucket_incarnation_id) .await } @@ -22443,6 +22454,61 @@ mod single_delete_namespace_owner_tests { .expect("seed a complete real object version"); } + #[tokio::test] + #[serial_test::serial] + async fn internal_metadata_mutations_do_not_reenter_pool_metadata() { + let (_dirs, store, _other_store) = + crate::services::rebalance::test_three_pool_stores_with_isolated_node_contexts(None).await; + let set = &store.pools[0].disk_set[0]; + let disks = set.disk_inventory().await.into_iter().flatten().collect::>(); + + for bucket in [RUSTFS_META_BUCKET, crate::disk::MIGRATING_META_BUCKET] { + for disk in &disks { + if let Err(err) = disk.make_volume(bucket).await { + assert_eq!(err, DiskError::VolumeExists, "prepare the internal metadata volume"); + } + } + let version = Uuid::new_v4(); + seed_version(set, bucket, "delete-under-pool-lock", version, b"metadata version").await; + let request = FileInfo { + name: "delete-under-pool-lock".to_string(), + version_id: Some(version), + mod_time: Some(OffsetDateTime::now_utc()), + ..Default::default() + }; + let mut reader = PutObjReader::from_vec(b"metadata under lock".to_vec()); + let opts = ObjectOptions { + write_completion: WriteCompletion::TailDrained, + ..Default::default() + }; + + // Pool metadata transactions retain this write guard while saving + // internal objects. None of these paths may look up a user bucket. + let _pool_meta_guard = store.pool_meta.write().await; + let limit = Duration::from_secs(30); + let (put, delete, partial) = tokio::join!( + tokio::time::timeout(limit, set.put_object(bucket, "put-under-pool-lock", &mut reader, &opts)), + tokio::time::timeout(limit, set.delete_object_version(bucket, &request.name, &request, false)), + tokio::time::timeout(limit, set.add_partial(bucket, "partial-under-pool-lock", "")), + ); + assert!( + matches!(&put, Ok(Ok(_))) && matches!(&delete, Ok(Ok(()))) && matches!(&partial, Ok(Ok(()))), + "internal operations must finish while pool metadata is locked: bucket={bucket}, put={put:?}, delete={delete:?}, partial={partial:?}" + ); + for disk in &disks { + let written = disk + .read_version("", bucket, "put-under-pool-lock", "", &ReadOptions::default()) + .await + .expect("the internal PUT must reach every disk"); + assert_eq!(written.size, b"metadata under lock".len() as i64); + let deleted = disk + .read_version("", bucket, &request.name, &version.to_string(), &ReadOptions::default()) + .await; + assert!(matches!(deleted, Err(DiskError::FileNotFound | DiskError::FileVersionNotFound))); + } + } + } + #[tokio::test] #[serial_test::serial(capacity_dirty_scope)] async fn single_delete_advances_namespace_generation_through_cleanup() { diff --git a/crates/ecstore/src/store/bucket.rs b/crates/ecstore/src/store/bucket.rs index 4695fee4a..1e67d852f 100644 --- a/crates/ecstore/src/store/bucket.rs +++ b/crates/ecstore/src/store/bucket.rs @@ -1872,8 +1872,12 @@ mod tests { assert_eq!(residue.diagnostic_bytes_read, 0); } - #[tokio::test] + #[tokio::test(start_paused = true)] async fn bucket_residue_scan_distinguishes_visible_and_tier_free_xlmeta() { + // Classification keeps the production budget, but real filesystem + // scheduling must not advance the clock before classification finishes. + let (clock_guard_tx, clock_guard_rx) = std::sync::mpsc::channel::<()>(); + let clock_guard = tokio::task::spawn_blocking(move || clock_guard_rx.recv()); let root = tempfile::tempdir().expect("temporary bucket root should be created"); let bucket_path = root.path().join("bucket"); let visible_path = bucket_path.join("visible").join(STORAGE_FORMAT_FILE); @@ -1895,6 +1899,7 @@ mod tests { let visible_scan = scan_metadata_less_residue(&bucket_path) .await .expect("visible xl.meta scan should succeed"); + assert!(!visible_scan.diagnostic_truncated, "{visible_scan:?}"); assert_eq!(visible_scan.xlmeta_blocker, Some(BucketDeleteBlockerKind::VisibleVersion)); tokio::fs::remove_dir_all(bucket_path.join("visible")) @@ -1931,6 +1936,7 @@ mod tests { let free_scan = scan_metadata_less_residue(&bucket_path) .await .expect("free-version xl.meta scan should succeed"); + assert!(!free_scan.diagnostic_truncated, "{free_scan:?}"); assert_eq!(free_scan.xlmeta_blocker, Some(BucketDeleteBlockerKind::TierFreeVersion)); tokio::fs::remove_dir_all(bucket_path.join("free")) @@ -1951,6 +1957,7 @@ mod tests { let exact_limit_scan = scan_metadata_less_residue(&bucket_path) .await .expect("exact-limit xl.meta scan should remain fail closed"); + assert!(!exact_limit_scan.diagnostic_truncated, "{exact_limit_scan:?}"); assert_eq!(exact_limit_scan.xlmeta_blocker, Some(BucketDeleteBlockerKind::UnknownXlMeta)); assert_eq!(exact_limit_scan.diagnostic_bytes_read, BUCKET_DELETE_XLMETA_DIAGNOSTIC_MAX_BYTES); tokio::fs::remove_dir_all(bucket_path.join("exact-limit")) @@ -1971,9 +1978,31 @@ mod tests { let oversized_scan = scan_metadata_less_residue(&bucket_path) .await .expect("oversized xl.meta scan should remain fail closed"); + assert!(!oversized_scan.diagnostic_truncated, "{oversized_scan:?}"); assert_eq!(oversized_scan.xlmeta_blocker, Some(BucketDeleteBlockerKind::UnknownXlMeta)); assert_eq!(oversized_scan.diagnostic_bytes_read, 0); assert!(oversized_scan.diagnostic_bytes_read <= BUCKET_DELETE_XLMETA_DIAGNOSTIC_MAX_BYTES); + + drop(clock_guard_tx); + let _ = clock_guard + .await + .expect("classification clock guard should exit after its sender is dropped"); + + // An exhausted diagnostic budget may return before observing xl.meta. + // It must report truncation rather than inventing a classification. + let first_io_started = Arc::new(AtomicBool::new(false)); + let mut budget = BucketDeleteDiagnosticBudget::new().with_first_io_delay( + BUCKET_DELETE_DIAGNOSTIC_MAX_ELAPSED + Duration::from_millis(100), + first_io_started.clone(), + ); + let truncated = scan_metadata_less_residue_with_budget(&bucket_path, &mut budget) + .await + .expect("a diagnostic timeout should return a partial scan"); + assert!(first_io_started.load(Ordering::SeqCst)); + assert!(truncated.diagnostic_truncated, "{truncated:?}"); + assert_eq!(truncated.xlmeta_blocker, None); + assert_eq!(truncated.entries_scanned, 0); + assert_eq!(truncated.diagnostic_bytes_read, 0); } #[tokio::test] diff --git a/crates/heal/src/heal/mrf_queue.rs b/crates/heal/src/heal/mrf_queue.rs index 956d89bd7..8a0484846 100644 --- a/crates/heal/src/heal/mrf_queue.rs +++ b/crates/heal/src/heal/mrf_queue.rs @@ -3194,25 +3194,25 @@ mod tests { replay_intent.kind = MrfKind::PartialWrite; replay_intent.version_id = None; let replay_payload = encoded_payload(&replay_intent); - let source_incarnation = storage + let source_bucket_incarnation_id = storage .mrf_bucket_incarnation_id(bucket) .await - .expect("read persisted replay source identity") - .expect("replay source bucket has a persisted incarnation"); + .expect("read the producer's bucket incarnation") + .expect("the fixture bucket must have a persisted incarnation"); let lifecycle_limit = config.journal_max_bytes.saturating_mul(4); let lifecycle = encode_mrf_lifecycle_checkpoint( replay_owner, 11, vec![ResponsibilityCheckpoint { - intent_digest: intent_digest(&replay_intent).expect("fixture intent has a canonical digest"), + intent_digest: intent_digest(&replay_intent).expect("encode the replay intent identity"), responsibility_id: Uuid::new_v4(), - source_bucket_incarnation_id: Some(source_incarnation), + source_bucket_incarnation_id: Some(source_bucket_incarnation_id), last_operator_acceptance: None, state: partial_write::ResponsibilityState::Active, }], lifecycle_limit, ) - .expect("encode generation-bound replay responsibility"); + .expect("encode the producer's source-bound responsibility"); snapshot::publish_committed_snapshot_with_companion( &disks, replay_owner, @@ -3222,7 +3222,7 @@ mod tests { Some((&MRF_LIFECYCLE_PATHS, &lifecycle, lifecycle_limit)), ) .await - .expect("publish committed replay checkpoint with its source identity"); + .expect("publish committed replay checkpoint with its producer identity"); let mut queue = MrfQueue::new(config.queue_capacity, config.journal_max_bytes); let mut backoff_until = None; diff --git a/crates/heal/tests/heal_concurrent_delete_test.rs b/crates/heal/tests/heal_concurrent_delete_test.rs index 5e86255c8..61657c4b7 100644 --- a/crates/heal/tests/heal_concurrent_delete_test.rs +++ b/crates/heal/tests/heal_concurrent_delete_test.rs @@ -513,14 +513,24 @@ async fn delete_during_root(case: DeleteCase) { fn run_delete_case(case: DeleteCase) { // Incarnation-bound repairs use a spawned owner; retain the server's stack - // budget while exercising the real EC2+2 encode/decode futures in debug. - let runtime = tokio::runtime::Builder::new_multi_thread() - .worker_threads(4) - .thread_stack_size(8 * 1024 * 1024) - .enable_all() - .build() - .expect("concurrent delete runtime"); - runtime.block_on(delete_during_root(case)); + // budget for the whole scenario, including its foreground DELETE. block_on + // alone would poll that future on libtest's smaller calling-thread stack. + let stack_size = 8 * 1024 * 1024; + std::thread::Builder::new() + .name("heal-concurrent-delete".to_owned()) + .stack_size(stack_size) + .spawn(move || { + let runtime = tokio::runtime::Builder::new_multi_thread() + .worker_threads(4) + .thread_stack_size(stack_size) + .enable_all() + .build() + .expect("concurrent delete runtime"); + runtime.block_on(delete_during_root(case)); + }) + .expect("concurrent delete scenario thread starts") + .join() + .expect("concurrent delete scenario must not panic"); } #[test] diff --git a/crates/heal/tests/mrf_pipeline_test.rs b/crates/heal/tests/mrf_pipeline_test.rs index a93e13bb2..f72a91b98 100644 --- a/crates/heal/tests/mrf_pipeline_test.rs +++ b/crates/heal/tests/mrf_pipeline_test.rs @@ -229,42 +229,10 @@ fn committed_manifest(owner: uuid::Uuid, sequence: u64, payload: &[u8]) -> Vec>(); - let lifecycle = serde_json::to_vec(&serde_json::json!({ - "format_version": 1, - "checkpoint_owner": owner, - "checkpoint_sequence": sequence, - "records": records, - })) - .expect("encode source-bound lifecycle fixture"); - let mut companion = b"RFMFLC01".to_vec(); - companion.push(1); - companion.extend_from_slice(&u64::try_from(lifecycle.len()).expect("fixture length fits").to_le_bytes()); - companion.extend_from_slice(&Sha256::digest(&lifecycle)); - companion.extend_from_slice(&lifecycle); - write_journal_path_to_disks(disk_paths, ".heal-mrf-lifecycle.0.bin", &companion); +fn write_committed_snapshot_to_disks(disk_paths: &[std::path::PathBuf], sequence: u64, payload: &[u8]) { + let manifest = committed_manifest(uuid::Uuid::new_v4(), sequence, payload); write_journal_path_to_disks(disk_paths, COMMITTED_PAYLOAD_REL, payload); - write_journal_path_to_disks(disk_paths, COMMITTED_MANIFEST_REL, &committed_manifest(owner, sequence, payload)); + write_journal_path_to_disks(disk_paths, COMMITTED_MANIFEST_REL, &manifest); } fn journal_exists_on_all_disks(disk_paths: &[std::path::PathBuf], relative_path: &str) -> bool { @@ -344,6 +312,87 @@ fn committed_checkpoint_matches_on_all_disks(disk_paths: &[PathBuf], sequence: u }) } +async fn assert_legacy_responsibilities_parked( + disk_paths: &[PathBuf], + expected_records: &[&[u8]], + observed_bucket_incarnation_id: uuid::Uuid, +) { + assert!(!observed_bucket_incarnation_id.is_nil(), "the observed bucket generation must be real"); + let checkpoint = mrf_queue::snapshot::inspect_local_committed_snapshot(rustfs_config::DEFAULT_HEAL_MRF_JOURNAL_MAX_BYTES) + .await + .expect("inspect the committed legacy responsibility checkpoint") + .expect("parked responsibilities must retain a committed checkpoint"); + assert_eq!( + checkpoint.payload().len(), + expected_records.iter().map(|record| record.len()).sum::() + ); + assert!( + expected_records + .iter() + .all(|record| { checkpoint.payload().windows(record.len()).any(|bytes| bytes == *record) }), + "the latest checkpoint must preserve each exact kind, scope, version and object identity" + ); + + // The production inspector selects and validates the committed slot. Read + // its private lifecycle companion without adding a test-only public API. + let companion_path = format!(".heal-mrf-lifecycle.{}.bin", checkpoint.slot()); + let companions: Vec<_> = disk_paths + .iter() + .map(|path| std::fs::read(path.join(META_BUCKET).join(&companion_path)).expect("read each parked lifecycle replica")) + .collect(); + let companion = companions.first().expect("the fixture must have journal disks"); + assert!( + companions.iter().all(|replica| replica == companion), + "every disk must retain the same lifecycle responsibility" + ); + assert!(companion.len() >= 49, "lifecycle envelope must be complete"); + assert_eq!(&companion[..8], b"RFMFLC01"); + assert_eq!(companion[8], 1); + assert_eq!( + u64::from_le_bytes(companion[9..17].try_into().expect("lifecycle payload length")), + u64::try_from(companion.len() - 49).expect("fixture lifecycle length fits") + ); + assert_eq!(&companion[17..49], Sha256::digest(&companion[49..]).as_slice()); + let payload: serde_json::Value = serde_json::from_slice(&companion[49..]).expect("decode committed lifecycle payload"); + assert_eq!(payload["format_version"], 1); + assert_eq!(payload["checkpoint_owner"], checkpoint.owner().to_string()); + assert_eq!(payload["checkpoint_sequence"], checkpoint.sequence()); + let records = payload["records"].as_array().expect("lifecycle responsibility records"); + assert_eq!(records.len(), expected_records.len()); + let mut responsibility_ids = std::collections::HashSet::new(); + for expected in expected_records { + let digest: [u8; 32] = Sha256::digest(expected).into(); + let matching: Vec<_> = records + .iter() + .filter(|record| record["intent_digest"] == serde_json::json!(digest)) + .collect(); + assert_eq!(matching.len(), 1, "each exact partial-write identity must own one parked responsibility"); + let record = matching[0]; + let id = uuid::Uuid::parse_str(record["responsibility_id"].as_str().expect("responsibility UUID")) + .expect("valid responsibility UUID"); + assert!( + !id.is_nil() && responsibility_ids.insert(id), + "scoped responsibilities must have distinct non-nil IDs" + ); + assert_eq!( + record.get("source_bucket_incarnation_id"), + Some(&serde_json::Value::Null), + "replay must not invent the unknown producer generation" + ); + assert_eq!( + record.get("last_operator_acceptance"), + Some(&serde_json::Value::Null), + "replay must not fabricate operator authorization" + ); + assert_eq!(record["state"]["state"], "legacy_generation_unknown"); + assert_eq!( + record["state"]["observed_bucket_incarnation_id"], + observed_bucket_incarnation_id.to_string() + ); + assert!(record["state"]["detected_at_ms"].as_u64().is_some_and(|time| time > 0)); + } +} + async fn wait_until(deadline: Duration, mut probe: F) -> bool where F: FnMut() -> Fut, @@ -393,12 +442,17 @@ async fn decode_failure_intent_maps_to_urgent_mrf_heal_request() { /// A journal left behind by a previous process must be replayed into the /// manager queue, and a torn tail must not block replay of the intact records. -/// The partial-write record keeps the legacy journal as the durable anchor -/// until an exact verified repair proof can discharge it. +/// A partial write with no producer generation remains durably parked rather +/// than being admitted against whichever bucket happens to exist at replay. #[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_with_bucket("replay-bucket").await; + let bucket_incarnation_id = storage + .mrf_bucket_incarnation_id("replay-bucket") + .await + .expect("read the current replay bucket generation") + .expect("the replay bucket must have a persisted generation"); // The journal reader resolves disks through the process-local disk map; // register the environment's disks the same way server startup does. @@ -418,7 +472,10 @@ async fn journal_replay_arms_intents_and_retains_unproven_partial_write_anchor() assert_eq!(replayed, 2, "the two intact records must be replayed"); let snapshot = manager.operations_snapshot().await; - assert_eq!(snapshot.queued_by_source.mrf, 2, "replayed intents must be attributed to the MRF source"); + assert_eq!( + snapshot.queued_by_source.mrf, 1, + "only the decode-failure intent may enter the MRF manager queue" + ); assert!( disk_paths @@ -441,7 +498,11 @@ async fn journal_replay_arms_intents_and_retains_unproven_partial_write_anchor() let snapshot = manager.operations_snapshot().await; assert_eq!(snapshot.queued_by_priority.urgent, 1, "the decode-failure record must replay as Urgent"); - assert!(snapshot.queued_by_priority.normal >= 1, "the partial-write record must replay as Normal"); + assert_eq!( + snapshot.queued_by_priority.normal, 0, + "an unknown producer generation must not enter the repair queue" + ); + assert_legacy_responsibilities_parked(&disk_paths, &[&successor], bucket_incarnation_id).await; } /// A committed checkpoint published by the new two-slot writer is the @@ -451,16 +512,16 @@ async fn journal_replay_arms_intents_and_retains_unproven_partial_write_anchor() #[serial] async fn committed_snapshot_replay_takes_precedence_over_stale_legacy_mirror() { let (disk_paths, storage) = heal_env_with_bucket("committed-bucket").await; + let bucket_incarnation_id = storage + .mrf_bucket_incarnation_id("committed-bucket") + .await + .expect("read the current committed bucket generation") + .expect("the committed bucket must have a persisted generation"); 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); let stale_legacy = journal_record(1, "legacy-bucket", "legacy-object", None, 0); - let incarnation = storage - .mrf_bucket_incarnation_id("committed-bucket") - .await - .expect("read committed source incarnation") - .expect("committed source bucket has a persisted incarnation"); - write_generation_bound_snapshot_to_disks(&disk_paths, 7, &committed, incarnation, &[&committed]); + write_committed_snapshot_to_disks(&disk_paths, 7, &committed); write_journal_path_to_disks(&disk_paths, SCOPED_JOURNAL_REL, &stale_legacy); write_journal_path_to_disks(&disk_paths, JOURNAL_REL, &stale_legacy); @@ -469,10 +530,10 @@ async fn committed_snapshot_replay_takes_precedence_over_stale_legacy_mirror() { assert_eq!(replayed, 1, "only the committed snapshot epoch may replay"); let snapshot = manager.operations_snapshot().await; - assert_eq!(snapshot.queued_by_source.mrf, 1); + assert_eq!(snapshot.queued_by_source.mrf, 0); assert_eq!( - snapshot.queued_by_priority.normal, 1, - "the committed partial-write record must replay instead of the stale legacy decode-failure" + snapshot.queued_by_priority.normal, 0, + "the committed partial write without a producer generation must remain parked" ); assert_eq!( snapshot.queued_by_priority.urgent, 0, @@ -480,8 +541,9 @@ async fn committed_snapshot_replay_takes_precedence_over_stale_legacy_mirror() { ); assert!( journal_exists_on_all_disks(&disk_paths, COMMITTED_MANIFEST_REL), - "the committed checkpoint remains until the accepted partial-write has proof" + "the committed checkpoint must retain the parked partial-write responsibility" ); + assert_legacy_responsibilities_parked(&disk_paths, &[&committed], bucket_incarnation_id).await; } /// A damaged committed checkpoint is ambiguous: replay must not fall back to @@ -588,6 +650,11 @@ async fn authoritative_journal_is_not_merged_with_legacy_mirror() { #[serial] async fn authoritative_journal_replay_preserves_kind_and_scope_identity() { let (disk_paths, storage) = heal_env_with_bucket("identity-bucket").await; + let bucket_incarnation_id = storage + .mrf_bucket_incarnation_id("identity-bucket") + .await + .expect("read the current identity bucket generation") + .expect("the identity bucket must have a persisted generation"); register_local_disks(&disk_paths, "mrf-authoritative-identity-test").await; let first_partial = scoped_journal_record(3, "identity-bucket", "same-object", None, 0, 3, 7); @@ -599,12 +666,6 @@ async fn authoritative_journal_replay_preserves_kind_and_scope_identity() { let stale_legacy = journal_record(3, "identity-bucket", "stale-legacy-object", None, 0); write_journal_path_to_disks(&disk_paths, SCOPED_JOURNAL_REL, &authoritative); write_journal_path_to_disks(&disk_paths, JOURNAL_REL, &stale_legacy); - let incarnation = storage - .mrf_bucket_incarnation_id("identity-bucket") - .await - .expect("read authoritative source incarnation") - .expect("authoritative source bucket has a persisted incarnation"); - write_generation_bound_snapshot_to_disks(&disk_paths, 1, &authoritative, incarnation, &[&first_partial, &second_partial]); let manager = make_manager(storage); let replayed = mrf_queue::replay_journal_once(&manager).await; @@ -615,12 +676,12 @@ async fn authoritative_journal_replay_preserves_kind_and_scope_identity() { let snapshot = manager.operations_snapshot().await; assert_eq!( - snapshot.queued_by_source.mrf, 4, - "same-object MRF replay must retain distinct kind and scope responsibilities" + snapshot.queued_by_source.mrf, 2, + "only metadata-corruption and decode-failure repairs may enter the manager" ); assert_eq!( - snapshot.queued_by_priority.normal, 2, - "the two scoped partial-write records must remain independently queued" + snapshot.queued_by_priority.normal, 0, + "neither scoped partial write may be repaired with an unknown producer generation" ); assert_eq!( snapshot.queued_by_priority.high, 1, @@ -640,6 +701,7 @@ async fn authoritative_journal_replay_preserves_kind_and_scope_identity() { journal_matches_on_all_disks(&disk_paths, JOURNAL_REL, &[]), "the legacy mirror must not misrepresent scoped-only partial-write responsibilities" ); + assert_legacy_responsibilities_parked(&disk_paths, &[&first_partial, &second_partial], bucket_incarnation_id).await; } /// If replay reaches a full heal-manager queue, the old journal remains the diff --git a/crates/scanner/src/scanner/backlog.rs b/crates/scanner/src/scanner/backlog.rs index 3981df07e..c172fc640 100644 --- a/crates/scanner/src/scanner/backlog.rs +++ b/crates/scanner/src/scanner/backlog.rs @@ -2402,18 +2402,57 @@ mod tests { assert!(pool_meta_lock.get_write_lock_quiet(Duration::from_millis(100)).await.is_err()); barrier.release(); tokio::time::timeout(Duration::from_secs(30), async { - loop { + let loaded = loop { let loaded = load_scanner_pause_backlog(Arc::clone(&store)) .await .expect("cancellation cannot erase the old authority"); assert_eq!(loaded.ledger, original); - let committed = loaded.authoritative_commit.expect("seed retains the old cohort proof"); + let committed = loaded + .authoritative_commit + .as_ref() + .expect("seed retains the old cohort proof"); assert_eq!(committed.replicas, vec![replica_id(0, 0), replica_id(0, 1)]); + assert_eq!(loaded.replica_count, 6); if loaded.healthy_replicas == loaded.replica_count { - break; + break loaded; } tokio::task::yield_now().await; - } + }; + let drained_fence = pool_meta_lock + .get_write_lock_quiet(Duration::from_secs(30)) + .await + .expect("the detached seed owner releases its membership fence after persistence"); + drop(drained_fence); + assert_eq!(loaded.ledger, original); + let committed = loaded + .authoritative_commit + .as_ref() + .expect("seed retains the old cohort proof"); + assert_eq!(committed.replicas, vec![replica_id(0, 0), replica_id(0, 1)]); + assert_eq!(loaded.replica_count, 6); + assert_eq!(loaded.healthy_replicas, 6); + let seeded = loaded + .replicas + .iter() + .find(|replica| replica.id == replica_id(2, 0)) + .expect("the canceled caller's admitted seed replica"); + assert!(matches!( + &seeded.state, + ScannerPauseBacklogReplicaState::Valid(record) + if record.stable.as_ref() == Some(&original) && record.committed.as_ref() == Some(committed) + )); + // The detached publication owner finishes every selected replica + // without advancing beyond the old committed authority. + let remaining = loaded + .replicas + .iter() + .find(|replica| replica.id == replica_id(2, 1)) + .expect("the remaining seed replica"); + assert!(matches!( + &remaining.state, + ScannerPauseBacklogReplicaState::Valid(record) + if record.stable.as_ref() == Some(&original) && record.committed.as_ref() == Some(committed) + )); }) .await .expect("detached native seed owners must drain without the canceled caller"); diff --git a/crates/scanner/src/scanner_folder/tests/mrf_ownership.rs b/crates/scanner/src/scanner_folder/tests/mrf_ownership.rs index 117541ac7..e83a73145 100644 --- a/crates/scanner/src/scanner_folder/tests/mrf_ownership.rs +++ b/crates/scanner/src/scanner_folder/tests/mrf_ownership.rs @@ -16,7 +16,7 @@ use super::*; use crate::storage_api::EcstoreHealResultItem as HealItem; use crate::storage_api::scanner_io::BucketInfo; use rustfs_common::mrf_channel::{ - MrfIngressResult, MrfKind, MrfScope, note_mrf_repaired, take_mrf_repaired_events_for, try_send_mrf_intent_typed, + MrfScope, note_mrf_repaired, persist_partial_write_intent_with_incarnation, take_mrf_repaired_events_for, }; use rustfs_heal::heal::{ manager::{HealConfig, HealManager}, @@ -379,30 +379,41 @@ async fn mrf_ownership_manager_completion_preserves_scanner_pending() { 1, 1, )); - let scope = Some(MrfScope { + let scope = MrfScope { pool_index: 0, set_index: 0, - }); - assert_eq!( - try_send_mrf_intent_typed(MrfKind::PartialWrite, &bucket, object, Some(version), scope), - MrfIngressResult::Enqueued - ); - // Re-admission establishes that the first terminal callback released - // its ingress lease. Statistics alone precede notice publication. + }; + let before = manager.get_statistics().await; + let completed_before = before.successful_tasks + before.failed_tasks; + persist_partial_write_intent_with_incarnation(&bucket, object, Some(version), scope, Some(storage.bucket_incarnation_id)) + .await + .expect("persist the partial write with its actual source incarnation"); + // The scheduler publishes terminal repair notices before updating these + // counters. Earlier objects can only enter their blocked retry here. tokio::time::timeout(Duration::from_secs(5), async { loop { - match try_send_mrf_intent_typed(MrfKind::PartialWrite, &bucket, object, Some(version), scope) { - MrfIngressResult::Enqueued => break, - MrfIngressResult::Coalesced => tokio::task::yield_now().await, - other => panic!("unexpected retry ingress result: {other:?}"), + let called = storage + .calls + .lock() + .expect("fixture calls") + .get(*object) + .copied() + .unwrap_or(0); + let stats = manager.get_statistics().await; + if called > 0 && stats.successful_tasks + stats.failed_tasks > completed_before { + break; } + tokio::task::yield_now().await; } }) .await - .expect("production terminal releases its ingress lease"); + .expect("production terminal publishes its repair notices and completion"); + persist_partial_write_intent_with_incarnation(&bucket, object, Some(version), scope, Some(storage.bucket_incarnation_id)) + .await + .expect("persist another generation for the same source incarnation"); tokio::time::timeout(Duration::from_secs(5), storage.retry_started.notified()) .await - .expect("the real consumer starts the second generation"); + .expect("the real consumer starts the second object call"); assert!( take_mrf_repaired_events_for(&bucket).is_empty(), "{object}: task completion must not emit an unproved repair" @@ -415,9 +426,9 @@ async fn mrf_ownership_manager_completion_preserves_scanner_pending() { ); } assert_eq!( - try_send_mrf_intent_typed(MrfKind::PartialWrite, &bucket, object, Some(version), scope), - MrfIngressResult::Coalesced, - "the in-flight retry retains its new ingress lease" + storage.calls.lock().expect("fixture calls").get(*object).copied(), + Some(2), + "the second object call remains blocked without a duplicate execution" ); note_mrf_repaired(&bucket, object, Some(*version.as_bytes())); let syncs_before_retry = scanner.pending_heal_sync_count; diff --git a/rustfs/src/admin/handlers/site_replication.rs b/rustfs/src/admin/handlers/site_replication.rs index 0935b9548..38f2828fa 100644 --- a/rustfs/src/admin/handlers/site_replication.rs +++ b/rustfs/src/admin/handlers/site_replication.rs @@ -8878,6 +8878,8 @@ mod tests { } /// Publish a ready IAM app context so `apply_iam_item` gets past its IAM guard. + /// Tests using this real storage fixture keep the `durable_iam_state_` prefix so + /// nextest reserves capacity while their internal races retain production deadlines. async fn publish_ready_iam_context() { use crate::admin::runtime_sources::{AppContext, publish_test_app_context}; use rustfs_iam::store::{Store as _, object::IAM_CONFIG_PREFIX}; @@ -9077,7 +9079,7 @@ mod tests { /// every gated item type. #[tokio::test] #[serial] - async fn apply_iam_item_applies_delayed_in_order_updates_through_the_receiver() { + async fn durable_iam_state_apply_iam_item_applies_delayed_in_order_updates_through_the_receiver() { publish_ready_iam_context().await; seed_two_peer_state_for_iam_apply().await; let iam = current_iam_handle().expect("test IAM"); @@ -9158,7 +9160,7 @@ mod tests { #[tokio::test] #[serial] - async fn apply_group_member_add_preserves_a_disabled_group_status() { + async fn durable_iam_state_apply_group_member_add_preserves_a_disabled_group_status() { publish_ready_iam_context().await; seed_two_peer_state_for_iam_apply().await; let iam = current_iam_handle().expect("test IAM"); @@ -9193,7 +9195,7 @@ mod tests { /// membership-only guard above exists to prevent. #[tokio::test] #[serial] - async fn apply_group_snapshot_propagates_a_disabled_status_with_members() { + async fn durable_iam_state_apply_group_snapshot_propagates_a_disabled_status_with_members() { publish_ready_iam_context().await; seed_two_peer_state_for_iam_apply().await; let iam = current_iam_handle().expect("test IAM"); @@ -9228,7 +9230,7 @@ mod tests { /// the write of one cannot interleave with the other's. #[tokio::test] #[serial] - async fn apply_iam_item_serializes_a_concurrent_older_grant_and_newer_revoke() { + async fn durable_iam_state_apply_iam_item_serializes_a_concurrent_older_grant_and_newer_revoke() { publish_ready_iam_context().await; seed_two_peer_state_for_iam_apply().await; let t1 = OffsetDateTime::now_utc() - time::Duration::hours(2); @@ -9267,7 +9269,7 @@ mod tests { /// applied the delete; a genuinely newer create still lands. #[tokio::test] #[serial] - async fn apply_iam_item_rejects_a_stale_recreate_after_a_replicated_delete() { + async fn durable_iam_state_apply_iam_item_rejects_a_stale_recreate_after_a_replicated_delete() { publish_ready_iam_context().await; seed_two_peer_state_for_iam_apply().await; let iam = current_iam_handle().expect("test IAM"); @@ -9312,7 +9314,7 @@ mod tests { /// then switched off, and carries the source stamp. #[tokio::test] #[serial] - async fn apply_iam_item_creates_a_replicated_service_account_with_its_status() { + async fn durable_iam_state_apply_iam_item_creates_a_replicated_service_account_with_its_status() { publish_ready_iam_context().await; seed_two_peer_state_for_iam_apply().await; let iam = current_iam_handle().expect("test IAM"); @@ -9341,7 +9343,7 @@ mod tests { #[tokio::test] #[serial] - async fn apply_iam_item_accepts_minio_sts_account_item_type() { + async fn durable_iam_state_apply_iam_item_accepts_minio_sts_account_item_type() { publish_ready_iam_context().await; // MinIO madmin-go sends `SRIAMItemSTSAcc = "sts-account"`. The bogus session token @@ -9363,7 +9365,7 @@ mod tests { #[tokio::test] #[serial] - async fn apply_iam_item_still_accepts_legacy_sts_credential_item_type() { + async fn durable_iam_state_apply_iam_item_still_accepts_legacy_sts_credential_item_type() { publish_ready_iam_context().await; // Older RustFS peers emit `sts-credential`; the alias stays accepted permanently @@ -15400,7 +15402,7 @@ mod tests { /// would delete the retry queue and every other field along with it. #[tokio::test(flavor = "multi_thread", worker_threads = 2)] #[serial] - async fn test_missed_pending_clear_must_not_rewrite_the_state_object() { + async fn durable_iam_state_test_missed_pending_clear_must_not_rewrite_the_state_object() { publish_ready_iam_context().await; // One peer, no pending records: exactly the shape the persist helper's clear @@ -15461,7 +15463,7 @@ mod tests { /// write before A finishes. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[serial] - async fn test_peer_join_admission_serializes_iam_apply_against_a_newer_join() { + async fn durable_iam_state_test_peer_join_admission_serializes_iam_apply_against_a_newer_join() { publish_ready_iam_context().await; // Whole-second timestamps so the RFC3339 round trip through the state @@ -15565,7 +15567,7 @@ mod tests { /// admission: the assertion on the empty IAM log turns red. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[serial] - async fn test_peer_join_admission_serializes_across_separate_nodes() { + async fn durable_iam_state_test_peer_join_admission_serializes_across_separate_nodes() { publish_ready_iam_context().await; let now = OffsetDateTime::now_utc().replace_nanosecond(0).expect("truncate nanos"); @@ -15662,7 +15664,7 @@ mod tests { /// acked rotation is cleared in the same transaction that reports true. #[tokio::test(flavor = "multi_thread", worker_threads = 2)] #[serial] - async fn test_finalize_pending_rotation_three_way_contract() { + async fn durable_iam_state_test_finalize_pending_rotation_three_way_contract() { publish_ready_iam_context().await; let local_peer = PeerInfo { @@ -15744,7 +15746,7 @@ mod tests { /// update. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[serial] - async fn test_state_object_lock_serializes_writers_from_separate_nodes() { + async fn durable_iam_state_test_state_object_lock_serializes_writers_from_separate_nodes() { publish_ready_iam_context().await; let seed = SiteReplicationState { @@ -15796,7 +15798,7 @@ mod tests { /// leave the receiver unable to tell which edit is newer. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[serial] - async fn test_peer_edit_generations_are_unique_across_nodes() { + async fn durable_iam_state_test_peer_edit_generations_are_unique_across_nodes() { publish_ready_iam_context().await; // A configured site: `persist_site_replication_state_no_lock` clears // the object once a site drops below two peers, and a cleared object @@ -15847,7 +15849,7 @@ mod tests { /// accepts the restarted counter instead of fencing it. #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[serial] - async fn test_recreated_state_object_allocates_over_the_previous_lifetimes_mark() { + async fn durable_iam_state_test_recreated_state_object_allocates_over_the_previous_lifetimes_mark() { publish_ready_iam_context().await; let seed = || SiteReplicationState { peers: ["site-a", "site-b"] @@ -15912,7 +15914,7 @@ mod tests { /// (the red-light commit pinned the exact interleaving). #[tokio::test(flavor = "multi_thread", worker_threads = 4)] #[serial] - async fn test_retry_event_persist_must_not_wipe_concurrent_locked_rmw() { + async fn durable_iam_state_test_retry_event_persist_must_not_wipe_concurrent_locked_rmw() { publish_ready_iam_context().await; const ROUNDS: usize = 8; diff --git a/rustfs/tests/connect_heartbeat.rs b/rustfs/tests/connect_heartbeat.rs index f02d2679d..b7fd033ed 100644 --- a/rustfs/tests/connect_heartbeat.rs +++ b/rustfs/tests/connect_heartbeat.rs @@ -41,6 +41,9 @@ use tokio::sync::watch; use tokio_rustls::TlsAcceptor; use tokio_util::sync::CancellationToken; +// Allow durable state and real TLS setup independently of transport and cancellation deadlines. +const STATUS_OBSERVATION_TIMEOUT: Duration = Duration::from_secs(10); + const ORGANIZATION_UID: &str = "0198f4b0-1a00-7c10-8d21-2e3f4a5b6c70"; const CLUSTER_UID: &str = "0198f4b0-2b00-7d20-9e31-3f4a5b6c7d81"; const DEVICE_UID: &str = "0198f4b0-3c00-7e30-8f41-4a5b6c7d8e92"; @@ -191,9 +194,23 @@ impl Reply { struct TestServer { endpoint: String, seen: Arc>>, + paths: Arc>>, task: tokio::task::JoinHandle<()>, } +impl TestServer { + fn request_diagnostics(&self) -> String { + match self.paths.try_lock() { + Ok(paths) => format!( + "request_count={} request_paths={paths:?} server_task_finished={}", + paths.len(), + self.task.is_finished() + ), + Err(error) => format!("request_paths_unavailable={error} server_task_finished={}", self.task.is_finished()), + } + } +} + impl Drop for TestServer { fn drop(&mut self) { self.task.abort(); @@ -207,16 +224,20 @@ async fn server(pki: &TestPki, replies: Vec) -> TestServer { let replies = Arc::new(Mutex::new(VecDeque::from(replies))); let seen = Arc::new(Mutex::new(Vec::new())); let captured = seen.clone(); + let paths = Arc::new(Mutex::new(Vec::new())); + let captured_paths = paths.clone(); let task = tokio::spawn(async move { while let Ok((stream, _)) = listener.accept().await { let acceptor = acceptor.clone(); let replies = replies.clone(); let seen = captured.clone(); + let paths = captured_paths.clone(); tokio::spawn(async move { let Ok(stream) = acceptor.accept(stream).await else { return }; let service = service_fn(move |request: Request| { let replies = replies.clone(); let seen = seen.clone(); + let paths = paths.clone(); async move { let path = request.uri().path().to_owned(); assert!( @@ -227,6 +248,7 @@ async fn server(pki: &TestPki, replies: Vec) -> TestServer { seen.lock() .expect("seen lock") .push(serde_json::from_slice(&body).expect("request JSON")); + paths.lock().expect("paths lock").push(path.clone()); let reply = if path.ends_with("/diagnosticExecutionReceipts") { Reply { status: StatusCode::OK, @@ -266,6 +288,7 @@ async fn server(pki: &TestPki, replies: Vec) -> TestServer { TestServer { endpoint: format!("https://localhost:{}/agent/", address.port()), seen, + paths, task, } } @@ -314,23 +337,33 @@ fn summary() -> CoarseNodeSummary { } async fn wait_for( + server: &TestServer, status: &mut watch::Receiver, predicate: impl Fn(&HeartbeatStatus) -> bool, ) -> HeartbeatStatus { - tokio::time::timeout(Duration::from_secs(3), async { + tokio::time::timeout(STATUS_OBSERVATION_TIMEOUT, async { loop { let current = status.borrow_and_update().clone(); if predicate(¤t) { return current; } - status - .changed() - .await - .unwrap_or_else(|error| panic!("heartbeat status channel closed: {error}; last status: {:?}", *status.borrow())); + status.changed().await.unwrap_or_else(|error| { + panic!( + "heartbeat status channel closed: {error}; last status: {:?}; {}", + *status.borrow(), + server.request_diagnostics() + ) + }); } }) .await - .unwrap_or_else(|error| panic!("heartbeat status timeout: {error}; last status: {:?}", *status.borrow())) + .unwrap_or_else(|error| { + panic!( + "heartbeat status timeout: {error}; last status: {:?}; {}", + *status.borrow(), + server.request_diagnostics() + ) + }) } async fn assert_credential_failure(config: HeartbeatConfig, server: &TestServer, expected: &str) { @@ -340,7 +373,7 @@ async fn assert_credential_failure(config: HeartbeatConfig, server: &TestServer, .expect("configured runtime"); let mut status = runtime.status(); assert!(matches!( - wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await, + wait_for(server, &mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await, HeartbeatStatus::Failed { reason } if reason.contains(expected) )); assert!(server.seen.lock().expect("seen lock").is_empty()); @@ -436,7 +469,7 @@ async fn corrupt_persisted_state_is_rejected_before_network_delivery() { let mut status = runtime.status(); assert!(matches!( - wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await, + wait_for(&server, &mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await, HeartbeatStatus::Failed { reason } if reason.contains("violates the protocol invariants") )); assert!(server.seen.lock().expect("seen lock").is_empty()); @@ -505,7 +538,7 @@ async fn sends_only_l0_fields_and_accepts_additive_response_fields() { let mut status = runtime.status(); assert_eq!( - wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await, + wait_for(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await, HeartbeatStatus::Online { server_time: "2038-01-19T03:14:07Z".to_owned() } @@ -599,7 +632,7 @@ async fn restart_replays_a_pending_heartbeat_from_before_environment_collection( .expect("configured runtime"); let mut status = runtime.status(); assert!(matches!( - wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await, + wait_for(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await, HeartbeatStatus::Online { .. } )); runtime.shutdown().await; @@ -620,7 +653,7 @@ async fn restart_replays_pending_request_then_advances_sequence() { .expect("start runtime") .expect("configured runtime"); let mut status = runtime.status(); - wait_for(&mut status, |status| matches!(status, HeartbeatStatus::BackingOff { .. })).await; + wait_for(&first_server, &mut status, |status| matches!(status, HeartbeatStatus::BackingOff { .. })).await; runtime.shutdown().await; let first = first_server.seen.lock().expect("seen lock")[0].clone(); drop(first_server); @@ -730,7 +763,7 @@ async fn restart_replays_exact_legacy_capabilities_with_and_without_jobs() { .expect("configured runtime"); let mut status = runtime.status(); assert!(matches!( - wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await, + wait_for(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await, HeartbeatStatus::Online { .. } )); runtime.shutdown().await; @@ -770,7 +803,7 @@ async fn legacy_compatibility_does_not_accept_changed_capabilities() { let mut status = runtime.status(); assert!( matches!( - wait_for(&mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await, + wait_for(&server, &mut status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await, HeartbeatStatus::Failed { reason } if reason.contains("violates the protocol invariants") ), "accepted altered persisted capabilities: {mutation}, jobs={job_capable}, service_memory={service_memory}" @@ -859,7 +892,7 @@ async fn retry_after_is_respected_with_the_local_upper_bound() { .expect("configured runtime"); let mut status = runtime.status(); assert_eq!( - wait_for(&mut status, |status| matches!(status, HeartbeatStatus::BackingOff { .. })).await, + wait_for(&server, &mut status, |status| matches!(status, HeartbeatStatus::BackingOff { .. })).await, HeartbeatStatus::BackingOff { delay: Duration::from_millis(80) } @@ -887,7 +920,7 @@ async fn disconnects_use_exponential_backoff_with_a_cap() { let mut status = runtime.status(); for delay in [20, 40, 80] { assert_eq!( - wait_for(&mut status, |status| { + wait_for(&server, &mut status, |status| { matches!(status, HeartbeatStatus::BackingOff { delay: observed } if *observed == Duration::from_millis(delay)) }) .await, @@ -912,7 +945,11 @@ async fn revoked_credential_stops_and_exposes_local_status() { .expect("configured runtime"); let mut status = runtime.status(); assert_eq!( - wait_for(&mut status, |status| matches!(status, HeartbeatStatus::AuthenticationStopped { .. })).await, + wait_for(&server, &mut status, |status| matches!( + status, + HeartbeatStatus::AuthenticationStopped { .. } + )) + .await, HeartbeatStatus::AuthenticationStopped { status: 401, reason: Some("CREDENTIAL_REVOKED".to_owned()) diff --git a/rustfs/tests/connect_perf_object.rs b/rustfs/tests/connect_perf_object.rs index 4ed6fb45a..be6fc6c53 100644 --- a/rustfs/tests/connect_perf_object.rs +++ b/rustfs/tests/connect_perf_object.rs @@ -38,6 +38,9 @@ use zeroize::Zeroizing; mod common; static TEST_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); +// Real storage setup shares the measurement window. Verify functionality with +// an integration budget; deadline behavior is checked against a stalled peer. +const REAL_SERVER_MEASUREMENT_BUDGET: Duration = Duration::from_secs(10); fn now() -> i64 { SystemTime::now().duration_since(UNIX_EPOCH).expect("current time").as_secs() as i64 @@ -305,6 +308,56 @@ async fn deadline_and_in_flight_cancellation_stop_a_stalled_probe() { assert_eq!(measurement.target.reason_code, ObjectTargetReasonCode::Cancelled); } +#[tokio::test] +async fn s3_request_deadline_stops_a_peer_that_never_responds() { + use tokio::io::AsyncReadExt as _; + + let _guard = TEST_LOCK.lock().await; + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.expect("response listener"); + let address = listener.local_addr().expect("listener address"); + let (received, request_received) = tokio::sync::oneshot::channel(); + let server = tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.expect("object connection"); + let mut request = [0_u8; 4096]; + assert!(socket.read(&mut request).await.expect("request headers") > 0); + received.send(()).expect("report received request"); + std::future::pending::<()>().await; + drop(socket); + }); + let probe = S3ObjectProbe::new( + &format!("http://{address}"), + None, + None, + Zeroizing::new("access".to_owned()), + Zeroizing::new("secret".to_owned()), + Zeroizing::new(String::new()), + Duration::from_secs(60), + ) + .expect("object probe"); + let mut request = request(ObjectOperation::GetObject); + request.duration = Duration::from_secs(30); + let measurement = tokio::spawn(async move { measure_object(&request, &probe, &CancellationToken::new()).await }); + + // Complete real TCP setup before pausing time. The 30s operation deadline + // must fire independently of the client's longer 60s transport timeout. + tokio::time::timeout(Duration::from_secs(10), request_received) + .await + .expect("peer received the request") + .expect("peer notification"); + tokio::time::pause(); + tokio::time::advance(Duration::from_secs(30)).await; + let measurement = tokio::time::timeout(Duration::from_secs(1), measurement) + .await + .expect("operation deadline must stop the request") + .expect("measurement task") + .expect("typed timeout"); + assert_eq!(measurement.result.outcome(), ObjectOutcome::Failed); + assert_eq!(measurement.target.reason_code, ObjectTargetReasonCode::TimedOut); + assert_eq!(measurement.target.completed_operations, 0); + server.abort(); + let _ = server.await; +} + #[tokio::test] async fn only_one_object_collector_can_run_at_a_time() { let _guard = TEST_LOCK.lock().await; @@ -519,6 +572,7 @@ fn real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put() { async fn real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put_body() { let _guard = TEST_LOCK.lock().await; + let startup = std::time::Instant::now(); let port = match find_available_port() { Ok(port) => port, Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return, @@ -531,6 +585,7 @@ async fn real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put_bod .build() .await .expect("start embedded server"); + eprintln!("embedded object fixture ready after {:?}", startup.elapsed()); let probe = S3ObjectProbe::new( &server.endpoint(), None, @@ -538,13 +593,13 @@ async fn real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put_bod Zeroizing::new(server.access_key().to_owned()), Zeroizing::new(server.secret_key().to_owned()), Zeroizing::new(String::new()), - Duration::from_secs(2), + REAL_SERVER_MEASUREMENT_BUDGET, ) .expect("object probe"); for operation in [ObjectOperation::GetObject, ObjectOperation::PutObject] { let mut request = request(operation); - request.duration = Duration::from_secs(2); + request.duration = REAL_SERVER_MEASUREMENT_BUDGET; if operation == ObjectOperation::PutObject { request.artifact_uid = "019e3ae0-0000-7000-8000-000000000016".to_owned(); } @@ -604,7 +659,7 @@ async fn real_rustfs_endpoint_and_production_cli_support_bounded_get_and_put_bod "--traffic-bytes", "65536", "--duration-millis", - "1000", + &REAL_SERVER_MEASUREMENT_BUDGET.as_millis().to_string(), "--acknowledge-l1", ]); let result = tokio::task::spawn_blocking(move || command.output()) diff --git a/rustfs/tests/connect_registration.rs b/rustfs/tests/connect_registration.rs index 2385dd960..e03a67200 100644 --- a/rustfs/tests/connect_registration.rs +++ b/rustfs/tests/connect_registration.rs @@ -59,6 +59,9 @@ const EXPECTED_CERTIFICATE_LIFETIME_SECONDS: i64 = 86_400; const EXPECTED_ROTATION_THRESHOLD_SECONDS: i64 = 120; #[cfg(not(feature = "connect-e2e-short-credentials"))] const EXPECTED_ROTATION_THRESHOLD_SECONDS: i64 = 8 * 60 * 60; +// Allow durable state and real TLS setup independently of transport and cancellation deadlines. +const STATUS_OBSERVATION_TIMEOUT: Duration = Duration::from_secs(10); + const ORGANIZATION_UID: &str = "0198f4b0-1a00-7c10-8d21-2e3f4a5b6c70"; const CLUSTER_UID: &str = "0198f4b0-2b00-7d20-9e31-3f4a5b6c7d81"; const DEVICE_UID: &str = "0198f4b0-3c00-7e30-8f41-4a5b6c7d8e92"; @@ -209,6 +212,19 @@ struct TestServer { task: tokio::task::JoinHandle<()>, } +impl TestServer { + fn request_diagnostics(&self) -> String { + match self.paths.try_lock() { + Ok(paths) => format!( + "request_count={} request_paths={paths:?} server_task_finished={}", + paths.len(), + self.task.is_finished() + ), + Err(error) => format!("request_paths_unavailable={error} server_task_finished={}", self.task.is_finished()), + } + } +} + impl Drop for TestServer { fn drop(&mut self) { self.task.abort(); @@ -678,7 +694,7 @@ async fn explicit_proxy_carries_registration_rotation_and_heartbeat_with_mtls() .expect("configured heartbeat runtime"); let mut status = runtime.status(); assert!(matches!( - wait_for_heartbeat_status(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await, + wait_for_heartbeat_status(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await, HeartbeatStatus::Online { .. } )); runtime.shutdown().await; @@ -872,22 +888,33 @@ async fn wait_for_requests(server: &TestServer, count: usize) { } async fn wait_for_heartbeat_status( + server: &TestServer, status: &mut watch::Receiver, predicate: impl Fn(&HeartbeatStatus) -> bool, ) -> HeartbeatStatus { - tokio::time::timeout(Duration::from_secs(3), async { + tokio::time::timeout(STATUS_OBSERVATION_TIMEOUT, async { loop { let current = status.borrow_and_update().clone(); if predicate(¤t) { return current; } status.changed().await.unwrap_or_else(|error| { - panic!("heartbeat status channel: {error}; last status: {:?}", *status.borrow()); + panic!( + "heartbeat status channel: {error}; last status: {:?}; {}", + *status.borrow(), + server.request_diagnostics() + ); }); } }) .await - .unwrap_or_else(|error| panic!("heartbeat status: {error}; last status: {:?}", *status.borrow())) + .unwrap_or_else(|error| { + panic!( + "heartbeat status: {error}; last status: {:?}; {}", + *status.borrow(), + server.request_diagnostics() + ) + }) } fn rotation_response(pki: &TestPki, identity: &rustfs::connect::DeviceIdentity, serial: u8) -> (Value, Value) { @@ -1301,7 +1328,25 @@ async fn heartbeat_runtime_respects_rotation_retry_after_without_blocking_heartb .expect("configured runtime"); let mut status = runtime.status(); - wait_for_heartbeat_status(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await; + // Start the progress bound after the durable rotation claim reaches the server. + // Heartbeats must still proceed before the five-second rotation Retry-After expires. + wait_for_requests(&server, 1).await; + assert!( + server.paths.lock().expect("paths lock")[0].ends_with(":rotateCredential"), + "rotation request must precede heartbeat progress" + ); + tokio::time::timeout( + Duration::from_secs(3), + wait_for_heartbeat_status(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })), + ) + .await + .unwrap_or_else(|error| { + panic!( + "heartbeat progress blocked after rotation: {error}; last status: {:?}; {}", + *status.borrow(), + server.request_diagnostics() + ) + }); wait_for_requests(&server, 4).await; runtime.shutdown().await; @@ -1353,7 +1398,7 @@ async fn heartbeat_runtime_skips_only_valid_pending_reenrollment() { let mut status = runtime.status(); assert!(matches!( - wait_for_heartbeat_status(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await, + wait_for_heartbeat_status(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await, HeartbeatStatus::Online { .. } )); runtime.shutdown().await; @@ -1393,7 +1438,8 @@ async fn heartbeat_runtime_skips_only_valid_pending_reenrollment() { .expect("configured runtime"); let mut corrupted_status = corrupted_runtime.status(); assert_eq!( - wait_for_heartbeat_status(&mut corrupted_status, |status| matches!(status, HeartbeatStatus::Failed { .. })).await, + wait_for_heartbeat_status(&server, &mut corrupted_status, |status| matches!(status, HeartbeatStatus::Failed { .. })) + .await, HeartbeatStatus::Failed { reason: HeartbeatError::StateConflict.to_string(), } @@ -1462,7 +1508,7 @@ async fn heartbeat_retries_with_the_new_credential_after_concurrent_rotation() { .expect("concurrent rotation") .expect("rotation due"); assert!(matches!( - wait_for_heartbeat_status(&mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await, + wait_for_heartbeat_status(&server, &mut status, |status| matches!(status, HeartbeatStatus::Online { .. })).await, HeartbeatStatus::Online { .. } )); runtime.shutdown().await; @@ -1547,7 +1593,7 @@ async fn inventory_retries_with_the_new_credential_after_concurrent_rotation() { .expect("concurrent rotation") .expect("rotation due"); assert!(matches!( - tokio::time::timeout(Duration::from_secs(3), async { + tokio::time::timeout(STATUS_OBSERVATION_TIMEOUT, async { loop { let current = status.borrow_and_update().clone(); if matches!(current, InventoryStatus::BackingOff { .. }) { @@ -1642,20 +1688,37 @@ async fn inventory_first_recovers_a_saved_reenrollment_before_telemetry() { .expect("start inventory runtime") .expect("configured inventory runtime"); let mut inventory_status = inventory.status(); + let startup_diagnostics = || { + format!( + "{}; pending_registration={} staged_identity={} completed_registration={}", + telemetry.request_diagnostics(), + temp.path().join("credential/registration.pending.json").exists(), + temp.path().join("identity/device.key.next").exists(), + temp.path().join("credential/registration.completed.json").exists() + ) + }; assert!(matches!( - tokio::time::timeout(Duration::from_secs(3), async { + tokio::time::timeout(STATUS_OBSERVATION_TIMEOUT, async { loop { let current = inventory_status.borrow_and_update().clone(); if matches!(current, InventoryStatus::Online { .. }) { break current; } inventory_status.changed().await.unwrap_or_else(|error| { - panic!("inventory status channel: {error}; last status: {:?}", *inventory_status.borrow()); + panic!( + "inventory status channel: {error}; last status: {:?}; {}", + *inventory_status.borrow(), + startup_diagnostics() + ); }); } }) .await - .unwrap_or_else(|error| panic!("inventory online status: {error}; last status: {:?}", *inventory_status.borrow())), + .unwrap_or_else(|error| panic!( + "inventory online status: {error}; last status: {:?}; {}", + *inventory_status.borrow(), + startup_diagnostics() + )), InventoryStatus::Online { .. } )); @@ -1684,7 +1747,11 @@ async fn inventory_first_recovers_a_saved_reenrollment_before_telemetry() { .expect("configured heartbeat runtime"); let mut heartbeat_status = heartbeat.status(); assert!(matches!( - wait_for_heartbeat_status(&mut heartbeat_status, |status| matches!(status, HeartbeatStatus::Online { .. })).await, + wait_for_heartbeat_status(&telemetry, &mut heartbeat_status, |status| matches!( + status, + HeartbeatStatus::Online { .. } + )) + .await, HeartbeatStatus::Online { .. } )); heartbeat.shutdown().await; @@ -1747,7 +1814,11 @@ async fn heartbeat_runtime_stops_when_rotation_reports_revocation() { let mut status = runtime.status(); assert_eq!( - wait_for_heartbeat_status(&mut status, |status| matches!(status, HeartbeatStatus::AuthenticationStopped { .. })).await, + wait_for_heartbeat_status(&server, &mut status, |status| matches!( + status, + HeartbeatStatus::AuthenticationStopped { .. } + )) + .await, HeartbeatStatus::AuthenticationStopped { status: 401, reason: Some("DEVICE_REVOKED".to_owned()),