From 590fab5c7eb1e45fb5dea2acd2b295c3e1a04119 Mon Sep 17 00:00:00 2001 From: overtrue Date: Tue, 8 Sep 2026 16:57:04 +0800 Subject: [PATCH] test: stabilize release recovery and transition gates --- .config/e2e-full-selection.txt | 4 +- .config/e2e-nightly-selection.txt | 4 +- .config/e2e-repl-nightly-selection.txt | 2 +- .github/workflows/e2e-distributed.yml | 4 +- crates/e2e_test/src/distributed/heal_test.rs | 68 ++++- .../src/heal_erasure_disk_rebuild_test.rs | 243 ++++++------------ .../src/inline_fast_path_cluster_test.rs | 70 ++++- crates/e2e_test/src/lib.rs | 3 + crates/e2e_test/src/scanner_heal_evidence.rs | 181 +++++++++++++ docs/testing/ci-gates.md | 4 +- 10 files changed, 405 insertions(+), 178 deletions(-) create mode 100644 crates/e2e_test/src/scanner_heal_evidence.rs diff --git a/.config/e2e-full-selection.txt b/.config/e2e-full-selection.txt index d519fff8e..395443805 100644 --- a/.config/e2e-full-selection.txt +++ b/.config/e2e-full-selection.txt @@ -1,2 +1,2 @@ -sha256-darwin=71d04825c143b85334802dbd2d67df6f08006f18dc2835df85bc21ea85eca02a -sha256-linux=baee6b0c8b38a5be139d5583bcd4fd645a27659a70b2ca8f81646d760c299395 +sha256-darwin=cca6d0bc1487f472dc354bbffcc2b5c7410edfccb499c6bce6638113fbe6bbec +sha256-linux=3b71936f6f4ea0cca3b5db6c2f387c990dddda7e96c982315eb42e3b2e6c92b8 diff --git a/.config/e2e-nightly-selection.txt b/.config/e2e-nightly-selection.txt index 163d3235b..bec86f799 100644 --- a/.config/e2e-nightly-selection.txt +++ b/.config/e2e-nightly-selection.txt @@ -1,2 +1,2 @@ -sha256-darwin=a5665318c9bdc0947514fb7008ba1b83b114b739fac775c3c446f207058b7c7a -sha256-linux=45d80e1723de5d25bb5b81f3ef5c82f583efc3e4f036a8cd2bb99e4f1eca9e51 +sha256-darwin=83a7dcaffd5a789517ae9f02a224f66a9713937885cff96fca2ad7e216f197ae +sha256-linux=626c10f8c964507ff987b6c86069e9019dc6d2ae7fb02db9be5df5aa8cc5145b diff --git a/.config/e2e-repl-nightly-selection.txt b/.config/e2e-repl-nightly-selection.txt index b879da5f6..6626a2217 100644 --- a/.config/e2e-repl-nightly-selection.txt +++ b/.config/e2e-repl-nightly-selection.txt @@ -1 +1 @@ -sha256=0fe8408874ccec3620262a9812d67920ddd72dc9edf0e36e0d0aed3f8bad026e +sha256=0e338d305260229e17ccfb2adc48a6212dbdfea36a9ebfb5a4e0d38658e6cc45 diff --git a/.github/workflows/e2e-distributed.yml b/.github/workflows/e2e-distributed.yml index 091d15c98..42842e880 100644 --- a/.github/workflows/e2e-distributed.yml +++ b/.github/workflows/e2e-distributed.yml @@ -19,7 +19,7 @@ # case is a two-site 4-node 1-drive pair or a 4-node upgrade). Membership is # `[profile.e2e-distributed]` in `.config/nextest.toml`. Storage-sensitive PRs, # nightly runs, and manual dispatches all execute the same fail-closed suite. -# Upgrade cases download the same pinned previous release as e2e-upgrade.yml. +# Upgrade cases use an independent 1.0.0-rc.2 pin defined below. # # Isolated pool filesystems: expand/decommission/rebalance cases require # independent `statfs` capacity. This job runs on GitHub-hosted @@ -87,7 +87,7 @@ jobs: NO_PROXY: 127.0.0.1,localhost HTTP_PROXY: "" HTTPS_PROXY: "" - # Pinned previous release used by distributed::upgrade_test (same pin as e2e-upgrade.yml). + # Independent 1.0.0-rc.2 source pin for distributed::upgrade_test. UPGRADE_SOURCE_VERSION: 1.0.0-rc.2 UPGRADE_SOURCE_ASSET: rustfs-linux-x86_64-gnu-v1.0.0-rc.2.zip UPGRADE_SOURCE_SHA256: 7c789386bf85278f865b8e0d359bf4edb84d5aa408cc3fa54a18c25ca74cd6e7 diff --git a/crates/e2e_test/src/distributed/heal_test.rs b/crates/e2e_test/src/distributed/heal_test.rs index 29b79a7d6..8ea2e54d5 100644 --- a/crates/e2e_test/src/distributed/heal_test.rs +++ b/crates/e2e_test/src/distributed/heal_test.rs @@ -14,9 +14,10 @@ use super::harness::{DistCluster, DistLayout, TestResult, assert_inventory, payload_for, put_object, unique_bucket, wait_until}; use crate::chaos::{ - VersionShardCensus, census_object_version_on_disk, signed_admin_post, wait_for_complete_physical_shard_on_disk, + VersionShardCensus, census_object_version_on_disk, sha256_hex, signed_admin_post, wait_for_complete_physical_shard_on_disk, }; -use crate::common::init_logging; +use crate::common::{init_logging, rustfs_binary_path}; +use crate::scanner_heal_evidence::{EvidenceTopology, RestartObservation, ScannerHealEvidenceCase, restart_evidence_run}; use aws_sdk_s3::Client; use aws_sdk_s3::primitives::ByteStream; use std::collections::{BTreeMap, HashSet}; @@ -89,6 +90,19 @@ async fn put_large_inventory(client: &Client, bucket: &str) -> TestResult TestResult { init_logging(); + let server_binary = rustfs_binary_path(); + let evidence_run = restart_evidence_run( + &server_binary, + ScannerHealEvidenceCase { + id: "ec84-target-drive-restart", + oracle: "ec84-target-drive-restart.json", + evidence: "process-restart", + unclean_shutdown_marker: false, + topology: EvidenceTopology::new(3, 4), + storage_class_standard: Some("EC:4"), + erasure_set_drive_count: Some("12"), + }, + )?; let mut dist = DistCluster::start_with_env( DistLayout::ThreeByFourEc84, &[ @@ -119,7 +133,17 @@ async fn three_node_four_drive_ec8_4_root_heal_rebuilds_replaced_drive_after_res let format_path = replaced_drive.join(".rustfs.sys").join("format.json"); let format_json = std::fs::read(&format_path)?; + let pid_before = dist.cluster.nodes[replaced_node] + .process + .as_ref() + .ok_or("target process is absent")? + .id(); dist.cluster.stop_node_gracefully(replaced_node).await?; + let unclean_shutdown_marker = Path::new(&dist.cluster.nodes[replaced_node].data_dir) + .join(".rustfs.sys") + .join("unclean-shutdown") + .is_file(); + assert!(!unclean_shutdown_marker, "graceful target shutdown must remove its unclean marker"); let retired_drive = PathBuf::from(format!("{}.retired", replaced_drive.display())); std::fs::rename(&replaced_drive, &retired_drive)?; std::fs::create_dir_all(format_path.parent().ok_or("replacement format path has no parent")?)?; @@ -170,6 +194,7 @@ async fn three_node_four_drive_ec8_4_root_heal_rebuilds_replaced_drive_after_res .chain(std::iter::once((outage_key.to_string(), outage_body.clone()))) .collect::>(); let expected_keys = inventory.keys().cloned().collect::>(); + let mut node_listings = Vec::new(); for node_index in 0..dist.cluster.nodes.len() { let client = dist.client(node_index)?; assert_inventory(&client, &bucket, &inventory).await?; @@ -180,6 +205,45 @@ async fn three_node_four_drive_ec8_4_root_heal_rebuilds_replaced_drive_after_res .filter_map(|object| object.key().map(str::to_owned)) .collect::>(); assert_eq!(observed, expected_keys, "node {node_index} listing diverged after EC8+4 heal"); + let mut keys = observed.into_iter().collect::>(); + keys.sort(); + node_listings.push(keys); + } + + if let Some(evidence_run) = evidence_run { + let target_client = dist.client(replaced_node)?; + let mut objects = Vec::with_capacity(inventory.len()); + for (key, body) in &inventory { + let response = target_client.get_object().bucket(&bucket).key(key).send().await?; + let actual = response.body.collect().await?.into_bytes(); + assert_eq!(actual.as_ref(), body.as_slice(), "object body changed for {key}"); + let physical = census_object_version_on_disk(&replaced_drive, &bucket, key, None)?; + assert_ec84_geometry(&physical, key)?; + let baseline = expected.iter().find(|item| item.key == *key).map(|item| &item.baseline); + objects.push(serde_json::json!({ + "key": key, "version_id": null, + "expected_bytes": body.len(), "actual_bytes": actual.len(), + "expected_sha256": sha256_hex(body), "actual_sha256": sha256_hex(&actual), + "expected_physical": baseline, "physical": physical, + })); + } + let pid_after = dist.cluster.nodes[replaced_node] + .process + .as_ref() + .ok_or("restarted target is absent")? + .id(); + evidence_run.write( + &server_binary, + RestartObservation { + nodes: dist.cluster.nodes.len(), + drives_per_node: dist.cluster.topology.drives_per_node, + pid_before, + pid_after, + unclean_shutdown_marker, + objects, + node_listings, + }, + )?; } Ok(()) diff --git a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs index 1a0eccfd8..a86f82473 100644 --- a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs +++ b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs @@ -21,16 +21,15 @@ mod tests { wait_for_complete_physical_shard_on_disk, }; use crate::common::{ - ClusterTopology, FAST_DATA_USAGE_SCANNER_ENV, RustFSTestClusterEnvironment, RustFSTestEnvironment, admin_request, - init_logging, rustfs_binary_path, + FAST_DATA_USAGE_SCANNER_ENV, RustFSTestClusterEnvironment, RustFSTestEnvironment, admin_request, init_logging, + rustfs_binary_path, }; + use crate::scanner_heal_evidence::{EvidenceTopology, RestartObservation, ScannerHealEvidenceCase, restart_evidence_run}; use crate::storage_api::RUSTFS_META_BUCKET; use aws_sdk_s3::primitives::ByteStream; use http::Method; - use sha2::{Digest, Sha256}; use std::collections::HashSet; use std::error::Error; - use std::io::{Read, Write}; use std::net::SocketAddr; use std::path::{Path, PathBuf}; use std::process::Command; @@ -42,52 +41,6 @@ mod tests { const POOL_METADATA_OBJECT: &str = "pool.bin"; - #[derive(serde::Deserialize)] - struct EvidenceBuild { - sha256: String, - } - - #[derive(serde::Deserialize)] - struct RestartEvidenceRun { - schema: u32, - run_id: String, - source_revision: String, - test_build: serde_json::Value, - binary: EvidenceBuild, - test_binary: EvidenceBuild, - } - - #[derive(Clone, Copy)] - struct ScannerHealEvidenceCase { - id: &'static str, - oracle: &'static str, - evidence: &'static str, - unclean_shutdown_marker: bool, - topology: EvidenceTopology, - storage_class_standard: Option<&'static str>, - erasure_set_drive_count: Option<&'static str>, - } - - #[derive(Clone, Copy)] - struct EvidenceTopology { - nodes: usize, - drives_per_node: usize, - } - - impl EvidenceTopology { - const fn new(nodes: usize, drives_per_node: usize) -> Self { - Self { nodes, drives_per_node } - } - - fn total_drives(self) -> usize { - self.nodes * self.drives_per_node - } - - fn cluster_topology(self) -> ClusterTopology { - ClusterTopology::single_pool_multidrive(self.nodes, self.drives_per_node) - } - } - const BACKGROUND_TARGET_RESTART_EVIDENCE: ScannerHealEvidenceCase = ScannerHealEvidenceCase { id: "background-target-restart", oracle: "background-target-restart.json", @@ -128,81 +81,6 @@ mod tests { erasure_set_drive_count: Some("12"), }; - struct RestartEvidenceContext { - directory: PathBuf, - run: RestartEvidenceRun, - case: ScannerHealEvidenceCase, - } - - fn file_sha256(path: &Path) -> Result> { - let mut file = std::fs::File::open(path)?; - let mut digest = Sha256::new(); - let mut buffer = [0_u8; 64 * 1024]; - loop { - let read = file.read(&mut buffer)?; - if read == 0 { - break; - } - digest.update(&buffer[..read]); - } - Ok(digest.finalize().iter().map(|byte| format!("{byte:02x}")).collect()) - } - - fn restart_evidence_run( - binary: &Path, - case: ScannerHealEvidenceCase, - ) -> Result, Box> { - let Some(directory) = std::env::var_os("RUSTFS_SCANNER_HEAL_RUN_DIR") else { - return Ok(None); - }; - if case.id.is_empty() - || case.oracle.is_empty() - || !case.oracle.ends_with(".json") - || case.oracle.contains('/') - || case.oracle.contains('\\') - || case.oracle.contains("..") - || !matches!(case.evidence, "process-restart" | "process-crash-restart") - || (case.evidence == "process-crash-restart") != case.unclean_shutdown_marker - { - return Err("invalid scanner/heal evidence case".into()); - } - let directory = PathBuf::from(directory); - let receipt = directory.join("run.json"); - if receipt.metadata()?.len() > 1024 * 1024 { - return Err("oversized scanner/heal execution receipt".into()); - } - let run: RestartEvidenceRun = serde_json::from_slice(&std::fs::read(receipt)?)?; - if run.schema != 1 || run.run_id.len() != 32 || run.source_revision.len() != 40 { - return Err("invalid scanner/heal execution identity".into()); - } - let built = compiled_test_identity(); - for key in ["source_revision", "dirty", "lock_blob", "features"] { - assert_eq!(built[key], run.test_build[key], "compiled test identity differs for {key}"); - } - assert_eq!(file_sha256(binary)?, run.binary.sha256, "server binary must match the run receipt"); - assert_eq!( - file_sha256(&std::env::current_exe()?)?, - run.test_binary.sha256, - "test executable must match the run receipt" - ); - if directory.join(case.oracle).exists() { - return Err("scanner/heal oracle already exists; create a new execution receipt".into()); - } - Ok(Some(RestartEvidenceContext { directory, run, case })) - } - - fn compiled_test_identity() -> serde_json::Value { - serde_json::json!({ - "source_revision": env!("RUSTFS_E2E_BUILD_COMMIT"), - "dirty": env!("RUSTFS_E2E_BUILD_DIRTY") != "false", - "lock_blob": env!("RUSTFS_E2E_BUILD_LOCK"), - "features": env!("RUSTFS_E2E_BUILD_FEATURES"), - "target": env!("RUSTFS_E2E_BUILD_TARGET"), - "profile": env!("RUSTFS_E2E_BUILD_PROFILE"), - "rustflags_hex": env!("RUSTFS_E2E_BUILD_RUSTFLAGS_HEX"), - }) - } - struct TcpPortBlackhole { port: u16, comment: String, @@ -943,7 +821,10 @@ mod tests { cluster: &RustFSTestClusterEnvironment, previous_cycle_end: u64, ) -> Result> { - let deadline = Instant::now() + Duration::from_secs(60); + let started = Instant::now(); + let mut deadline = started + Duration::from_secs(60); + let catch_up_deadline = deadline + Duration::from_secs(300); + let mut catch_up_wait_observed = false; loop { let mut latest_cycle_end = 0; let mut versions_observed = false; @@ -971,24 +852,61 @@ mod tests { let versions_scanned = metrics["versions_scanned"] .as_u64() .ok_or("scanner status is missing its version-coverage counter")?; - latest_cycle_end = latest_cycle_end.max(cycle_end); + let cycle_result = metrics["last_cycle_result"] + .as_str() + .ok_or("scanner status is missing its cycle result")?; + if cycle_result == "success" { + latest_cycle_end = latest_cycle_end.max(cycle_end); + } versions_observed |= versions_scanned > 0; + let backlog = &status["pause_backlog"]; + if !catch_up_wait_observed + && backlog["persistence_state"].as_str() == Some("healthy") + && backlog["durable"].as_bool() == Some(true) + && backlog["phase"].as_str() == Some("catching_up") + && backlog["rate_limited"].as_bool() == Some(true) + && backlog["retry_exhausted"].as_bool() == Some(false) + { + let next_attempt = backlog["next_attempt_at_unix_secs"] + .as_u64() + .ok_or("rate-limited scanner backlog is missing its next attempt")?; + let interval = backlog["thresholds"]["catch_up_min_interval_seconds"] + .as_u64() + .ok_or("rate-limited scanner backlog is missing its catch-up interval")?; + let now = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH)?.as_secs(); + let remaining = next_attempt.saturating_sub(now); + if remaining > 0 { + if interval > 300 || remaining > interval { + return Err( + format!("scanner catch-up schedule exceeds the bounded recovery budget: {backlog}").into() + ); + } + // The durable catch-up interval overrides SCANNER_CYCLE=1. + // Honor one observed retry without restarting the deadline on every poll. + deadline = deadline + .max(Instant::now() + Duration::from_secs(remaining + 60)) + .min(catch_up_deadline); + catch_up_wait_observed = true; + } + } observations.push(format!( - "node{node_index}: end={cycle_end}, versions={versions_scanned}, cycle={}, active={}, leader={}, result={}", + "node{node_index}: end={cycle_end}, versions={versions_scanned}, cycle={}, active={}, leader={}, result={}, backlog={}", metrics["current_cycle"], metrics["current_cycle_active"], metrics["leader_lock_state"], metrics["last_cycle_result"], + backlog, )); } - // The coordinator records cycle completion, but remote workers - // record scanned versions. Both witnesses need not share a node. + // Only a successful coordinator cycle counts as completion; deferred + // and superseded attempts also advance its end timestamp. Remote + // workers record version coverage, so the witnesses can span nodes. if latest_cycle_end > previous_cycle_end && versions_observed { return Ok(latest_cycle_end); } if Instant::now() >= deadline { return Err(format!( - "enabled scanner did not complete an object-scanning cycle after {previous_cycle_end}: {observations:?}" + "enabled scanner did not complete a successful object-scanning cycle after {previous_cycle_end}: {observations:?}" ) .into()); } @@ -1008,7 +926,7 @@ mod tests { async fn test_cluster_root_heal_recovers_remote_shards_after_background_target_restart() -> Result<(), Box> { timeout( - Duration::from_secs(420), + Duration::from_secs(720), run_cluster_root_heal_interruption(InterruptionScenario::BackgroundTargetRestart), ) .await? @@ -1018,7 +936,7 @@ mod tests { async fn test_cluster_root_heal_recovers_remote_shards_after_background_target_crash() -> Result<(), Box> { timeout( - Duration::from_secs(420), + Duration::from_secs(720), run_cluster_root_heal_interruption(InterruptionScenario::BackgroundTargetCrash), ) .await? @@ -1028,7 +946,7 @@ mod tests { async fn test_cluster_root_heal_recovers_ec84_shards_after_background_target_restart() -> Result<(), Box> { timeout( - Duration::from_secs(420), + Duration::from_secs(720), run_cluster_root_heal_interruption(InterruptionScenario::BackgroundTargetRestartEc84), ) .await? @@ -1038,7 +956,7 @@ mod tests { async fn test_cluster_root_heal_recovers_ec84_shards_after_background_target_crash() -> Result<(), Box> { timeout( - Duration::from_secs(420), + Duration::from_secs(720), run_cluster_root_heal_interruption(InterruptionScenario::BackgroundTargetCrashEc84), ) .await? @@ -1048,7 +966,7 @@ mod tests { async fn test_cluster_root_heal_recovers_remote_shards_after_coordinator_restart() -> Result<(), Box> { timeout( - Duration::from_secs(420), + Duration::from_secs(720), run_cluster_root_heal_interruption(InterruptionScenario::BackgroundCoordinatorRestart), ) .await? @@ -1155,10 +1073,19 @@ mod tests { let server_rust_log = std::env::var("RUSTFS_HEAL_CHAOS_SERVER_RUST_LOG") .unwrap_or_else(|_| "rustfs::heal::task=info,rustfs=error".to_string()); cluster.set_env("RUST_LOG", server_rust_log); - let log_dir = std::env::var("RUSTFS_HEAL_CHAOS_LOG_DIR").unwrap_or_else(|_| format!("{}/logs", cluster.temp_dir)); + let log_dir = if let Some(directory) = std::env::var_os("RUSTFS_HEAL_CHAOS_LOG_DIR") { + PathBuf::from(directory) + } else if let Some(directory) = std::env::var_os("RUSTFS_E2E_LOG_DIR") { + let cluster_name = Path::new(&cluster.temp_dir) + .file_name() + .ok_or("cluster directory has no name")?; + PathBuf::from(directory).join(cluster_name).join("heal") + } else { + PathBuf::from(&cluster.temp_dir).join("logs") + }; std::fs::create_dir_all(&log_dir)?; for node_index in 0..cluster.nodes.len() { - cluster.set_node_capture_log_path(node_index, format!("{log_dir}/node{node_index}.log"))?; + cluster.set_node_capture_log_path(node_index, log_dir.join(format!("node{node_index}.log")).to_string_lossy())?; } cluster.start_with_binary(&server_binary).await?; let clients = cluster.create_all_clients()?; @@ -1435,7 +1362,7 @@ mod tests { let pre_interrupt_status: serde_json::Value = serde_json::from_str(&pre_interrupt_status_body) .map_err(|err| format!("pre-interrupt background heal status is not JSON ({err}): {pre_interrupt_status_body}"))?; let pre_interrupt_replacement = replacement_recovery_status(&cluster).await?; - let coordinator_log = std::fs::read_to_string(format!("{log_dir}/node0.log"))?; + let coordinator_log = std::fs::read_to_string(log_dir.join("node0.log"))?; assert!( coordinator_log .lines() @@ -1797,32 +1724,18 @@ mod tests { if let Some(evidence_context) = evidence_run { let restarted_pid = cluster.nodes[1].process.as_ref().ok_or("restarted target is absent")?.id(); assert_ne!(target_pid, restarted_pid, "target must be a new process"); - assert_eq!( - file_sha256(&server_binary)?, - evidence_context.run.binary.sha256, - "server build changed during restart" - ); - let evidence = serde_json::json!({ - "schema": 1, "case": evidence_context.case.id, "evidence": evidence_context.case.evidence, - "run_id": evidence_context.run.run_id, "source_revision": evidence_context.run.source_revision, - "test_build": compiled_test_identity(), - "binary_sha256": evidence_context.run.binary.sha256, - "test_binary_sha256": evidence_context.run.test_binary.sha256, - "topology": {"nodes": cluster.nodes.len(), "drives_per_node": cluster.nodes[0].data_dirs.len()}, - "pid_before": target_pid, "pid_after": restarted_pid, - "unclean_shutdown_marker": unclean_shutdown_marker_observed.unwrap_or(false), - "objects": evidence_objects, "node_listings": node_listings, - }); - let data = serde_json::to_vec(&evidence)?; - if data.len() > 1024 * 1024 { - return Err("scanner/heal oracle exceeds the 1 MiB artifact budget".into()); - } - let mut output = std::fs::OpenOptions::new() - .write(true) - .create_new(true) - .open(evidence_context.directory.join(evidence_context.case.oracle))?; - output.write_all(&data)?; - output.sync_all()?; + evidence_context.write( + &server_binary, + RestartObservation { + nodes: cluster.nodes.len(), + drives_per_node: cluster.nodes[0].data_dirs.len(), + pid_before: target_pid, + pid_after: restarted_pid, + unclean_shutdown_marker: unclean_shutdown_marker_observed.ok_or("missing shutdown marker observation")?, + objects: evidence_objects, + node_listings, + }, + )?; } Ok(()) diff --git a/crates/e2e_test/src/inline_fast_path_cluster_test.rs b/crates/e2e_test/src/inline_fast_path_cluster_test.rs index b568364d7..18e38cfa3 100644 --- a/crates/e2e_test/src/inline_fast_path_cluster_test.rs +++ b/crates/e2e_test/src/inline_fast_path_cluster_test.rs @@ -21,7 +21,7 @@ //! One S3 GET can select readers on multiple EC nodes, so the counter tracks //! distributed reader selection rather than HTTP request count. -use crate::common::{RustFSTestClusterEnvironment, RustFSTestEnvironment, init_logging}; +use crate::common::{RustFSTestClusterEnvironment, RustFSTestEnvironment, init_logging, signal_process}; use aws_sdk_s3::Client; use aws_sdk_s3::primitives::ByteStream; use aws_sdk_s3::types::{ @@ -2207,6 +2207,33 @@ async fn four_node_manual_transition_job_status_survives_node_restart() -> TestR Ok(()) } +struct SuspendedTransitionTarget<'a> { + // Keep the owned child borrowed until it is resumed so its PID cannot be reused. + child: &'a std::process::Child, + suspended: bool, +} + +impl<'a> SuspendedTransitionTarget<'a> { + fn suspend(child: &'a std::process::Child) -> TestResult { + signal_process(child.id(), "STOP")?; + Ok(Self { child, suspended: true }) + } + + fn resume(&mut self) -> TestResult { + signal_process(self.child.id(), "CONT")?; + self.suspended = false; + Ok(()) + } +} + +impl Drop for SuspendedTransitionTarget<'_> { + fn drop(&mut self) { + if self.suspended { + let _ = signal_process(self.child.id(), "CONT"); + } + } +} + #[tokio::test] async fn four_node_manual_transition_distributed_admission_conflict_reports_status_and_backpressure() -> TestResult { init_logging(); @@ -2244,7 +2271,21 @@ async fn four_node_manual_transition_distributed_admission_conflict_reports_stat .send() .await?; } - put_lifecycle_with_transition_retry(&hot_client, &bucket, &tier_name).await?; + // Lifecycle PUT starts its own backfill. Keep its first page on a separate + // node and stop it at queue backpressure before it reaches the tested prefix: + // one active worker, one queued item, then the first rejected item. + for index in 0u8..3 { + hot_client + .put_object() + .bucket(&bucket) + .key(format!("transition/automatic-admission/object-{index:02}.bin")) + .body(ByteStream::from(payload(KIB, index))) + .send() + .await?; + } + let mut suspended_cold = SuspendedTransitionTarget::suspend(cold.process.as_ref().ok_or("cold-tier process missing")?)?; + let lifecycle_client = hot.create_s3_client(2)?; + put_lifecycle_with_transition_retry(&lifecycle_client, &bucket, &tier_name).await?; let (node0, node1) = tokio::join!( start_manual_transition_job_on_node(&hot, 0, &bucket, prefix, &tier_name, false, 64), @@ -2304,6 +2345,31 @@ async fn four_node_manual_transition_distributed_admission_conflict_reports_stat assert_eq!(status["job_id"].as_str(), Some(job_id)); assert_eq!(status["status_endpoint"].as_str(), Some(status_endpoint)); + let deadline = Instant::now() + Duration::from_secs(30); + loop { + let status = read_manual_transition_job_status_endpoint(&hot, accepted.0, status_endpoint).await?; + assert_eq!( + status["status"].as_str(), + Some("running"), + "blocked cold tier must keep the admitted job running: {status}" + ); + if status["report"]["skipped_queue_full"].as_u64().is_some_and(|count| count > 0) { + assert!( + status["report"]["enqueued"].as_u64().is_some_and(|count| count > 0), + "the job must own pending transitions while the cold tier is suspended: {status}" + ); + break; + } + if Instant::now() >= deadline { + return Err(format!( + "manual transition job did not reach queue backpressure while the cold tier was suspended: {status}" + ) + .into()); + } + sleep(Duration::from_millis(50)).await; + } + + suspended_cold.resume()?; let terminal = wait_for_manual_transition_job_terminal(&hot, conflict.0, job_id, false).await?; assert_eq!(terminal["job_id"].as_str(), Some(job_id)); assert_eq!(terminal["bucket"].as_str(), Some(bucket.as_str())); diff --git a/crates/e2e_test/src/lib.rs b/crates/e2e_test/src/lib.rs index 72f5121fb..61495cafd 100644 --- a/crates/e2e_test/src/lib.rs +++ b/crates/e2e_test/src/lib.rs @@ -23,6 +23,9 @@ pub mod common; #[cfg(test)] pub mod chaos; +#[cfg(test)] +mod scanner_heal_evidence; + // Programmable S3 target for replication failure-path tests (backlog#1147 repl-8) // and on-demand-migration source scenarios (backlog#2151). #[cfg(test)] diff --git a/crates/e2e_test/src/scanner_heal_evidence.rs b/crates/e2e_test/src/scanner_heal_evidence.rs new file mode 100644 index 000000000..780eac9df --- /dev/null +++ b/crates/e2e_test/src/scanner_heal_evidence.rs @@ -0,0 +1,181 @@ +// Copyright 2026 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +//! Build-bound evidence for scanner and heal restart tests. + +use crate::common::ClusterTopology; +use sha2::{Digest, Sha256}; +use std::error::Error; +use std::io::{Read, Write}; +use std::path::{Path, PathBuf}; + +#[derive(serde::Deserialize)] +struct EvidenceBuild { + sha256: String, +} + +#[derive(serde::Deserialize)] +struct RestartEvidenceRun { + schema: u32, + run_id: String, + source_revision: String, + test_build: serde_json::Value, + binary: EvidenceBuild, + test_binary: EvidenceBuild, +} + +#[derive(Clone, Copy)] +pub(crate) struct ScannerHealEvidenceCase { + pub(crate) id: &'static str, + pub(crate) oracle: &'static str, + pub(crate) evidence: &'static str, + pub(crate) unclean_shutdown_marker: bool, + pub(crate) topology: EvidenceTopology, + pub(crate) storage_class_standard: Option<&'static str>, + pub(crate) erasure_set_drive_count: Option<&'static str>, +} + +#[derive(Clone, Copy)] +pub(crate) struct EvidenceTopology { + pub(crate) nodes: usize, + pub(crate) drives_per_node: usize, +} + +impl EvidenceTopology { + pub(crate) const fn new(nodes: usize, drives_per_node: usize) -> Self { + Self { nodes, drives_per_node } + } + + pub(crate) fn total_drives(self) -> usize { + self.nodes * self.drives_per_node + } + + pub(crate) fn cluster_topology(self) -> ClusterTopology { + ClusterTopology::single_pool_multidrive(self.nodes, self.drives_per_node) + } +} + +pub(crate) struct RestartEvidenceContext { + directory: PathBuf, + run: RestartEvidenceRun, + case: ScannerHealEvidenceCase, +} + +fn file_sha256(path: &Path) -> Result> { + let mut file = std::fs::File::open(path)?; + let mut digest = Sha256::new(); + let mut buffer = [0_u8; 64 * 1024]; + loop { + let read = file.read(&mut buffer)?; + if read == 0 { + break; + } + digest.update(&buffer[..read]); + } + Ok(digest.finalize().iter().map(|byte| format!("{byte:02x}")).collect()) +} + +pub(crate) fn restart_evidence_run( + binary: &Path, + case: ScannerHealEvidenceCase, +) -> Result, Box> { + let Some(directory) = std::env::var_os("RUSTFS_SCANNER_HEAL_RUN_DIR") else { + return Ok(None); + }; + if case.id.is_empty() + || case.oracle.is_empty() + || !case.oracle.ends_with(".json") + || case.oracle.contains('/') + || case.oracle.contains('\\') + || case.oracle.contains("..") + || !matches!(case.evidence, "process-restart" | "process-crash-restart") + || (case.evidence == "process-crash-restart") != case.unclean_shutdown_marker + { + return Err("invalid scanner/heal evidence case".into()); + } + let directory = PathBuf::from(directory); + let receipt = directory.join("run.json"); + if receipt.metadata()?.len() > 1024 * 1024 { + return Err("oversized scanner/heal execution receipt".into()); + } + let run: RestartEvidenceRun = serde_json::from_slice(&std::fs::read(receipt)?)?; + if run.schema != 1 || run.run_id.len() != 32 || run.source_revision.len() != 40 { + return Err("invalid scanner/heal execution identity".into()); + } + let built = compiled_test_identity(); + for key in ["source_revision", "dirty", "lock_blob", "features"] { + assert_eq!(built[key], run.test_build[key], "compiled test identity differs for {key}"); + } + assert_eq!(file_sha256(binary)?, run.binary.sha256, "server binary must match the run receipt"); + assert_eq!( + file_sha256(&std::env::current_exe()?)?, + run.test_binary.sha256, + "test executable must match the run receipt" + ); + if directory.join(case.oracle).exists() { + return Err("scanner/heal oracle already exists; create a new execution receipt".into()); + } + Ok(Some(RestartEvidenceContext { directory, run, case })) +} + +fn compiled_test_identity() -> serde_json::Value { + serde_json::json!({ + "source_revision": env!("RUSTFS_E2E_BUILD_COMMIT"), + "dirty": env!("RUSTFS_E2E_BUILD_DIRTY") != "false", + "lock_blob": env!("RUSTFS_E2E_BUILD_LOCK"), + "features": env!("RUSTFS_E2E_BUILD_FEATURES"), + "target": env!("RUSTFS_E2E_BUILD_TARGET"), + "profile": env!("RUSTFS_E2E_BUILD_PROFILE"), + "rustflags_hex": env!("RUSTFS_E2E_BUILD_RUSTFLAGS_HEX"), + }) +} + +pub(crate) struct RestartObservation { + pub(crate) nodes: usize, + pub(crate) drives_per_node: usize, + pub(crate) pid_before: u32, + pub(crate) pid_after: u32, + pub(crate) unclean_shutdown_marker: bool, + pub(crate) objects: Vec, + pub(crate) node_listings: Vec>, +} + +impl RestartEvidenceContext { + pub(crate) fn write(self, binary: &Path, observed: RestartObservation) -> Result<(), Box> { + assert_ne!(observed.pid_before, observed.pid_after, "target must be a new process"); + assert_eq!(file_sha256(binary)?, self.run.binary.sha256, "server build changed during restart"); + let evidence = serde_json::json!({ + "schema": 1, "case": self.case.id, "evidence": self.case.evidence, + "run_id": self.run.run_id, "source_revision": self.run.source_revision, + "test_build": compiled_test_identity(), + "binary_sha256": self.run.binary.sha256, + "test_binary_sha256": self.run.test_binary.sha256, + "topology": {"nodes": observed.nodes, "drives_per_node": observed.drives_per_node}, + "pid_before": observed.pid_before, "pid_after": observed.pid_after, + "unclean_shutdown_marker": observed.unclean_shutdown_marker, + "objects": observed.objects, "node_listings": observed.node_listings, + }); + let data = serde_json::to_vec(&evidence)?; + if data.len() > 1024 * 1024 { + return Err("scanner/heal oracle exceeds the 1 MiB artifact budget".into()); + } + let mut output = std::fs::OpenOptions::new() + .write(true) + .create_new(true) + .open(self.directory.join(self.case.oracle))?; + output.write_all(&data)?; + output.sync_all()?; + Ok(()) + } +} diff --git a/docs/testing/ci-gates.md b/docs/testing/ci-gates.md index a7087d736..20e9a3c0e 100644 --- a/docs/testing/ci-gates.md +++ b/docs/testing/ci-gates.md @@ -50,7 +50,7 @@ Promotion rule: never promote a report-only lane to required from one green run. | 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 | | PR touching `paths` in `coverage.yml` | `Workspace line coverage` | `coverage.yml` `coverage` | Report-only | `make coverage`; `python3 scripts/check_security_coverage.py target/llvm-cov/coverage.json` | -| PR touching `paths` in `e2e-upgrade.yml` | `Direct upgrade from the previous release`, `Mixed-version rolling upgrade from the previous release`, `Bucket configuration survives the upgrade`, `Rollback reads current bucket metadata` | `e2e-upgrade.yml` `upgrade` matrix | Report-only | the `cargo test --locked -p e2e_test` command in the job with `RUSTFS_UPGRADE_SOURCE_BINARY` pointing at the pinned previous release (`UPGRADE_SOURCE_VERSION`) | +| PR touching `paths` in `e2e-upgrade.yml` | `Direct upgrade from the previous release`, `Mixed-version rolling upgrade from the previous release`, `Bucket configuration survives the upgrade`, `Rollback reads current bucket metadata`, `ODM configuration recovery after rc.5 rollback`, `Multipart layouts survive the rc.5 upgrade`, `rc.5 multipart replication baseline` | `e2e-upgrade.yml` `upgrade` matrix | Report-only | the `cargo test --locked -p e2e_test` command in the job with `RUSTFS_UPGRADE_SOURCE_BINARY` pointing at the pinned previous release (`UPGRADE_SOURCE_VERSION`) | | PR touching `paths` in `oidc-keycloak.yml` | `OIDC Keycloak live gate` | `oidc-keycloak.yml` `oidc-keycloak-live` | Report-only | `cargo build --locked -p rustfs --bin rustfs`, then `bash scripts/test/oidc_keycloak_live.sh ./target/debug/rustfs` | | PR touching `paths` in `targets-integration.yml` | `PostgreSQL, MySQL, AMQP, and NATS` | `targets-integration.yml` `targets-live` | Report-only | start the containers as in the job, export the `RUSTFS_TEST_*` DSNs, then the job's `cargo test --locked -p rustfs-targets --test -- --ignored --test-threads=1` commands | | PR limited to main-CI-excluded paths | `Quick Checks`, `Test and Lint` | `ci-docs-only.yml` `quick-checks`, `test-and-lint` | Required | `git diff --check`; `make doc-paths-check`; `scripts/check_no_planning_docs.sh` | @@ -86,7 +86,7 @@ Scheduled lanes never block a PR. Their workflow-local gate fails the run, sched | `mint.yml` (weekly) | `mint` | report-only by design; per-suite PASS/FAIL/NA and raw `log.json` | yes | pinned Docker sequence in the workflow | | `coverage.yml` (weekly) | `coverage` | report-only trend; lcov and JSON artifact | yes | `make coverage` | | `runner-hygiene.yml` (monthly) | `check-ephemerality` | runner ephemerality | yes | dispatch | -| `e2e-upgrade.yml` (weekly) | `upgrade` (4-case matrix) | upgrade and rollback gate; server logs | no | see the PR row | +| `e2e-upgrade.yml` (weekly) | `upgrade` (7-case matrix) | upgrade and rollback gate; server logs | no | see the PR row | | `oidc-keycloak.yml` (weekly) | `oidc-keycloak-live` | live OIDC gate | no | see the PR row | | `targets-integration.yml` (nightly) | `targets-live` | live target gate; container logs | no | see the PR row | | `scheduled-validation-freshness.yml` (nightly) | `check-freshness` | fails on a never-created or stale schedule | n/a | dispatch |