Merge branch 'main' into overtrue/fix-1905-activation-fence

This commit is contained in:
houseme
2026-08-23 12:36:14 +08:00
committed by GitHub
23 changed files with 889 additions and 191 deletions
+20 -8
View File
@@ -988,14 +988,11 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for Sets {
}
}
#[async_trait::async_trait]
impl crate::storage_api_contracts::heal::HealOperations for Sets {
type Error = Error;
type HealResultItem = HealResultItem;
type HealOptions = HealOpts;
#[tracing::instrument(skip(self))]
async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
impl Sets {
pub(crate) async fn heal_format_with_fence<F>(&self, dry_run: bool, fence_lost: F) -> Result<(HealResultItem, Option<Error>)>
where
F: Fn() -> bool + Send + Sync,
{
let (disks, init_errs) = init_storage_disks_with_errors(
&self.endpoints.endpoints,
&DiskOption {
@@ -1068,6 +1065,9 @@ impl crate::storage_api_contracts::heal::HealOperations for Sets {
// Save new formats `format.json` on unformatted disks.
for (index, (fm, disk)) in tmp_new_formats.iter_mut().zip(disks.iter()).enumerate() {
if fm.is_some() && disk.is_some() {
if fence_lost() {
return Ok((res, Some(StorageError::SlowDown)));
}
if let Err(err) = save_format_file(disk, fm).await {
if let Some(disk) = disk.as_ref() {
let _ = disk.close().await;
@@ -1101,6 +1101,18 @@ impl crate::storage_api_contracts::heal::HealOperations for Sets {
}
Ok((res, None))
}
}
#[async_trait::async_trait]
impl crate::storage_api_contracts::heal::HealOperations for Sets {
type Error = Error;
type HealResultItem = HealResultItem;
type HealOptions = HealOpts;
#[tracing::instrument(skip(self))]
async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option<Error>)> {
self.heal_format_with_fence(dry_run, || false).await
}
#[tracing::instrument(skip(self))]
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
let mut result = HealResultItem {
+38 -6
View File
@@ -784,6 +784,24 @@ pub(crate) fn create_deferred_bitrot_reader_with_stripe_handle(
///
/// # Returns
/// A Result containing the BitrotWriterWrapper or an error
/// Size hint handed to `DiskAPI::create_file` for a bitrot-wrapped shard.
///
/// A known length is grown by one checksum per shard so the on-disk file size
/// matches what the bitrot writer emits. A negative length is the
/// unknown-size sentinel (`HashReader::SIZE_PRESERVE_LAYER`, used by SSE and
/// compression) and must be preserved: `RemoteDisk::create_file` forwards it
/// in the `put_file_stream` query, and the receiver only treats `size > 0` as
/// a fixed body length when locating the authenticated trailer. Clamping it
/// to `0` would claim an empty body and misframe the stream. `0` stays `0`
/// because a genuinely empty object still means an empty body.
fn bitrot_create_file_size(length: i64, shard_size: usize, checksum_algo: &HashAlgorithm) -> i64 {
if length <= 0 {
return length;
}
let length = length as usize;
(length.div_ceil(shard_size) * checksum_algo.size() + length) as i64
}
pub async fn create_bitrot_writer(
is_inline_buffer: bool,
disk: Option<&DiskStore>,
@@ -796,12 +814,7 @@ pub async fn create_bitrot_writer(
let writer = if is_inline_buffer {
CustomWriter::new_inline_buffer()
} else if let Some(disk) = disk {
let length = if length > 0 {
let length = length as usize;
(length.div_ceil(shard_size) * checksum_algo.size() + length) as i64
} else {
0
};
let length = bitrot_create_file_size(length, shard_size, &checksum_algo);
let file = disk.create_file("", volume, path, length).await?;
#[cfg(feature = "hotpath")]
@@ -820,6 +833,25 @@ mod tests {
use rustfs_rio::ChunkReader;
use std::collections::VecDeque;
#[test]
fn bitrot_create_file_size_grows_known_length_by_checksums() {
// 10 bytes over 4-byte shards = 3 shards, each followed by a 32-byte hash.
assert_eq!(bitrot_create_file_size(10, 4, &HashAlgorithm::HighwayHash256), 10 + 3 * 32);
assert_eq!(bitrot_create_file_size(10, 4, &HashAlgorithm::None), 10);
}
#[test]
fn bitrot_create_file_size_keeps_empty_and_unknown_distinct() {
assert_eq!(bitrot_create_file_size(0, 4, &HashAlgorithm::HighwayHash256), 0);
// SSE/compression streams advertise SIZE_PRESERVE_LAYER (-1); the remote
// put_file_stream receiver relies on a non-positive size to parse the auth
// trailer from the stream tail, so the sentinel must survive untouched.
assert_eq!(
bitrot_create_file_size(rustfs_rio::HashReader::SIZE_PRESERVE_LAYER, 4, &HashAlgorithm::HighwayHash256),
rustfs_rio::HashReader::SIZE_PRESERVE_LAYER
);
}
struct TestChunkReader {
chunks: VecDeque<Bytes>,
}
+1 -14
View File
@@ -2124,26 +2124,13 @@ impl SetDisks {
let put_object_size = known_put_object_storage_size(data.size());
let shard_file_size_raw = erasure.shard_file_size(put_object_size);
let is_inline_buffer =
storage_class_config.should_inline(shard_file_size_raw, erasure.data_shards, opts.versioned);
let is_inline_buffer = storage_class_config.should_inline(shard_file_size_raw, erasure.data_shards, opts.versioned);
let collect_stage_timing = rustfs_io_metrics::put_stage_metrics_enabled() || issue3031_diag_enabled();
let shard_file_size = shard_file_size_raw;
let shard_size = erasure.shard_size();
let write_path = classify_put_write_path(is_inline_buffer, put_object_size, fi.erasure.block_size);
let direct_inline_commit = matches!(write_path, SmallWritePath::Inline);
{
use std::io::Write;
let msg = format!(
"INLINE_DEBUG: bucket={} obj={} size={} shard_fs={} ds={} bs={} inline={} direct={} path={} iblock={} ver={}\n",
bucket, object, put_object_size, shard_file_size_raw, erasure.data_shards, fi.erasure.block_size,
is_inline_buffer, direct_inline_commit, write_path.metric_label(), storage_class_config.inline_block(), opts.versioned
);
if let Ok(mut f) = std::fs::OpenOptions::new().create(true).append(true).open("/tmp/rustfs_inline_debug.log") {
let _ = f.write_all(msg.as_bytes());
}
let _ = std::io::stderr().write_all(msg.as_bytes());
}
rustfs_io_metrics::record_put_object_path(write_path.metric_label());
let writer_setup_stage_start = collect_stage_timing.then(Instant::now);
let (mut writers, errors) = if direct_inline_commit {
+332 -3
View File
@@ -13,7 +13,12 @@
// limitations under the License.
use super::*;
use crate::core::pools::POOL_META_NAME;
use crate::services::rebalance::{REBAL_META_NAME, RebalStatus};
use crate::set_disk::get_lock_acquire_timeout;
use crate::storage_api_contracts::heal::HealOperations as _;
use crate::storage_api_contracts::namespace::NamespaceLocking as _;
use rustfs_lock::NamespaceLockGuard;
use tracing::trace;
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
@@ -30,7 +35,119 @@ fn invalid_heal_pool_index(pool_idx: usize, pool_count: usize) -> Error {
)
}
#[derive(Debug, Clone, Copy)]
enum HealFormatPoolSkip {
Completed,
Retryable,
}
fn classify_heal_format_pool(
pool_idx: usize,
pool_cmd_line: &str,
pool_meta: &PoolMeta,
rebalance_meta: Option<&RebalanceMeta>,
) -> Option<HealFormatPoolSkip> {
let Some(pool) = pool_meta.pools.get(pool_idx) else {
return Some(HealFormatPoolSkip::Retryable);
};
if pool.id != pool_idx || pool_cmd_line.is_empty() || pool.cmd_line.is_empty() || pool.cmd_line != pool_cmd_line {
return Some(HealFormatPoolSkip::Retryable);
}
if let Some(decommission) = pool.decommission.as_ref() {
if decommission.complete {
return Some(HealFormatPoolSkip::Completed);
}
if decommission.failed || decommission.canceled || decommission.queued || pool_meta.is_suspended(pool_idx) {
return Some(HealFormatPoolSkip::Retryable);
}
}
if let Some(meta) = rebalance_meta {
let Some(pool_stats) = meta.pool_stats.get(pool_idx) else {
return Some(HealFormatPoolSkip::Retryable);
};
if pool_stats.info.stopping || (pool_stats.participating && pool_stats.info.status == RebalStatus::Started) {
return Some(HealFormatPoolSkip::Retryable);
}
}
None
}
fn heal_format_pool_skip_error(skip: HealFormatPoolSkip) -> Error {
match skip {
HealFormatPoolSkip::Completed => StorageError::NoHealRequired,
HealFormatPoolSkip::Retryable => StorageError::SlowDown,
}
}
fn heal_format_fence_lost_error() -> Error {
StorageError::SlowDown
}
impl ECStore {
async fn acquire_heal_format_fence(
&self,
) -> Result<(NamespaceLockGuard, NamespaceLockGuard, PoolMeta, Option<RebalanceMeta>)> {
let metadata_pool = self
.pools
.first()
.cloned()
.ok_or_else(|| Error::other("heal format requires at least one storage pool"))?;
// Metadata fence order is part of the decommission/rebalance protocol:
// pool.bin must always be acquired before rebalance.bin.
let pool_lock = metadata_pool.new_ns_lock(RUSTFS_META_BUCKET, POOL_META_NAME).await?;
let pool_guard = pool_lock.get_write_lock(get_lock_acquire_timeout()).await?;
let rebalance_lock = metadata_pool.new_ns_lock(RUSTFS_META_BUCKET, REBAL_META_NAME).await?;
let rebalance_guard = rebalance_lock.get_write_lock(get_lock_acquire_timeout()).await?;
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
return Err(heal_format_fence_lost_error());
}
let mut pool_meta = PoolMeta::default();
pool_meta.load_no_lock(metadata_pool.clone()).await?;
if pool_meta.pools.len() != self.pools.len()
|| pool_meta.pools.iter().enumerate().any(|(pool_idx, pool)| {
pool.id != pool_idx || pool.cmd_line.is_empty() || pool.cmd_line != self.pools[pool_idx].endpoints.cmd_line
})
{
return Err(heal_format_fence_lost_error());
}
let mut rebalance_meta = RebalanceMeta::new();
let rebalance_meta = match rebalance_meta
.load_with_opts(
metadata_pool,
ObjectOptions {
no_lock: true,
..Default::default()
},
)
.await
{
Ok(()) => Some(rebalance_meta),
Err(Error::ConfigNotFound) => None,
Err(err) => return Err(err),
};
if rebalance_meta
.as_ref()
.is_some_and(|meta| meta.pool_stats.len() != self.pools.len())
{
return Err(heal_format_fence_lost_error());
}
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
return Err(heal_format_fence_lost_error());
}
Ok((pool_guard, rebalance_guard, pool_meta, rebalance_meta))
}
fn get_pools_for_heal_object(&self, opts: &HealOpts) -> Result<Vec<Arc<Sets>>> {
match opts.pool {
Some(pool_idx) => Ok(vec![
@@ -52,9 +169,26 @@ impl ECStore {
};
let mut count_no_heal = 0;
let mut count_completed = 0;
let mut first_error = None;
for pool in self.pools.iter() {
let (mut result, err) = pool.heal_format(dry_run).await?;
for (pool_idx, pool) in self.pools.iter().enumerate() {
let (pool_guard, rebalance_guard, pool_meta, rebalance_meta) = self.acquire_heal_format_fence().await?;
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
first_error.get_or_insert(heal_format_fence_lost_error());
break;
}
if let Some(skip) = classify_heal_format_pool(pool_idx, &pool.endpoints.cmd_line, &pool_meta, rebalance_meta.as_ref())
{
if matches!(skip, HealFormatPoolSkip::Completed) {
count_completed += 1;
} else {
first_error.get_or_insert(heal_format_pool_skip_error(skip));
}
continue;
}
let fence_lost = || pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost();
let (mut result, err) = pool.heal_format_with_fence(dry_run, fence_lost).await?;
if let Some(err) = err {
match err {
StorageError::NoHealRequired => {
@@ -69,11 +203,18 @@ impl ECStore {
r.set_count += result.set_count;
r.before.drives.append(&mut result.before.drives);
r.after.drives.append(&mut result.after.drives);
// A lease can be lost after the final write; fail closed before
// reporting the pool as successfully healed.
if pool_guard.is_lock_lost() || rebalance_guard.is_lock_lost() {
first_error.get_or_insert(heal_format_fence_lost_error());
break;
}
}
if let Some(err) = first_error {
return Ok((r, Some(err)));
}
if count_no_heal == self.pools.len() {
if count_no_heal + count_completed == self.pools.len() {
info!(
event = EVENT_HEAL_FORMAT_COMPLETED,
component = LOG_COMPONENT_ECSTORE,
@@ -302,6 +443,7 @@ mod tests {
use crate::disk::{DeleteOptions, DiskOption, format::FormatV3, new_disk};
use crate::layout::endpoints::{EndpointServerPools, Endpoints, PoolEndpoints};
use crate::runtime::instance::InstanceContext;
use crate::services::rebalance::{RebalanceInfo, RebalanceStats};
use crate::storage_api_contracts::bucket::{BucketOperations, MakeBucketOptions};
use crate::storage_api_contracts::object::{ObjectIO as _, ObjectOperations};
use crate::store::init_format::{load_format_erasure, save_format_file};
@@ -353,6 +495,164 @@ mod tests {
}
}
fn pool_meta_with_decommission(info: PoolDecommissionInfo) -> PoolMeta {
PoolMeta {
pools: vec![PoolStatus {
id: 0,
cmd_line: "pool-0".to_string(),
last_update: OffsetDateTime::UNIX_EPOCH,
decommission: Some(info),
}],
..Default::default()
}
}
#[test]
fn heal_format_pool_state_barriers_are_classified() {
let active = pool_meta_with_decommission(PoolDecommissionInfo {
start_time: Some(OffsetDateTime::UNIX_EPOCH),
..Default::default()
});
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &active, None),
Some(HealFormatPoolSkip::Retryable)
));
for info in [
PoolDecommissionInfo {
failed: true,
..Default::default()
},
PoolDecommissionInfo {
canceled: true,
..Default::default()
},
] {
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &pool_meta_with_decommission(info), None),
Some(HealFormatPoolSkip::Retryable)
));
}
let completed = pool_meta_with_decommission(PoolDecommissionInfo {
complete: true,
..Default::default()
});
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &completed, None),
Some(HealFormatPoolSkip::Completed)
));
}
#[test]
fn heal_format_pool_rebalance_barriers_and_identity_are_fail_closed() {
let identity_meta = pool_meta_with_decommission(PoolDecommissionInfo::default());
let rebalance = RebalanceMeta {
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Started,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&rebalance)),
Some(HealFormatPoolSkip::Retryable)
));
let stopping = RebalanceMeta {
pool_stats: vec![RebalanceStats {
info: RebalanceInfo {
stopping: true,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopping)),
Some(HealFormatPoolSkip::Retryable)
));
let identity = pool_meta_with_decommission(PoolDecommissionInfo::default());
assert!(matches!(
classify_heal_format_pool(0, "pool-new", &identity, None),
Some(HealFormatPoolSkip::Retryable)
));
let identity_without_decommission = PoolMeta {
pools: vec![PoolStatus {
id: 0,
cmd_line: "pool-0".to_string(),
last_update: OffsetDateTime::UNIX_EPOCH,
decommission: None,
}],
..Default::default()
};
assert!(matches!(
classify_heal_format_pool(0, "pool-new", &identity_without_decommission, None),
Some(HealFormatPoolSkip::Retryable)
));
assert!(matches!(
classify_heal_format_pool(0, "", &identity_meta, None),
Some(HealFormatPoolSkip::Retryable)
));
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &PoolMeta::default(), None),
Some(HealFormatPoolSkip::Retryable)
));
let stopped = RebalanceMeta {
stopped_at: Some(OffsetDateTime::UNIX_EPOCH),
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Stopped,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
assert!(classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopped)).is_none());
let stopping_after_stop = RebalanceMeta {
stopped_at: Some(OffsetDateTime::UNIX_EPOCH),
pool_stats: vec![RebalanceStats {
participating: true,
info: RebalanceInfo {
status: RebalStatus::Started,
stopping: true,
..Default::default()
},
..Default::default()
}],
..Default::default()
};
assert!(matches!(
classify_heal_format_pool(0, "pool-0", &identity_meta, Some(&stopping_after_stop)),
Some(HealFormatPoolSkip::Retryable)
));
}
#[test]
fn skipped_heal_format_pool_is_never_reported_as_success() {
assert!(matches!(
heal_format_pool_skip_error(HealFormatPoolSkip::Retryable),
StorageError::SlowDown
));
assert!(matches!(
heal_format_pool_skip_error(HealFormatPoolSkip::Completed),
StorageError::NoHealRequired
));
}
async fn multi_pool_heal_store() -> (tempfile::TempDir, Arc<ECStore>, CancellationToken) {
let temp_dir = tempfile::tempdir().expect("multi-pool heal test directory should be created");
let mut pool_endpoints = Vec::new();
@@ -889,6 +1189,18 @@ mod tests {
bucket_fence_registry: std::sync::Arc::default(),
};
let err = store
.handle_heal_format(false)
.await
.expect_err("missing pool metadata must fail closed before format writes");
assert!(matches!(err, StorageError::SlowDown));
let pool_meta = PoolMeta::new(&store.pools, &PoolMeta::default());
pool_meta
.save(store.pools.clone())
.await
.expect("pool metadata should be persisted before format heal");
let (result, err) = store
.handle_heal_format(false)
.await
@@ -902,5 +1214,22 @@ mod tests {
.await
.expect("the later pool should be healed despite the first pool error");
assert_eq!(healed.erasure.this, recoverable_format.erasure.sets[0][2]);
let mut completed_meta = PoolMeta::new(&store.pools, &PoolMeta::default());
for status in &mut completed_meta.pools {
status.decommission = Some(PoolDecommissionInfo {
complete: true,
..Default::default()
});
}
completed_meta
.save(store.pools.clone())
.await
.expect("completed pool metadata should be persisted");
let (_, err) = store
.handle_heal_format(false)
.await
.expect("completed pools should be reported as a no-op");
assert!(matches!(err, Some(StorageError::NoHealRequired)));
}
}