fix(ecstore): recover interrupted pool metadata writes (#7387)

* fix(ecstore): recover interrupted pool metadata writes

* chore(ci): update error format ratchet baseline

Co-Authored-By: heihutu <heihutu@gmail.com>

Co-Authored-By: zhi22915 <qiuzgang@gmail.com>

---------

Co-authored-by: houseme <housemecn@gmail.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
Co-authored-by: Zhengchao An <anzhengchao@gmail.com>
This commit is contained in:
cxymds
2026-09-07 22:17:51 +08:00
committed by GitHub
parent bc0f608431
commit fd92853ac4
11 changed files with 1370 additions and 98 deletions
+2 -2
View File
@@ -416,8 +416,8 @@ pub mod disk {
pub mod error {
pub use crate::error::{
Error, Result, StorageError, classify_system_path_failure_reason, is_err_bucket_not_found, is_err_object_not_found,
is_err_version_not_found,
Error, PoolMetadataError, PoolMetadataFailure, Result, StorageError, classify_system_path_failure_reason,
is_err_bucket_not_found, is_err_object_not_found, is_err_version_not_found,
};
}
File diff suppressed because it is too large Load Diff
+74 -4
View File
@@ -23,17 +23,59 @@ use s3s::S3ErrorCode;
pub type Error = StorageError;
pub type Result<T> = core::result::Result<T, Error>;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PoolMetadataFailure {
ReadUnavailable,
RecoveryRequired,
TransactionUnknown,
FenceLost,
}
impl PoolMetadataFailure {
fn recovery_hint(self) -> &'static str {
match self {
Self::ReadUnavailable => "read unavailable; retry after the replicas are readable",
Self::TransactionUnknown => "writes remain blocked pending fenced transaction recovery",
Self::RecoveryRequired | Self::FenceLost => {
"writes remain blocked after a recovery-required replica state; restart after all replicas are readable and consistent, with compatible formats"
}
}
}
pub fn as_str(self) -> &'static str {
match self {
Self::ReadUnavailable => "read_unavailable",
Self::RecoveryRequired => "recovery_required",
Self::TransactionUnknown => "transaction_unknown",
Self::FenceLost => "fence_lost",
}
}
}
/// Local control-plane context. Keep the existing storage error wire codes;
/// the HTTP boundary recognizes this typed source, not an error-message prefix.
#[derive(Debug, Clone, thiserror::Error)]
#[error("{operation}: pool metadata {hint} ({reason}, {phase}): {detail}", hint = kind.recovery_hint(), reason = kind.as_str(), detail = source.as_ref().map(ToString::to_string).unwrap_or_default())]
pub struct PoolMetadataError {
pub kind: PoolMetadataFailure,
pub operation: String,
pub phase: &'static str,
pub since: time::OffsetDateTime,
#[source]
pub source: Option<std::sync::Arc<StorageError>>,
}
/// Keeps high-cardinality diagnostic detail in the error source while making
/// the rendered `io::Error` stable for quorum aggregation.
#[derive(Debug)]
struct StableIoContextError {
message: &'static str,
message: std::borrow::Cow<'static, str>,
source: Box<dyn std::error::Error + Send + Sync>,
}
impl std::fmt::Display for StableIoContextError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str(self.message)
formatter.write_str(&self.message)
}
}
@@ -48,7 +90,7 @@ where
E: Into<Box<dyn std::error::Error + Send + Sync>>,
{
std::io::Error::other(StableIoContextError {
message,
message: message.into(),
source: source.into(),
})
}
@@ -300,6 +342,22 @@ impl From<crate::erasure::coding::ErasureConstructionError> for StorageError {
}
impl StorageError {
pub fn pool_metadata_failure(&self) -> Option<&PoolMetadataError> {
let mut current: Option<&(dyn std::error::Error + 'static)> = Some(self);
while let Some(error) = current {
if let Some(context) = error.downcast_ref::<PoolMetadataError>() {
return Some(context);
}
// io::Error::source skips its boxed context itself.
current = if let Some(io) = error.downcast_ref::<std::io::Error>() {
io.get_ref().map(|inner| inner as &(dyn std::error::Error + 'static))
} else {
error.source()
};
}
None
}
pub fn other<E>(error: E) -> Self
where
E: Into<Box<dyn std::error::Error + Send + Sync>>,
@@ -548,7 +606,19 @@ impl PartialEq for StorageError {
impl Clone for StorageError {
fn clone(&self) -> Self {
match self {
StorageError::Io(e) => StorageError::Io(std::io::Error::new(e.kind(), e.to_string())),
StorageError::Io(e) => {
if let Some(context) = self.pool_metadata_failure() {
Self::Io(std::io::Error::new(
e.kind(),
StableIoContextError {
message: e.to_string().into(),
source: Box::new(context.clone()),
},
))
} else {
StorageError::Io(std::io::Error::new(e.kind(), e.to_string()))
}
}
StorageError::FaultyDisk => StorageError::FaultyDisk,
StorageError::DiskFull => StorageError::DiskFull,
StorageError::VolumeNotFound => StorageError::VolumeNotFound,
+12 -7
View File
@@ -2353,14 +2353,19 @@ mod tests {
.await
.expect("quorum boundary heal should return a mapped result");
*store.pools[0].disk_set[0].disks.write().await = original_quorum_disks;
let quorum_err_text = quorum_err.as_ref().map(ToString::to_string);
assert!(
quorum_err_text.as_deref().is_some_and(|err| {
err.contains("target capacity admission failed")
&& err.contains("pool metadata update cannot overwrite an unreadable replica")
}),
"heal must fail closed when capacity admission cannot verify pool metadata, got {quorum_err:?}"
let quorum_err = quorum_err
.as_ref()
.expect("heal must fail closed when capacity admission cannot verify pool metadata");
let quorum_failure = quorum_err
.pool_metadata_failure()
.expect("capacity admission failure should preserve typed pool metadata context");
assert_eq!(
quorum_failure.kind,
crate::error::PoolMetadataFailure::ReadUnavailable,
"read-only capacity admission failure must remain retryable"
);
assert_eq!(quorum_failure.operation, "target capacity admission failed");
assert_eq!(quorum_failure.phase, "pool_read");
assert!(
store.pool_meta_writes_ready().await,
"read-only capacity admission failure must not latch the pool metadata writer"
+131 -7
View File
@@ -14,8 +14,8 @@
use super::*;
use crate::core::pools::{
PoolMetaBootstrapAuthority, PoolMetaReplicaState, PoolMetaWriteState, load_pool_meta_identity_observing,
local_decommission_queue_prefix, persist_pool_meta_identity_for_startup, pool_meta_has_active_decommission,
PoolMetaBootstrapAuthority, PoolMetaReplicaState, PoolMetaWriteState, local_decommission_queue_prefix,
persist_pool_meta_identity_for_startup, pool_meta_has_active_decommission,
};
use crate::runtime::instance::InstanceContext;
use crate::runtime::sources as runtime_sources;
@@ -153,14 +153,11 @@ async fn load_pool_meta_for_startup<S>(
where
S: EcstoreObjectIO,
{
load_pool_meta_identity_observing(pools.clone(), write_state)
.await
.map_err(|err| Error::other(format!("store init failed during load_pool_meta_identity: {err}")))?;
let mut meta = PoolMeta::default();
let replica_state = meta
.load_no_lock_from_replicas_observing(pools, write_state)
.load_for_startup_observing(pools, write_state)
.await
.map_err(|err| Error::other(format!("store init failed during load_pool_meta: {err}")))?;
.map_err(|err| Error::other_with_context("store init failed during load_pool_meta", err))?;
write_state.observe_replicas(replica_state);
write_state
.ensure_missing_metadata_can_initialize()
@@ -769,6 +766,33 @@ impl ECStore {
});
}
let recovery_store = self.clone();
let recovery_rx = rx.clone();
tokio::spawn(async move {
let mut delay = std::time::Duration::from_secs(5);
loop {
tokio::select! {
_ = recovery_rx.cancelled() => return,
_ = tokio::time::sleep(delay) => {}
}
let result = tokio::select! {
_ = recovery_rx.cancelled() => return,
result = tokio::time::timeout(std::time::Duration::from_secs(30), recovery_store.recover_pool_meta_transaction()) => result,
};
delay = match result {
Ok(Ok(_)) => std::time::Duration::from_secs(5),
failure => {
let error = match failure {
Ok(Err(error)) => error,
_ => Error::Timeout,
};
recovery_store.record_pool_meta_recovery_failure(error);
(delay * 2).min(std::time::Duration::from_secs(60))
}
};
}
});
runtime_sources::init_bucket_monitor_for_current_endpoints();
crate::bucket::bucket_target_sys::BucketTargetSys::get().start_heartbeat();
@@ -2493,6 +2517,106 @@ mod tests {
shutdown.cancel();
}
#[tokio::test]
#[serial_test::serial(storage_class_env)]
async fn pool_metadata_preflight_recovery_preserves_single_and_multi_pool_public_mutations() {
for layout in [vec![4], vec![4, 4]] {
let temp_dir = tempfile::tempdir().unwrap();
let (_ctx, store, shutdown) =
without_storage_class_env(build_isolated_test_store(temp_dir.path(), "pool-meta-retry", &layout)).await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(Arc::clone(&store), Vec::new()).await;
let bucket = format!("pool-meta-retry-{}", Uuid::new_v4());
store.make_bucket(&bucket, &MakeBucketOptions::default()).await.unwrap();
let mut saved_disks = Vec::new();
for set in &store.pools[0].disk_set {
let mut disks = set.disks.write().await;
let count = disks.len();
saved_disks.push((set.clone(), std::mem::replace(&mut *disks, vec![None; count])));
}
let indices = (0..layout.len()).collect::<Vec<_>>();
let err = store.save_current_pool_meta_for_test(&indices).await.unwrap_err();
assert_eq!(
err.pool_metadata_failure().unwrap().kind,
crate::error::PoolMetadataFailure::ReadUnavailable
);
for (set, disks) in saved_disks {
*set.disks.write().await = disks;
}
store.save_current_pool_meta_for_test(&indices).await.unwrap();
assert!(store.pool_meta_writes_ready().await);
let payload = b"pool metadata recovery payload".to_vec();
store
.put_object(&bucket, "put", &mut PutObjReader::from_vec(payload.clone()), &ObjectOptions::default())
.await
.unwrap();
let mut reader = store
.get_object_reader(&bucket, "put", None, HeaderMap::new(), &ObjectOptions::default())
.await
.unwrap();
let mut actual = Vec::new();
reader.stream.read_to_end(&mut actual).await.unwrap();
assert_eq!(actual, payload);
drop(reader);
store.delete_object(&bucket, "put", ObjectOptions::default()).await.unwrap();
assert!(crate::error::is_err_object_not_found(
&store
.get_object_info(&bucket, "put", &ObjectOptions::default())
.await
.unwrap_err()
));
let upload = store
.new_multipart_upload(&bucket, "multipart", &ObjectOptions::default())
.await
.unwrap();
let part = store
.put_object_part(
&bucket,
"multipart",
&upload.upload_id,
1,
&mut PutObjReader::from_vec(payload.clone()),
&ObjectOptions::default(),
)
.await
.unwrap();
store
.clone()
.complete_multipart_upload(
&bucket,
"multipart",
&upload.upload_id,
vec![crate::storage_api_contracts::multipart::CompletePart {
part_num: part.part_num,
etag: part.etag,
..Default::default()
}],
&ObjectOptions::default(),
)
.await
.unwrap();
let mut reader = store
.get_object_reader(&bucket, "multipart", None, HeaderMap::new(), &ObjectOptions::default())
.await
.unwrap();
actual.clear();
reader.stream.read_to_end(&mut actual).await.unwrap();
assert_eq!(actual, payload);
drop(reader);
let upload = store
.new_multipart_upload(&bucket, "abort", &ObjectOptions::default())
.await
.unwrap();
store
.abort_multipart_upload(&bucket, "abort", &upload.upload_id, &ObjectOptions::default())
.await
.unwrap();
assert!(store.pool_meta_writes_ready().await);
shutdown.cancel();
}
}
#[cfg(feature = "test-util")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
#[serial_test::serial(storage_class_env)]