Compare commits

..

2 Commits

Author SHA1 Message Date
overtrue 360f773f69 Merge remote-tracking branch 'origin/main' into overtrue/fix-1912-heal-bucket-fence
# Conflicts:
#	crates/ecstore/src/store/heal.rs
2026-08-23 04:40:48 +08:00
overtrue c422425431 fix(ecstore): fence bucket heal during decommission 2026-08-23 04:22:12 +08:00
19 changed files with 669 additions and 107 deletions
+2 -2
View File
@@ -1,2 +1,2 @@
sha256-darwin=c88d762eb265fd90f467e57c948a7e6806d5918dc75798651afe120af85c9d07
sha256-linux=9b694893afb2d750b37f08e2a4fe229b37cc9d21332985a3727daf5ebd77f1a9
sha256-darwin=9f767b37ed8b1c82da62ea441462d75487785c8086e56f08fb6f6cd89c6e2e52
sha256-linux=fbdaf42b220958d4b1e8880e0f8b5a7992d38e21051bb60596dd4538424757d6
Generated
+27 -22
View File
@@ -1858,9 +1858,9 @@ dependencies = [
[[package]]
name = "cc"
version = "1.4.4"
version = "1.4.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0ad534f4357a5264cce5019c989cf66a4f0dc4e0d1b1d15f8aacec0ff7360273"
checksum = "509591b7bcd67f4ef775afad7662703b4935daaa6ec0e5605cfb1090b32a2b6d"
dependencies = [
"find-msvc-tools",
"jobserver",
@@ -2522,6 +2522,12 @@ dependencies = [
"subtle",
]
[[package]]
name = "cty"
version = "0.2.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b365fabc795046672053e29c954733ec3b05e4be654ab130fe8f1f94d7051f35"
[[package]]
name = "curve25519-dalek"
version = "4.1.3"
@@ -5982,6 +5988,15 @@ version = "0.2.16"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b6d2cec3eae94f9f509c767b45932f1ada8350c4bdb85af2fcab4a3c14807981"
[[package]]
name = "libmimalloc-sys"
version = "0.1.49"
source = "git+https://github.com/xonatius/mimalloc_rust.git?rev=6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11#6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11"
dependencies = [
"cc",
"cty",
]
[[package]]
name = "libredox"
version = "0.1.20"
@@ -6382,6 +6397,14 @@ dependencies = [
"synstructure 0.13.2",
]
[[package]]
name = "mimalloc"
version = "0.1.52"
source = "git+https://github.com/xonatius/mimalloc_rust.git?rev=6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11#6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11"
dependencies = [
"libmimalloc-sys",
]
[[package]]
name = "mime"
version = "0.3.17"
@@ -9139,11 +9162,13 @@ dependencies = [
"insta",
"jiff",
"libc",
"libmimalloc-sys",
"libsystemd",
"matchit 0.9.2",
"md-5 0.11.0",
"metrics",
"metrics-util",
"mimalloc",
"mime_guess",
"opentelemetry",
"opentelemetry_sdk",
@@ -9179,8 +9204,6 @@ dependencies = [
"rustfs-lock",
"rustfs-log-analyzer",
"rustfs-madmin",
"rustfs-mimalloc",
"rustfs-mimalloc-sys",
"rustfs-notify",
"rustfs-object-capacity",
"rustfs-object-data-cache",
@@ -9852,24 +9875,6 @@ dependencies = [
"tokio",
]
[[package]]
name = "rustfs-mimalloc"
version = "0.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a406f4aa07084301d485beec873af6dccc8e3f8762da244743df92038b1db1a6"
dependencies = [
"rustfs-mimalloc-sys",
]
[[package]]
name = "rustfs-mimalloc-sys"
version = "0.5.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c3051b819175f58445d4c369a72f0ab88149f3885ba8bea2aff3be01f53fe7cd"
dependencies = [
"cc",
]
[[package]]
name = "rustfs-notify"
version = "1.0.0-rc.3"
+2 -2
View File
@@ -350,8 +350,8 @@ russh-sftp = "2.4.0"
dav-server = "0.11.0"
# Performance Analysis and Memory Profiling
rustfs-mimalloc = { version = "0.5.0" }
rustfs-mimalloc-sys = { version = "0.5.0" }
mimalloc = { version = "0.1.52", git = "https://github.com/xonatius/mimalloc_rust.git", rev = "6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11" }
libmimalloc-sys = { version = "0.1.49", git = "https://github.com/xonatius/mimalloc_rust.git", rev = "6d4c41bb10c6d9da1d1b6f07b38c4cc051667f11", features = ["extended"] }
hotpath = { version = "0.23.3", default-features = false }
# Snapshot testing for output format regression detection
insta = { version = "1.48" }
+14 -4
View File
@@ -59,6 +59,7 @@ where
/// Regression test for data usage accuracy (issue #1012).
/// Launches rustfs, writes 1000 objects, then asserts admin data usage reports the full count.
#[tokio::test(flavor = "multi_thread")]
#[ignore = "Starts a rustfs server and requires awscurl; enable when running full E2E"]
async fn data_usage_reports_all_objects() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
init_logging();
@@ -85,20 +86,28 @@ async fn data_usage_reports_all_objects() -> Result<(), Box<dyn std::error::Erro
usage
.buckets_usage
.get(TEST_BUCKET)
.map(|bucket_usage| usage.objects_total_count == 1000 && bucket_usage.objects_count == 1000)
.map(|bucket_usage| usage.objects_total_count >= 1000 && bucket_usage.objects_count >= 1000)
.unwrap_or(false)
})
.await?;
// Assert total object count and per-bucket count are exact.
// Assert total object count and per-bucket count are not truncated
let bucket_usage = usage
.buckets_usage
.get(TEST_BUCKET)
.cloned()
.expect("bucket usage should exist");
assert_eq!(usage.objects_total_count, 1000, "total object count should be exact");
assert_eq!(bucket_usage.objects_count, 1000, "bucket object count should be exact");
assert!(
usage.objects_total_count >= 1000,
"total object count should be at least 1000, got {}",
usage.objects_total_count
);
assert!(
bucket_usage.objects_count >= 1000,
"bucket object count should be at least 1000, got {}",
bucket_usage.objects_count
);
env.stop_server();
Ok(())
@@ -107,6 +116,7 @@ async fn data_usage_reports_all_objects() -> Result<(), Box<dyn std::error::Erro
/// Regression test for issue #3898.
/// Versioned buckets should expose versions and delete markers through admin data usage.
#[tokio::test(flavor = "multi_thread")]
#[ignore = "Starts a rustfs server and requires awscurl; enable when running full E2E"]
async fn data_usage_reports_versioned_objects_and_delete_markers() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
init_logging();
+8 -7
View File
@@ -454,13 +454,14 @@ pub mod rpc {
AuthenticatedChannel, KMS_SIGNAL_SUBSYSTEM, LocalPeerS3Client, PEER_RESTDRY_RUN, PEER_RESTSIGNAL, PEER_RESTSUB_SYS,
PeerRestClient, PeerS3Client, S3PeerSys, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC,
ScannerBucketListing, ScannerPeerActivity, TONIC_RPC_PREFIX, TonicInterceptor, build_put_file_auth_trailer,
check_and_record_signed_rpc_nonce, gen_signature_headers, gen_tonic_replay_scope_headers, gen_tonic_signature_headers,
gen_tonic_signature_interceptor, node_service_time_out_client, node_service_time_out_client_no_auth,
normalize_tonic_rpc_audience, set_tonic_canonical_body_digest, sign_ns_scanner_capability, sign_put_file_capability,
sign_tonic_rpc_response_proof, tonic_boot_epoch_challenge, tonic_boot_epoch_response_headers,
tonic_rpc_auth_failure_reason, verify_put_file_auth_trailer, verify_put_file_capability, verify_rpc_signature,
verify_tonic_boot_epoch_response, verify_tonic_canonical_body_digest, verify_tonic_mutation_body_digest,
verify_tonic_rpc_response_proof, verify_tonic_rpc_signature, verify_tonic_rpc_signature_with_bootstrap,
check_and_record_signed_rpc_nonce, decode_heal_bucket_rpc_options, encode_heal_bucket_rpc_options, gen_signature_headers,
gen_tonic_replay_scope_headers, gen_tonic_signature_headers, gen_tonic_signature_interceptor,
node_service_time_out_client, node_service_time_out_client_no_auth, normalize_tonic_rpc_audience,
set_tonic_canonical_body_digest, sign_ns_scanner_capability, sign_put_file_capability, sign_tonic_rpc_response_proof,
tonic_boot_epoch_challenge, tonic_boot_epoch_response_headers, tonic_rpc_auth_failure_reason,
verify_put_file_auth_trailer, verify_put_file_capability, verify_rpc_signature, verify_tonic_boot_epoch_response,
verify_tonic_canonical_body_digest, verify_tonic_mutation_body_digest, verify_tonic_rpc_response_proof,
verify_tonic_rpc_signature, verify_tonic_rpc_signature_with_bootstrap,
};
}
+4 -1
View File
@@ -50,6 +50,9 @@ pub use peer_rest_client::{
SERVICE_SIGNAL_RELOAD_DYNAMIC, ScannerPeerActivity,
};
pub(crate) use peer_s3_client::heal_bucket_local_on_disks;
pub use peer_s3_client::{LocalPeerS3Client, PeerS3Client, S3PeerSys, ScannerBucketListing, ScannerSetBucketListing};
pub use peer_s3_client::{
LocalPeerS3Client, PeerS3Client, S3PeerSys, ScannerBucketListing, ScannerSetBucketListing, decode_heal_bucket_rpc_options,
encode_heal_bucket_rpc_options,
};
pub use remote_disk::RemoteDisk;
pub use remote_locker::RemoteClient;
+414 -23
View File
@@ -18,6 +18,7 @@ use crate::cluster::rpc::client::{
node_service_time_out_client,
};
use crate::cluster::rpc::set_tonic_mutation_body_digest;
use crate::core::pools::PoolMeta;
use crate::disk::error::DiskError;
use crate::disk::error::{Error, Result};
use crate::disk::error_reduce::{BUCKET_OP_IGNORED_ERRS, is_all_buckets_not_found, reduce_write_quorum_errs};
@@ -46,7 +47,12 @@ use std::sync::{
Mutex as StdMutex,
atomic::{AtomicBool, Ordering},
};
use std::{collections::HashMap, fmt::Debug, sync::Arc, time::Duration};
use std::{
collections::{BTreeSet, HashMap},
fmt::Debug,
sync::Arc,
time::Duration,
};
#[cfg(test)]
use tokio::sync::Notify;
use tokio::{net::TcpStream, sync::RwLock, time};
@@ -99,6 +105,9 @@ impl DeleteBucketEmptyScanBarrier {
#[cfg(test)]
static DELETE_BUCKET_EMPTY_SCAN_BARRIER: StdMutex<Option<Arc<DeleteBucketEmptyScanBarrier>>> = StdMutex::new(None);
#[cfg(test)]
static HEAL_BUCKET_PRE_MUTATION_BARRIER: StdMutex<Option<Arc<DeleteBucketEmptyScanBarrier>>> = StdMutex::new(None);
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
enum HealBucketOperation {
Make,
@@ -171,6 +180,15 @@ pub(crate) fn install_delete_bucket_empty_scan_barrier() -> Arc<DeleteBucketEmpt
barrier
}
#[cfg(test)]
fn install_heal_bucket_pre_mutation_barrier() -> Arc<DeleteBucketEmptyScanBarrier> {
let barrier = Arc::new(DeleteBucketEmptyScanBarrier::default());
*HEAL_BUCKET_PRE_MUTATION_BARRIER
.lock()
.expect("heal bucket mutation barrier lock should not be poisoned") = Some(barrier.clone());
barrier
}
#[cfg(test)]
async fn pause_after_delete_bucket_empty_scan() {
let barrier = DELETE_BUCKET_EMPTY_SCAN_BARRIER
@@ -182,6 +200,20 @@ async fn pause_after_delete_bucket_empty_scan() {
}
}
#[cfg(test)]
async fn pause_before_heal_bucket_volume_mutation() {
let barrier = HEAL_BUCKET_PRE_MUTATION_BARRIER
.lock()
.expect("heal bucket mutation barrier lock should not be poisoned")
.take();
if let Some(barrier) = barrier {
barrier.pause().await;
}
}
#[cfg(not(test))]
async fn pause_before_heal_bucket_volume_mutation() {}
#[derive(Clone, Debug)]
pub struct ScannerBucketListing {
pub buckets: Vec<BucketInfo>,
@@ -253,9 +285,41 @@ fn resolve_heal_bucket_mode(opts: &mut HealOpts, pool_errs: &[Option<Error>]) ->
Ok(())
}
#[derive(serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct HealBucketRpcEnvelope {
options: HealOpts,
#[serde(rename = "fencedPools", default)]
fenced_pools: Vec<usize>,
}
pub fn encode_heal_bucket_rpc_options(opts: HealOpts, fenced_pools: &[usize]) -> Result<String> {
serde_json::to_string(&HealBucketRpcEnvelope {
options: opts,
fenced_pools: fenced_pools.to_vec(),
})
.map_err(Into::into)
}
pub fn decode_heal_bucket_rpc_options(payload: &str) -> Result<(HealOpts, Vec<usize>)> {
match serde_json::from_str::<HealBucketRpcEnvelope>(payload) {
Ok(envelope) => Ok((envelope.options, envelope.fenced_pools)),
Err(envelope_err) => serde_json::from_str::<HealOpts>(payload)
.map(|options| (options, Vec::new()))
.map_err(|legacy_err| {
Error::other(format!(
"decode heal bucket RPC options failed: envelope={envelope_err}; legacy={legacy_err}"
))
}),
}
}
#[async_trait]
pub trait PeerS3Client: Debug + Sync + Send + 'static {
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem>;
async fn heal_bucket_with_fence(&self, bucket: &str, opts: &HealOpts, _fenced_pools: &[usize]) -> Result<HealResultItem> {
self.heal_bucket(bucket, opts).await
}
async fn make_bucket(&self, bucket: &str, opts: &MakeBucketOptions) -> Result<()>;
async fn list_bucket(&self, opts: &BucketOptions) -> Result<Vec<BucketInfo>>;
async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()>;
@@ -309,6 +373,10 @@ impl S3PeerSys {
impl S3PeerSys {
pub async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
self.heal_bucket_with_fence(bucket, opts, &[]).await
}
pub async fn heal_bucket_with_fence(&self, bucket: &str, opts: &HealOpts, fenced_pools: &[usize]) -> Result<HealResultItem> {
let mut opts = *opts;
let mut futures = Vec::with_capacity(self.clients.len());
for client in self.clients.iter() {
@@ -331,7 +399,7 @@ impl S3PeerSys {
let opts_clone = opts;
let heal_bucket_results_clone = heal_bucket_results.clone();
futures.push(async move {
match client.heal_bucket(bucket, &opts_clone).await {
match client.heal_bucket_with_fence(bucket, &opts_clone, fenced_pools).await {
Ok(res) => {
heal_bucket_results_clone.write().await[idx] = res;
None
@@ -635,8 +703,18 @@ impl PeerS3Client for LocalPeerS3Client {
}
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
self.heal_bucket_with_fence(bucket, opts, &[]).await
}
async fn heal_bucket_with_fence(&self, bucket: &str, opts: &HealOpts, fenced_pools: &[usize]) -> Result<HealResultItem> {
let disks = self.local_disks_for_pools().await.into_iter().map(Some).collect();
heal_bucket_local_on_disks(bucket, opts, disks).await
let store = runtime_sources::object_store_handle().filter(|store| Arc::ptr_eq(&store.ctx, &self.instance_ctx));
#[cfg(not(test))]
if store.is_none() {
return Err(Error::other("bucket heal refused: pool metadata is unavailable for this instance"));
}
heal_bucket_local_on_disks_with_pool_meta(bucket, opts, disks, store.as_ref().map(|store| &store.pool_meta), fenced_pools)
.await
}
async fn list_bucket(&self, _opts: &BucketOptions) -> Result<Vec<BucketInfo>> {
@@ -1079,9 +1157,13 @@ impl PeerS3Client for RemotePeerS3Client {
}
async fn heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
self.heal_bucket_with_fence(bucket, opts, &[]).await
}
async fn heal_bucket_with_fence(&self, bucket: &str, opts: &HealOpts, fenced_pools: &[usize]) -> Result<HealResultItem> {
self.execute_with_timeout(
|| async {
let options: String = serde_json::to_string(opts)?;
let options = encode_heal_bucket_rpc_options(*opts, fenced_pools)?;
let mut client = self.get_client().await?;
let mut request = Request::new(HealBucketRequest {
bucket: bucket.to_string(),
@@ -1229,6 +1311,117 @@ pub(crate) async fn heal_bucket_local_on_disks(
opts: &HealOpts,
disks: Vec<Option<DiskStore>>,
) -> Result<HealResultItem> {
if let Some(store) = runtime_sources::object_store_handle() {
return heal_bucket_local_on_disks_with_pool_meta(bucket, opts, disks, Some(&store.pool_meta), &[]).await;
}
#[cfg(test)]
return heal_bucket_local_on_disks_with_pool_meta(bucket, opts, disks, None, &[]).await;
#[cfg(not(test))]
Err(Error::other("bucket heal refused: pool metadata is unavailable"))
}
fn disk_pool_index(disk: &DiskStore) -> Result<usize> {
usize::try_from(disk.endpoint().pool_idx)
.map_err(|_| Error::other(format!("invalid bucket-heal pool index {}", disk.endpoint().pool_idx)))
}
fn fenced_decommission_drive_state() -> DriveState {
DriveState::Unknown("skipped-decommission-suspended".to_string())
}
fn heal_bucket_fence_detail(fenced_pools: &BTreeSet<usize>) -> Option<String> {
if fenced_pools.is_empty() {
return None;
}
let pools = fenced_pools.iter().map(usize::to_string).collect::<Vec<_>>().join(", ");
Some(format!("skipped: bucket-volume heal fenced on decommission-suspended pool(s): {pools}"))
}
async fn snapshot_heal_bucket_fence(
disks: &[Option<DiskStore>],
pool_meta: Option<&RwLock<PoolMeta>>,
dispatch_fenced_pools: &[usize],
) -> Result<(Vec<bool>, BTreeSet<usize>)> {
let mut fenced_disks = vec![false; disks.len()];
let mut fenced_pools = dispatch_fenced_pools.iter().copied().collect::<BTreeSet<_>>();
let pool_meta = match pool_meta {
Some(pool_meta) => Some(pool_meta.read().await),
None => None,
};
if let Some(pool_meta) = pool_meta.as_ref()
&& let Some(pool_idx) = fenced_pools.iter().find(|pool_idx| **pool_idx >= pool_meta.pools.len())
{
return Err(Error::other(format!(
"bucket-heal dispatch fence pool index {pool_idx} is absent from {} pool metadata entries",
pool_meta.pools.len()
)));
}
for (disk_index, disk) in disks.iter().enumerate() {
let Some(disk) = disk else {
continue;
};
let pool_idx = disk_pool_index(disk)?;
if let Some(pool_meta) = pool_meta.as_ref() {
if pool_idx >= pool_meta.pools.len() {
return Err(Error::other(format!(
"bucket-heal pool index {pool_idx} is absent from {} pool metadata entries",
pool_meta.pools.len()
)));
}
if pool_meta.is_suspended(pool_idx) {
fenced_pools.insert(pool_idx);
}
}
if fenced_pools.contains(&pool_idx) {
fenced_disks[disk_index] = true;
}
}
Ok((fenced_disks, fenced_pools))
}
async fn run_heal_bucket_volume_mutation<F, Fut>(
disk: &DiskStore,
pool_meta: Option<&RwLock<PoolMeta>>,
operation: F,
) -> Result<Option<usize>>
where
F: FnOnce() -> Fut,
Fut: std::future::Future<Output = Result<()>>,
{
let Some(pool_meta) = pool_meta else {
operation().await?;
return Ok(None);
};
let pool_idx = disk_pool_index(disk)?;
let pool_meta = pool_meta.read().await;
if pool_idx >= pool_meta.pools.len() {
return Err(Error::other(format!(
"bucket-heal pool index {pool_idx} is absent from {} pool metadata entries",
pool_meta.pools.len()
)));
}
if pool_meta.is_suspended(pool_idx) {
return Ok(Some(pool_idx));
}
// Keep the metadata read guard through the disk mutation so a decommission
// transition cannot pass between this state check and the destructive action.
operation().await?;
Ok(None)
}
async fn heal_bucket_local_on_disks_with_pool_meta(
bucket: &str,
opts: &HealOpts,
disks: Vec<Option<DiskStore>>,
pool_meta: Option<&RwLock<PoolMeta>>,
dispatch_fenced_pools: &[usize],
) -> Result<HealResultItem> {
let (fenced_disks, mut fenced_pool_idxs) = snapshot_heal_bucket_fence(&disks, pool_meta, dispatch_fenced_pools).await?;
let fenced_disks = Arc::new(fenced_disks);
let before_state = Arc::new(RwLock::new(vec![String::new(); disks.len()]));
let after_state = Arc::new(RwLock::new(vec![String::new(); disks.len()]));
@@ -1238,7 +1431,14 @@ pub(crate) async fn heal_bucket_local_on_disks(
let bucket = bucket.to_string();
let bs_clone = before_state.clone();
let as_clone = after_state.clone();
let fenced_disks = fenced_disks.clone();
futures.push(async move {
if fenced_disks[index] {
let skipped = fenced_decommission_drive_state().to_string();
bs_clone.write().await[index] = skipped.clone();
as_clone.write().await[index] = skipped;
return None;
}
let disk = match disk {
Some(disk) => disk,
None => {
@@ -1301,9 +1501,14 @@ pub(crate) async fn heal_bucket_local_on_disks(
state: state.to_string(),
});
}
if let Some(detail) = heal_bucket_fence_detail(&fenced_pool_idxs) {
res.detail = detail;
}
return Ok(res);
}
pause_before_heal_bucket_volume_mutation().await;
let mut operation_error = errs
.iter()
.filter_map(|err| match err {
@@ -1315,26 +1520,35 @@ pub(crate) async fn heal_bucket_local_on_disks(
if opts.remove && !bucket.starts_with(disk::RUSTFS_META_BUCKET) && !is_all_buckets_not_found(&errs) {
let mut futures = Vec::new();
for (index, disk) in disks.iter().enumerate() {
if matches!(errs[index].as_ref(), Some(Error::DiskNotFound | Error::VolumeNotFound)) {
if fenced_disks[index] || matches!(errs[index].as_ref(), Some(Error::DiskNotFound | Error::VolumeNotFound)) {
continue;
}
let Some(disk) = disk.clone() else {
continue;
};
let bucket = bucket.to_string();
let mutation_disk = disk.clone();
futures.push(async move {
if let Some(err) = injected_heal_bucket_operation_error(&bucket, index, HealBucketOperation::Delete) {
return (index, Err(err));
}
(index, disk.delete_volume(&bucket, false).await)
let result = run_heal_bucket_volume_mutation(&disk, pool_meta, || async move {
if let Some(err) = injected_heal_bucket_operation_error(&bucket, index, HealBucketOperation::Delete) {
return Err(err);
}
mutation_disk.delete_volume(&bucket, false).await
})
.await;
(index, result)
});
}
for (index, result) in join_all(futures).await {
match result {
Ok(()) | Err(Error::VolumeNotFound) => {
Ok(None) | Err(Error::VolumeNotFound) => {
after_state.write().await[index] = DriveState::Missing.to_string();
}
Ok(Some(pool_idx)) => {
fenced_pool_idxs.insert(pool_idx);
after_state.write().await[index] = fenced_decommission_drive_state().to_string();
}
Err(Error::VolumeNotEmpty) => {
warn!(
bucket,
@@ -1365,30 +1579,38 @@ pub(crate) async fn heal_bucket_local_on_disks(
let bs_clone = before_state.clone();
futures.push(async move {
if bs_clone.read().await[idx] == DriveState::Missing.to_string() {
let Some(disk) = disk.as_ref() else {
return (idx, Some(Error::DiskNotFound));
let Some(disk) = disk else {
return (idx, Err(Error::DiskNotFound));
};
if let Some(err) = injected_heal_bucket_operation_error(&bucket, idx, HealBucketOperation::Make) {
return (idx, Some(err));
}
match disk.make_volume(&bucket).await {
Ok(()) | Err(Error::VolumeExists) => return (idx, None),
Err(err) => return (idx, Some(err)),
}
let mutation_disk = disk.clone();
let result = run_heal_bucket_volume_mutation(&disk, pool_meta, || async move {
if let Some(err) = injected_heal_bucket_operation_error(&bucket, idx, HealBucketOperation::Make) {
return Err(err);
}
match mutation_disk.make_volume(&bucket).await {
Ok(()) | Err(Error::VolumeExists) => Ok(()),
Err(err) => Err(err),
}
})
.await;
return (idx, result);
}
(idx, None)
(idx, Ok(None))
});
}
for (index, result) in join_all(futures).await {
match result {
None => {
Ok(None) => {
if before_state.read().await[index] == DriveState::Missing.to_string() {
after_state.write().await[index] = DriveState::Ok.to_string();
}
}
Some(err) => {
Ok(Some(pool_idx)) => {
fenced_pool_idxs.insert(pool_idx);
after_state.write().await[index] = fenced_decommission_drive_state().to_string();
}
Err(err) => {
after_state.write().await[index] = match &err {
Error::DiskNotFound => DriveState::Offline.to_string(),
_ => DriveState::Corrupt.to_string(),
@@ -1409,6 +1631,10 @@ pub(crate) async fn heal_bucket_local_on_disks(
});
}
if let Some(detail) = heal_bucket_fence_detail(&fenced_pool_idxs) {
res.detail = detail;
}
match operation_error {
Some(err) => Err(err),
None => Ok(res),
@@ -1426,6 +1652,7 @@ async fn clone_drives() -> Vec<Option<DiskStore>> {
#[cfg(test)]
mod tests {
use super::*;
use crate::core::pools::{PoolDecommissionInfo, PoolStatus};
use crate::disk::WalkDirOptions;
use crate::disk::disk_store::LocalDiskWrapper;
use crate::disk::endpoint::Endpoint;
@@ -1599,6 +1826,23 @@ mod tests {
disks
}
fn heal_bucket_pool_meta(suspended_pool: Option<usize>) -> PoolMeta {
PoolMeta {
pools: (0..2)
.map(|pool_idx| PoolStatus {
id: pool_idx,
cmd_line: format!("pool-{pool_idx}"),
last_update: ::time::OffsetDateTime::UNIX_EPOCH,
decommission: (suspended_pool == Some(pool_idx)).then(|| PoolDecommissionInfo {
start_time: Some(::time::OffsetDateTime::UNIX_EPOCH),
..Default::default()
}),
})
.collect(),
..Default::default()
}
}
fn test_remote_peer(addr: &str) -> RemotePeerS3Client {
RemotePeerS3Client {
pools: Some(vec![0]),
@@ -1910,6 +2154,127 @@ mod tests {
reset_local_disk_test_state().await;
}
#[tokio::test]
#[serial]
async fn heal_bucket_rechecks_decommission_before_recreating_volume() {
reset_local_disk_test_state().await;
let temp_dir = TempDir::new().expect("create temp dir for bucket-heal fence regression");
let disks = init_test_local_disks_for_pools(&temp_dir, &[(0, 1), (1, 1)], "heal-bucket-mutation-fence").await;
let bucket = "fenced-recreate-bucket";
disks[0]
.make_volume(bucket)
.await
.expect("active pool should start with the bucket volume");
let pool_meta = Arc::new(RwLock::new(heal_bucket_pool_meta(None)));
let barrier = install_heal_bucket_pre_mutation_barrier();
let heal = tokio::spawn({
let disks = disks.clone();
let pool_meta = pool_meta.clone();
async move {
heal_bucket_local_on_disks_with_pool_meta(
bucket,
&HealOpts {
recreate: true,
..Default::default()
},
disks.into_iter().map(Some).collect(),
Some(pool_meta.as_ref()),
&[],
)
.await
}
});
barrier.wait_until_paused().await;
pool_meta.write().await.pools[1].decommission = Some(PoolDecommissionInfo {
start_time: Some(::time::OffsetDateTime::UNIX_EPOCH),
..Default::default()
});
barrier.release();
let result = heal
.await
.expect("bucket-heal task should join")
.expect("suspended pool should be reported as skipped");
assert!(result.detail.contains("skipped") && result.detail.contains('1'));
assert!(matches!(disks[1].stat_volume(bucket).await, Err(Error::VolumeNotFound)));
reset_local_disk_test_state().await;
}
#[tokio::test]
#[serial]
async fn heal_bucket_dispatch_fence_blocks_stale_active_peer_state() {
reset_local_disk_test_state().await;
let temp_dir = TempDir::new().expect("create temp dir for stale bucket-heal peer regression");
let disks = init_test_local_disks_for_pools(&temp_dir, &[(0, 1), (1, 1)], "heal-bucket-dispatch-fence").await;
let bucket = "dispatch-fenced-bucket";
disks[0]
.make_volume(bucket)
.await
.expect("active pool should start with the bucket volume");
let stale_pool_meta = RwLock::new(heal_bucket_pool_meta(None));
let result = heal_bucket_local_on_disks_with_pool_meta(
bucket,
&HealOpts {
recreate: true,
..Default::default()
},
disks.iter().cloned().map(Some).collect(),
Some(&stale_pool_meta),
&[1],
)
.await
.expect("dispatch fence should override stale active peer metadata");
assert!(result.detail.contains("skipped") && result.detail.contains('1'));
assert!(matches!(disks[1].stat_volume(bucket).await, Err(Error::VolumeNotFound)));
reset_local_disk_test_state().await;
}
#[tokio::test]
#[serial]
async fn heal_bucket_keeps_suspended_pool_volume_on_remove() {
reset_local_disk_test_state().await;
let temp_dir = TempDir::new().expect("create temp dir for bucket-heal delete fence regression");
let disks = init_test_local_disks_for_pools(&temp_dir, &[(0, 1), (1, 1)], "heal-bucket-delete-fence").await;
let bucket = "fenced-remove-bucket";
for disk in &disks {
disk.make_volume(bucket)
.await
.expect("bucket volume should exist before heal");
}
let pool_meta = RwLock::new(heal_bucket_pool_meta(Some(1)));
let result = heal_bucket_local_on_disks_with_pool_meta(
bucket,
&HealOpts {
remove: true,
..Default::default()
},
disks.iter().cloned().map(Some).collect(),
Some(&pool_meta),
&[],
)
.await
.expect("suspended pool should be skipped during bucket-volume removal");
assert!(result.detail.contains("skipped") && result.detail.contains('1'));
assert!(matches!(disks[0].stat_volume(bucket).await, Err(Error::VolumeNotFound)));
disks[1]
.stat_volume(bucket)
.await
.expect("suspended pool bucket volume must not be deleted");
reset_local_disk_test_state().await;
}
#[tokio::test]
#[serial]
async fn heal_bucket_local_dry_run_reports_discovered_drive_states() {
@@ -2123,6 +2488,32 @@ mod tests {
assert!(partial.recreate);
}
#[test]
fn heal_bucket_rpc_envelope_preserves_legacy_compatibility_fail_closed() {
let opts = HealOpts {
recreate: true,
pool: Some(2),
..Default::default()
};
let encoded = encode_heal_bucket_rpc_options(opts, &[1, 2]).expect("encode bucket-heal RPC envelope");
assert!(
serde_json::from_str::<HealOpts>(&encoded).is_err(),
"an old peer must reject the nested request instead of ignoring its dispatch fence"
);
let (decoded, fenced_pools) =
decode_heal_bucket_rpc_options(&encoded).expect("new peer should decode bucket-heal RPC envelope");
assert!(decoded.recreate);
assert_eq!(decoded.pool, Some(2));
assert_eq!(fenced_pools, vec![1, 2]);
let legacy = serde_json::to_string(&opts).expect("encode legacy HealOpts");
let (decoded, fenced_pools) = decode_heal_bucket_rpc_options(&legacy).expect("new peer should accept a legacy request");
assert!(decoded.recreate);
assert_eq!(decoded.pool, Some(2));
assert!(fenced_pools.is_empty());
}
#[tokio::test]
async fn test_make_bucket_reduces_quorum_by_pool_participants() {
let peer_sys = S3PeerSys {
+14 -1
View File
@@ -2124,13 +2124,26 @@ 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 {
+98 -1
View File
@@ -19,6 +19,7 @@ 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 std::collections::BTreeSet;
use tracing::trace;
const LOG_COMPONENT_ECSTORE: &str = "ecstore";
@@ -291,7 +292,45 @@ impl ECStore {
#[instrument(skip(self))]
pub(super) async fn handle_heal_bucket(&self, bucket: &str, opts: &HealOpts) -> Result<HealResultItem> {
let res = self.peer_sys.heal_bucket(bucket, opts).await?;
let mut fenced_pools = BTreeSet::new();
{
let pool_meta = self.pool_meta.read().await;
fenced_pools.extend((0..pool_meta.pools.len()).filter(|pool_idx| pool_meta.is_suspended(*pool_idx)));
if let Some(pool_idx) = opts.pool {
if pool_idx >= pool_meta.pools.len() {
return Err(invalid_heal_pool_index(pool_idx, pool_meta.pools.len()));
}
if pool_meta.is_suspended(pool_idx) {
let complete = pool_meta.pools[pool_idx]
.decommission
.as_ref()
.is_some_and(|decommission| decommission.complete);
return Err(if complete {
StorageError::InvalidArgument(
"heal".to_string(),
"pool".to_string(),
format!("heal pool {pool_idx} has completed decommission"),
)
} else {
Error::SlowDown
});
}
}
}
let dispatch_fenced_pools = fenced_pools.iter().copied().collect::<Vec<_>>();
let mut res = self
.peer_sys
.heal_bucket_with_fence(bucket, opts, &dispatch_fenced_pools)
.await?;
{
let pool_meta = self.pool_meta.read().await;
fenced_pools.extend((0..pool_meta.pools.len()).filter(|pool_idx| pool_meta.is_suspended(*pool_idx)));
}
if !fenced_pools.is_empty() {
let pools = fenced_pools.iter().map(usize::to_string).collect::<Vec<_>>().join(", ");
res.detail = format!("skipped: bucket-volume heal fenced on decommission-suspended pool(s): {pools}");
}
Ok(res)
}
@@ -857,6 +896,64 @@ mod tests {
}
}
#[tokio::test]
async fn scoped_heal_bucket_blocks_before_dispatch_when_pool_is_suspended() {
let mut store = minimal_heal_store().await;
store.pool_meta = RwLock::new(PoolMeta {
pools: vec![
PoolStatus {
id: 0,
cmd_line: "pool-0".to_string(),
last_update: OffsetDateTime::UNIX_EPOCH,
decommission: None,
},
PoolStatus {
id: 1,
cmd_line: "pool-1".to_string(),
last_update: OffsetDateTime::UNIX_EPOCH,
decommission: Some(PoolDecommissionInfo {
start_time: Some(OffsetDateTime::UNIX_EPOCH),
..Default::default()
}),
},
],
..Default::default()
});
let err = store
.handle_heal_bucket(
"bucket",
&HealOpts {
pool: Some(1),
..Default::default()
},
)
.await
.expect_err("suspended pool must be blocked before bucket-heal fan-out");
assert_eq!(err, Error::SlowDown);
store.pool_meta.write().await.pools[1]
.decommission
.as_mut()
.expect("decommission state should exist")
.complete = true;
let err = store
.handle_heal_bucket(
"bucket",
&HealOpts {
pool: Some(1),
..Default::default()
},
)
.await
.expect_err("completed pool must remain fenced from bucket heal");
assert!(
matches!(err, StorageError::InvalidArgument(_, ref field, ref reason)
if field == "pool" && reason.contains("completed decommission")),
"unexpected completed-pool error: {err:?}"
);
}
#[tokio::test]
#[serial_test::serial]
async fn unscoped_heal_object_suspended_owner_semantics() {
+1 -1
View File
@@ -3194,7 +3194,7 @@ impl ECStore {
// Default return value
let mut del_objects = vec![DeletedObject::default(); objects.len()];
let accounting = vec![None; objects.len()];
let mut accounting = vec![None; objects.len()];
let mut del_errs = Vec::with_capacity(objects.len());
for _ in 0..objects.len() {
@@ -271,7 +271,7 @@ pub(super) fn resolve_latest_object_info_candidates(
.filter(|candidate| latest_candidate_mod_time(candidate) == Some(latest_mod_time))
.collect::<Vec<_>>();
latest_candidates.sort_by_key(|candidate| std::cmp::Reverse(candidate.idx));
latest_candidates.sort_by(|left, right| right.idx.cmp(&left.idx));
let Some(winner) = latest_candidates.first() else {
return Err(Error::ErasureReadQuorum);
+3
View File
@@ -43,6 +43,9 @@ allow-git = [
# RustFS fork carrying presigned expiry and constant-time authentication fixes.
# owner: rustfs-maintainers review: 2026-10
"https://github.com/rustfs/s3s.git",
# MiMalloc fork pinned for hotpath allocation counting support.
# owner: houseme review: 2026-10
"https://github.com/xonatius/mimalloc_rust.git",
]
[bans]
+1 -2
View File
@@ -41,7 +41,6 @@
| copy_object_version_restore_test | 2 | |
| copy_source_invalid_date_test | 1 | ✅ |
| create_bucket_region_test | 2 | ✅ |
| data_usage_test | 2 | |
| degraded_read_eof_regression_test | 3 | |
| delete_marker_migration_semantics_test | 2 | ✅ |
| delete_object_no_content_length_test | 1 | |
@@ -100,4 +99,4 @@
| tls_hot_reload_test | 1 | ✅ |
| version_id_regression_test | 10 | ✅ |
**Total listed: 577 tests across 83 modules · PR smoke: 163 tests / 36 modules · merge/main full: 455 tests / 74 modules · nightly replication: 55 tests · nightly cluster faults: 28 tests / 7 modules · nightly protocols: 16 tests** · updated 2026-08-23.
**Total listed: 575 tests across 82 modules · PR smoke: 163 tests / 36 modules · merge/main full: 453 tests / 73 modules · nightly replication: 55 tests · nightly cluster faults: 28 tests / 7 modules · nightly protocols: 16 tests** · updated 2026-08-23.
+2 -2
View File
@@ -336,13 +336,13 @@ opentelemetry = { workspace = true }
tracing-opentelemetry = { workspace = true }
# Data structures
hashbrown = { workspace = true, features = ["serde", "rayon"] }
rustfs-mimalloc = { workspace = true }
mimalloc = { workspace = true }
[target.'cfg(target_os = "linux")'.dependencies]
libsystemd.workspace = true
[target.'cfg(not(target_os = "windows"))'.dependencies]
rustfs-mimalloc-sys.workspace = true
libmimalloc-sys.workspace = true
[dev-dependencies]
uuid = { workspace = true, features = ["v4", "v5", "fast-rng", "macro-diagnostics"] }
+7 -1
View File
@@ -369,8 +369,14 @@ pub fn allocator_reclaim_controller_snapshot(ctx: &CancellationToken) -> Allocat
}
#[cfg(not(target_os = "windows"))]
#[allow(unsafe_code)]
fn collect_allocator_memory(force: bool) -> Result<(), String> {
rustfs_mimalloc::MiMalloc::collect(force);
// SAFETY: `mi_collect` is provided by the active global allocator backend
// on this target family. It is explicitly intended to reclaim retained
// pages/segments and does not require additional invariants from the caller.
unsafe {
libmimalloc_sys::mi_collect(force);
}
Ok(())
}
+8 -10
View File
@@ -26,22 +26,22 @@ struct MiMallocAllocator;
unsafe impl GlobalAlloc for MiMallocAllocator {
unsafe fn alloc(&self, layout: Layout) -> *mut u8 {
// SAFETY: the caller upholds GlobalAlloc's contract for layout.
unsafe { rustfs_mimalloc::MiMalloc.alloc(layout) }
unsafe { mimalloc::MiMalloc.alloc(layout) }
}
unsafe fn alloc_zeroed(&self, layout: Layout) -> *mut u8 {
// SAFETY: the caller upholds GlobalAlloc's contract for layout.
unsafe { rustfs_mimalloc::MiMalloc.alloc_zeroed(layout) }
unsafe { mimalloc::MiMalloc.alloc_zeroed(layout) }
}
unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) {
// SAFETY: ptr and layout came from this allocator and are forwarded unchanged.
unsafe { rustfs_mimalloc::MiMalloc.dealloc(ptr, layout) }
unsafe { mimalloc::MiMalloc.dealloc(ptr, layout) }
}
unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, new_size: usize) -> *mut u8 {
// SAFETY: ptr and layout came from this allocator and are forwarded unchanged.
unsafe { rustfs_mimalloc::MiMalloc.realloc(ptr, layout, new_size) }
unsafe { mimalloc::MiMalloc.realloc(ptr, layout, new_size) }
}
}
@@ -51,7 +51,7 @@ static GLOBAL: hotpath::CountingAllocator<MiMallocAllocator> = hotpath::Counting
#[cfg(not(all(feature = "hotpath", feature = "hotpath-alloc")))]
#[global_allocator]
static GLOBAL: rustfs_mimalloc::MiMalloc = rustfs_mimalloc::MiMalloc;
static GLOBAL: mimalloc::MiMalloc = mimalloc::MiMalloc;
fn main() {
let _hotpath_guard = hotpath::HotpathGuardBuilder::new("main").build();
@@ -71,9 +71,8 @@ mod tests {
allocation.extend_from_slice(&[7_u8; 64]);
assert_eq!(allocation.len(), 64);
let heap = rustfs_mimalloc::heap::Heap::main();
// SAFETY: the live Vec pointer is valid to inspect for heap ownership.
assert!(unsafe { heap.contains(allocation.as_ptr()) });
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(allocation.as_ptr().cast()) });
}
#[test]
@@ -86,13 +85,12 @@ mod tests {
let layout = Layout::from_size_align(32, 8).expect("valid test allocation layout");
let grown_layout = Layout::from_size_align(64, 8).expect("valid grown test allocation layout");
let allocator = super::MiMallocAllocator;
let heap = rustfs_mimalloc::heap::Heap::main();
// SAFETY: The pointer is checked for null before use and later released
// through the same allocator with the corresponding layout.
let ptr = unsafe { allocator.alloc_zeroed(layout) };
assert!(!ptr.is_null());
assert!(unsafe { heap.contains(ptr) });
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(ptr.cast()) });
assert!(unsafe { std::slice::from_raw_parts(ptr, 32).iter().all(|byte| *byte == 0) });
// SAFETY: `ptr` was allocated by `allocator` with `layout`; on failure
@@ -104,7 +102,7 @@ mod tests {
panic!("mimalloc realloc failed in allocator smoke test");
}
assert!(unsafe { heap.contains(grown_ptr) });
assert!(unsafe { libmimalloc_sys::mi_is_in_heap_region(grown_ptr.cast()) });
// SAFETY: `grown_ptr` was reallocated by `allocator` and is released
// with the matching grown layout.
unsafe { allocator.dealloc(grown_ptr, grown_layout) };
+36 -19
View File
@@ -17,7 +17,10 @@ use rustfs_io_metrics::{
record_cpu_usage, record_memory_usage, record_process_memory_split,
};
use serde::Serialize;
#[cfg(any(test, not(target_os = "windows")))]
use serde_json::Value;
#[cfg(not(target_os = "windows"))]
use std::ffi::CStr;
use std::path::Path;
use std::sync::{Arc, Mutex, OnceLock};
use std::time::Duration;
@@ -228,18 +231,7 @@ fn read_cgroup_memory_snapshot() -> Option<CgroupMemorySnapshot> {
read_cgroup_v2().or_else(read_cgroup_v1)
}
fn read_allocator_memory_snapshot() -> Option<AllocatorMemorySnapshot> {
let json = rustfs_mimalloc::MiMalloc::stats_json();
if json.is_empty() {
return None;
}
let observation = parse_mimalloc_stats_json(&json)?;
Some(AllocatorMemorySnapshot {
backend: crate::allocator_reclaim::allocator_backend(),
observation,
})
}
#[cfg(any(test, not(target_os = "windows")))]
fn numeric_json_value(value: &Value) -> Option<u64> {
match value {
Value::Number(number) => number
@@ -250,6 +242,7 @@ fn numeric_json_value(value: &Value) -> Option<u64> {
}
}
#[cfg(any(test, not(target_os = "windows")))]
fn numeric_json_field(value: &Value, field: &str) -> Option<u64> {
match value {
Value::Object(fields) => fields
@@ -261,6 +254,7 @@ fn numeric_json_field(value: &Value, field: &str) -> Option<u64> {
}
}
#[cfg(any(test, not(target_os = "windows")))]
fn mimalloc_stat_field(value: &Value, metric: &str, field: &str) -> Option<u64> {
match value {
Value::Object(fields) => {
@@ -277,10 +271,12 @@ fn mimalloc_stat_field(value: &Value, metric: &str, field: &str) -> Option<u64>
}
}
#[cfg(any(test, not(target_os = "windows")))]
fn mimalloc_stat_current(value: &Value, metric: &str) -> Option<u64> {
mimalloc_stat_field(value, metric, "current")
}
#[cfg(any(test, not(target_os = "windows")))]
fn mimalloc_stat_sum(value: &Value, metrics: &[&str], field: &str) -> Option<u64> {
metrics
.iter()
@@ -289,6 +285,7 @@ fn mimalloc_stat_sum(value: &Value, metrics: &[&str], field: &str) -> Option<u64
.filter(|value| *value > 0)
}
#[cfg(any(test, not(target_os = "windows")))]
fn parse_mimalloc_stats_json(stats_json: &str) -> Option<AllocatorMemoryObservation> {
let value = serde_json::from_str::<Value>(stats_json).ok()?;
let malloc_metrics = ["malloc_normal", "malloc_huge"];
@@ -315,6 +312,33 @@ fn parse_mimalloc_stats_json(stats_json: &str) -> Option<AllocatorMemoryObservat
}
}
#[cfg(not(target_os = "windows"))]
#[allow(unsafe_code)]
fn read_allocator_memory_snapshot() -> Option<AllocatorMemorySnapshot> {
// SAFETY: `mi_stats_get_json` returns a null-terminated JSON buffer owned by
// mimalloc when called with a null input buffer. The mimalloc API requires
// freeing that buffer with `mi_free`; parsing finishes before the buffer is freed.
let observation = unsafe {
let stats_ptr = libmimalloc_sys::mi_stats_get_json(0, std::ptr::null_mut());
if stats_ptr.is_null() {
return None;
}
let observation = CStr::from_ptr(stats_ptr).to_str().ok().and_then(parse_mimalloc_stats_json);
libmimalloc_sys::mi_free(stats_ptr.cast());
observation?
};
Some(AllocatorMemorySnapshot {
backend: crate::allocator_reclaim::allocator_backend(),
observation,
})
}
#[cfg(target_os = "windows")]
fn read_allocator_memory_snapshot() -> Option<AllocatorMemorySnapshot> {
None
}
fn configured_memory_observability_interval_secs() -> u64 {
rustfs_utils::get_env_u64(ENV_MEMORY_OBSERVABILITY_INTERVAL_SECS, DEFAULT_MEMORY_OBSERVABILITY_INTERVAL_SECS).max(1)
}
@@ -542,13 +566,6 @@ mod tests {
assert_eq!(parse_mimalloc_stats_json(r#"{ "allocator": "unknown" }"#), None);
}
#[test]
fn read_allocator_memory_snapshot_uses_mimalloc_stats_json() {
let snapshot = super::read_allocator_memory_snapshot();
#[cfg(not(target_os = "windows"))]
assert!(snapshot.is_some(), "allocator snapshot should be available on non-Windows");
}
#[test]
fn memory_observability_snapshot_reports_disabled_when_metrics_are_disabled() {
let snapshot = build_memory_observability_status_snapshot(false, 15, false);
@@ -17,9 +17,8 @@ use crate::storage::storage_api::rpc_consumer::node_service::contract::bucket::{
BucketOptions, DeleteBucketOptions, MakeBucketOptions,
};
use crate::storage::storage_api::rpc_consumer::node_service::{
DiskError, StoragePeerS3ClientExt as _, reload_bucket_metadata, remove_bucket_metadata,
DiskError, StoragePeerS3ClientExt as _, decode_heal_bucket_rpc_options, reload_bucket_metadata, remove_bucket_metadata,
};
use rustfs_common::heal_channel::HealOpts;
use rustfs_protos::proto_gen::node_service::*;
use tonic::{Request, Response, Status};
use tracing::debug;
@@ -239,7 +238,7 @@ impl NodeService {
) -> Result<Response<HealBucketResponse>, Status> {
debug!("heal bucket");
let request = request.into_inner();
let options = match serde_json::from_str::<HealOpts>(&request.options) {
let (options, fenced_pools) = match decode_heal_bucket_rpc_options(&request.options) {
Ok(options) => options,
Err(err) => {
return Ok(Response::new(HealBucketResponse {
@@ -249,7 +248,11 @@ impl NodeService {
}
};
match self.local_peer.heal_bucket(&request.bucket, &options).await {
match self
.local_peer
.heal_bucket_with_fence(&request.bucket, &options, &fenced_pools)
.await
{
Ok(_) => Ok(Response::new(HealBucketResponse {
success: true,
error: None,
+20 -4
View File
@@ -239,6 +239,7 @@ pub(crate) mod rpc_consumer {
}
pub(crate) mod node_service {
pub(crate) use super::super::ecstore_rpc::decode_heal_bucket_rpc_options;
pub(crate) use super::super::storage_contracts::{
SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION, SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION,
};
@@ -520,10 +521,10 @@ pub(crate) mod ecstore_rpc {
pub(crate) use rustfs_ecstore::api::rpc::{
KMS_SIGNAL_SUBSYSTEM, LocalPeerS3Client, PEER_RESTDRY_RUN, PEER_RESTSIGNAL, PEER_RESTSUB_SYS, PeerRestClient,
PeerS3Client, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, TONIC_RPC_PREFIX,
check_and_record_signed_rpc_nonce, normalize_tonic_rpc_audience, sign_ns_scanner_capability, sign_put_file_capability,
sign_tonic_rpc_response_proof, tonic_boot_epoch_challenge, tonic_boot_epoch_response_headers,
tonic_rpc_auth_failure_reason, verify_put_file_auth_trailer, verify_rpc_signature, verify_tonic_canonical_body_digest,
verify_tonic_mutation_body_digest, verify_tonic_rpc_signature_with_bootstrap,
check_and_record_signed_rpc_nonce, decode_heal_bucket_rpc_options, normalize_tonic_rpc_audience,
sign_ns_scanner_capability, sign_put_file_capability, sign_tonic_rpc_response_proof, tonic_boot_epoch_challenge,
tonic_boot_epoch_response_headers, tonic_rpc_auth_failure_reason, verify_put_file_auth_trailer, verify_rpc_signature,
verify_tonic_canonical_body_digest, verify_tonic_mutation_body_digest, verify_tonic_rpc_signature_with_bootstrap,
};
#[cfg(test)]
pub(crate) use rustfs_ecstore::api::rpc::{
@@ -1412,6 +1413,12 @@ pub(crate) trait StoragePeerS3ClientExt {
bucket: &str,
opts: &rustfs_common::heal_channel::HealOpts,
) -> DiskResult<rustfs_madmin::heal_commands::HealResultItem>;
async fn heal_bucket_with_fence(
&self,
bucket: &str,
opts: &rustfs_common::heal_channel::HealOpts,
fenced_pools: &[usize],
) -> DiskResult<rustfs_madmin::heal_commands::HealResultItem>;
async fn make_bucket(&self, bucket: &str, opts: &contract::bucket::MakeBucketOptions) -> DiskResult<()>;
async fn list_bucket(&self, opts: &contract::bucket::BucketOptions) -> DiskResult<Vec<contract::bucket::BucketInfo>>;
async fn delete_bucket(&self, bucket: &str, opts: &contract::bucket::DeleteBucketOptions) -> DiskResult<()>;
@@ -1431,6 +1438,15 @@ impl StoragePeerS3ClientExt for LocalPeerS3Client {
ecstore_rpc::PeerS3Client::heal_bucket(self, bucket, opts).await
}
async fn heal_bucket_with_fence(
&self,
bucket: &str,
opts: &rustfs_common::heal_channel::HealOpts,
fenced_pools: &[usize],
) -> DiskResult<rustfs_madmin::heal_commands::HealResultItem> {
ecstore_rpc::PeerS3Client::heal_bucket_with_fence(self, bucket, opts, fenced_pools).await
}
async fn make_bucket(&self, bucket: &str, opts: &contract::bucket::MakeBucketOptions) -> DiskResult<()> {
ecstore_rpc::PeerS3Client::make_bucket(self, bucket, opts).await
}