Compare commits

..

4 Commits

Author SHA1 Message Date
马登山 82fb0a8843 fix(scanner): correct observational usage arguments 2026-08-23 10:31:42 +08:00
马登山 307510749e style: format usage freshness changes 2026-08-23 10:28:47 +08:00
马登山 5496e14960 fix(ecstore): preserve quota baseline across restart 2026-08-23 10:20:39 +08:00
马登山 ec3b7a7dc6 fix(scanner): publish partial usage observations 2026-08-23 09:50:16 +08:00
27 changed files with 820 additions and 405 deletions
+1 -2
View File
@@ -30,8 +30,7 @@ make build-docker BUILD_OS=ubuntu22.04
- Crate membership: `Cargo.toml` `[workspace].members`
- Architecture, layering, crate map: [ARCHITECTURE.md](ARCHITECTURE.md)
- Migration guardrails & readiness contracts: [docs/architecture/](docs/architecture/README.md)
- CI workflow steps: `.github/workflows/`; event, timeout, and required-status
matrix: [docs/testing/ci-gates.md](docs/testing/ci-gates.md)
- CI gates: `.github/workflows/ci.yml` (source of truth; never copy its steps into docs)
- Test-layer taxonomy, per-layer entry commands, serial/nextest rules, flake
policy: [docs/testing/README.md](docs/testing/README.md)
- Tier/ILM transition debugging (xl.meta inspection, versionId tracing):
-2
View File
@@ -70,8 +70,6 @@ make pre-pr
> For the full test-layer taxonomy (unit / ecstore black-box / e2e / s3s-e2e / S3 compatibility / chaos / fuzz / bench), each layer's entry command, the naming conventions the migration gate depends on, and the serial/nextest rules, see [docs/testing/README.md](docs/testing/README.md).
> For the event, timeout, required-status, and local reproduction matrix, see [docs/testing/ci-gates.md](docs/testing/ci-gates.md).
### 🔒 Automated Pre-commit Hooks
#### What `make pre-commit` and `make pre-pr` actually run
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" }
+108 -1
View File
@@ -236,6 +236,15 @@ pub struct DataUsageInfo {
/// without relying on synchronized clocks.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub usage_snapshot_authoritative_baseline: Option<DataUsageSnapshotIdentity>,
/// Per-set freshness for an observational aggregate. A set entry is
/// never sufficient to make the aggregate authoritative; it only records
/// which last-known-good generation contributed to the view.
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub usage_snapshot_set_states: Vec<DataUsageSnapshotSetState>,
/// An observational view may contain only the sets that completed this
/// cycle (or retained a compatible last-known-good cache).
#[serde(default)]
pub usage_snapshot_partial: bool,
/// Deprecated kept here for backward compatibility reasons
pub bucket_sizes: HashMap<String, u64>,
/// Per-disk snapshot information when available
@@ -252,6 +261,22 @@ pub struct DataUsageSnapshotIdentity {
pub scanner_epoch: Option<u64>,
}
#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
pub struct DataUsageSnapshotSetState {
pub pool_index: u64,
pub set_index: u64,
#[serde(default)]
pub scanner_cycle: Option<u64>,
#[serde(default)]
pub scanner_epoch: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub scan_plan_digest: Option<[u8; 32]>,
#[serde(default)]
pub complete: bool,
#[serde(default)]
pub tombstone: bool,
}
impl DataUsageInfo {
pub fn snapshot_identity(&self) -> DataUsageSnapshotIdentity {
DataUsageSnapshotIdentity {
@@ -291,7 +316,7 @@ pub fn data_usage_snapshot_is_newer(candidate: &DataUsageInfo, baseline: &DataUs
/// rollback delete/recreate fences the previous bucket incarnation too.
pub fn observed_data_usage_is_newer(observed: &DataUsageInfo, authoritative: &DataUsageInfo) -> bool {
observed.usage_snapshot_converged == Some(false)
&& observed.is_complete_bucket_usage_snapshot()
&& (observed.is_complete_bucket_usage_snapshot() || observed.is_valid_partial_snapshot())
&& observed.usage_snapshot_authoritative_baseline.as_ref() == Some(&authoritative.snapshot_identity())
&& data_usage_snapshot_is_newer(observed, authoritative)
}
@@ -1436,6 +1461,39 @@ impl DataUsageInfo {
&& u64::try_from(self.buckets_usage.len()).ok() == Some(self.buckets_count)
}
/// Validate provenance before an observational view can be selected for
/// admin display. Partial data is accepted only with unique set states,
/// a plan digest for every state, and at least one usable generation.
pub fn is_valid_partial_snapshot(&self) -> bool {
if !self.usage_snapshot_partial
|| self.usage_snapshot_converged != Some(false)
|| self.last_update.is_none()
|| self.scanner_cycle.is_none()
|| self.scanner_epoch.is_none()
|| self.usage_snapshot_set_states.is_empty()
|| u64::try_from(self.buckets_usage.len()).ok() != Some(self.buckets_count)
{
return false;
}
let mut previous = None;
let mut plan_digest = None;
let mut has_source = false;
for state in &self.usage_snapshot_set_states {
if state.scan_plan_digest.is_none()
|| plan_digest.is_some_and(|digest| Some(digest) != state.scan_plan_digest)
|| state.scanner_cycle.is_some() != state.scanner_epoch.is_some()
|| previous.is_some_and(|(pool, set)| (pool, set) >= (state.pool_index, state.set_index))
{
return false;
}
previous = Some((state.pool_index, state.set_index));
plan_digest = state.scan_plan_digest;
has_source |= state.scanner_cycle.is_some() && !state.tombstone;
}
has_source
}
/// Add object metadata to data usage statistics
pub fn add_object(&mut self, object_path: &str, meta_object: &rustfs_filemeta::MetaObject) {
// This method is kept for backward compatibility
@@ -2263,6 +2321,55 @@ mod tests {
assert!(!observed_data_usage_is_newer(&candidate(2, 9, Some(false), true), &authoritative));
assert!(!observed_data_usage_is_newer(&candidate(2, 11, Some(true), true), &authoritative));
assert!(!observed_data_usage_is_newer(&candidate(2, 11, Some(false), false), &authoritative));
let mut partial = candidate(2, 11, Some(false), false);
partial.usage_snapshot_partial = true;
partial.usage_snapshot_set_states = vec![DataUsageSnapshotSetState {
pool_index: 0,
set_index: 0,
scanner_cycle: Some(10),
scanner_epoch: Some(2),
scan_plan_digest: Some([1; 32]),
complete: false,
tombstone: false,
}];
assert!(observed_data_usage_is_newer(&partial, &authoritative));
}
#[test]
fn mixed_topology_snapshot_is_rejected() {
let mut partial = DataUsageInfo {
last_update: Some(SystemTime::UNIX_EPOCH + Duration::from_secs(2)),
scanner_cycle: Some(11),
scanner_epoch: Some(2),
buckets_count: 0,
usage_snapshot_converged: Some(false),
usage_snapshot_partial: true,
usage_snapshot_set_states: vec![
DataUsageSnapshotSetState {
pool_index: 0,
set_index: 0,
scanner_cycle: Some(11),
scanner_epoch: Some(2),
scan_plan_digest: Some([1; 32]),
complete: true,
tombstone: false,
},
DataUsageSnapshotSetState {
pool_index: 1,
set_index: 0,
scanner_cycle: Some(10),
scanner_epoch: Some(2),
scan_plan_digest: Some([2; 32]),
complete: false,
tombstone: false,
},
],
..Default::default()
};
assert!(!partial.is_valid_partial_snapshot());
partial.usage_snapshot_set_states[1].scan_plan_digest = Some([1; 32]);
assert!(partial.is_valid_partial_snapshot());
}
#[test]
@@ -15,14 +15,12 @@
use crate::common::RustFSTestClusterEnvironment;
use aws_sdk_s3::Client;
use aws_sdk_s3::error::SdkError;
use aws_sdk_s3::types::{CorsConfiguration, CorsRule};
use bytes::Bytes;
use std::sync::Arc;
use tokio::sync::Barrier;
use tracing::{info, warn};
const BUCKET: &str = "conditional-put-race-bucket";
const BUCKET_METADATA_RELOAD_BUCKET: &str = "bucket-metadata-reload-barrier";
async fn cleanup_object(client: &Client, key: &str) {
if let Err(e) = client.delete_object().bucket(BUCKET).key(key).send().await {
@@ -30,16 +28,6 @@ async fn cleanup_object(client: &Client, key: &str) {
}
}
async fn assert_bucket_cors_missing(client: &Client) {
let result = client.get_bucket_cors().bucket(BUCKET_METADATA_RELOAD_BUCKET).send().await;
match result {
Err(SdkError::ServiceError(error)) => {
assert_eq!(error.err().meta().code(), Some("NoSuchCORSConfiguration"));
}
result => panic!("expected the peer to report a missing CORS configuration: {result:?}"),
}
}
async fn conditional_put(
client: &Client,
key: &str,
@@ -248,48 +236,3 @@ async fn test_conditional_put_basic_cluster() -> Result<(), Box<dyn std::error::
cleanup_object(&client, test_key).await;
Ok(())
}
#[tokio::test]
async fn test_bucket_cors_write_is_visible_on_peer_before_response() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
crate::common::init_logging();
let mut cluster = RustFSTestClusterEnvironment::new(2).await?;
cluster.start().await?;
cluster.create_test_bucket(BUCKET_METADATA_RELOAD_BUCKET).await?;
let writer = cluster.create_s3_client(0)?;
let reader = cluster.create_s3_client(1)?;
assert_bucket_cors_missing(&reader).await;
let rule = CorsRule::builder()
.allowed_methods("GET")
.allowed_origins("https://example.com")
.build()?;
let configuration = CorsConfiguration::builder().cors_rules(rule).build()?;
writer
.put_bucket_cors()
.bucket(BUCKET_METADATA_RELOAD_BUCKET)
.cors_configuration(configuration)
.send()
.await?;
let response = reader.get_bucket_cors().bucket(BUCKET_METADATA_RELOAD_BUCKET).send().await?;
let rules = response.cors_rules();
assert_eq!(
rules.len(),
1,
"peer should observe the committed CORS rule before the write response returns"
);
assert_eq!(rules[0].allowed_methods(), ["GET"]);
assert_eq!(rules[0].allowed_origins(), ["https://example.com"]);
writer
.delete_bucket_cors()
.bucket(BUCKET_METADATA_RELOAD_BUCKET)
.send()
.await?;
assert_bucket_cors_missing(&reader).await;
writer.delete_bucket().bucket(BUCKET_METADATA_RELOAD_BUCKET).send().await?;
Ok(())
}
@@ -85,7 +85,6 @@ const PEER_REST_RECOVERY_MAX_ATTEMPTS: u32 = 60;
const PEER_REST_RECOVERY_MAX_BACKOFF: Duration = Duration::from_secs(30);
const SCANNER_ACTIVITY_MAX_MESSAGE_SIZE: usize = 1024;
const REPLICATION_STATS_MAX_MESSAGE_SIZE: usize = 8 * 1024 * 1024;
const BUCKET_METADATA_RELOAD_TIMEOUT: Duration = Duration::from_secs(5);
/// Error for a peer that reported `success = false` without an `error_info` payload.
///
@@ -1329,38 +1328,27 @@ impl PeerRestClient {
}
pub async fn load_bucket_metadata(&self, bucket: &str, scanner_maintenance_change: bool) -> Result<()> {
let result = tokio::time::timeout(BUCKET_METADATA_RELOAD_TIMEOUT, async {
let result = self.load_bucket_metadata_once(bucket, scanner_maintenance_change).await;
if let Err(err) = &result
&& Self::is_network_like_error(err)
{
self.prepare_retry().await;
return self.load_bucket_metadata_once(bucket, scanner_maintenance_change).await;
self.finalize_result(
async {
let mut client = self.get_client().await?;
let mut request = Request::new(LoadBucketMetadataRequest {
bucket: bucket.to_string(),
scanner_maintenance_change,
});
set_tonic_mutation_body_digest(&mut request)?;
let response = client.load_bucket_metadata(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("load_bucket_metadata", Some(bucket)));
}
Ok(())
}
result
})
.await,
)
.await
.unwrap_or_else(|_| Err(Error::other(format!("load_bucket_metadata({bucket}) timed out"))));
self.finalize_result(result).await
}
async fn load_bucket_metadata_once(&self, bucket: &str, scanner_maintenance_change: bool) -> Result<()> {
let mut client = self.get_client().await?;
let mut request = Request::new(LoadBucketMetadataRequest {
bucket: bucket.to_string(),
scanner_maintenance_change,
});
set_tonic_mutation_body_digest(&mut request)?;
request.set_timeout(BUCKET_METADATA_RELOAD_TIMEOUT);
let response = client.load_bucket_metadata(request).await?.into_inner();
if !response.success {
if let Some(msg) = response.error_info {
return Err(Error::other(msg));
}
return Err(peer_failure_without_details("load_bucket_metadata", Some(bucket)));
}
Ok(())
}
pub async fn delete_bucket_metadata(&self, bucket: &str) -> Result<()> {
+109 -3
View File
@@ -73,6 +73,16 @@ struct CachedBucketUsage {
// mutation. A strictly later generation is required before the mutation
// evidence can be discarded.
pending_scanner_position: Option<(u64, u64)>,
// Deletes are visible to admin immediately, but quota admission keeps
// them pending until a complete scanner generation reconciles the set.
// This marker intentionally remains process-local: the delete request
// updates this overlay before the scanner writes a durable snapshot. If
// the process restarts first, loading the persisted complete snapshot
// restores the pre-reconciliation (larger) baseline, which is
// conservative for quota admission. A persisted post-delete snapshot is
// necessarily a complete scanner reconciliation and therefore creates a
// fresh cache entry with no pending hold.
pending_negative_delta: u64,
}
type UsageMemoryCache = Arc<RwLock<HashMap<String, CachedBucketUsage>>>;
@@ -948,7 +958,12 @@ async fn load_observed_data_usage_snapshot(store: Arc<ECStore>) -> Option<DataUs
};
match parse_usage_snapshot(&data) {
Ok(info) if info.usage_snapshot_converged == Some(false) && info.is_complete_bucket_usage_snapshot() => Some(info),
Ok(info)
if info.usage_snapshot_converged == Some(false)
&& (info.is_complete_bucket_usage_snapshot() || info.is_valid_partial_snapshot()) =>
{
Some(info)
}
Ok(_) => {
error!(
event = "data_usage_snapshot_load_failed",
@@ -993,7 +1008,7 @@ async fn load_admin_data_usage_from_backend(store: Arc<ECStore>) -> Result<DataU
}
fn discard_incomplete_bucket_usage(data_usage_info: &mut DataUsageInfo) {
if !data_usage_info.is_complete_bucket_usage_snapshot() {
if !data_usage_info.is_complete_bucket_usage_snapshot() && !data_usage_info.usage_snapshot_partial {
data_usage_info.usage_snapshot_complete = false;
data_usage_info.buckets_usage.clear();
data_usage_info.bucket_sizes.clear();
@@ -1643,6 +1658,7 @@ fn cached_bucket_usage_from_backend(usage: BucketUsageInfo, updated_at: SystemTi
dirty: false,
stale_snapshot_pending: false,
pending_scanner_position: None,
pending_negative_delta: 0,
}
}
@@ -1656,6 +1672,7 @@ fn cached_bucket_usage_now(usage: BucketUsageInfo) -> CachedBucketUsage {
dirty: false,
stale_snapshot_pending: false,
pending_scanner_position: None,
pending_negative_delta: 0,
}
}
@@ -1808,6 +1825,7 @@ pub async fn record_bucket_object_delete_memory(bucket: &str, deleted_size: u64,
.or_insert_with(|| cached_bucket_usage_now(BucketUsageInfo::default()));
entry.usage.size = entry.usage.size.saturating_sub(deleted_size);
entry.pending_negative_delta = entry.pending_negative_delta.saturating_add(deleted_size);
if removed_current_object {
entry.usage.objects_count = entry.usage.objects_count.saturating_sub(1);
entry.usage.versions_count = entry.usage.versions_count.saturating_sub(1);
@@ -1863,7 +1881,7 @@ pub async fn get_bucket_usage_memory(bucket: &str) -> Option<u64> {
cache
.get(bucket)
.filter(|cached| cached.authoritative)
.map(|cached| cached.usage.size)
.map(|cached| cached.usage.size.saturating_add(cached.pending_negative_delta))
}
async fn update_usage_cache_if_needed() {
@@ -2943,6 +2961,45 @@ mod tests {
assert_eq!(selected.usage_snapshot_converged, Some(true));
}
#[test]
fn persisted_authoritative_stalls_but_memory_overlay_remains_visible() {
let authoritative = DataUsageInfo {
last_update: Some(SystemTime::UNIX_EPOCH),
scanner_epoch: Some(4),
scanner_cycle: Some(10),
usage_snapshot_complete: true,
..Default::default()
};
let mut partial = authoritative.clone();
partial.last_update = Some(SystemTime::UNIX_EPOCH + Duration::from_secs(1));
partial.scanner_cycle = Some(11);
partial.usage_snapshot_complete = false;
partial.usage_snapshot_partial = true;
partial.usage_snapshot_converged = Some(false);
partial.usage_snapshot_authoritative_baseline = Some(authoritative.snapshot_identity());
partial.usage_snapshot_set_states = vec![rustfs_data_usage::DataUsageSnapshotSetState {
pool_index: 0,
set_index: 0,
scanner_cycle: Some(10),
scanner_epoch: Some(4),
scan_plan_digest: Some([1; 32]),
complete: false,
tombstone: false,
}];
partial.buckets_usage.insert(
"bucket".to_string(),
BucketUsageInfo {
size: 100,
..Default::default()
},
);
partial.buckets_count = 1;
let (selected, _) = select_admin_data_usage_snapshot(authoritative, true, Some(partial));
assert!(selected.usage_snapshot_partial);
assert_eq!(selected.buckets_usage.get("bucket").map(|usage| usage.size), Some(100));
}
#[tokio::test]
async fn authoritative_save_cleanup_removes_observed_snapshot_best_effort() {
let store = UsageCasStore::default();
@@ -4665,6 +4722,55 @@ mod tests {
);
}
#[tokio::test]
#[serial]
async fn partial_usage_is_observational_not_authoritative_for_quota() {
clear_usage_memory_cache_for_test().await;
let mut partial = data_usage_info_for_test("bucket-a", 10, 100, SystemTime::now());
partial.usage_snapshot_complete = false;
partial.usage_snapshot_partial = true;
replace_bucket_usage_memory_from_info(&partial).await;
assert_eq!(get_bucket_usage_memory("bucket-a").await, None);
}
#[tokio::test]
#[serial]
async fn stale_quota_uses_complete_baseline_plus_positive_deltas() {
clear_usage_memory_cache_for_test().await;
let baseline = data_usage_info_for_test("bucket-a", 1, 100, SystemTime::now());
replace_bucket_usage_memory_from_info(&baseline).await;
record_bucket_object_write_memory("bucket-a", None, 25).await;
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(125));
}
#[tokio::test]
#[serial]
async fn negative_delta_waits_for_set_reconciliation() {
clear_usage_memory_cache_for_test().await;
let baseline = data_usage_info_for_test("bucket-a", 1, 100, SystemTime::UNIX_EPOCH + Duration::from_secs(100));
replace_bucket_usage_memory_from_info(&baseline).await;
record_bucket_object_delete_memory("bucket-a", 25, true).await;
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(100));
// Simulate a process restart: the request-path overlay is gone, but
// the persisted authoritative snapshot is still the pre-reconciliation
// baseline. Quota must remain conservative until a complete scanner
// result proves the delete.
clear_usage_memory_cache_for_test().await;
replace_bucket_usage_memory_from_info(&baseline).await;
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(100));
let reconciled = data_usage_info_for_test("bucket-a", 0, 75, SystemTime::UNIX_EPOCH + Duration::from_secs(101));
replace_bucket_usage_memory_from_info(&reconciled).await;
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(75));
}
#[tokio::test]
#[serial]
async fn memory_overlay_counts_versioned_overwrite_as_new_version() {
+6 -38
View File
@@ -784,24 +784,6 @@ 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>,
@@ -814,7 +796,12 @@ pub async fn create_bitrot_writer(
let writer = if is_inline_buffer {
CustomWriter::new_inline_buffer()
} else if let Some(disk) = disk {
let length = bitrot_create_file_size(length, shard_size, &checksum_algo);
let length = if length > 0 {
let length = length as usize;
(length.div_ceil(shard_size) * checksum_algo.size() + length) as i64
} else {
0
};
let file = disk.create_file("", volume, path, length).await?;
#[cfg(feature = "hotpath")]
@@ -833,25 +820,6 @@ 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>,
}
+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 {
+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);
+21 -3
View File
@@ -28,8 +28,9 @@ use rustfs_common::heal_channel::HealScanMode;
use rustfs_config::ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS;
pub use rustfs_data_usage::{
AllTierStats, BucketTargetUsageInfo, BucketUsageInfo, DATA_USAGE_OBJECT_NAME, DATA_USAGE_OBSERVED_OBJECT_NAME,
DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageInfo, LEGACY_DATA_USAGE_OBJECT_NAME, PrefixUsageEntry,
PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeSummary, TierStats, hash_path, prefix_usage_in_cache,
DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageInfo, DataUsageSnapshotSetState, LEGACY_DATA_USAGE_OBJECT_NAME,
PrefixUsageEntry, PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeSummary, TierStats, hash_path,
prefix_usage_in_cache,
};
use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf};
use tokio::time::{Duration, Instant, sleep, timeout};
@@ -344,6 +345,18 @@ pub struct DataUsageCacheInfo {
pub scan_plan_digest: Option<DataUsageScanPlanDigest>,
#[serde(default)]
pub cache_key_format: u16,
/// Whether the entries retained while a set scan was incomplete come
/// from a prior complete set snapshot. This is observational input only.
#[serde(default)]
pub lkg_snapshot_complete: bool,
#[serde(default)]
pub lkg_next_cycle: Option<u64>,
#[serde(default)]
pub lkg_last_update: Option<SystemTime>,
#[serde(default)]
pub lkg_leader_epoch: Option<u64>,
#[serde(default)]
pub lkg_scan_plan_digest: Option<DataUsageScanPlanDigest>,
}
impl Serialize for DataUsageCacheInfo {
@@ -353,7 +366,7 @@ impl Serialize for DataUsageCacheInfo {
{
// Keep this metadata map-encoded so older readers can ignore fields
// appended by newer scanner versions during rolling upgrades.
let mut state = serializer.serialize_map(Some(16))?;
let mut state = serializer.serialize_map(Some(21))?;
state.serialize_entry("name", &self.name)?;
state.serialize_entry("next_cycle", &self.next_cycle)?;
state.serialize_entry("leader_epoch", &self.leader_epoch)?;
@@ -370,6 +383,11 @@ impl Serialize for DataUsageCacheInfo {
state.serialize_entry("snapshot_complete", &self.snapshot_complete)?;
state.serialize_entry("scan_plan_digest", &self.scan_plan_digest)?;
state.serialize_entry("cache_key_format", &self.cache_key_format)?;
state.serialize_entry("lkg_snapshot_complete", &self.lkg_snapshot_complete)?;
state.serialize_entry("lkg_next_cycle", &self.lkg_next_cycle)?;
state.serialize_entry("lkg_last_update", &self.lkg_last_update)?;
state.serialize_entry("lkg_leader_epoch", &self.lkg_leader_epoch)?;
state.serialize_entry("lkg_scan_plan_digest", &self.lkg_scan_plan_digest)?;
state.end()
}
}
+2 -3
View File
@@ -2274,9 +2274,8 @@ async fn final_data_usage_publication_defer_reason(
}
}
ScannerCycleStatus::Deferred(reason) => Some(reason),
// Incomplete cycles do not publish a usage snapshot. Keep the
// decision permissive so existing partial-cycle handling remains
// unchanged if a future scanner path emits a bookkeeping update.
// Incomplete cycles may publish a non-authoritative observational
// snapshot when at least one set has a usable current/LKG view.
ScannerCycleStatus::Incomplete => None,
}
}
+1 -1
View File
@@ -198,7 +198,7 @@ where
data_usage_info.usage_snapshot_authoritative_baseline = Some(authoritative.snapshot_identity());
}
if !data_usage_info.is_complete_bucket_usage_snapshot() {
if !data_usage_info.is_complete_bucket_usage_snapshot() && !data_usage_info.usage_snapshot_partial {
error!(
target: "rustfs::scanner",
event = EVENT_SCANNER_PERSIST_STATE,
+13 -2
View File
@@ -18,8 +18,8 @@ use crate::scanner_folder::{ScannerItem, scan_data_folder};
use crate::sleeper::SCANNER_SLEEPER;
use crate::{
DATA_USAGE_CACHE_NAME, DATA_USAGE_ROOT, DataUsageCache, DataUsageCacheInfo, DataUsageCachePrepareOutcome,
DataUsageCacheSource, DataUsageEntry, DataUsageEntryInfo, DataUsageInfo, DataUsageScanPlanDigest, ScannerError, SizeSummary,
TierStats,
DataUsageCacheSource, DataUsageEntry, DataUsageEntryInfo, DataUsageInfo, DataUsageScanPlanDigest, DataUsageSnapshotSetState,
ScannerError, SizeSummary, TierStats,
};
use futures::future::join_all;
use metrics::counter;
@@ -278,6 +278,17 @@ async fn publish_usage_snapshot(
Ok(true)
}
async fn publish_observational_snapshot(
updates: &mpsc::Sender<DataUsageInfo>,
mut data_usage_info: DataUsageInfo,
) -> Result<bool> {
data_usage_info.usage_snapshot_complete = false;
data_usage_info.usage_snapshot_partial = true;
data_usage_info.usage_snapshot_converged = Some(false);
send_data_usage_update(updates, data_usage_info).await?;
Ok(true)
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum ScannerCycleActivityStatus {
Unchanged,
+145 -2
View File
@@ -188,7 +188,7 @@ pub(super) fn completed_data_usage_info(
}
let mut total = DataUsageEntry::default();
let mut buckets_usage = HashMap::with_capacity(all_buckets.len());
let mut bucket_entries = HashMap::with_capacity(all_buckets.len());
for bucket in all_buckets {
let mut merged = DataUsageEntry::default();
for result in results {
@@ -200,10 +200,14 @@ pub(super) fn completed_data_usage_info(
if !total.checked_merge(&merged) {
return None;
}
buckets_usage.insert(bucket.clone(), checked_bucket_usage_info(&merged)?);
bucket_entries.insert(bucket.clone(), merged);
}
let merged_last_update = results.iter().filter_map(|result| result.info.last_update).max()?;
let buckets_usage = bucket_entries
.iter()
.map(|(bucket, entry)| Some((bucket.clone(), checked_bucket_usage_info(entry)?)))
.collect::<Option<HashMap<_, _>>>()?;
let bucket_sizes = buckets_usage
.iter()
.map(|(bucket, usage)| (bucket.clone(), usage.size))
@@ -225,6 +229,145 @@ pub(super) fn completed_data_usage_info(
Some((data_usage_info, merged_last_update))
}
/// Build a non-authoritative view from the set snapshots that completed this
/// cycle plus compatible per-set last-known-good caches. The caller must
/// persist this result only on the observational object; a missing set is
/// intentionally represented by an incomplete state and is never treated as
/// an empty set.
pub(super) fn observational_data_usage_info(
results: &[DataUsageCache],
expected_sources: &HashSet<DataUsageCacheSource>,
all_buckets: &[String],
expected_plan_digest: DataUsageScanPlanDigest,
scanner_cycle: u64,
leader_epoch: u64,
) -> Option<(DataUsageInfo, SystemTime)> {
let mut by_source = HashMap::with_capacity(results.len());
for result in results {
let source = result.info.source?;
if !expected_sources.contains(&source) || by_source.insert(source, result).is_some() {
return None;
}
}
let mut usable = Vec::new();
let mut set_states = Vec::with_capacity(expected_sources.len());
let mut sources = expected_sources.iter().copied().collect::<Vec<_>>();
sources.sort_by_key(|source| (source.pool_index, source.set_index));
for source in sources {
let result = by_source.get(&source).copied();
let current = result.filter(|result| {
result.info.snapshot_complete
&& result.info.next_cycle == scanner_cycle
&& result.info.leader_epoch == leader_epoch
&& result.info.scan_plan_digest == Some(expected_plan_digest)
});
let lkg = result.filter(|result| {
!result.info.snapshot_complete
&& result.info.lkg_snapshot_complete
&& result.info.lkg_scan_plan_digest == Some(expected_plan_digest)
&& result.info.lkg_leader_epoch.is_some_and(|epoch| {
epoch < leader_epoch
|| (epoch == leader_epoch && result.info.lkg_next_cycle.is_some_and(|cycle| cycle <= scanner_cycle))
})
});
let current_snapshot = current.is_some();
let selected = current.or(lkg);
if let Some(selected) = selected {
let (cycle, epoch, digest, last_update, complete) = if current_snapshot {
(
Some(selected.info.next_cycle),
Some(selected.info.leader_epoch),
selected.info.scan_plan_digest.map(|digest| digest.0),
selected.info.last_update,
true,
)
} else {
(
selected.info.lkg_next_cycle,
selected.info.lkg_leader_epoch,
selected.info.lkg_scan_plan_digest.map(|digest| digest.0),
selected.info.lkg_last_update,
false,
)
};
set_states.push(DataUsageSnapshotSetState {
pool_index: u64::try_from(source.pool_index).ok()?,
set_index: u64::try_from(source.set_index).ok()?,
scanner_cycle: cycle,
scanner_epoch: epoch,
scan_plan_digest: digest,
complete,
tombstone: false,
});
usable.push((selected, last_update));
} else {
set_states.push(DataUsageSnapshotSetState {
pool_index: u64::try_from(source.pool_index).ok()?,
set_index: u64::try_from(source.set_index).ok()?,
scanner_cycle: None,
scanner_epoch: None,
scan_plan_digest: Some(expected_plan_digest.0),
complete: false,
tombstone: false,
});
}
}
if usable.is_empty() {
return None;
}
let mut total = DataUsageEntry::default();
let mut bucket_entries = HashMap::with_capacity(all_buckets.len());
let mut merged_last_update = None;
for (result, last_update) in usable {
if let Some(update) = last_update {
merged_last_update = Some(merged_last_update.map_or(update, |current: SystemTime| current.max(update)));
}
for bucket in all_buckets {
let Some(entry) = result.checked_flatten(bucket) else {
continue;
};
let bucket_entry = bucket_entries.entry(bucket.clone()).or_insert_with(DataUsageEntry::default);
if !bucket_entry.checked_merge(&entry) {
return None;
}
if !total.checked_merge(&entry) {
return None;
}
}
}
let merged_last_update = merged_last_update?;
let buckets_usage = bucket_entries
.iter()
.map(|(bucket, entry)| Some((bucket.clone(), checked_bucket_usage_info(entry)?)))
.collect::<Option<HashMap<_, _>>>()?;
Some((
DataUsageInfo {
last_update: Some(merged_last_update),
scanner_cycle: Some(scanner_cycle),
scanner_epoch: Some(leader_epoch),
objects_total_count: u64::try_from(total.objects).ok()?,
versions_total_count: u64::try_from(total.versions).ok()?,
delete_markers_total_count: u64::try_from(total.delete_markers).ok()?,
objects_total_size: u64::try_from(total.size).ok()?,
tier_stats: total.all_tier_stats.filter(|tiers| !tiers.is_empty()),
buckets_count: u64::try_from(buckets_usage.len()).ok()?,
bucket_sizes: buckets_usage
.iter()
.map(|(bucket, usage)| (bucket.clone(), usage.size))
.collect(),
buckets_usage,
usage_snapshot_complete: false,
usage_snapshot_partial: true,
usage_snapshot_converged: Some(false),
usage_snapshot_set_states: set_states,
..Default::default()
},
merged_last_update,
))
}
pub(super) async fn send_cache_root_entry_info(
bucket_result_tx: &mpsc::Sender<DataUsageEntryInfo>,
cache: &DataUsageCache,
+89 -30
View File
@@ -40,6 +40,21 @@ impl ScannerIOCache for SetDisks {
let set_label = self.set_index.to_string();
let source = DataUsageCacheSource::new(self.pool_index, self.set_index);
let mut old_cache = DataUsageCache::default();
if let Err(e) = old_cache.load(self.clone(), DATA_USAGE_CACHE_NAME).await {
warn!(
target: "rustfs::scanner::io",
event = EVENT_SCANNER_CACHE_PERSIST_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_IO,
pool = self.pool_index,
set = self.set_index,
cache_name = DATA_USAGE_CACHE_NAME,
state = "old_cache_load_failed",
error = %e,
"Scanner old data usage cache load failed; rebuilding from bucket caches"
);
}
if buckets.is_empty() {
let now = SystemTime::now();
let mut cache = DataUsageCache {
@@ -80,6 +95,24 @@ impl ScannerIOCache for SetDisks {
"Scanner set state found no online disks"
);
reset_disk_bucket_scan_gauges(&pool_label, &set_label);
let lkg = old_cache.info.snapshot_complete.then(|| old_cache.clone());
let mut incomplete_scope = lkg.clone().unwrap_or_default();
incomplete_scope.info.name = DATA_USAGE_ROOT.to_string();
incomplete_scope.info.next_cycle = want_cycle;
incomplete_scope.info.last_update = None;
incomplete_scope.info.leader_epoch = leader_epoch;
incomplete_scope.info.source = Some(source);
incomplete_scope.info.snapshot_complete = false;
incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest);
incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT;
if let Some(lkg) = lkg {
incomplete_scope.info.lkg_snapshot_complete = true;
incomplete_scope.info.lkg_next_cycle = Some(lkg.info.next_cycle);
incomplete_scope.info.lkg_last_update = lkg.info.last_update;
incomplete_scope.info.lkg_leader_epoch = Some(lkg.info.leader_epoch);
incomplete_scope.info.lkg_scan_plan_digest = lkg.info.scan_plan_digest;
}
let _ = updates.send(incomplete_scope).await;
return Ok(());
}
// Preserve the original set topology across capability filtering. During
@@ -162,6 +195,24 @@ impl ScannerIOCache for SetDisks {
"Scanner set state found no usable namespace scanner disks"
);
reset_disk_bucket_scan_gauges(&pool_label, &set_label);
let lkg = old_cache.info.snapshot_complete.then(|| old_cache.clone());
let mut incomplete_scope = lkg.clone().unwrap_or_default();
incomplete_scope.info.name = DATA_USAGE_ROOT.to_string();
incomplete_scope.info.next_cycle = want_cycle;
incomplete_scope.info.last_update = None;
incomplete_scope.info.leader_epoch = leader_epoch;
incomplete_scope.info.source = Some(source);
incomplete_scope.info.snapshot_complete = false;
incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest);
incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT;
if let Some(lkg) = lkg {
incomplete_scope.info.lkg_snapshot_complete = true;
incomplete_scope.info.lkg_next_cycle = Some(lkg.info.next_cycle);
incomplete_scope.info.lkg_last_update = lkg.info.last_update;
incomplete_scope.info.lkg_leader_epoch = Some(lkg.info.leader_epoch);
incomplete_scope.info.lkg_scan_plan_digest = lkg.info.scan_plan_digest;
}
let _ = updates.send(incomplete_scope).await;
return Ok(());
}
let set_disk_inventory = Arc::new(scanner_set_disk_inventory(self.as_ref()).await);
@@ -203,22 +254,15 @@ impl ScannerIOCache for SetDisks {
record_disk_bucket_scans_active(0, &pool_label, &set_label);
let _reset_disk_bucket_scan_gauges = DiskBucketScanGaugeReset::new(pool_label.clone(), set_label.clone());
let mut old_cache = DataUsageCache::default();
if let Err(e) = old_cache.load(self.clone(), DATA_USAGE_CACHE_NAME).await {
warn!(
target: "rustfs::scanner::io",
event = EVENT_SCANNER_CACHE_PERSIST_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_IO,
pool = self.pool_index,
set = self.set_index,
cache_name = DATA_USAGE_CACHE_NAME,
state = "old_cache_load_failed",
error = %e,
"Scanner old data usage cache load failed; rebuilding from bucket caches"
);
}
match old_cache.prepare_for_scan(
let old_lkg = old_cache.info.snapshot_complete.then(|| {
(
old_cache.info.next_cycle,
old_cache.info.last_update,
old_cache.info.leader_epoch,
old_cache.info.scan_plan_digest,
)
});
let prepare_outcome = match old_cache.prepare_for_scan(
DATA_USAGE_ROOT,
want_cycle,
leader_epoch,
@@ -259,7 +303,16 @@ impl ScannerIOCache for SetDisks {
);
return Ok(());
}
DataUsageCachePrepareOutcome::Reused | DataUsageCachePrepareOutcome::Reset => {}
outcome => outcome,
};
if matches!(prepare_outcome, DataUsageCachePrepareOutcome::Reused)
&& let Some((cycle, last_update, epoch, digest)) = old_lkg
{
old_cache.info.lkg_snapshot_complete = true;
old_cache.info.lkg_next_cycle = Some(cycle);
old_cache.info.lkg_last_update = last_update;
old_cache.info.lkg_leader_epoch = Some(epoch);
old_cache.info.lkg_scan_plan_digest = digest;
}
let mut cache = DataUsageCache {
@@ -1099,23 +1152,29 @@ impl ScannerIOCache for SetDisks {
cache.info.next_cycle = want_cycle;
cache.info.last_update.get_or_insert_with(SystemTime::now);
cache.info.snapshot_complete = true;
cache.info.lkg_snapshot_complete = false;
cache.info.lkg_next_cycle = None;
cache.info.lkg_last_update = None;
cache.info.lkg_leader_epoch = None;
cache.info.lkg_scan_plan_digest = None;
cache.clone()
};
let _ = persist_and_publish_cache_snapshot(self.clone(), &updates, cache_snapshot, cache_cycle_floor.as_ref()).await;
} else {
let incomplete_scope = DataUsageCache {
info: DataUsageCacheInfo {
name: DATA_USAGE_ROOT.to_string(),
next_cycle: want_cycle,
leader_epoch,
source: Some(source),
snapshot_complete: false,
scan_plan_digest: Some(scan_plan_digest),
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
..Default::default()
},
cache: HashMap::new(),
};
let mut incomplete_scope = cache_mutex.lock().await.clone();
incomplete_scope.info.name = DATA_USAGE_ROOT.to_string();
incomplete_scope.info.next_cycle = want_cycle;
incomplete_scope.info.last_update = None;
incomplete_scope.info.leader_epoch = leader_epoch;
incomplete_scope.info.source = Some(source);
incomplete_scope.info.snapshot_complete = false;
incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest);
incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT;
incomplete_scope.info.lkg_snapshot_complete = old_cache.info.lkg_snapshot_complete;
incomplete_scope.info.lkg_next_cycle = old_cache.info.lkg_next_cycle;
incomplete_scope.info.lkg_last_update = old_cache.info.lkg_last_update;
incomplete_scope.info.lkg_leader_epoch = old_cache.info.lkg_leader_epoch;
incomplete_scope.info.lkg_scan_plan_digest = old_cache.info.lkg_scan_plan_digest;
if let Err(e) = updates.send(incomplete_scope).await {
error!(
target: "rustfs::scanner::io",
+33
View File
@@ -234,6 +234,7 @@ impl ScannerIOCycle for ECStore {
let active_set_scans_clone = active_set_scans.clone();
let (tx, mut rx) = mpsc::channel::<DataUsageCache>(1);
let failed_scope_tx = tx.clone();
// Spawn task to receive and store results
let receiver_fut = tokio::spawn(async move {
@@ -314,6 +315,21 @@ impl ScannerIOCycle for ECStore {
state = "set_scan_failed",
"Scanner set scan failed; continuing cycle"
);
let _ = failed_scope_tx
.send(DataUsageCache {
info: DataUsageCacheInfo {
name: DATA_USAGE_ROOT.to_string(),
next_cycle: want_cycle_clone,
leader_epoch,
source: Some(source),
snapshot_complete: false,
scan_plan_digest: Some(scan_plan_digest),
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
..Default::default()
},
cache: HashMap::new(),
})
.await;
let mut first_err = first_err_mutex_clone.lock().await;
record_set_scan_failure(&mut first_err, e);
}
@@ -370,6 +386,19 @@ impl ScannerIOCycle for ECStore {
budget_elapsed,
ctx.is_cancelled(),
);
let observational_usage = completed_usage
.is_none()
.then(|| {
observational_data_usage_info(
&results,
&expected_sources,
&all_bucket_names,
scan_plan_digest,
want_cycle,
leader_epoch,
)
})
.flatten();
let structurally_complete_snapshot = result.is_ok() && completed_all_sets && completed_usage.is_some();
let cycle_status = classify_nsscanner_cycle(
structurally_complete_snapshot,
@@ -381,6 +410,10 @@ impl ScannerIOCycle for ECStore {
);
if let Some((data_usage_info, _)) = completed_usage {
publish_usage_snapshot(&updates, cycle_status, data_usage_info).await?;
} else if !ctx.is_cancelled()
&& let Some((data_usage_info, _)) = observational_usage
{
publish_observational_snapshot(&updates, data_usage_info).await?;
}
let dirty_usage_clear = should_clear_dirty_usage_snapshot(
result.is_ok(),
@@ -105,6 +105,160 @@ fn completed_data_usage_info_for_test(
completed_data_usage_info(results, &expected_sources, all_buckets, true, budget_elapsed, cancelled)
}
fn lkg_root_cache(bucket: &str, objects: usize, source: DataUsageCacheSource) -> DataUsageCache {
let mut cache = completed_root_cache(bucket, objects, 10, source);
cache.info.snapshot_complete = false;
cache.info.next_cycle = 8;
cache.info.leader_epoch = 3;
cache.info.lkg_snapshot_complete = true;
cache.info.lkg_next_cycle = Some(7);
cache.info.lkg_last_update = cache.info.last_update;
cache.info.lkg_leader_epoch = Some(3);
cache.info.lkg_scan_plan_digest = Some(TEST_PLAN_DIGEST);
cache
}
#[test]
fn partial_usage_is_observational_not_authoritative_for_quota() {
let all_buckets = vec!["bucket".to_string()];
let current_source = DataUsageCacheSource::new(0, 0);
let stalled_source = DataUsageCacheSource::new(1, 0);
let mut current = completed_root_cache("bucket", 2, 20, current_source);
current.info.next_cycle = 8;
current.info.leader_epoch = 3;
let stalled = lkg_root_cache("bucket", 1, stalled_source);
let expected = HashSet::from([current_source, stalled_source]);
assert!(
completed_data_usage_info(&[current.clone(), stalled.clone()], &expected, &all_buckets, true, false, false).is_none()
);
let (observed, _) = observational_data_usage_info(&[current, stalled], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
.expect("a completed set should produce an observational view");
assert!(observed.usage_snapshot_partial);
assert!(!observed.usage_snapshot_complete);
assert_eq!(observed.usage_snapshot_converged, Some(false));
assert_eq!(observed.usage_snapshot_set_states.len(), 2);
}
#[test]
fn lkg_scope_does_not_count_as_current_cycle_completion() {
let source = DataUsageCacheSource::new(0, 0);
let mut lkg = lkg_root_cache("bucket", 1, source);
lkg.info.last_update = None;
let expected = HashSet::from([source]);
assert!(!scanner_results_form_complete_snapshot(&[lkg], &expected));
}
#[test]
fn stale_quota_uses_complete_baseline_plus_positive_deltas() {
let all_buckets = vec!["bucket".to_string()];
let source = DataUsageCacheSource::new(0, 0);
let mut current = completed_root_cache("bucket", 3, 20, source);
current.info.next_cycle = 8;
current.info.leader_epoch = 3;
let expected = HashSet::from([source]);
let (observed, _) = observational_data_usage_info(&[current], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
.expect("complete set data is a valid observational baseline");
assert_eq!(observed.objects_total_size, 30);
assert_eq!(observed.usage_snapshot_set_states[0].complete, true);
}
#[test]
fn negative_delta_waits_for_set_reconciliation() {
let all_buckets = vec!["bucket".to_string()];
let source = DataUsageCacheSource::new(0, 0);
let mut stalled = lkg_root_cache("bucket", 4, source);
stalled.info.lkg_scan_plan_digest = Some(DataUsageScanPlanDigest([9; 32]));
let expected = HashSet::from([source]);
assert!(observational_data_usage_info(&[stalled], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3).is_none());
}
#[test]
fn set_membership_add_remove_uses_generation_and_tombstone() {
let state = DataUsageSnapshotSetState {
pool_index: 1,
set_index: 2,
scanner_cycle: Some(9),
scanner_epoch: Some(4),
scan_plan_digest: Some(TEST_PLAN_DIGEST.0),
complete: false,
tombstone: true,
};
let encoded = serde_json::to_vec(&state).expect("set state should serialize");
let decoded: DataUsageSnapshotSetState = serde_json::from_slice(&encoded).expect("set state should deserialize");
assert_eq!(decoded, state);
let snapshot = DataUsageInfo {
last_update: Some(SystemTime::UNIX_EPOCH + Duration::from_secs(10)),
scanner_cycle: Some(9),
scanner_epoch: Some(4),
buckets_count: 0,
usage_snapshot_converged: Some(false),
usage_snapshot_partial: true,
usage_snapshot_set_states: vec![
DataUsageSnapshotSetState {
pool_index: 0,
set_index: 0,
scanner_cycle: Some(9),
scanner_epoch: Some(4),
scan_plan_digest: Some(TEST_PLAN_DIGEST.0),
complete: true,
tombstone: false,
},
state,
],
..Default::default()
};
assert!(snapshot.is_valid_partial_snapshot());
}
#[test]
fn old_set_completion_cannot_overwrite_new_aggregate() {
let all_buckets = vec!["bucket".to_string()];
let source = DataUsageCacheSource::new(0, 0);
let mut old = completed_root_cache("bucket", 1, 20, source);
old.info.next_cycle = 7;
old.info.leader_epoch = 2;
let expected = HashSet::from([source]);
assert!(observational_data_usage_info(&[old], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3).is_none());
}
#[test]
fn usage_aggregate_survives_restart_and_leader_failover() {
let all_buckets = vec!["bucket".to_string()];
let source = DataUsageCacheSource::new(0, 0);
let mut lkg = lkg_root_cache("bucket", 5, source);
lkg.info.lkg_leader_epoch = Some(4);
lkg.info.lkg_next_cycle = Some(9);
let expected = HashSet::from([source]);
let (observed, _) = observational_data_usage_info(&[lkg], &expected, &all_buckets, TEST_PLAN_DIGEST, 10, 5)
.expect("compatible LKG should survive a leader change");
assert_eq!(observed.usage_snapshot_set_states[0].scanner_epoch, Some(4));
assert_eq!(observed.objects_total_size, 50);
}
#[test]
fn usage_aggregate_cost_is_linear_in_set_count() {
let all_buckets = vec!["bucket".to_string()];
let mut results = Vec::new();
let mut expected = HashSet::new();
for index in 0..32 {
let source = DataUsageCacheSource::new(index, 0);
expected.insert(source);
let mut cache = completed_root_cache("bucket", 1, 20, source);
cache.info.next_cycle = 8;
cache.info.leader_epoch = 3;
results.push(cache);
}
let (observed, _) = observational_data_usage_info(&results, &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
.expect("all set snapshots should aggregate");
assert_eq!(observed.objects_total_count, 32);
let reversed = results.iter().rev().cloned().collect::<Vec<_>>();
let (reversed_observed, _) = observational_data_usage_info(&reversed, &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
.expect("reordered set snapshots should aggregate");
assert_eq!(observed.usage_snapshot_set_states, reversed_observed.usage_snapshot_set_states);
}
#[test]
fn completed_data_usage_info_publishes_tier_stats_across_sets() {
let all_buckets = vec!["bucket-a".to_string(), "bucket-b".to_string()];
+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]
-149
View File
@@ -1,149 +0,0 @@
# CI gate matrix
This file is the source of truth for which validation runs on each event, its
configured wall-clock budget, and whether it can block a merge. Test taxonomy,
naming, and nextest serialization rules remain in [README.md](README.md); e2e
membership and counts remain in
[e2e-suite-inventory.md](e2e-suite-inventory.md).
The distinction between **required** and **report-only** is load-bearing:
a failing job blocks a merge only when its exact check name is present in the
live `main` ruleset. A workflow name, a `merge_group` trigger, or a red PR check
does not make a job required by itself.
## Required merge checks
The live `main` ruleset (`6436880`) currently requires exactly these contexts:
| Required context | Producer | Validation |
|---|---|---|
| `CLA Check` | `.github/workflows/cla.yml` | Contributor agreement |
| `Quick Checks` | `.github/workflows/ci.yml` | Formatting and repository guard scripts |
| `Test and Lint` | `.github/workflows/ci.yml` | Clippy, workspace nextest excluding `e2e_test`, doctests, and migration proofs |
For pull requests limited to the paths excluded by the main CI workflow,
`.github/workflows/ci-docs-only.yml` reports `Quick Checks` and
`Test and Lint` under the same names. It runs the real quick checks and the
planning-document guard; it does not claim that Rust compilation or runtime
tests ran. Despite the workflow name, these paths also include selected deploy,
workflow, and lock files.
Verify the live rule rather than trusting this snapshot before changing merge
policy:
```bash
gh api repos/rustfs/rustfs/rulesets/6436880 \
--jq '.rules[] | select(.type == "required_status_checks") | .parameters'
```
The ruleset currently has `strict_required_status_checks_policy=false`.
`Continuous Integration` accepts `merge_group` events and runs `e2e-full` for
them, but `End-to-End Tests (full merge gate)` is not currently a required
context. Therefore the repository is prepared to test a merge-queue SHA, but
the workflow alone does not prove that every merge passed that lane.
## Pull request and merge matrix
Budgets below are job `timeout-minutes`, not typical runtimes. “Report-only”
means the result is visible and actionable but is not in the live required
context list.
| Event | Validation | Budget | Merge status | Reproduction |
|---|---|---:|---|---|
| PR, non-doc change | `Quick Checks` | 10 min | Required | `make pre-commit` (broader local umbrella) |
| PR, non-doc change | `Test and Lint` | 90 min | Required | `cargo nextest run --profile ci --all --exclude e2e_test` |
| PR, non-doc change | `Typos` | 10 min | Report-only | `typos` |
| PR, non-doc change | `ILM Integration (serial)` | 90 min | Report-only | Use the exact command in `.github/workflows/ci.yml` |
| PR, non-doc change | rio-v2 / swift / sftp test-and-lint variants | 90 min each | Report-only | `cargo nextest run` with the workflow's feature set |
| PR, non-doc change | `Build RustFS Debug Binary` | 30 min | Report-only; prerequisite for black-box lanes | `cargo build -p rustfs --bins` |
| PR, non-doc change | `io_uring Integration (real)` | 30 min | Report-only | `cargo test -p rustfs-ecstore --lib uring_ -- --test-threads=1 --nocapture` |
| PR, non-doc change | `End-to-End Tests` (`e2e-smoke` plus `s3s-e2e`) | 30 min | Report-only | `cargo nextest run --profile e2e-smoke -p e2e_test`; then `./scripts/e2e-run.sh ./target/debug/rustfs <data-dir>` |
| PR, non-doc change | `S3 Implemented Tests` | 60 min | Report-only | Build `rustfs`, then run `scripts/s3-tests/run.sh` with `DEPLOY_MODE=binary`, `TEST_MODE=single`, and `MAXFAIL=0` |
| PR, non-doc change | `S3 Lifecycle Behavior Tests` | 30 min | Report-only | Use the accelerated scanner environment in `.github/workflows/ci.yml` with `scripts/s3-tests/run.sh` |
| PR touching dependency or workflow inputs | Cargo Deny / Workflow Pin Report / Dependency Review | 20 / 5 / 30 min | Report-only | `cargo deny check`; `scripts/security/check_workflow_pins.sh` |
| PR touching architecture rules or architecture docs | `Architecture Migration Rules` | 10 min | Report-only | `scripts/check_architecture_migration_rules.sh` |
| PR touching Nix or workspace manifests | `Nix Build & Check` | 60 min | Report-only | `nix flake check` |
| PR limited to main-CI-excluded paths | companion `Quick Checks` and `Test and Lint` | 10 min each | Required | `git diff --check`; `make doc-paths-check` when documentation paths changed |
| `merge_group` | Standard CI plus `e2e-full` | 55 min for `e2e-full` | Standard required contexts only; `e2e-full` report-only | `cargo nextest run --profile e2e-full -p e2e_test` |
| Push to `main` | Standard CI plus `e2e-full` | 55 min for `e2e-full` | Post-merge detection | Same as `merge_group` |
| PR touching fuzz inputs or harness paths | Build plus five 60-second fuzz smoke targets | 60 min build; 30 min per target | Report-only | `MAX_TOTAL_TIME=60 ./scripts/fuzz/run.sh` |
| PR touching selected ecstore disk/format paths | `Rename Safety` on Windows | 60 min | Report-only | Run the four `cargo test -p rustfs-ecstore --lib <filter>` commands in `windows-filesystem.yml` on Windows |
The authoritative e2e filters live in `.config/nextest.toml`; extend a profile
instead of adding a second ad-hoc selector. Before a profile runs,
`scripts/check_test_wiring.py` compares its exact membership to the committed
digest so a silent test drop fails closed.
## Scheduled and manual validation
Scheduled lanes are independent fault domains. They do not block a pull
request, but their workflow-local gate can fail the run and scheduled failures
are routed to the shared failure-issue action. The scheduled-validation
watchdog and freshness workflow separately detect incomplete runs and missing
schedules.
| Cadence (UTC unless noted) | Workflow / validation | Budget | Verdict and artifacts | Reproduction |
|---|---|---:|---|---|
| Daily 02:17 | Fuzz: five nightly corpus targets | 60 min build; 60 min per target | Gate; corpus/crash artifacts, scheduled failure alert | `MAX_TOTAL_TIME=<seconds> ./scripts/fuzz/run.sh` |
| Daily 03:17 | MinIO interop (EC + SSE read parity) | 40 min | Gate; scheduled failure alert | Dispatch `minio-interop.yml` or follow its pinned Docker fixture steps |
| Daily 04:29 | Replication / cluster-fault / protocol e2e | 45 / 90 / 90 min | Three independent gates; JUnit, membership, and server logs | `cargo nextest run --profile e2e-repl-nightly -p e2e_test`; `--profile e2e-nightly`; `-j 1 --profile e2e-protocols` |
| Daily 06:31 | Warp performance A/B | 180 min | Regression budget gate; A/B summaries and server logs | `bash scripts/run_hotpath_warp_abba.sh --help` |
| Daily 00:07 Asia/Shanghai (16:07 UTC previous day) | Nightly GNU build and Vault lanes | 150 / 90 / 60 min | Build, live Vault, and HA failover gates | Use the commands and pinned Vault images in `nightly-gnu.yml` |
| Daily 03:23 | Security Audit | 20 / 5 min, plus 30 min on PR dependency review | Cargo Deny and workflow-pin gates; scheduled failure alert | `cargo deny check`; `scripts/security/check_workflow_pins.sh` |
| Daily 23:47 | Scheduled Validation Freshness | 10 min | Fails when a critical schedule was never created or is stale | Dispatch `scheduled-validation-freshness.yml` |
| Sunday 00:11 | Full `Continuous Integration` matrix | Per-job budgets above | Weekly variant coverage, including dormant rio-v2 binary/e2e lanes | Dispatch `ci.yml` |
| Sunday 01:13 | Seven-platform build matrix | 150 min per platform | Build/package integrity; scheduled failure alert | Dispatch `build.yml` with an exact platform set |
| Sunday 02:19 | Ceph s3-tests full sweep: single and real four-node, four shards each | 180 min per shard | Compatibility gate; report, JUnit, exact node IDs, and server logs | `scripts/s3-tests/run.sh` against an existing single or distributed target |
| Sunday 06:41 | Mint | 120 min | **Report-only by design**; per-suite PASS/FAIL/NA and raw `log.json` | Reproduce the pinned Docker sequence in `mint.yml` or dispatch it |
| Sunday 07:43 | Workspace line coverage | 120 min | Report-only trend; lcov and JSON retained 90 days | `make coverage` |
| Monthly, day 1 06:37 | Runner Hygiene | 15 min | Validates runner ephemerality; scheduled failure alert | Dispatch `runner-hygiene.yml` |
Manual `workflow_dispatch` exists for the scheduled workflows above. Manual
runs are debugging evidence and intentionally do not open scheduled-failure
issues. A manual performance run may explicitly allow a known regression; that
override must not be treated as an ordinary passing baseline.
## Release validation
Release validation is post-merge and tag-driven; it does not substitute for a
pull-request gate.
| Event | Validation | Budget | Result |
|---|---|---:|---|
| Push to `main` or weekly schedule | `Build and Release` platform matrix | 150 min per platform | Build artifacts for all selected targets; no release publication on a main push |
| Valid release or preview tag | `Build and Release` plus asset checks | 150 min per platform | Draft release, checksummed assets, and publish step |
| Successful non-preview release-tag build | Docker image build and image scan | 60 min build; 30 min scan | Multi-architecture images plus vulnerability report |
| Successful release-tag build | DEB/RPM packaging | 30 min per architecture | Packages and checksum files uploaded to the release |
| Successful non-preview release-tag build | Helm template test and package | 30 min build; 30 min publish | Versioned chart and repository index |
Use an exact preview tag for end-to-end release rehearsal. Manual dispatches
are backfill/debug paths and do not prove the automatic `workflow_run` chain.
## Evidence requirements
A green check is useful only when it proves the intended behavior ran:
- Record the exact commit SHA and run URL.
- Separate product failure from runner prerequisites, service readiness, and
cancellation. Repair the precondition, then rerun the exact workload.
- Preserve membership manifests, JUnit, raw compatibility logs, seeds, and
server logs where the workflow provides them.
- For a bug fix or a new fault checker, provide sensitivity evidence: the old
behavior or an intentional mutation must fail the new oracle, and the fixed
behavior must pass it.
- Never promote a report-only lane to required from one green run. Require at
least 14 days and 30 representative pull requests with at least 99% complete
execution, then update the ruleset and this table together.
## Change checklist
Update this file in the same pull request when any of these change:
- workflow triggers, job names, timeouts, or nextest profile ownership;
- required status contexts or strict/merge-queue policy;
- scheduled cadence, alert routing, artifact contract, or local reproduction;
- report-only versus gating semantics.
Do not copy per-module test counts here. Update
[e2e-suite-inventory.md](e2e-suite-inventory.md) and its enforced membership
digest instead.
+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(())
}
+18 -22
View File
@@ -513,15 +513,13 @@ fn sr_bucket_meta_item(bucket: String, item_type: &str) -> SRBucketMeta {
}
}
async fn notify_bucket_metadata_reload(
fn notify_bucket_metadata_reload(
bucket: String,
operation: &'static str,
request_context: Option<request_context::RequestContext>,
scanner_maintenance_change: bool,
) {
record_local_scanner_maintenance_reload(&bucket, scanner_maintenance_change);
// Keep reload detached across request cancellation, but wait before a healthy peer can serve the previous config.
let (completed_tx, completed_rx) = tokio::sync::oneshot::channel();
spawn_background_with_context(request_context, async move {
if let Some(notification_sys) = current_notification_system() {
let result = if scanner_maintenance_change {
@@ -533,9 +531,7 @@ async fn notify_bucket_metadata_reload(
warn!(bucket = %bucket, error = %err, "failed to notify peers after {operation}");
}
}
let _ = completed_tx.send(());
});
let _ = completed_rx.await;
}
fn record_local_scanner_maintenance_reload(bucket: &str, scanner_maintenance_change: bool) {
@@ -1480,7 +1476,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "delete bucket encryption", request_context, false).await;
notify_bucket_metadata_reload(bucket.clone(), "delete bucket encryption", request_context, false);
let item = sr_bucket_meta_item(bucket.clone(), "sse-config");
if let Err(err) = site_replication_bucket_meta_hook(item).await {
@@ -1512,7 +1508,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "delete bucket cors", request_context, false).await;
notify_bucket_metadata_reload(bucket.clone(), "delete bucket cors", request_context, false);
let item = sr_bucket_meta_item(bucket.clone(), "cors-config");
if let Err(err) = site_replication_bucket_meta_hook(item).await {
@@ -1544,7 +1540,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "delete bucket lifecycle", request_context, true).await;
notify_bucket_metadata_reload(bucket.clone(), "delete bucket lifecycle", request_context, true);
let item = sr_bucket_meta_item(bucket.clone(), "lc-config");
if let Err(err) = site_replication_bucket_meta_hook(item).await {
@@ -1576,7 +1572,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "delete bucket policy", request_context, false).await;
notify_bucket_metadata_reload(bucket.clone(), "delete bucket policy", request_context, false);
let item = sr_bucket_meta_item(bucket.clone(), "policy");
if let Err(err) = site_replication_bucket_meta_hook(item).await {
@@ -1634,7 +1630,7 @@ impl DefaultBucketUsecase {
}
drop(targets_guard);
notify_bucket_metadata_reload(bucket.clone(), "delete bucket replication", request_context, true).await;
notify_bucket_metadata_reload(bucket.clone(), "delete bucket replication", request_context, true);
let item = sr_bucket_meta_item(bucket.clone(), "replication-config");
if let Err(err) = site_replication_bucket_meta_hook(item).await {
@@ -1659,7 +1655,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "delete bucket tagging", request_context, false).await;
notify_bucket_metadata_reload(bucket.clone(), "delete bucket tagging", request_context, false);
let item = sr_bucket_meta_item(bucket.clone(), "tags");
if let Err(err) = site_replication_bucket_meta_hook(item).await {
@@ -1692,7 +1688,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "delete public access block", request_context, false).await;
notify_bucket_metadata_reload(bucket.clone(), "delete public access block", request_context, false);
Ok(S3Response::with_status(DeletePublicAccessBlockOutput::default(), StatusCode::NO_CONTENT))
}
@@ -2147,7 +2143,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "put bucket encryption", request_context, false).await;
notify_bucket_metadata_reload(bucket.clone(), "put bucket encryption", request_context, false);
let mut item = sr_bucket_meta_item(bucket.clone(), "sse-config");
item.sse_config = Some(
@@ -2226,7 +2222,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "put bucket lifecycle", request_context, true).await;
notify_bucket_metadata_reload(bucket.clone(), "put bucket lifecycle", request_context, true);
let mut item = sr_bucket_meta_item(bucket.clone(), "lc-config");
item.expiry_lc_config =
@@ -2311,7 +2307,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "put bucket notification", request_context, false).await;
notify_bucket_metadata_reload(bucket.clone(), "put bucket notification", request_context, false);
let region = resolve_notification_region(self.global_region(), request_region);
let notify = current_notify_interface_for_context(self.context.as_deref());
@@ -2416,7 +2412,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "put bucket policy", request_context, false).await;
notify_bucket_metadata_reload(bucket.clone(), "put bucket policy", request_context, false);
let mut item = sr_bucket_meta_item(bucket.clone(), "policy");
item.policy = Some(serde_json::from_str(&policy).map_err(|e| s3_error!(InvalidArgument, "parse policy failed {:?}", e))?);
@@ -2451,7 +2447,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "put bucket cors", request_context, false).await;
notify_bucket_metadata_reload(bucket.clone(), "put bucket cors", request_context, false);
let mut item = sr_bucket_meta_item(bucket.clone(), "cors-config");
item.cors =
@@ -2495,7 +2491,7 @@ impl DefaultBucketUsecase {
.map_err(ApiError::from)?;
drop(targets_guard);
notify_bucket_metadata_reload(bucket.clone(), "put bucket replication", request_context, true).await;
notify_bucket_metadata_reload(bucket.clone(), "put bucket replication", request_context, true);
let mut item = sr_bucket_meta_item(bucket.clone(), "replication-config");
item.replication_config = Some(
@@ -2535,7 +2531,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "put public access block", request_context, false).await;
notify_bucket_metadata_reload(bucket.clone(), "put public access block", request_context, false);
Ok(S3Response::new(PutPublicAccessBlockOutput::default()))
}
@@ -2564,7 +2560,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "put bucket tagging", request_context, false).await;
notify_bucket_metadata_reload(bucket.clone(), "put bucket tagging", request_context, false);
let mut item = sr_bucket_meta_item(bucket.clone(), "tags");
item.tags = Some(serialize_config(&tagging).and_then(|bytes| String::from_utf8(bytes).map_err(to_internal_error))?);
@@ -2597,7 +2593,7 @@ impl DefaultBucketUsecase {
.await
.map_err(ApiError::from)?;
notify_bucket_metadata_reload(bucket.clone(), "put bucket versioning", request_context, false).await;
notify_bucket_metadata_reload(bucket.clone(), "put bucket versioning", request_context, false);
let mut item = sr_bucket_meta_item(bucket.clone(), "version-config");
item.versioning = Some(
@@ -3048,7 +3044,7 @@ mod tests {
"{method} should identify the bucket metadata operation in reload logs"
);
let expected_reload = format!(
"notify_bucket_metadata_reload(bucket.clone(), \"{operation}\", request_context, {scanner_maintenance_change}).await;"
"notify_bucket_metadata_reload(bucket.clone(), \"{operation}\", request_context, {scanner_maintenance_change});"
);
assert!(
body.contains(&expected_reload),
+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);