mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-06 12:09:12 +00:00
Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 3340a755c1 | |||
| b4e0838b01 |
@@ -745,4 +745,593 @@ mod tests {
|
||||
sleep(Duration::from_millis(25)).await;
|
||||
}
|
||||
}
|
||||
|
||||
/// A stopped, initialized pool loses only its pool.bin object replicas.
|
||||
/// Restart must repair both pools through their actual conditional writes.
|
||||
#[tokio::test]
|
||||
async fn test_two_pool_restart_repairs_missing_pool_metadata_via_bootstrap_cas() -> TestResult {
|
||||
use crate::common::ClusterTopology;
|
||||
use futures::FutureExt;
|
||||
use std::path::PathBuf;
|
||||
|
||||
init_logging();
|
||||
let run = uuid::Uuid::new_v4().to_string();
|
||||
let artifact = std::env::var_os("RUSTFS_E2E_STARTUP_CAS_ARTIFACT_DIR")
|
||||
.map(PathBuf::from)
|
||||
.unwrap_or_else(std::env::temp_dir)
|
||||
.join(format!("repair-startup-cas-{run}"));
|
||||
std::fs::create_dir_all(&artifact)?;
|
||||
let binary_dir = tempfile::tempdir()?;
|
||||
let binary = prepare_startup_cas_binary(binary_dir.path(), &artifact)?;
|
||||
probe_startup_cas_binary(&binary, &run, &artifact).await?;
|
||||
let mut owned = RepairStartupCluster(Some(
|
||||
RustFSTestClusterEnvironment::with_topology(ClusterTopology::per_node_pools(2, vec![vec![0], vec![1]])).await?,
|
||||
));
|
||||
let cluster = owned.0.as_mut().expect("fixture owns its cluster");
|
||||
assert_eq!(cluster.nodes.len(), 2);
|
||||
for (pool, node) in cluster.nodes.iter().enumerate() {
|
||||
assert_eq!(node.pool_idx, pool);
|
||||
assert_eq!(node.data_dirs.len(), 2);
|
||||
}
|
||||
let volumes = cluster.rustfs_volumes_arg();
|
||||
assert_eq!(volumes.split_whitespace().count(), 2);
|
||||
cluster.set_env("RUSTFS_OBS_LOG_DIRECTORY", "");
|
||||
cluster.set_env("RUSTFS_OBS_LOG_STDOUT_ENABLED", "true");
|
||||
cluster.set_env("RUST_LOG", "rustfs=info,rustfs_ecstore=trace");
|
||||
for key in [
|
||||
"HTTP_PROXY",
|
||||
"HTTPS_PROXY",
|
||||
"ALL_PROXY",
|
||||
"http_proxy",
|
||||
"https_proxy",
|
||||
"all_proxy",
|
||||
] {
|
||||
cluster.set_env(key, "");
|
||||
}
|
||||
for key in ["NO_PROXY", "no_proxy"] {
|
||||
cluster.set_env(key, "127.0.0.1,localhost");
|
||||
}
|
||||
let attempt = tokio::time::timeout(
|
||||
Duration::from_secs(600),
|
||||
std::panic::AssertUnwindSafe(async {
|
||||
let mut first = RepairStartupPhase::new(cluster, &artifact, "initial")?;
|
||||
let previous = run_repair_startup_phase(cluster, &binary, &mut first, None).await?;
|
||||
let bucket = format!("repair-cas-{run}");
|
||||
let key = "preserved/body";
|
||||
let body = vec![0x6bu8; 256 * 1024];
|
||||
tokio::time::timeout(Duration::from_secs(20), async {
|
||||
cluster.create_test_bucket(&bucket).await?;
|
||||
cluster.create_all_clients()?[0]
|
||||
.put_object()
|
||||
.bucket(&bucket)
|
||||
.key(key)
|
||||
.body(ByteStream::from(body.clone()))
|
||||
.send()
|
||||
.await?;
|
||||
repair_full_get(cluster, &bucket, key, &body).await
|
||||
})
|
||||
.await
|
||||
.map_err(std::io::Error::other)??;
|
||||
|
||||
// Retain each Child until wait confirms exit; only then mutate
|
||||
// the stopped fixture's exact pool1 object subtrees.
|
||||
let stopped = stop_repair_cluster(cluster)?;
|
||||
assert_eq!(stopped.len(), 2, "both initial children must be reaped before disk mutation");
|
||||
std::fs::write(artifact.join("initial-stopped.json"), serde_json::to_vec_pretty(&stopped)?)?;
|
||||
let mut before = Vec::new();
|
||||
for (pool, node) in cluster.nodes.iter().enumerate() {
|
||||
for disk in &node.data_dirs {
|
||||
let root = std::fs::canonicalize(disk)?;
|
||||
assert!(root.starts_with(std::fs::canonicalize(&cluster.temp_dir)?));
|
||||
let snapshot = repair_disk_snapshot(&root)?;
|
||||
assert!(snapshot.contains_key(std::path::Path::new(".rustfs.sys/format.json")));
|
||||
assert!(snapshot.contains_key(std::path::Path::new(".rustfs.sys/pool.bin.identity/xl.meta")));
|
||||
assert!(snapshot.contains_key(std::path::Path::new(".rustfs.sys/pool.bin/xl.meta")));
|
||||
before.push((pool, root, snapshot));
|
||||
}
|
||||
}
|
||||
std::fs::write(artifact.join("before-removal.json"), serde_json::to_vec_pretty(&before)?)?;
|
||||
for (pool, root, _) in &before {
|
||||
if *pool == 1 {
|
||||
let object = root.join(".rustfs.sys/pool.bin");
|
||||
assert_eq!(std::fs::canonicalize(&object)?, object, "no aliased deletion target");
|
||||
std::fs::remove_dir_all(&object)?;
|
||||
assert!(!object.try_exists()?);
|
||||
}
|
||||
}
|
||||
for (pool, root, snapshot) in &before {
|
||||
let mut expected = snapshot.clone();
|
||||
if *pool == 1 {
|
||||
expected.retain(|path, _| !path.starts_with(".rustfs.sys/pool.bin"));
|
||||
}
|
||||
assert_eq!(repair_disk_snapshot(root)?, expected, "only pool1's complete pool.bin objects may change");
|
||||
}
|
||||
assert_eq!(cluster.rustfs_volumes_arg(), volumes, "repair reuses identical topology, ports and roots");
|
||||
let mut restart = RepairStartupPhase::new(cluster, &artifact, "repair")?;
|
||||
assert_ne!(restart.nonce, first.nonce);
|
||||
let repaired = run_repair_startup_phase(cluster, &binary, &mut restart, Some(&previous)).await?;
|
||||
tokio::time::timeout(Duration::from_secs(20), repair_full_get(cluster, &bucket, key, &body))
|
||||
.await
|
||||
.map_err(std::io::Error::other)??;
|
||||
std::fs::write(
|
||||
artifact.join("repair-proof.json"),
|
||||
serde_json::to_vec_pretty(&serde_json::json!({
|
||||
"initial": previous, "repair": repaired, "volumes": volumes,
|
||||
"negative_control": "NOT_RUN", "body_length": body.len(), "full_get_nodes": [0, 1],
|
||||
}))?,
|
||||
)?;
|
||||
Ok::<(), Box<dyn Error + Send + Sync>>(())
|
||||
})
|
||||
.catch_unwind(),
|
||||
)
|
||||
.await;
|
||||
// Borrowed startup futures and their release guards have dropped here.
|
||||
// Attempt every child even when an earlier wait or assertion failed.
|
||||
let stopped = stop_repair_cluster(cluster);
|
||||
eprintln!("two-pool repair CAS evidence: {}", artifact.display());
|
||||
match attempt {
|
||||
Ok(Ok(result)) => {
|
||||
result?;
|
||||
std::fs::write(artifact.join("repair-stopped.json"), serde_json::to_vec_pretty(&stopped?)?)?;
|
||||
Ok(())
|
||||
}
|
||||
Ok(Err(panic)) => std::panic::resume_unwind(panic),
|
||||
Err(error) => Err(std::io::Error::other(format!("two-pool repair fixture deadline: {error}")).into()),
|
||||
}
|
||||
}
|
||||
|
||||
// Unlike the general harness stop_node, keep a failed Child wait attached
|
||||
// so Drop can retry without deleting a possibly active disk directory.
|
||||
struct RepairStartupCluster(Option<RustFSTestClusterEnvironment>);
|
||||
|
||||
impl Drop for RepairStartupCluster {
|
||||
fn drop(&mut self) {
|
||||
if let Some(mut cluster) = self.0.take() {
|
||||
if let Err(error) = stop_repair_cluster(&mut cluster) {
|
||||
eprintln!("repair child cleanup failed; preserving {}: {error}", cluster.temp_dir);
|
||||
}
|
||||
if cluster.nodes.iter().any(|node| node.process.is_some()) {
|
||||
std::mem::forget(cluster);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn stop_repair_cluster(cluster: &mut RustFSTestClusterEnvironment) -> std::io::Result<Vec<serde_json::Value>> {
|
||||
let mut stopped = Vec::new();
|
||||
let mut failure = None;
|
||||
for (node, state) in cluster.nodes.iter_mut().enumerate() {
|
||||
let Some(child) = state.process.as_mut() else { continue };
|
||||
let pid = child.id();
|
||||
let deadline = std::time::Instant::now() + Duration::from_secs(5);
|
||||
let waited = (|| {
|
||||
if child.try_wait()?.is_none() {
|
||||
// Exit can race kill; the actual wait below decides whether
|
||||
// the child is gone, rather than a successful signal alone.
|
||||
let _ = child.kill();
|
||||
}
|
||||
loop {
|
||||
if child.try_wait()?.is_some() {
|
||||
return child.wait();
|
||||
}
|
||||
if std::time::Instant::now() >= deadline {
|
||||
return Err(std::io::Error::other(format!("node {node} pid {pid} did not exit")));
|
||||
}
|
||||
std::thread::sleep(Duration::from_millis(10));
|
||||
}
|
||||
})();
|
||||
match waited {
|
||||
Ok(status) => {
|
||||
stopped.push(serde_json::json!({"node": node, "pid": pid, "status": status.to_string()}));
|
||||
state.process = None;
|
||||
}
|
||||
Err(error) => {
|
||||
failure.get_or_insert(error);
|
||||
}
|
||||
}
|
||||
}
|
||||
failure.map_or(Ok(stopped), Err)
|
||||
}
|
||||
|
||||
fn repair_disk_snapshot(
|
||||
root: &std::path::Path,
|
||||
) -> std::io::Result<std::collections::BTreeMap<std::path::PathBuf, (u64, String)>> {
|
||||
let mut snapshot = std::collections::BTreeMap::new();
|
||||
let mut pending = vec![std::path::PathBuf::new()];
|
||||
while let Some(relative) = pending.pop() {
|
||||
if relative.components().count() > 64 || snapshot.len() > 10_000 {
|
||||
return Err(std::io::Error::other("repair fixture snapshot exceeds its bound"));
|
||||
}
|
||||
for entry in std::fs::read_dir(root.join(&relative))? {
|
||||
let entry = entry?;
|
||||
let path = relative.join(entry.file_name());
|
||||
let metadata = std::fs::symlink_metadata(entry.path())?;
|
||||
if metadata.is_dir() {
|
||||
snapshot.insert(path.clone(), (0, "directory".to_owned()));
|
||||
pending.push(path);
|
||||
} else if metadata.is_file() {
|
||||
snapshot.insert(path, (metadata.len(), startup_cas_sha256(&entry.path())?));
|
||||
} else {
|
||||
return Err(std::io::Error::other("repair fixture contains a symlink or special file"));
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(snapshot)
|
||||
}
|
||||
|
||||
async fn repair_full_get(cluster: &RustFSTestClusterEnvironment, bucket: &str, key: &str, body: &[u8]) -> TestResult {
|
||||
for (node, client) in cluster.create_all_clients()?.iter().enumerate() {
|
||||
let received = client
|
||||
.get_object()
|
||||
.bucket(bucket)
|
||||
.key(key)
|
||||
.send()
|
||||
.await?
|
||||
.body
|
||||
.collect()
|
||||
.await?
|
||||
.into_bytes();
|
||||
assert_eq!(received.as_ref(), body, "node {node} must return every preserved object byte");
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
struct RepairStartupPhase {
|
||||
nonce: String,
|
||||
artifact: std::path::PathBuf,
|
||||
logs: Vec<std::path::PathBuf>,
|
||||
disks: Vec<Vec<std::path::PathBuf>>,
|
||||
endpoints: Vec<Vec<String>>,
|
||||
releases: StartupCasReleases,
|
||||
}
|
||||
|
||||
impl RepairStartupPhase {
|
||||
fn new(
|
||||
cluster: &mut RustFSTestClusterEnvironment,
|
||||
artifact: &std::path::Path,
|
||||
phase: &str,
|
||||
) -> Result<Self, Box<dyn Error + Send + Sync>> {
|
||||
let artifact = artifact.join(phase);
|
||||
std::fs::create_dir(&artifact)?;
|
||||
let nonce = uuid::Uuid::new_v4().to_string();
|
||||
cluster.set_env("RUSTFS_E2E_STARTUP_CAS_NONCE", &nonce);
|
||||
let mut result = Self {
|
||||
nonce,
|
||||
artifact,
|
||||
logs: Vec::new(),
|
||||
disks: Vec::new(),
|
||||
endpoints: Vec::new(),
|
||||
releases: StartupCasReleases(Vec::new()),
|
||||
};
|
||||
for node in 0..2 {
|
||||
let log = result.artifact.join(format!("node-{node}.log"));
|
||||
let release = result.artifact.join(format!("release-{node}"));
|
||||
assert!(!release.try_exists()?, "each startup has a new unreleased gate");
|
||||
cluster.set_node_capture_log_path(node, log.to_string_lossy())?;
|
||||
cluster.set_node_env(node, "RUSTFS_E2E_STARTUP_CAS_RELEASE", release.to_string_lossy())?;
|
||||
result
|
||||
.disks
|
||||
.push(cluster.nodes[node].data_dirs.iter().map(std::path::PathBuf::from).collect());
|
||||
result.endpoints.push(
|
||||
cluster.nodes[node]
|
||||
.data_dirs
|
||||
.iter()
|
||||
.map(|disk| format!("http://{}{disk}", cluster.nodes[node].address))
|
||||
.collect(),
|
||||
);
|
||||
result.logs.push(log);
|
||||
result.releases.0.push(release);
|
||||
}
|
||||
Ok(result)
|
||||
}
|
||||
}
|
||||
|
||||
async fn run_repair_startup_phase(
|
||||
cluster: &mut RustFSTestClusterEnvironment,
|
||||
binary: &std::path::Path,
|
||||
phase: &mut RepairStartupPhase,
|
||||
previous: Option<&serde_json::Value>,
|
||||
) -> Result<serde_json::Value, Box<dyn Error + Send + Sync>> {
|
||||
let mut startup = Box::pin(cluster.start_with_binary(binary));
|
||||
let mut observer = Box::pin(wait_two_pool_startup_cas(phase, previous));
|
||||
let (observed, finished) = tokio::select! {
|
||||
observed = tokio::time::timeout(Duration::from_secs(120), &mut observer) => {
|
||||
(observed.map_err(std::io::Error::other).and_then(|result| result), false)
|
||||
}
|
||||
result = &mut startup => {
|
||||
(Err(std::io::Error::other(format!("startup ended before both unreleased repair gates: {result:?}"))), true)
|
||||
}
|
||||
};
|
||||
drop(observer);
|
||||
let released = phase.releases.release();
|
||||
let drained = if finished {
|
||||
Ok(())
|
||||
} else {
|
||||
tokio::time::timeout(Duration::from_secs(if observed.is_ok() { 60 } else { 5 }), &mut startup)
|
||||
.await
|
||||
.map_err(std::io::Error::other)
|
||||
.map_err(|error| -> Box<dyn Error + Send + Sync> { error.into() })
|
||||
.and_then(|result| result)
|
||||
};
|
||||
drop(startup);
|
||||
let observed = observed?;
|
||||
released?;
|
||||
drained?;
|
||||
for (node, process) in cluster.nodes.iter().enumerate() {
|
||||
assert_eq!(observed["pids"][node], process.process.as_ref().expect("live phase child").id());
|
||||
}
|
||||
Ok(observed)
|
||||
}
|
||||
|
||||
async fn wait_two_pool_startup_cas(
|
||||
phase: &RepairStartupPhase,
|
||||
previous: Option<&serde_json::Value>,
|
||||
) -> std::io::Result<serde_json::Value> {
|
||||
use serde_json::{Value, json};
|
||||
loop {
|
||||
let logs = phase
|
||||
.logs
|
||||
.iter()
|
||||
.map(|path| startup_cas_log(path))
|
||||
.collect::<std::io::Result<Vec<_>>>()?;
|
||||
let events: Vec<Vec<&Value>> = logs
|
||||
.iter()
|
||||
.map(|log| log.iter().filter(|event| event["nonce"] == phase.nonce).collect())
|
||||
.collect();
|
||||
if events
|
||||
.iter()
|
||||
.flat_map(|events| events.iter())
|
||||
.any(|event| event["kind"] == "cas" && event["ok"] == false)
|
||||
{
|
||||
return Err(std::io::Error::other("normal two-pool startup produced a failed CAS; inspect phase logs"));
|
||||
}
|
||||
if !events.iter().all(|events| {
|
||||
events
|
||||
.iter()
|
||||
.any(|event| event["kind"] == "gate" && event["slot_installed"] == false)
|
||||
}) {
|
||||
sleep(Duration::from_millis(25)).await;
|
||||
continue;
|
||||
}
|
||||
let mut pids = Vec::new();
|
||||
for records in &events {
|
||||
let ready: Vec<_> = records.iter().filter(|event| event["kind"] == "observer-ready").collect();
|
||||
assert_eq!(ready.len(), 1, "one real process per phase log");
|
||||
assert!(ready[0]["pid"].as_u64().is_some());
|
||||
assert!(records.iter().all(|event| event["pid"] == ready[0]["pid"]));
|
||||
pids.push(ready[0]["pid"].clone());
|
||||
}
|
||||
assert_ne!(pids[0], pids[1]);
|
||||
let source = &events[0];
|
||||
let classified: Vec<_> = source
|
||||
.iter()
|
||||
.copied()
|
||||
.filter(|event| event["kind"] == "startup-classifier")
|
||||
.collect();
|
||||
assert_eq!(classified.len(), 1, "the normal startup must classify once, without retrying failures");
|
||||
let classified = classified[0];
|
||||
let attempt = &classified["attempt"];
|
||||
uuid::Uuid::parse_str(attempt.as_str().expect("real init attempt UUID")).expect("valid init attempt UUID");
|
||||
assert_eq!(classified["elected_writer"], true);
|
||||
assert_eq!(classified["needs_repair"], true);
|
||||
assert_eq!(classified["repair_write_safe"], true);
|
||||
assert_eq!(classified["topology_update"], previous.is_none());
|
||||
assert!(
|
||||
events[1]
|
||||
.iter()
|
||||
.all(|event| event["kind"] != "cas" || event["object"] != "pool.bin"),
|
||||
"only elected pool0 may persist pool.bin"
|
||||
);
|
||||
let initial: Vec<_> = source
|
||||
.iter()
|
||||
.copied()
|
||||
.filter(|event| {
|
||||
event["kind"] == "replica-read" && event["startup_phase"] == "load" && event["attempt"] == *attempt
|
||||
})
|
||||
.collect();
|
||||
assert_eq!(initial.len(), 2, "complete initial reads of both real pools");
|
||||
assert_eq!(initial[0]["batch"], initial[1]["batch"]);
|
||||
let mut prepares = Vec::new();
|
||||
let mut commits = Vec::new();
|
||||
for pool in 0..2 {
|
||||
let read: Vec<_> = initial.iter().filter(|event| event["pool"] == pool).collect();
|
||||
assert_eq!(read.len(), 1);
|
||||
let read = *read[0];
|
||||
if let Some(previous) = previous.filter(|_| pool == 0) {
|
||||
assert_eq!(read["state"], "valid");
|
||||
assert_eq!(read["committed"], true);
|
||||
for field in [
|
||||
"payload_sha256",
|
||||
"raw_sha256",
|
||||
"version",
|
||||
"cluster_id",
|
||||
"epoch",
|
||||
"generation",
|
||||
"transaction_id",
|
||||
"etag",
|
||||
] {
|
||||
assert_eq!(read[field], previous["replicas"][0][field], "pool0 baseline {field} is retained");
|
||||
}
|
||||
} else {
|
||||
assert_eq!(read["state"], "missing");
|
||||
assert_eq!(read["cas"], "missing");
|
||||
}
|
||||
let mut pair = Vec::new();
|
||||
for stage in ["prepare_cas", "commit_cas"] {
|
||||
let matching: Vec<_> = source
|
||||
.iter()
|
||||
.copied()
|
||||
.filter(|event| {
|
||||
event["kind"] == "cas"
|
||||
&& event["object"] == "pool.bin"
|
||||
&& event["phase"] == stage
|
||||
&& event["pool"] == pool
|
||||
&& event["attempt"] == *attempt
|
||||
&& event["startup_phase"] == "persist"
|
||||
})
|
||||
.collect();
|
||||
assert_eq!(matching.len(), 1, "exactly one successful {stage} on actual pool {pool}");
|
||||
let cas = matching[0];
|
||||
assert_eq!(cas["ok"], true);
|
||||
assert_eq!(cas["tail_drained"], true);
|
||||
assert_eq!(cas["no_lock"], true);
|
||||
assert!(cas["etag"].as_str().is_some_and(|value| !value.is_empty()));
|
||||
assert!(cas["mod_time"].as_str().is_some_and(|value| !value.is_empty()));
|
||||
pair.push(cas);
|
||||
}
|
||||
if previous.is_some() && pool == 0 {
|
||||
assert_eq!(read["cas"], "existing");
|
||||
assert_eq!(pair[0]["if_match"], read["etag"]);
|
||||
assert!(pair[0]["if_none_match"].is_null());
|
||||
} else {
|
||||
assert_eq!(pair[0]["if_none_match"], "*");
|
||||
assert!(pair[0]["if_match"].is_null());
|
||||
}
|
||||
assert_eq!(pair[1]["if_match"], pair[0]["etag"]);
|
||||
assert!(pair[1]["if_none_match"].is_null());
|
||||
assert_ne!(pair[0]["etag"], pair[1]["etag"]);
|
||||
assert_ne!(pair[0]["payload_sha256"], pair[1]["payload_sha256"]);
|
||||
prepares.push(pair[0]);
|
||||
commits.push(pair[1]);
|
||||
}
|
||||
assert_eq!(
|
||||
source
|
||||
.iter()
|
||||
.filter(|event| event["kind"] == "cas" && event["object"] == "pool.bin")
|
||||
.count(),
|
||||
4
|
||||
);
|
||||
assert_eq!(prepares[0]["payload_sha256"], prepares[1]["payload_sha256"]);
|
||||
assert_eq!(commits[0]["payload_sha256"], commits[1]["payload_sha256"]);
|
||||
let mut replicas = Vec::new();
|
||||
for pool in 0..2 {
|
||||
let matching: Vec<_> = source
|
||||
.iter()
|
||||
.copied()
|
||||
.filter(|event| {
|
||||
event["kind"] == "replica-read"
|
||||
&& event["startup_phase"] == "persist"
|
||||
&& event["attempt"] == *attempt
|
||||
&& event["pool"] == pool
|
||||
&& event["payload_sha256"] == commits[pool]["payload_sha256"]
|
||||
})
|
||||
.collect();
|
||||
assert_eq!(matching.len(), 1, "actual final complete decoded read for pool {pool}");
|
||||
let replica = matching[0];
|
||||
assert_eq!(replica["state"], "valid");
|
||||
assert_eq!(replica["committed"], true);
|
||||
assert_eq!(replica["pool_count"], 2);
|
||||
assert_eq!(replica["raw_sha256"], replica["payload_sha256"]);
|
||||
assert_eq!(replica["etag"], commits[pool]["etag"]);
|
||||
assert_eq!(replica["cas"], "existing");
|
||||
assert!(
|
||||
replica["cluster_id"]
|
||||
.as_str()
|
||||
.is_some_and(|value| uuid::Uuid::parse_str(value).is_ok())
|
||||
);
|
||||
assert!(
|
||||
replica["transaction_id"]
|
||||
.as_str()
|
||||
.is_some_and(|value| uuid::Uuid::parse_str(value).is_ok())
|
||||
);
|
||||
assert!(replica["epoch"].as_u64().is_some_and(|value| value > 0));
|
||||
let expected_generation = match previous {
|
||||
Some(previous) => {
|
||||
previous["replicas"][0]["generation"]
|
||||
.as_u64()
|
||||
.expect("previous decoded generation")
|
||||
+ 1
|
||||
}
|
||||
None => 1,
|
||||
};
|
||||
assert_eq!(replica["generation"], expected_generation);
|
||||
if let Some(previous) = previous {
|
||||
assert_eq!(replica["cluster_id"], previous["replicas"][0]["cluster_id"]);
|
||||
assert_eq!(replica["epoch"], previous["replicas"][0]["epoch"]);
|
||||
assert_ne!(replica["transaction_id"], previous["replicas"][0]["transaction_id"]);
|
||||
}
|
||||
replicas.push(replica);
|
||||
}
|
||||
for field in [
|
||||
"batch",
|
||||
"version",
|
||||
"cluster_id",
|
||||
"epoch",
|
||||
"generation",
|
||||
"transaction_id",
|
||||
"payload_sha256",
|
||||
] {
|
||||
assert_eq!(replicas[0][field], replicas[1][field], "same decoded final revision: {field}");
|
||||
}
|
||||
let confirmed: Vec<_> = source
|
||||
.iter()
|
||||
.filter(|event| {
|
||||
event["kind"] == "confirmed"
|
||||
&& event["attempt"] == *attempt
|
||||
&& event["payload_sha256"] == replicas[0]["payload_sha256"]
|
||||
&& event["generation"] == replicas[0]["generation"]
|
||||
&& event["transaction_id"] == replicas[0]["transaction_id"]
|
||||
})
|
||||
.collect();
|
||||
assert_eq!(confirmed.len(), 1);
|
||||
let mut receivers = Vec::new();
|
||||
for drive in 0..2 {
|
||||
for cas in [prepares[1], commits[1]] {
|
||||
let matching: Vec<_> = events[1]
|
||||
.iter()
|
||||
.copied()
|
||||
.filter(|event| {
|
||||
event["kind"] == "receiver"
|
||||
&& event["disk"] == phase.endpoints[1][drive]
|
||||
&& event["dst_volume"] == ".rustfs.sys"
|
||||
&& event["dst_path"] == "pool.bin"
|
||||
&& event["etag"] == cas["etag"]
|
||||
})
|
||||
.collect();
|
||||
assert!(matching.len() <= 1, "one actual receiver per disk/CAS");
|
||||
if let Some(receiver) = matching.first() {
|
||||
assert_eq!(receiver["ok"], true);
|
||||
assert_eq!(receiver["target"], "bootstrap");
|
||||
assert_eq!(receiver["mod_time"], cas["mod_time"]);
|
||||
assert!(receiver["body_sha256"].as_str().is_some_and(|hash| hash.len() == 64));
|
||||
if startup_cas_remote_matches(&logs[0], receiver) == 1 {
|
||||
receivers.push(*receiver);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if receivers.len() != 4 {
|
||||
// Remote started tracing may flush after direct receiver JSON.
|
||||
sleep(Duration::from_millis(25)).await;
|
||||
continue;
|
||||
}
|
||||
for (pool, disks) in phase.disks.iter().enumerate() {
|
||||
for (drive, disk) in disks.iter().enumerate() {
|
||||
let raw = std::fs::read(disk.join(".rustfs.sys/pool.bin/xl.meta"))?;
|
||||
let latest = rustfs_filemeta::get_file_info(
|
||||
&raw,
|
||||
".rustfs.sys",
|
||||
"pool.bin",
|
||||
"",
|
||||
rustfs_filemeta::FileInfoOpts {
|
||||
data: false,
|
||||
include_free_versions: false,
|
||||
include_part_checksums: false,
|
||||
},
|
||||
)
|
||||
.map_err(std::io::Error::other)?;
|
||||
assert!(!raw.is_empty());
|
||||
assert_eq!(latest.metadata.get("etag").map(String::as_str), commits[pool]["etag"].as_str());
|
||||
assert_eq!(
|
||||
latest.mod_time.map(|time| time.unix_timestamp_nanos().to_string()).as_deref(),
|
||||
commits[pool]["mod_time"].as_str()
|
||||
);
|
||||
std::fs::write(phase.artifact.join(format!("pool-{pool}-drive-{drive}-latest-xl.meta")), raw)?;
|
||||
}
|
||||
}
|
||||
let proof = json!({"nonce": phase.nonce, "pids": pids, "classifier": classified, "initial": initial, "prepare": prepares, "commit": commits, "replicas": replicas, "receivers": receivers, "confirmed": confirmed});
|
||||
std::fs::write(phase.artifact.join("cas-proof.json"), serde_json::to_vec_pretty(&proof)?)?;
|
||||
return Ok(proof);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5108,7 +5108,49 @@ async fn read_pool_meta_replicas<S>(pools: Vec<Arc<S>>, no_lock: bool) -> Vec<Po
|
||||
where
|
||||
S: EcstoreObjectIO,
|
||||
{
|
||||
join_all(pools.into_iter().map(|pool| read_pool_meta_replica(pool, no_lock))).await
|
||||
let reads = join_all(pools.into_iter().map(|pool| read_pool_meta_replica(pool, no_lock))).await;
|
||||
#[cfg(feature = "e2e-test-hooks")]
|
||||
if STARTUP_CAS_OBSERVATION.try_with(|_| ()).is_ok() {
|
||||
let batch = uuid::Uuid::new_v4();
|
||||
for (pool, read) in reads.iter().enumerate() {
|
||||
let mut observation = serde_json::json!({
|
||||
"kind": "replica-read", "object": POOL_META_NAME, "batch": batch, "pool": pool,
|
||||
"cas": match &read.cas {
|
||||
PoolMetaCasToken::Missing => "missing",
|
||||
PoolMetaCasToken::Existing(_) => "existing",
|
||||
PoolMetaCasToken::Unsafe => "unsafe",
|
||||
},
|
||||
"etag": match &read.cas { PoolMetaCasToken::Existing(etag) => Some(etag), _ => None },
|
||||
});
|
||||
match &read.replica {
|
||||
PoolMetaReplica::Valid {
|
||||
raw,
|
||||
canonical,
|
||||
meta,
|
||||
revision,
|
||||
committed,
|
||||
..
|
||||
} => {
|
||||
observation["state"] = serde_json::json!("valid");
|
||||
observation["committed"] = serde_json::json!(committed);
|
||||
observation["version"] = serde_json::json!(revision.version);
|
||||
observation["cluster_id"] = serde_json::json!(revision.cluster_id);
|
||||
observation["epoch"] = serde_json::json!(revision.epoch);
|
||||
observation["generation"] = serde_json::json!(revision.generation);
|
||||
observation["transaction_id"] = serde_json::json!(revision.transaction_id);
|
||||
observation["pool_count"] = serde_json::json!(meta.pools.len());
|
||||
observation["payload_sha256"] = serde_json::json!(rustfs_utils::crypto::hex(Sha256::digest(canonical)));
|
||||
observation["raw_sha256"] = serde_json::json!(rustfs_utils::crypto::hex(Sha256::digest(raw)));
|
||||
}
|
||||
PoolMetaReplica::Missing => observation["state"] = serde_json::json!("missing"),
|
||||
PoolMetaReplica::Corrupt(_) => observation["state"] = serde_json::json!("corrupt"),
|
||||
PoolMetaReplica::Incompatible(_) => observation["state"] = serde_json::json!("incompatible"),
|
||||
PoolMetaReplica::Unreadable(_) => observation["state"] = serde_json::json!("unreadable"),
|
||||
}
|
||||
startup_cas_test_observe(observation);
|
||||
}
|
||||
}
|
||||
reads
|
||||
}
|
||||
|
||||
fn select_pool_meta_replicas_observing<R>(write_state: &mut PoolMetaWriteState, replicas: Vec<R>) -> Result<PoolMetaSelection>
|
||||
@@ -5480,9 +5522,44 @@ fn pool_meta_cas_preconditions(token: &PoolMetaCasToken, object: &str) -> Result
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "e2e-test-hooks")]
|
||||
struct StartupCasObservation {
|
||||
attempt: uuid::Uuid,
|
||||
phase: &'static str,
|
||||
pools: Vec<usize>,
|
||||
}
|
||||
|
||||
#[cfg(feature = "e2e-test-hooks")]
|
||||
tokio::task_local! {
|
||||
static STARTUP_CAS_OBSERVATION: StartupCasObservation;
|
||||
}
|
||||
|
||||
// This scope follows only the directly polled startup future. Spawned work
|
||||
// does not inherit it; receiver evidence retains its existing RPC tuple.
|
||||
#[cfg(feature = "e2e-test-hooks")]
|
||||
pub(crate) async fn startup_cas_test_scope<S, F: std::future::Future>(
|
||||
attempt: uuid::Uuid,
|
||||
phase: &'static str,
|
||||
pools: &[Arc<S>],
|
||||
future: F,
|
||||
) -> F::Output {
|
||||
STARTUP_CAS_OBSERVATION
|
||||
.scope(
|
||||
StartupCasObservation {
|
||||
attempt,
|
||||
phase,
|
||||
// These identities are never dereferenced or logged. The
|
||||
// caller and operation keep the same pool Arcs alive.
|
||||
pools: pools.iter().map(|pool| Arc::as_ptr(pool) as usize).collect(),
|
||||
},
|
||||
future,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
// Direct JSON diagnostics are independent of the startup tracing subscriber.
|
||||
#[cfg(feature = "e2e-test-hooks")]
|
||||
fn startup_cas_test_observe(mut observation: serde_json::Value) {
|
||||
pub(crate) fn startup_cas_test_observe(mut observation: serde_json::Value) {
|
||||
let Some(nonce) = std::env::var("RUSTFS_E2E_STARTUP_CAS_NONCE")
|
||||
.ok()
|
||||
.and_then(|value| uuid::Uuid::parse_str(&value).ok())
|
||||
@@ -5491,6 +5568,10 @@ fn startup_cas_test_observe(mut observation: serde_json::Value) {
|
||||
};
|
||||
observation["nonce"] = serde_json::json!(nonce);
|
||||
observation["pid"] = serde_json::json!(std::process::id());
|
||||
let _ = STARTUP_CAS_OBSERVATION.try_with(|scope| {
|
||||
observation["attempt"] = serde_json::json!(scope.attempt);
|
||||
observation["startup_phase"] = serde_json::json!(scope.phase);
|
||||
});
|
||||
let line = format!("RUSTFS_E2E_STARTUP_CAS {observation}\n");
|
||||
let _ = std::io::Write::write_all(&mut std::io::stderr().lock(), line.as_bytes());
|
||||
}
|
||||
@@ -5519,6 +5600,9 @@ where
|
||||
let observation = std::env::var_os("RUSTFS_E2E_STARTUP_CAS_NONCE").map(|_| {
|
||||
serde_json::json!({
|
||||
"kind": "cas", "object": object, "phase": phase,
|
||||
"pool": STARTUP_CAS_OBSERVATION.try_with(|scope| {
|
||||
scope.pools.iter().position(|identity| *identity == Arc::as_ptr(&pool) as usize)
|
||||
}).ok().flatten(),
|
||||
"payload_sha256": rustfs_utils::crypto::hex(Sha256::digest(&data)),
|
||||
"if_match": opts.http_preconditions.as_ref().and_then(|p| p.if_match.as_deref()),
|
||||
"if_none_match": opts.http_preconditions.as_ref().and_then(|p| p.if_none_match.as_deref()),
|
||||
@@ -5538,6 +5622,13 @@ where
|
||||
if let Some(mut observation) = observation {
|
||||
observation["ok"] = serde_json::json!(result.is_ok());
|
||||
observation["etag"] = serde_json::json!(result.as_ref().ok().and_then(|info| info.etag.as_deref()));
|
||||
observation["mod_time"] = serde_json::json!(
|
||||
result
|
||||
.as_ref()
|
||||
.ok()
|
||||
.and_then(|info| info.mod_time)
|
||||
.map(|time| time.unix_timestamp_nanos().to_string())
|
||||
);
|
||||
observation["error"] = serde_json::json!(result.as_ref().err().map(ToString::to_string));
|
||||
startup_cas_test_observe(observation);
|
||||
}
|
||||
|
||||
@@ -630,14 +630,27 @@ impl ECStore {
|
||||
.pools
|
||||
.first()
|
||||
.is_some_and(|pool| pool_first_endpoint_is_local(&pool.endpoints));
|
||||
#[cfg(feature = "e2e-test-hooks")]
|
||||
let startup_attempt = uuid::Uuid::new_v4();
|
||||
let (meta, pool_meta_replica_state) = {
|
||||
let mut write_state = self.pool_meta_save_gate.lock().await;
|
||||
establish_pool_meta_bootstrap_identity_if_proven(self.pools.clone(), &mut write_state, should_persist_pool_meta)
|
||||
.await
|
||||
.map_err(|err| Error::other(format!("store init failed during establish_pool_meta_bootstrap_identity: {err}")))?;
|
||||
load_pool_meta_for_startup(self.pools.clone(), &mut write_state).await?
|
||||
let load = load_pool_meta_for_startup(self.pools.clone(), &mut write_state);
|
||||
#[cfg(feature = "e2e-test-hooks")]
|
||||
let load = crate::core::pools::startup_cas_test_scope(startup_attempt, "load", &self.pools, load);
|
||||
load.await?
|
||||
};
|
||||
let update = meta.validate(self.pools.clone())?;
|
||||
#[cfg(feature = "e2e-test-hooks")]
|
||||
crate::core::pools::startup_cas_test_observe(serde_json::json!({
|
||||
"kind": "startup-classifier", "attempt": startup_attempt,
|
||||
"elected_writer": should_persist_pool_meta,
|
||||
"needs_repair": pool_meta_replica_state.needs_repair,
|
||||
"repair_write_safe": pool_meta_replica_state.repair_write_safe,
|
||||
"topology_update": update,
|
||||
}));
|
||||
let endpoints = runtime_sources::endpoint_pools_or_default();
|
||||
|
||||
let mut installed_pool_meta = if update {
|
||||
@@ -649,15 +662,17 @@ impl ECStore {
|
||||
// distributed startup can race on the same lock and replay the prior init bug.
|
||||
{
|
||||
let mut write_state = self.pool_meta_save_gate.lock().await;
|
||||
installed_pool_meta = persist_pool_meta_for_startup_if_safe(
|
||||
let persist = persist_pool_meta_for_startup_if_safe(
|
||||
&installed_pool_meta,
|
||||
self.pools.clone(),
|
||||
pool_meta_replica_state,
|
||||
&mut write_state,
|
||||
update,
|
||||
should_persist_pool_meta,
|
||||
)
|
||||
.await?;
|
||||
);
|
||||
#[cfg(feature = "e2e-test-hooks")]
|
||||
let persist = crate::core::pools::startup_cas_test_scope(startup_attempt, "persist", &self.pools, persist);
|
||||
installed_pool_meta = persist.await?;
|
||||
}
|
||||
|
||||
{
|
||||
|
||||
Reference in New Issue
Block a user