Compare commits

..

1 Commits

Author SHA1 Message Date
overtrue 950cce6f57 fix(admin): preserve decommission readiness error 2026-09-06 17:03:24 +08:00
23 changed files with 187 additions and 1214 deletions
Generated
+34 -34
View File
@@ -2527,18 +2527,18 @@ checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b"
[[package]]
name = "crossbeam-channel"
version = "0.5.17"
version = "0.5.16"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "98b0cc327b5bc766e7fda9c9260cc0fa81b43a8e240440422dff70788e3f9ef1"
checksum = "d85363c37faeca707aef026efa9f3b34d077bce547e48f770770625c6013679e"
dependencies = [
"crossbeam-utils",
]
[[package]]
name = "crossbeam-deque"
version = "0.8.8"
version = "0.8.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "622f3fc73690be383c7214310406f28a90e6edeadc3cea882f9d71e495b9711a"
checksum = "5181e0de7b61eb03a81e347d6dd8797bae9da5146707b51077e2d71a54ec0ceb"
dependencies = [
"crossbeam-epoch",
"crossbeam-utils",
@@ -2546,27 +2546,27 @@ dependencies = [
[[package]]
name = "crossbeam-epoch"
version = "0.9.21"
version = "0.9.20"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "dc74980687109a3b14c72fd458107bf0baa1da1a1a805e178d15501ba9b86d9d"
checksum = "2d6914041f254d6e9176c01941b21115dcfb7089e55135a35411081bd106ef3f"
dependencies = [
"crossbeam-utils",
]
[[package]]
name = "crossbeam-queue"
version = "0.3.14"
version = "0.3.13"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "03e8bd762f7479489c70ed6c768ddca99d7296857de437a68dcb2a94365b3fae"
checksum = "803d13fb3b09d88be9f4dbc29062c66b19bf7170867ceb746d2a8689bf6c7a26"
dependencies = [
"crossbeam-utils",
]
[[package]]
name = "crossbeam-utils"
version = "0.8.23"
version = "0.8.22"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a31eee39dddec8330830986fcd7625edb5a24ec90ea038215273bbc3adb08ac6"
checksum = "61803da095bee82a81bb1a452ecc25d3b2f1416d1897eb86430c6159ef717c17"
[[package]]
name = "crunchy"
@@ -3673,9 +3673,9 @@ dependencies = [
[[package]]
name = "der"
version = "0.8.2"
version = "0.8.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a878c850e9e421b20262e9b41f9c860e4785fa07541c266b62ff9d1ef998a80a"
checksum = "a69dedd701da44b0536442edf09c81a64b0ab97a7a4a5e3d1971f00027cbc63d"
dependencies = [
"const-oid 0.10.2",
"pem-rfc7468 1.0.0",
@@ -3946,7 +3946,7 @@ dependencies = [
"libc",
"option-ext",
"redox_users 0.5.2",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -4090,7 +4090,7 @@ version = "0.17.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c0681a4fc24c767085329728d8dfba959af91228aa4610cca4f8ce317ba46ae0"
dependencies = [
"der 0.8.2",
"der 0.8.1",
"digest 0.11.3",
"elliptic-curve 0.14.1",
"rfc6979 0.6.0",
@@ -4295,7 +4295,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb"
dependencies = [
"libc",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -5693,7 +5693,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "20fd6de4ccfcc187e38bc21cfa543cb5a302cb86a8b114eb7f0bf0dc9f8ac00f"
dependencies = [
"io-lifetimes 3.0.1",
"windows-sys 0.59.0",
"windows-sys 0.60.2",
]
[[package]]
@@ -5734,9 +5734,9 @@ dependencies = [
[[package]]
name = "ipnet"
version = "2.12.2"
version = "2.12.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "791930b43c0d5973160d90a8f3894509f2b273430f5c5c73b668636d0287c5c0"
checksum = "6a756c3fac73139e83f14c2d742155dd2b78d3ee56597b419a0579b7bdd6dd78"
dependencies = [
"serde",
]
@@ -5758,7 +5758,7 @@ checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46"
dependencies = [
"hermit-abi",
"libc",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -6955,7 +6955,7 @@ version = "0.50.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5"
dependencies = [
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -7933,7 +7933,7 @@ version = "0.8.0-rc.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "986d2e952779af96ea048f160fd9194e1751b4faea78bcf3ceb456efe008088e"
dependencies = [
"der 0.8.2",
"der 0.8.1",
"spki 0.8.0",
]
@@ -7976,7 +7976,7 @@ dependencies = [
"aes 0.9.3",
"aes-gcm",
"cbc 0.2.1",
"der 0.8.2",
"der 0.8.1",
"pbkdf2 0.13.0",
"rand_core 0.10.1",
"scrypt 0.12.0",
@@ -8000,7 +8000,7 @@ version = "0.11.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "451913da69c775a56034ea8d9003d27ee8948e12443eae7c038ba100a4f21cb7"
dependencies = [
"der 0.8.2",
"der 0.8.1",
"pkcs5 0.8.1",
"rand_core 0.10.1",
"spki 0.8.0",
@@ -8706,7 +8706,7 @@ dependencies = [
"once_cell",
"socket2",
"tracing",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -8942,9 +8942,9 @@ dependencies = [
[[package]]
name = "redis"
version = "1.7.0"
version = "1.6.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2acbc41a996f7652b2ddd9dfd98cc4ff602cfd742ae35382f07f608405ab50ed"
checksum = "e37a4ca5c6ca42aa3e6df2fd32b987a65d32a4c2159a6f3fe0fd1df306a2658f"
dependencies = [
"arc-swap",
"arcstr",
@@ -9326,7 +9326,7 @@ dependencies = [
"curve25519-dalek 5.0.0",
"data-encoding",
"delegate",
"der 0.8.2",
"der 0.8.1",
"digest 0.11.3",
"ecdsa 0.17.0",
"ed25519-dalek 3.0.0",
@@ -11066,7 +11066,7 @@ dependencies = [
"errno",
"libc",
"linux-raw-sys",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -11149,7 +11149,7 @@ dependencies = [
"security-framework",
"security-framework-sys",
"webpki-root-certs",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -11420,7 +11420,7 @@ checksum = "d56d437c2f19203ce5f7122e507831de96f3d2d4d3be5af44a0b0a09d8a80e4d"
dependencies = [
"base16ct 1.0.0",
"ctutils",
"der 0.8.2",
"der 0.8.1",
"hybrid-array",
"subtle",
"zeroize",
@@ -11994,7 +11994,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1d9efca8738c78ee9484207732f728b1ef517bbb1833d6fc0879ca898a522f6f"
dependencies = [
"base64ct",
"der 0.8.2",
"der 0.8.1",
]
[[package]]
@@ -12406,10 +12406,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd"
dependencies = [
"fastrand",
"getrandom 0.3.4",
"getrandom 0.4.3",
"once_cell",
"rustix",
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
@@ -13528,7 +13528,7 @@ version = "0.1.11"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22"
dependencies = [
"windows-sys 0.59.0",
"windows-sys 0.61.2",
]
[[package]]
+5 -5
View File
@@ -256,10 +256,10 @@ clap = { version = "4.6.6" }
const-str = { version = "1.1.0" }
convert_case = "0.12.0"
criterion = { version = "0.8" }
crossbeam-queue = "0.3.14"
crossbeam-channel = "0.5.17"
crossbeam-deque = "0.8.8"
crossbeam-utils = "0.8.23"
crossbeam-queue = "0.3.13"
crossbeam-channel = "0.5.16"
crossbeam-deque = "0.8.7"
crossbeam-utils = "0.8.22"
datafusion = { default-features = false, version = "55.0.0" }
derive_builder = "0.20.2"
enumset = "1.1.14"
@@ -306,7 +306,7 @@ rustfs-erasure-codec = { version = "8.0.2" }
reed-solomon-simd = "3.1.0"
regex = { version = "1.13.1" }
rumqttc = { package = "rumqttc-next", version = "0.34.0" }
redis = { version = "1.7.0" }
redis = { version = "1.6.0" }
rustify = { version = "0.7", default-features = false }
rustix = { version = "1.1.4" }
rust-embed = { version = "8.12.0" }
+1 -1
View File
@@ -293,7 +293,7 @@ pub mod cache {
pub mod capacity {
pub use crate::core::pools::{
DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free,
path2_bucket_object, path2_bucket_object_with_base_path,
is_pool_activation_fleet_proof_error, path2_bucket_object, path2_bucket_object_with_base_path,
};
pub use crate::store::utils::is_reserved_or_invalid_bucket;
}
+1 -1
View File
@@ -3385,7 +3385,7 @@ pub(crate) async fn acquire_pool_activation_fleet_proof(
.ok_or_else(|| Error::other(POOL_ACTIVATION_FLEET_PROOF_REQUIRED))
}
pub(crate) fn is_pool_activation_fleet_proof_error(err: &Error) -> bool {
pub fn is_pool_activation_fleet_proof_error(err: &Error) -> bool {
// Save-stage helpers add context by formatting the original error, so the
// marker may be nested in the display string. Restrict matching to the
// `Error::other` I/O shape used by this activation path.
+1 -17
View File
@@ -16291,23 +16291,7 @@ mod transition_upload_integrity_tests {
let bucket = format!("transition-real-bitrot-{}", position.label());
let object = format!("{}-corrupt.bin", position.label());
let payload = vec![0x41; 2 * 1024 * 1024];
// Corruption edits physical shards, so every healthy rename must finish first.
for disk in &disk_stores {
disk.make_volume(&bucket).await.expect("bucket volume should be created");
}
let mut reader = PutObjReader::from_vec(payload.to_vec());
let original = set_disks
.put_object(
&bucket,
&object,
&mut reader,
&ObjectOptions {
write_completion: WriteCompletion::TailDrained,
..Default::default()
},
)
.await
.expect("source object should be written");
let original = write_source(&set_disks, &disk_stores, &bucket, &object, &payload).await;
let source = set_disks
.get_object_fileinfo(
&bucket,
+11 -54
View File
@@ -1730,7 +1730,6 @@ where
};
let baseline_publication_epoch = baseline_publication_guard.epoch();
let usage_persist_baseline_result = read_data_usage_persist_baseline(storeapi.clone()).await;
let observed_usage_candidate_result = read_config(storeapi.clone(), DATA_USAGE_OBSERVED_OBJ_NAME_PATH.as_str()).await;
drop(baseline_publication_guard);
let usage_persist_baseline = match usage_persist_baseline_result {
Ok(baseline) => baseline,
@@ -1751,24 +1750,6 @@ where
return ScannerCycleOutcome::Failed;
}
};
let observed_usage_candidate = match observed_usage_candidate_result {
Ok(candidate) => Some(Bytes::from(candidate)),
Err(EcstoreError::ConfigNotFound) => None,
Err(err) => {
debug!(
target: "rustfs::scanner",
event = EVENT_SCANNER_PERSIST_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_RUNTIME,
cycle = cycle_info.current,
path = %DATA_USAGE_OBSERVED_OBJ_NAME_PATH.as_str(),
state = "observed_candidate_load_failed",
error = %err,
"Scanner skipped an unavailable observed usage candidate for scoped refresh"
);
None
}
};
let (sender, receiver) = mpsc::channel::<DataUsageInfo>(1);
let done_cycle = Metrics::time(Metric::ScanCycle);
@@ -1783,7 +1764,6 @@ where
scan_mode,
scan_scope: crate::scanner_io::ScannerBucketScanScope::default(),
persisted_usage_baseline: usage_persist_baseline.data.clone(),
observed_usage_candidate,
requires_full_scan: scheduling.requires_full_scan,
service_cohort: scheduling.service_cohort,
#[cfg(test)]
@@ -1910,10 +1890,10 @@ where
.as_ref()
.map(|(notification_system, grants)| (Arc::clone(notification_system), grants.clone()));
let remote_lease_release_safe = Arc::new(AtomicBool::new(true));
let mut usage_publication_result = match publication_defer_reason {
let mut usage_persist_outcome = match publication_defer_reason {
Some(reason) => {
drop(receiver);
DataUsagePublicationResult::from(DataUsagePersistOutcome::Deferred(reason))
DataUsagePersistOutcome::Deferred(reason)
}
None => {
// ScannerIO emits its complete or observational update only after
@@ -1928,11 +1908,6 @@ where
.as_ref()
.map(|(_, grants)| grants.iter().map(|grant| grant.lease.token).collect())
.unwrap_or_default();
let ack_expectation = scan_result
.as_ref()
.ok()
.filter(|result| result.has_dirty_usage_to_acknowledge())
.and_then(ScannerCycleResult::publication_expectation);
let mut usage_persist_task = AbortOnDropHandle::new(tokio::spawn(async move {
store_data_usage_in_backend_with_outcome_for_epoch_and_baseline_and_route_probe_for_publication_epoch_and_lease_fence(
ctx_clone,
@@ -1945,7 +1920,6 @@ where
remote_lease_deadline,
remote_lease_fence,
)
.with_ack_expectation(ack_expectation)
.with_remote_lease_tokens(remote_lease_tokens)
.with_lease_release_flag(remote_lease_release_safe_for_task),
move || {
@@ -1979,7 +1953,7 @@ where
error = %err,
"Scanner data usage persistence task failed"
);
DataUsagePublicationResult::from(DataUsagePersistOutcome::Failed)
DataUsagePersistOutcome::Failed
}
DataUsagePersistTaskResult::Cancelled => {
debug!(
@@ -1991,7 +1965,7 @@ where
state = "usage_persist_task_cancelled",
"Scanner data usage persistence task cancelled"
);
DataUsagePublicationResult::from(DataUsagePersistOutcome::Failed)
DataUsagePersistOutcome::Failed
}
DataUsagePersistTaskResult::TimedOut => {
error!(
@@ -2004,12 +1978,11 @@ where
state = "usage_persist_task_timed_out",
"Scanner data usage persistence task timed out"
);
DataUsagePublicationResult::from(DataUsagePersistOutcome::Failed)
DataUsagePersistOutcome::Failed
}
}
}
};
let mut usage_persist_outcome = usage_publication_result.outcome();
let lease_expired = remote_publication_leases
.as_ref()
.is_some_and(|(_, grants)| grants.iter().any(|grant| !grant.lease.is_valid()));
@@ -2209,9 +2182,8 @@ where
};
}
usage_publication_result.restrict_outcome(usage_persist_outcome);
let (completion_outcome, scanner_pending_maintenance_work, remote_dirty_usage_acknowledgements) =
finalize_scanner_cycle_result(scan_cycle_result, usage_publication_result);
finalize_scanner_cycle_result(scan_cycle_result, usage_persist_outcome);
let remote_dirty_usage_pending = if remote_dirty_usage_acknowledgements.is_empty() {
false
} else if let Some(notification_system) = storeapi.scanner_notification_system() {
@@ -3445,35 +3417,21 @@ fn scanner_cycle_completion_outcome(
fn finalize_scanner_cycle_result(
scan_cycle_result: crate::scanner_io::ScannerCycleResult,
publication: DataUsagePublicationResult,
usage_persist_outcome: DataUsagePersistOutcome,
) -> (ScannerCycleOutcome, bool, Vec<ScannerDirtyUsageAcknowledgement>) {
let (usage_persist_outcome, proof) = publication.into_parts();
let completion_outcome = scanner_cycle_completion_outcome_for_result(&scan_cycle_result, usage_persist_outcome);
let pending_maintenance_work = scan_cycle_result.has_pending_maintenance_work();
let durable_complete_snapshot = scan_cycle_result.status == ScannerCycleStatus::Complete
&& matches!(
usage_persist_outcome,
DataUsagePersistOutcome::Saved | DataUsagePersistOutcome::AlreadyDurable
)
&& scan_cycle_result.publication_expectation().as_ref().is_some_and(|expected| {
proof
.as_ref()
.is_some_and(|proof| proof.verified_version_for(expected).is_some())
});
let pending_maintenance_work = scan_cycle_result.has_pending_maintenance_work()
|| (scan_cycle_result.has_dirty_usage_to_acknowledge() && !durable_complete_snapshot);
);
let remote_dirty_usage_acknowledgements = if durable_complete_snapshot {
match proof {
Some(proof) => scan_cycle_result.acknowledge_durable_usage(&proof),
None => Vec::new(),
}
scan_cycle_result.acknowledge_durable_usage()
} else {
Vec::new()
};
(
completion_outcome,
pending_maintenance_work || crate::scanner_io::dirty_usage_buckets_pending(),
remote_dirty_usage_acknowledgements,
)
(completion_outcome, pending_maintenance_work, remote_dirty_usage_acknowledgements)
}
fn scanner_cycle_completion_outcome_for_result(
@@ -3573,7 +3531,6 @@ use activity::*;
use backlog::*;
use cycle_state::*;
use leadership::*;
pub(crate) use usage_store::RootPublicationProof;
use usage_store::*;
pub use activity::scanner_topology_digest;
+15 -22
View File
@@ -37,7 +37,6 @@ use tokio::time::{Duration, advance};
const TEST_DEFAULT_SCANNER_CYCLE_SECS: u64 = 24 * 60 * 60;
mod recovery_control;
mod scoped_ack_publication;
async fn setup_scanner_cycle_store() -> (tempfile::TempDir, Arc<ECStore>) {
setup_scanner_cycle_store_with_usage_baseline(true).await
@@ -6185,7 +6184,7 @@ async fn coordinator_classifies_an_expired_publication_lease() {
.await;
assert_eq!(
outcome.outcome(),
outcome,
DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::PublicationLeaseDeadlineExceeded)
);
assert!(store.put_counts.lock().await.is_empty(), "expired lease must prevent a PUT");
@@ -7326,7 +7325,7 @@ fn scanner_cycle_cache_floor_stays_pending_during_deferred_usage_publication() {
#[test]
#[serial]
fn finalizing_a_saved_enum_without_proof_keeps_dirty_pending() {
fn finalizing_a_saved_cycle_acknowledges_its_exact_dirty_snapshot() {
crate::scanner_io::clear_dirty_usage_bucket("photos");
crate::scanner_io::record_dirty_usage_bucket("photos");
let dirty_snapshot = crate::scanner_io::dirty_usage_buckets_for_tests();
@@ -7338,19 +7337,17 @@ fn finalizing_a_saved_enum_without_proof_keeps_dirty_pending() {
};
let unsaved = crate::scanner_io::ScannerCycleResult::new(ScannerCycleStatus::Complete, Some(dirty_snapshot.clone()))
.with_remote_dirty_usage_acknowledgements(vec![remote_acknowledgement.clone()]);
let (outcome, _, acknowledgements) = finalize_scanner_cycle_result(unsaved, DataUsagePersistOutcome::NoUpdate.into());
let (outcome, _, acknowledgements) = finalize_scanner_cycle_result(unsaved, DataUsagePersistOutcome::NoUpdate);
assert_eq!(outcome, ScannerCycleOutcome::Failed);
assert!(acknowledgements.is_empty());
assert!(crate::scanner_io::dirty_usage_buckets_pending());
let saved = crate::scanner_io::ScannerCycleResult::new(ScannerCycleStatus::Complete, Some(dirty_snapshot))
.with_remote_dirty_usage_acknowledgements(vec![remote_acknowledgement]);
let (outcome, pending, acknowledgements) = finalize_scanner_cycle_result(saved, DataUsagePersistOutcome::Saved.into());
.with_remote_dirty_usage_acknowledgements(vec![remote_acknowledgement.clone()]);
let (outcome, _, acknowledgements) = finalize_scanner_cycle_result(saved, DataUsagePersistOutcome::Saved);
assert_eq!(outcome, ScannerCycleOutcome::Completed);
assert!(acknowledgements.is_empty());
assert!(pending);
assert!(crate::scanner_io::dirty_usage_buckets_pending());
crate::scanner_io::clear_dirty_usage_bucket("photos");
assert_eq!(acknowledgements, vec![remote_acknowledgement]);
assert!(!crate::scanner_io::dirty_usage_buckets_pending());
}
#[test]
@@ -7362,7 +7359,7 @@ fn finalizing_a_deferred_usage_save_keeps_dirty_work_pending() {
let deferred = crate::scanner_io::ScannerCycleResult::new(ScannerCycleStatus::Complete, Some(dirty_snapshot));
let (outcome, _, acknowledgements) =
finalize_scanner_cycle_result(deferred, DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement).into());
finalize_scanner_cycle_result(deferred, DataUsagePersistOutcome::Deferred(ScannerCycleDeferReason::DataMovement));
assert_eq!(outcome, ScannerCycleOutcome::Deferred(ScannerCycleDeferReason::DataMovement));
assert!(acknowledgements.is_empty());
@@ -7382,7 +7379,7 @@ fn finalizing_post_scan_observation_advances_partially_without_dirty_ack() {
)
.with_observational_snapshot_published(true);
let (outcome, _, acknowledgements) = finalize_scanner_cycle_result(observed, DataUsagePersistOutcome::Saved.into());
let (outcome, _, acknowledgements) = finalize_scanner_cycle_result(observed, DataUsagePersistOutcome::Saved);
assert_eq!(outcome, ScannerCycleOutcome::Partial);
assert!(acknowledgements.is_empty());
@@ -7418,20 +7415,17 @@ async fn scanner_cycle_keeps_remote_pending_acknowledgement() {
#[test]
#[serial]
fn finalizing_an_already_durable_enum_without_proof_keeps_dirty_pending() {
fn finalizing_an_already_durable_cycle_acknowledges_its_exact_dirty_snapshot() {
crate::scanner_io::clear_dirty_usage_bucket("photos");
crate::scanner_io::record_dirty_usage_bucket("photos");
let dirty_snapshot = crate::scanner_io::dirty_usage_buckets_for_tests();
let durable = crate::scanner_io::ScannerCycleResult::new(ScannerCycleStatus::Complete, Some(dirty_snapshot));
let (outcome, pending, acknowledgements) =
finalize_scanner_cycle_result(durable, DataUsagePersistOutcome::AlreadyDurable.into());
let (outcome, _, acknowledgements) = finalize_scanner_cycle_result(durable, DataUsagePersistOutcome::AlreadyDurable);
assert_eq!(outcome, ScannerCycleOutcome::Completed);
assert!(acknowledgements.is_empty());
assert!(pending);
assert!(crate::scanner_io::dirty_usage_buckets_pending());
crate::scanner_io::clear_dirty_usage_bucket("photos");
assert!(!crate::scanner_io::dirty_usage_buckets_pending());
}
#[test]
@@ -7442,8 +7436,7 @@ fn finalizing_a_prior_same_cycle_snapshot_keeps_new_dirty_work_pending() {
let dirty_snapshot = crate::scanner_io::dirty_usage_buckets_for_tests();
let durable = crate::scanner_io::ScannerCycleResult::new(ScannerCycleStatus::Complete, Some(dirty_snapshot));
let (outcome, _, acknowledgements) =
finalize_scanner_cycle_result(durable, DataUsagePersistOutcome::PriorCycleDurable.into());
let (outcome, _, acknowledgements) = finalize_scanner_cycle_result(durable, DataUsagePersistOutcome::PriorCycleDurable);
assert_eq!(outcome, ScannerCycleOutcome::Completed);
assert!(acknowledgements.is_empty());
@@ -7459,7 +7452,7 @@ fn finalizing_a_durable_superseded_snapshot_keeps_dirty_work_pending() {
let dirty_snapshot = crate::scanner_io::dirty_usage_buckets_for_tests();
let superseded = crate::scanner_io::ScannerCycleResult::new(ScannerCycleStatus::Superseded, Some(dirty_snapshot));
let (outcome, _, acknowledgements) = finalize_scanner_cycle_result(superseded, DataUsagePersistOutcome::Saved.into());
let (outcome, _, acknowledgements) = finalize_scanner_cycle_result(superseded, DataUsagePersistOutcome::Saved);
assert_eq!(outcome, ScannerCycleOutcome::Superseded);
assert!(acknowledgements.is_empty());
@@ -8888,7 +8881,7 @@ fn post_lease_activity_proof_rejects_a_put_tail_that_finished_before_lease_acqui
]);
let (outcome, _, acknowledgements) = finalize_scanner_cycle_result(
result,
DataUsagePersistOutcome::Deferred(reason.expect("changed namespace should defer publication")).into(),
DataUsagePersistOutcome::Deferred(reason.expect("changed namespace should defer publication")),
);
assert_eq!(
outcome,
@@ -1,490 +0,0 @@
// Copyright 2026 RustFS Team
// Licensed under the Apache License, Version 2.0.
use super::super::usage_store::DataUsagePublicationResult;
use super::*;
use crate::scanner_io::ScannerBucketScanScope;
use rustfs_utils::path::path_join_buf;
use sha2::Digest;
use std::time::SystemTime;
const PROOF_BUCKET: &str = "publication-proof-bucket";
const PROOF_EPOCH: u64 = 7;
const PROOF_CYCLE: u64 = 11;
async fn settle_namespace_commits(store: &ECStore) {
tokio::time::timeout(Duration::from_secs(30), async {
while store.scanner_data_usage_publication_blocked().await {
tokio::time::sleep(Duration::from_millis(1)).await;
}
})
.await
.expect("fixture namespace commits must settle before collecting complete coverage");
}
async fn complete_candidate(store: &Arc<ECStore>, cycle: u64) -> (crate::scanner_io::ScannerCycleResult, DataUsageInfo) {
settle_namespace_commits(store).await;
let ctx = CancellationToken::new();
let budget = ScannerCycleBudget::new_with_progress_tracking(
&ctx,
ScannerCycleBudgetConfig {
max_objects: Some(8),
..Default::default()
},
);
let (updates, mut receiver) = mpsc::channel(1);
let result = crate::scanner_io::nsscanner_with_storage_status_scoped(
store.as_ref(),
crate::scanner_io::ScannerCycleRequest {
ctx,
budget,
updates,
want_cycle: cycle,
leader_epoch: PROOF_EPOCH,
scan_mode: HealScanMode::Normal,
scan_scope: ScannerBucketScanScope::default(),
persisted_usage_baseline: None,
observed_usage_candidate: None,
requires_full_scan: true,
service_cohort: None,
resolved_scope_observer: None,
},
)
.await
.expect("real scanner must produce the fixture candidate");
assert_eq!(result.status, ScannerCycleStatus::Complete);
let candidate = receiver.recv().await.expect("complete scanner snapshot");
assert!(candidate.usage_snapshot_complete);
assert_eq!(candidate.usage_snapshot_converged, Some(true));
assert_eq!(candidate.scanner_cycle, Some(cycle));
assert_eq!(candidate.scanner_epoch, Some(PROOF_EPOCH));
(result, candidate)
}
async fn candidate_store() -> (tempfile::TempDir, Arc<ECStore>) {
crate::scanner_io::clear_dirty_usage_buckets_for_tests();
let (directory, store) = setup_scanner_cycle_store_with_usage_baseline(false).await;
store
.make_bucket(PROOF_BUCKET, &crate::storage_api::scan::MakeBucketOptions::default())
.await
.expect("create proof fixture bucket through the owner");
let mut reader = PutObjReader::from_vec(b"proof".to_vec());
store.pools[0].disk_set[0]
.put_object(PROOF_BUCKET, "initial", &mut reader, &ObjectOptions::default())
.await
.expect("persist fixture object through the owner");
crate::scanner_io::record_dirty_usage_bucket(PROOF_BUCKET);
settle_namespace_commits(&store).await;
(directory, store)
}
async fn read_root(store: &Arc<ECStore>) -> (Option<Vec<u8>>, DataUsageCacheRevision) {
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("read actual v2 root bytes and revision")
}
async fn publish_candidate(
store: &Arc<ECStore>,
scan: &crate::scanner_io::ScannerCycleResult,
candidate: DataUsageInfo,
baseline: Option<DataUsagePersistBaseline>,
) -> DataUsagePublicationResult {
let expectation = scan.publication_expectation();
assert!(expectation.is_some(), "only a real complete scan may supply the expectation");
let (sender, receiver) = mpsc::channel(1);
sender.send(candidate).await.expect("enqueue the real scan candidate");
drop(sender);
store_data_usage_in_backend_with_outcome_for_epoch_and_baseline_and_route_probe_for_publication_epoch_and_lease_fence(
CancellationToken::new(),
store.clone(),
receiver,
Some(PROOF_EPOCH),
baseline,
ScannerPublicationFence::new(scan.publication_epoch(), None, None).with_ack_expectation(expectation),
|| async { None },
)
.await
}
#[tokio::test]
#[serial]
async fn scoped_ack_publication_companion_only_does_not_authorize_root_ack() {
for companion in [
format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()),
LEGACY_DATA_USAGE_OBJ_NAME_PATH.to_string(),
format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()),
] {
let (_directory, store) = candidate_store().await;
let (scan, candidate) = complete_candidate(&store, PROOF_CYCLE).await;
let bytes = serde_json::to_vec(&candidate).expect("actual candidate JSON");
save_config(store.clone(), &companion, bytes.clone())
.await
.expect("persist the companion on real disks");
let baseline = read_data_usage_persist_baseline(store.clone())
.await
.expect("companion fallback baseline");
assert_eq!(baseline.data.as_deref(), Some(bytes.as_slice()));
assert_eq!(baseline.revision, DataUsageCacheRevision::Missing);
assert_eq!(read_root(&store).await.0, None);
let dirty = crate::scanner_io::dirty_usage_buckets_for_tests();
let publication = publish_candidate(&store, &scan, candidate, Some(baseline)).await;
assert_eq!(publication.outcome(), DataUsagePersistOutcome::AlreadyDurable);
let (_, pending, acknowledgements) = finalize_scanner_cycle_result(scan, publication);
assert!(pending, "unacknowledged durable companion work must remain pending");
assert!(acknowledgements.is_empty());
assert_eq!(
crate::scanner_io::dirty_usage_buckets_for_tests(),
dirty,
"a companion is not the v2 root target"
);
assert_eq!(read_root(&store).await, (None, DataUsageCacheRevision::Missing));
assert_eq!(read_config(store.clone(), &companion).await.expect("companion retained"), bytes);
}
crate::scanner_io::clear_dirty_usage_buckets_for_tests();
}
#[tokio::test]
#[serial]
async fn scoped_ack_publication_actual_root_readback_accepts_semantic_json_equivalence() {
let (_directory, store) = candidate_store().await;
let (scan, candidate) = complete_candidate(&store, PROOF_CYCLE).await;
let canonical = serde_json::to_vec(&candidate).expect("candidate encoding");
let mut value = serde_json::to_value(&candidate).expect("candidate value");
value
.as_object_mut()
.expect("usage object")
.insert("fixture_unknown_field".into(), serde_json::json!({"retained": true}));
let different_bytes = serde_json::to_vec_pretty(&value).expect("noncanonical primary JSON");
assert_ne!(different_bytes, canonical);
assert_eq!(
serde_json::from_slice::<DataUsageInfo>(&different_bytes).expect("semantic primary"),
candidate
);
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), different_bytes.clone())
.await
.expect("persist actual primary representation");
let before = read_root(&store).await;
assert!(matches!(&before.1, DataUsageCacheRevision::Etag(etag) if !etag.is_empty()));
let baseline = read_data_usage_persist_baseline(store.clone())
.await
.expect("real primary revision");
assert!(crate::scanner_io::dirty_usage_buckets_pending());
let publication = publish_candidate(&store, &scan, candidate.clone(), Some(baseline)).await;
assert_eq!(publication.outcome(), DataUsagePersistOutcome::AlreadyDurable);
let (_, proof) = publication.into_parts();
let proof = proof.expect("actual primary readback must produce its own root proof");
let expected = scan.publication_expectation().expect("real scan expectation");
let (etag, raw_digest) = proof.verified_version_for(&expected).expect("proof must bind this candidate");
let DataUsageCacheRevision::Etag(expected_etag) = &before.1 else { panic!("actual root ETag") };
assert_eq!(etag, expected_etag);
let expected_digest: [u8; 32] = sha2::Sha256::digest(&different_bytes).into();
assert_eq!(
*raw_digest, expected_digest,
"proof must record actual bytes, not reserialized candidate bytes"
);
// Obtain another proof through the same real readback path rather than
// fabricating a publication result from the inspected proof above.
let publication = publish_candidate(&store, &scan, candidate, None).await;
let (outcome, _, acknowledgements) = finalize_scanner_cycle_result(scan, publication);
assert_eq!(outcome, ScannerCycleOutcome::Completed);
assert!(acknowledgements.is_empty(), "the single-node fixture has no remote targets");
assert!(
!crate::scanner_io::dirty_usage_buckets_pending(),
"actual root bytes plus a real revision authorize this scan"
);
assert_eq!(read_root(&store).await, before, "readback must not rewrite unknown fields or whitespace");
}
#[tokio::test]
#[serial]
async fn scoped_ack_publication_successful_root_cas_authorizes_its_scan() {
let (_directory, store) = candidate_store().await;
let (scan, candidate) = complete_candidate(&store, PROOF_CYCLE).await;
let baseline = read_data_usage_persist_baseline(store.clone())
.await
.expect("initial root revision");
assert_eq!(baseline.revision, DataUsageCacheRevision::Missing);
assert!(crate::scanner_io::dirty_usage_buckets_pending());
let publication = publish_candidate(&store, &scan, candidate.clone(), Some(baseline)).await;
assert_eq!(publication.outcome(), DataUsagePersistOutcome::Saved);
let (bytes, revision) = read_root(&store).await;
assert!(matches!(revision, DataUsageCacheRevision::Etag(etag) if !etag.is_empty()));
assert_eq!(
serde_json::from_slice::<DataUsageInfo>(&bytes.expect("actual saved root")).expect("root JSON"),
candidate
);
let (outcome, pending, acknowledgements) = finalize_scanner_cycle_result(scan, publication);
assert_eq!(outcome, ScannerCycleOutcome::Completed);
assert!(!pending);
assert!(acknowledgements.is_empty(), "the single-node fixture has no remote targets");
assert!(
!crate::scanner_io::dirty_usage_buckets_pending(),
"the real root CAS must authorize its matching scan"
);
}
#[tokio::test]
#[serial]
async fn scoped_ack_publication_observed_candidate_reuse_requires_a_new_root_proof() {
let (_directory, store) = candidate_store().await;
let bootstrap = scanner_usage_bootstrap_marker(SystemTime::now(), Some(PROOF_EPOCH));
save_config(
store.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
serde_json::to_vec(&bootstrap).expect("bootstrap root encoding"),
)
.await
.expect("persist authoritative bootstrap root");
let (prior_scan, mut observed_candidate) = complete_candidate(&store, PROOF_CYCLE).await;
// Seed a complete but unconverged observation from real scanner coverage;
// the production writer attaches its authoritative baseline identity.
observed_candidate.usage_snapshot_converged = Some(false);
let observation = publish_candidate(&store, &prior_scan, observed_candidate, None).await;
let (outcome, proof) = observation.into_parts();
assert_eq!(outcome, DataUsagePersistOutcome::Saved);
assert!(proof.is_none(), "an observational write cannot authorize a root ACK");
let (root_before, revision_before) = read_root(&store).await;
let observed = read_config(store.clone(), DATA_USAGE_OBSERVED_OBJ_NAME_PATH.as_str())
.await
.expect("read real persisted observation");
let ctx = CancellationToken::new();
let budget = ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default());
let (updates, mut receiver) = mpsc::channel(1);
let (observer, selected) = tokio::sync::oneshot::channel();
let scan = crate::scanner_io::nsscanner_with_storage_status_scoped(
store.as_ref(),
crate::scanner_io::ScannerCycleRequest {
ctx,
budget,
updates,
want_cycle: PROOF_CYCLE + 1,
leader_epoch: PROOF_EPOCH,
scan_mode: HealScanMode::Normal,
scan_scope: ScannerBucketScanScope::default(),
persisted_usage_baseline: root_before.clone().map(Bytes::from),
observed_usage_candidate: Some(Bytes::from(observed)),
requires_full_scan: false,
service_cohort: None,
resolved_scope_observer: Some(observer),
},
)
.await
.expect("observation-backed scope must run through the real scanner");
let scope = selected.await.expect("production resolver decision");
assert_eq!(scope.selected_buckets_for_tests(), Some(&HashSet::from([PROOF_BUCKET.to_string()])));
assert_eq!(scan.status, ScannerCycleStatus::Complete);
let expectation = scan.publication_expectation().expect("reused coverage must be revalidated");
assert!(
!expectation.same_candidate(&prior_scan.publication_expectation().expect("prior real candidate")),
"the observation cannot transfer the previous scan's expectation"
);
assert_eq!(read_root(&store).await, (root_before, revision_before));
assert!(crate::scanner_io::dirty_usage_buckets_pending());
let candidate = receiver.recv().await.expect("new validated root candidate");
assert_eq!(candidate.scanner_cycle, Some(PROOF_CYCLE + 1));
assert_eq!(candidate.usage_snapshot_converged, Some(true));
let publication = publish_candidate(&store, &scan, candidate, None).await;
assert_eq!(publication.outcome(), DataUsagePersistOutcome::Saved);
let (outcome, pending, acknowledgements) = finalize_scanner_cycle_result(scan, publication);
assert_eq!(outcome, ScannerCycleOutcome::Completed);
assert!(!pending);
assert!(acknowledgements.is_empty());
assert!(!crate::scanner_io::dirty_usage_buckets_pending());
}
#[tokio::test]
#[serial]
async fn scoped_ack_publication_stale_root_cas_keeps_dirty_after_bucket_save() {
let (_directory, store) = candidate_store().await;
let (scan, candidate) = complete_candidate(&store, PROOF_CYCLE).await;
let mut bucket_cache = DataUsageCache::default();
bucket_cache
.load(store.pools[0].disk_set[0].clone(), &path_join_buf(&[PROOF_BUCKET, DATA_USAGE_CACHE_NAME]))
.await
.expect("real bucket checkpoint must be persisted before root publication");
assert!(bucket_cache.info.snapshot_complete);
assert_eq!(
bucket_cache
.checked_flatten(PROOF_BUCKET)
.expect("persisted bucket root")
.objects,
1
);
let stale_baseline = read_data_usage_persist_baseline(store.clone())
.await
.expect("missing root revision");
assert_eq!(stale_baseline.revision, DataUsageCacheRevision::Missing);
let mut competing = candidate.clone();
competing.scanner_epoch = Some(PROOF_EPOCH + 1);
competing.scanner_cycle = Some(PROOF_CYCLE + 1);
for state in &mut competing.usage_snapshot_set_states {
state.scanner_epoch = Some(PROOF_EPOCH + 1);
state.scanner_cycle = Some(PROOF_CYCLE + 1);
}
let competing_bytes = serde_json::to_vec(&competing).expect("competing root");
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), competing_bytes.clone())
.await
.expect("another publisher wins the actual root slot");
let before = read_root(&store).await;
let dirty = crate::scanner_io::dirty_usage_buckets_for_tests();
let publication = publish_candidate(&store, &scan, candidate, Some(stale_baseline)).await;
assert_eq!(
publication.outcome(),
DataUsagePersistOutcome::Current,
"the old missing revision loses CAS and reconciles the newer root"
);
let (_, _, acknowledgements) = finalize_scanner_cycle_result(scan, publication);
assert!(acknowledgements.is_empty());
assert_eq!(crate::scanner_io::dirty_usage_buckets_for_tests(), dirty);
assert_eq!(
read_root(&store).await,
before,
"bucket durability must not authorize replacing the winning root"
);
crate::scanner_io::clear_dirty_usage_buckets_for_tests();
}
#[tokio::test]
#[serial]
async fn scoped_ack_publication_cannot_transfer_proof_between_real_scan_results() {
let (_directory, store) = candidate_store().await;
let (first_scan, first_candidate) = complete_candidate(&store, PROOF_CYCLE).await;
let (second_scan, second_candidate) = complete_candidate(&store, PROOF_CYCLE).await;
assert_eq!(first_candidate.scanner_epoch, second_candidate.scanner_epoch);
assert_eq!(first_candidate.scanner_cycle, second_candidate.scanner_cycle);
assert_eq!(first_candidate.objects_total_count, second_candidate.objects_total_count);
let baseline = read_data_usage_persist_baseline(store.clone())
.await
.expect("initial root revision");
let dirty = crate::scanner_io::dirty_usage_buckets_for_tests();
let publication = publish_candidate(&store, &first_scan, first_candidate, Some(baseline)).await;
assert_eq!(publication.outcome(), DataUsagePersistOutcome::Saved);
assert!(read_root(&store).await.0.is_some(), "the first scan really published its root");
let (_, pending, acknowledgements) = finalize_scanner_cycle_result(second_scan, publication);
assert!(pending, "another scan's publication must not finish this scan's dirty maintenance work");
assert!(acknowledgements.is_empty());
assert_eq!(
crate::scanner_io::dirty_usage_buckets_for_tests(),
dirty,
"same counters and cycle cannot transfer another scan's proof"
);
crate::scanner_io::clear_dirty_usage_buckets_for_tests();
}
#[tokio::test]
#[serial]
async fn scoped_ack_publication_stale_baseline_cannot_prove_a_replaced_root() {
let (_directory, store) = candidate_store().await;
let (first_scan, first_candidate) = complete_candidate(&store, PROOF_CYCLE).await;
save_config(
store.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
serde_json::to_vec(&first_candidate).expect("first candidate"),
)
.await
.expect("persist the first candidate on real disks");
let stale_baseline = read_data_usage_persist_baseline(store.clone())
.await
.expect("capture the genuine first root revision");
let dirty = crate::scanner_io::dirty_usage_buckets_for_tests();
let mut reader = PutObjReader::from_vec(b"second".to_vec());
store.pools[0].disk_set[0]
.put_object(PROOF_BUCKET, "second", &mut reader, &ObjectOptions::default())
.await
.expect("commit a real namespace change");
assert_eq!(
crate::scanner_io::dirty_usage_buckets_for_tests(),
dirty,
"direct storage writes leave this fixture's scanner hint generation unchanged"
);
let (_, replacement) = complete_candidate(&store, PROOF_CYCLE).await;
assert_eq!(first_candidate.scanner_epoch, replacement.scanner_epoch);
assert_eq!(first_candidate.scanner_cycle, replacement.scanner_cycle);
assert_eq!((first_candidate.objects_total_count, replacement.objects_total_count), (1, 2));
save_config(
store.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
serde_json::to_vec(&replacement).expect("replacement candidate"),
)
.await
.expect("publish the replacement root");
let current = read_root(&store).await;
assert_ne!(current.1, stale_baseline.revision);
// The supplied baseline still equals candidate A, but the actual target
// now contains B. Compatibility's AlreadyDurable outcome is not proof.
let publication = publish_candidate(&store, &first_scan, first_candidate, Some(stale_baseline)).await;
assert_eq!(publication.outcome(), DataUsagePersistOutcome::AlreadyDurable);
let (_, pending, acknowledgements) = finalize_scanner_cycle_result(first_scan, publication);
assert!(pending);
assert!(acknowledgements.is_empty());
assert_eq!(crate::scanner_io::dirty_usage_buckets_for_tests(), dirty);
assert_eq!(read_root(&store).await, current);
crate::scanner_io::clear_dirty_usage_buckets_for_tests();
}
#[tokio::test]
#[serial]
async fn scoped_ack_publication_rejects_builder_mutation_after_real_root_publish() {
for mutation in ["remote_ack_target", "publication_epoch", "remote_lease_targets"] {
let (_directory, store) = candidate_store().await;
let (scan, candidate) = complete_candidate(&store, PROOF_CYCLE).await;
let baseline = read_data_usage_persist_baseline(store.clone())
.await
.expect("initial root revision");
let dirty = crate::scanner_io::dirty_usage_buckets_for_tests();
let changed_generation = dirty
.get(PROOF_BUCKET)
.expect("the real scan has dirty work")
.checked_add(1)
.expect("bounded fixture generation");
let changed_epoch = scan
.publication_epoch()
.expect("real scan publication epoch")
.checked_add(1)
.expect("bounded fixture epoch");
let publication = publish_candidate(&store, &scan, candidate.clone(), Some(baseline)).await;
assert_eq!(publication.outcome(), DataUsagePersistOutcome::Saved, "{mutation}");
let root_before = read_root(&store).await;
assert_eq!(
serde_json::from_slice::<DataUsageInfo>(root_before.0.as_deref().expect("actual saved root"))
.expect("persisted root JSON"),
candidate,
"{mutation}: the original candidate really reached root storage"
);
let changed = match mutation {
"remote_ack_target" => scan.with_remote_dirty_usage_acknowledgements(vec![ScannerDirtyUsageAcknowledgement {
host: "proof-peer:9000".to_string(),
instance_id: crate::scanner_activity_epoch().to_string(),
generation: changed_generation,
}]),
"publication_epoch" => scan.with_publication_epoch(Some(changed_epoch)),
"remote_lease_targets" => scan.with_remote_publication_lease_targets(vec![(
"proof-peer:9000".to_string(),
crate::scanner_activity_epoch().to_string(),
changed_generation,
)]),
_ => unreachable!("fixed mutation cases"),
};
let (_, pending, acknowledgements) = finalize_scanner_cycle_result(changed, publication);
assert!(
acknowledgements.is_empty(),
"{mutation}: the old root proof must not authorize changed ACK work"
);
assert!(pending, "{mutation}: changed maintenance work must remain pending");
assert_eq!(
crate::scanner_io::dirty_usage_buckets_for_tests(),
dirty,
"{mutation}: the changed scan must not clear local dirty work"
);
assert_eq!(read_root(&store).await, root_before, "{mutation}: the durable original root is retained");
}
crate::scanner_io::clear_dirty_usage_buckets_for_tests();
}
+15 -192
View File
@@ -34,139 +34,6 @@ pub(super) enum DataUsagePersistOutcome {
Failed,
}
#[derive(Debug)]
pub(crate) struct RootPublicationProof {
candidate: crate::scanner_io::ScannerPublicationExpectation,
root_version: (String, [u8; 32]),
}
impl RootPublicationProof {
pub(crate) fn verified_version_for(
&self,
expected: &crate::scanner_io::ScannerPublicationExpectation,
) -> Option<&(String, [u8; 32])> {
self.candidate.same_candidate(expected).then_some(&self.root_version)
}
}
#[derive(Debug)]
pub(super) struct DataUsagePublicationResult {
outcome: DataUsagePersistOutcome,
proof: Option<RootPublicationProof>,
}
impl From<DataUsagePersistOutcome> for DataUsagePublicationResult {
fn from(outcome: DataUsagePersistOutcome) -> Self {
Self { outcome, proof: None }
}
}
impl DataUsagePublicationResult {
pub(super) fn outcome(&self) -> DataUsagePersistOutcome {
self.outcome
}
pub(super) fn restrict_outcome(&mut self, outcome: DataUsagePersistOutcome) {
if outcome != self.outcome {
self.proof = None;
}
self.outcome = outcome;
}
pub(super) fn into_parts(self) -> (DataUsagePersistOutcome, Option<RootPublicationProof>) {
(self.outcome, self.proof)
}
}
fn root_ack_write_is_confirmed<T, E>(
result: &std::result::Result<T, E>,
state: Option<ScannerPublicationCommitState>,
written_etag: Option<&str>,
) -> bool {
result.is_ok() && state == Some(ScannerPublicationCommitState::Committed) && written_etag.is_some_and(|etag| !etag.is_empty())
}
#[cfg(test)]
mod root_publication_confirmation_tests {
use super::*;
#[test]
fn root_publication_confirmation_requires_committed_state_and_write_revision() {
let saved = Ok::<(), ()>(());
for state in [
None,
Some(ScannerPublicationCommitState::Admitted),
Some(ScannerPublicationCommitState::InFlight),
Some(ScannerPublicationCommitState::AbortedBeforeCommit),
Some(ScannerPublicationCommitState::Indeterminate),
] {
assert!(!root_ack_write_is_confirmed(&saved, state, Some("revision")));
}
for etag in [None, Some("")] {
assert!(!root_ack_write_is_confirmed(&saved, Some(ScannerPublicationCommitState::Committed), etag));
}
assert!(root_ack_write_is_confirmed(
&saved,
Some(ScannerPublicationCommitState::Committed),
Some("revision")
));
}
#[test]
fn root_publication_confirmation_does_not_carry_state_across_cas_attempts() {
let attempts = [
(Err(()), Some(ScannerPublicationCommitState::Committed), Some("first")),
(Ok(()), Some(ScannerPublicationCommitState::AbortedBeforeCommit), Some("second")),
(Ok(()), None, Some("legacy")),
(Ok(()), Some(ScannerPublicationCommitState::Committed), Some("confirmed")),
];
let confirmations = attempts
.iter()
.map(|(result, state, etag)| root_ack_write_is_confirmed(result, *state, *etag))
.collect::<Vec<_>>();
assert_eq!(confirmations, [false, false, false, true]);
}
}
async fn read_root_publication_proof<S: ScannerObjectIO + ScannerConfigObjectDelete>(
store: Arc<S>,
ctx: &CancellationToken,
deadline: tokio::time::Instant,
epoch: u64,
expected: &crate::scanner_io::ScannerPublicationExpectation,
candidate: &DataUsageInfo,
written_etag: Option<&str>,
) -> Option<RootPublicationProof> {
let read = async {
let _admission = scanner_publication_admission_for_epoch(store.clone(), epoch).await?;
let (bytes, revision) = read_config_with_revision(store, DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.ok()?;
let bytes = bytes?;
let DataUsageCacheRevision::Etag(etag) = revision else {
return None;
};
if etag.is_empty() || written_etag.is_some_and(|written| written != etag) {
return None;
}
let persisted: DataUsageInfo = serde_json::from_slice(&bytes).ok()?;
if &persisted != candidate {
return None;
}
let root_digest = Sha256::digest(&bytes).into();
if ctx.is_cancelled() || tokio::time::Instant::now() >= deadline {
return None;
}
Some(RootPublicationProof {
candidate: expected.clone(),
root_version: (etag, root_digest),
})
};
tokio::select! {
biased;
_ = ctx.cancelled() => None,
result = tokio::time::timeout_at(deadline, read) => result.ok().flatten(),
}
}
fn remote_lease_expired(deadline: Option<std::time::Instant>) -> bool {
deadline.is_some_and(|deadline| std::time::Instant::now() >= deadline)
}
@@ -299,7 +166,6 @@ pub(super) struct ScannerPublicationFence {
pub(super) scanner_publication_lease_fence: Option<String>,
pub(super) remote_lease_tokens: Vec<Uuid>,
pub(super) lease_release_safe: Arc<AtomicBool>,
pub(super) ack_expectation: Option<crate::scanner_io::ScannerPublicationExpectation>,
}
impl ScannerPublicationFence {
@@ -314,7 +180,6 @@ impl ScannerPublicationFence {
scanner_publication_lease_fence,
remote_lease_tokens: Vec::new(),
lease_release_safe: Arc::new(AtomicBool::new(true)),
ack_expectation: None,
}
}
@@ -327,26 +192,21 @@ impl ScannerPublicationFence {
self.lease_release_safe = lease_release_safe;
self
}
pub(super) fn with_ack_expectation(mut self, expected: Option<crate::scanner_io::ScannerPublicationExpectation>) -> Self {
self.ack_expectation = expected;
self
}
}
#[derive(Debug)]
pub(super) enum DataUsagePersistTaskResult<T = DataUsagePersistOutcome> {
Completed(T),
pub(super) enum DataUsagePersistTaskResult {
Completed(DataUsagePersistOutcome),
Cancelled,
TimedOut,
JoinFailed(tokio::task::JoinError),
}
pub(super) async fn wait_for_data_usage_persist_task<T>(
pub(super) async fn wait_for_data_usage_persist_task(
ctx: &CancellationToken,
task: &mut AbortOnDropHandle<T>,
task: &mut AbortOnDropHandle<DataUsagePersistOutcome>,
timeout: Duration,
) -> DataUsagePersistTaskResult<T> {
) -> DataUsagePersistTaskResult {
tokio::select! {
biased;
result = &mut *task => match result {
@@ -460,7 +320,6 @@ where
route_probe,
)
.await
.outcome()
}
pub(super) async fn store_data_usage_in_backend_with_outcome_for_epoch_and_baseline_and_route_probe_for_publication_epoch_and_lease_fence<
@@ -474,7 +333,7 @@ pub(super) async fn store_data_usage_in_backend_with_outcome_for_epoch_and_basel
initial_baseline: Option<DataUsagePersistBaseline>,
publication_fence: ScannerPublicationFence,
route_probe: F,
) -> DataUsagePublicationResult
) -> DataUsagePersistOutcome
where
F: Fn() -> Fut + Send + Sync,
Fut: Future<Output = Option<ScannerCycleDeferReason>> + Send,
@@ -485,15 +344,11 @@ where
scanner_publication_lease_fence,
remote_lease_tokens,
lease_release_safe,
ack_expectation,
} = publication_fence;
let ack_deadline = scanner_publication_scope_deadline(data_usage_persist_timeout(), remote_lease_deadline);
let mut outcome = DataUsagePersistOutcome::NoUpdate;
let mut proof = None;
let mut next_baseline = initial_baseline;
'updates: while let Some(mut data_usage_info) = receiver.recv().await {
proof = None;
let _activity_guard = ScannerActivityGuard::new();
if ctx.is_cancelled() {
break;
@@ -668,14 +523,10 @@ where
continue;
}
};
let data_digest: [u8; 32] = Sha256::digest(&data).into();
let sha256hex = (!data.is_empty()).then(|| hex_simd::encode_to_string(data_digest, hex_simd::AsciiCase::Lower));
let sha256hex = (!data.is_empty()).then(|| hex_simd::encode_to_string(Sha256::digest(&data), hex_simd::AsciiCase::Lower));
let data = Bytes::from(data);
let backup_due = !observational && data_usage_backup_due(&data_usage_info);
let mut cas_retry = 0usize;
let mut ack_epoch = None;
let mut write_confirmed = false;
let mut written_etag = None;
let save_outcome = loop {
if ctx.is_cancelled() {
break 'updates;
@@ -706,7 +557,6 @@ where
} else {
None
};
ack_epoch = Some(publication_epoch_for_save);
let (existing_data, revision) = match baseline {
Some(baseline) => (baseline.data, baseline.revision),
None => match read_config_with_revision(storeapi.clone(), target_path).await {
@@ -795,7 +645,7 @@ where
}
let done_save = Metrics::time(Metric::SaveUsage);
let (save_result, commit_state) = {
let save_result = {
let publication_scope = storeapi
.scanner_data_usage_publication_commit_scope_with_release_flag(
publication_epoch_for_save,
@@ -831,33 +681,24 @@ where
.await;
drop(legacy_publication_admission);
if let Some(scope) = publication_scope {
let state = scope.wait_for_completion().await;
let result = match state {
ScannerPublicationCommitState::Committed => save_result,
ScannerPublicationCommitState::AbortedBeforeCommit => save_result,
match scope.wait_for_completion().await {
ScannerPublicationCommitState::Committed | ScannerPublicationCommitState::AbortedBeforeCommit => {
save_result
}
ScannerPublicationCommitState::Indeterminate
| ScannerPublicationCommitState::Admitted
| ScannerPublicationCommitState::InFlight => Err(EcstoreError::other(
"scanner publication commit scope did not reach a safe terminal state",
)),
};
(result, Some(state))
}
} else {
(save_result, None)
save_result
}
};
done_save();
let attempt_confirmed = root_ack_write_is_confirmed(
&save_result,
commit_state,
save_result.as_ref().ok().and_then(|info| info.etag.as_deref()),
);
match save_result {
Ok(object_info) => {
write_confirmed = attempt_confirmed;
written_etag = object_info.etag.as_ref().filter(|etag| !etag.is_empty()).cloned();
if !observational {
next_baseline = object_info
.etag
@@ -1068,27 +909,9 @@ where
break 'updates;
}
}
if !observational
&& data_usage_info.usage_snapshot_converged == Some(true)
&& matches!(outcome, DataUsagePersistOutcome::Saved | DataUsagePersistOutcome::AlreadyDurable)
&& (outcome == DataUsagePersistOutcome::AlreadyDurable || write_confirmed)
&& let (Some(expected), Some(epoch)) = (ack_expectation.as_ref(), ack_epoch)
&& expected.matches_encoded_candidate(&data_digest)
{
proof = read_root_publication_proof(
storeapi.clone(),
&ctx,
ack_deadline,
epoch,
expected,
&data_usage_info,
written_etag.as_deref(),
)
.await;
}
}
DataUsagePublicationResult { outcome, proof }
outcome
}
async fn cleanup_observed_data_usage_snapshot_for_epoch_and_lease(
+14 -82
View File
@@ -27,7 +27,7 @@ use metrics::counter;
use rand::seq::SliceRandom as _;
#[cfg(test)]
use rustfs_config::{ENV_SCANNER_MAX_CONCURRENT_DISK_SCANS, ENV_SCANNER_MAX_CONCURRENT_SET_SCANS};
use rustfs_data_usage::{BucketTargetUsageInfo, BucketUsageInfo, observed_data_usage_is_newer};
use rustfs_data_usage::{BucketTargetUsageInfo, BucketUsageInfo};
use rustfs_filemeta::FileMeta;
use rustfs_heal_contracts::heal_channel::HealScanMode;
use rustfs_lock::{LockError, NamespaceLockGuard};
@@ -114,11 +114,6 @@ pub(crate) struct ScannerBucketScanScope {
}
impl ScannerBucketScanScope {
#[cfg(test)]
pub(crate) fn selected_buckets_for_tests(&self) -> Option<&HashSet<String>> {
self.selected_buckets.as_deref()
}
fn is_default(&self) -> bool {
self.selected_buckets.is_none() && self.baseline_scan_plan_digest.is_none()
}
@@ -133,8 +128,7 @@ impl ScannerBucketScanScope {
#[derive(Clone, Copy)]
pub(super) struct ScannerCacheBaselineProof<'a> {
pub(super) authoritative_data: Option<&'a Bytes>,
pub(super) observed_candidate_data: Option<&'a Bytes>,
pub(super) data: Option<&'a Bytes>,
pub(super) expected_sources: &'a HashSet<DataUsageCacheSource>,
pub(super) leader_epoch: u64,
pub(super) want_cycle: u64,
@@ -177,23 +171,21 @@ fn verified_remote_dirty_usage_buckets(
(received_peers.len() == expected_peers.len()).then_some(dirty_buckets)
}
fn complete_scanner_cache_snapshot_plan_digest(
snapshot: &DataUsageInfo,
proof: ScannerCacheBaselineProof<'_>,
expected_converged: bool,
) -> Option<DataUsageScanPlanDigest> {
if !snapshot.is_complete_bucket_usage_snapshot()
|| snapshot.usage_snapshot_partial
|| snapshot.usage_snapshot_converged != Some(expected_converged)
|| snapshot.scanner_epoch != Some(proof.leader_epoch)
|| snapshot.usage_snapshot_set_states.len() != proof.expected_sources.len()
fn complete_scanner_cache_baseline_plan_digest(proof: ScannerCacheBaselineProof<'_>) -> Option<DataUsageScanPlanDigest> {
let data = proof.data?;
let baseline = serde_json::from_slice::<DataUsageInfo>(data).ok()?;
if !baseline.is_complete_bucket_usage_snapshot()
|| baseline.usage_snapshot_partial
|| baseline.usage_snapshot_converged != Some(true)
|| baseline.scanner_epoch != Some(proof.leader_epoch)
|| baseline.usage_snapshot_set_states.len() != proof.expected_sources.len()
{
return None;
}
// Completed maintenance also covers ordinary usage. Keep its exact stored
// proof for cache reuse, and reject mixtures of different set work proofs.
let baseline_plan_digest = DataUsageScanPlanDigest(snapshot.usage_snapshot_set_states.first()?.scan_plan_digest?);
let baseline_plan_digest = DataUsageScanPlanDigest(baseline.usage_snapshot_set_states.first()?.scan_plan_digest?);
if ![
proof.scan_plan_digest,
scanner_bucket_work_digest(proof.scan_plan_digest, HealScanMode::Normal, true),
@@ -203,8 +195,8 @@ fn complete_scanner_cache_snapshot_plan_digest(
{
return None;
}
let mut states = HashSet::with_capacity(snapshot.usage_snapshot_set_states.len());
for state in &snapshot.usage_snapshot_set_states {
let mut states = HashSet::with_capacity(baseline.usage_snapshot_set_states.len());
for state in &baseline.usage_snapshot_set_states {
let source = DataUsageCacheSource::new(usize::try_from(state.pool_index).ok()?, usize::try_from(state.set_index).ok()?);
if !proof.expected_sources.contains(&source)
|| !states.insert(source)
@@ -221,30 +213,6 @@ fn complete_scanner_cache_snapshot_plan_digest(
(states == *proof.expected_sources).then_some(baseline_plan_digest)
}
fn complete_scanner_cache_baseline_plan_digest(proof: ScannerCacheBaselineProof<'_>) -> Option<DataUsageScanPlanDigest> {
let authoritative = serde_json::from_slice::<DataUsageInfo>(proof.authoritative_data?).ok()?;
if let Some(validated_digest) = complete_scanner_cache_snapshot_plan_digest(&authoritative, proof, true) {
return Some(validated_digest);
}
// A complete but superseded observation may reuse its per-set cache only
// when it was explicitly tied to the durable authoritative baseline. It
// remains observational: this proof grants bucket-scope reuse only and
// never changes authoritative usage publication or dirty acknowledgement.
let authoritative_has_identity = (crate::scanner::data_usage_info_has_persisted_baseline_identity(&authoritative)
&& authoritative.usage_snapshot_converged != Some(false))
|| crate::scanner::data_usage_info_is_bootstrap_pending(&authoritative);
if !authoritative_has_identity {
return None;
}
let observed = serde_json::from_slice::<DataUsageInfo>(proof.observed_candidate_data?).ok()?;
if !observed_data_usage_is_newer(&observed, &authoritative) {
return None;
}
complete_scanner_cache_snapshot_plan_digest(&observed, proof, false)
}
fn scoped_scan_scope_from_dirty_buckets(
requested_scope: ScannerBucketScanScope,
dirty_buckets: HashSet<String>,
@@ -901,7 +869,6 @@ pub(crate) struct ScannerCycleResult {
failed_dirty_usage: bool,
pending_maintenance_work: bool,
required_cycle_floor: Option<u64>,
publication_expectation: Option<ScannerPublicationExpectation>,
}
impl ScannerCycleResult {
@@ -917,12 +884,10 @@ impl ScannerCycleResult {
failed_dirty_usage: false,
pending_maintenance_work: false,
required_cycle_floor: None,
publication_expectation: None,
}
}
pub(crate) fn with_publication_epoch(mut self, publication_epoch: Option<u64>) -> Self {
self.publication_expectation = None;
self.publication_epoch = publication_epoch;
self
}
@@ -932,7 +897,6 @@ impl ScannerCycleResult {
}
fn with_activity_digest(mut self, activity_digest: [u8; 32]) -> Self {
self.publication_expectation = None;
self.activity_digest = Some(activity_digest);
self
}
@@ -942,7 +906,6 @@ impl ScannerCycleResult {
}
pub(crate) fn with_observational_snapshot_published(mut self, published: bool) -> Self {
self.publication_expectation = None;
self.observational_snapshot_published = published;
self
}
@@ -952,19 +915,16 @@ impl ScannerCycleResult {
}
fn with_failed_dirty_usage(mut self, failed_dirty_usage: bool) -> Self {
self.publication_expectation = None;
self.failed_dirty_usage = failed_dirty_usage;
self
}
fn with_pending_maintenance_work(mut self, pending_maintenance_work: bool) -> Self {
self.publication_expectation = None;
self.pending_maintenance_work = pending_maintenance_work;
self
}
fn with_required_cycle_floor(mut self, required_cycle_floor: Option<u64>) -> Self {
self.publication_expectation = None;
self.required_cycle_floor = required_cycle_floor;
self
}
@@ -973,13 +933,11 @@ impl ScannerCycleResult {
mut self,
acknowledgements: Vec<crate::scanner::ScannerDirtyUsageAcknowledgement>,
) -> Self {
self.publication_expectation = None;
self.remote_dirty_usage_acknowledgements = acknowledgements;
self
}
pub(crate) fn with_remote_publication_lease_targets(mut self, targets: Vec<(String, String, u64)>) -> Self {
self.publication_expectation = None;
self.remote_publication_lease_targets = targets;
self
}
@@ -988,32 +946,7 @@ impl ScannerCycleResult {
&self.remote_publication_lease_targets
}
pub(crate) fn publication_expectation(&self) -> Option<ScannerPublicationExpectation> {
self.publication_expectation.clone()
}
fn with_publication_expectation(mut self, expectation: Option<ScannerPublicationExpectation>) -> Self {
// Seal only after all coverage and acknowledgement inputs are final.
self.publication_expectation = expectation;
self
}
pub(crate) fn acknowledge_durable_usage(
self,
proof: &crate::scanner::RootPublicationProof,
) -> Vec<crate::scanner::ScannerDirtyUsageAcknowledgement> {
if self.status != ScannerCycleStatus::Complete
|| self
.publication_expectation
.as_ref()
.is_none_or(|expected| proof.verified_version_for(expected).is_none())
{
return Vec::new();
}
self.clear_verified_usage()
}
fn clear_verified_usage(self) -> Vec<crate::scanner::ScannerDirtyUsageAcknowledgement> {
pub(crate) fn acknowledge_durable_usage(self) -> Vec<crate::scanner::ScannerDirtyUsageAcknowledgement> {
if let Some(snapshot) = self.dirty_usage_clear {
clear_dirty_usage_buckets(&snapshot);
}
@@ -1053,7 +986,6 @@ mod publish_gate_tests;
#[cfg(test)]
mod tests;
pub(crate) use cache::ScannerPublicationExpectation;
use cache::*;
use dirty_usage::*;
use guards::*;
+1 -98
View File
@@ -308,86 +308,6 @@ impl<'a> ValidatedScannerSnapshot<'a> {
}
}
#[derive(Clone, Debug)]
pub(crate) struct ScannerPublicationExpectation {
candidate: Arc<([u8; 32], DataUsageScanPlanDigest)>,
}
impl ScannerPublicationExpectation {
pub(crate) fn matches_encoded_candidate(&self, digest: &[u8; 32]) -> bool {
&self.candidate.0 == digest
}
pub(crate) fn same_candidate(&self, other: &Self) -> bool {
Arc::ptr_eq(&self.candidate, &other.candidate) && self.candidate.1 == other.candidate.1
}
}
pub(super) struct ValidatedUsageCandidate {
data: DataUsageInfo,
#[cfg(test)]
last_update: SystemTime,
coverage_digest: DataUsageScanPlanDigest,
}
pub(super) fn empty_namespace_usage_candidate(
all_buckets: &[BucketInfo],
sources: &HashSet<DataUsageCacheSource>,
buckets_by_source: &HashMap<DataUsageCacheSource, Vec<BucketInfo>>,
identity: ScannerSnapshotIdentity,
) -> Option<ValidatedUsageCandidate> {
if !all_buckets.is_empty()
|| sources.is_empty()
|| sources.len() != buckets_by_source.len()
|| sources
.iter()
.any(|source| buckets_by_source.get(source).is_none_or(|buckets| !buckets.is_empty()))
{
return None;
}
let last_update = SystemTime::now();
Some(ValidatedUsageCandidate {
data: DataUsageInfo {
last_update: Some(last_update),
scanner_cycle: Some(identity.cycle),
scanner_epoch: Some(identity.leader_epoch),
usage_snapshot_complete: true,
..Default::default()
},
#[cfg(test)]
last_update,
coverage_digest: identity.coverage_digest,
})
}
impl ValidatedUsageCandidate {
pub(super) fn prepare(mut self, status: ScannerCycleStatus) -> (DataUsageInfo, Option<ScannerPublicationExpectation>) {
self.data.usage_snapshot_converged = Some(status == ScannerCycleStatus::Complete);
let expectation = if status == ScannerCycleStatus::Complete {
struct DigestWriter(Sha256);
impl std::io::Write for DigestWriter {
fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
self.0.update(bytes);
Ok(bytes.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
let mut writer = DigestWriter(Sha256::new());
serde_json::to_writer(&mut writer, &self.data)
.ok()
.map(|()| ScannerPublicationExpectation {
candidate: Arc::new((writer.0.finalize().into(), self.coverage_digest)),
})
} else {
None
};
(self.data, expectation)
}
}
#[cfg(test)]
pub(super) fn completed_data_usage_info(
results: &[DataUsageCache],
scope: &ScannerSnapshotScope<'_>,
@@ -396,18 +316,6 @@ pub(super) fn completed_data_usage_info(
budget_elapsed: bool,
cancelled: bool,
) -> Option<(DataUsageInfo, SystemTime)> {
completed_usage_candidate(results, scope, tier_registry_names, bucket_plan_complete, budget_elapsed, cancelled)
.map(|candidate| (candidate.data, candidate.last_update))
}
pub(super) fn completed_usage_candidate(
results: &[DataUsageCache],
scope: &ScannerSnapshotScope<'_>,
tier_registry_names: &[String],
bucket_plan_complete: bool,
budget_elapsed: bool,
cancelled: bool,
) -> Option<ValidatedUsageCandidate> {
if !bucket_plan_complete {
return None;
}
@@ -485,12 +393,7 @@ pub(super) fn completed_usage_candidate(
usage_snapshot_set_states,
..Default::default()
};
Some(ValidatedUsageCandidate {
data: data_usage_info,
#[cfg(test)]
last_update: merged_last_update,
coverage_digest: scope.identity.coverage_digest,
})
Some((data_usage_info, merged_last_update))
}
fn tier_accounting_proof_is_publishable(
+10 -28
View File
@@ -72,7 +72,6 @@ where
scan_mode,
scan_scope: ScannerBucketScanScope::default(),
persisted_usage_baseline: None,
observed_usage_candidate: None,
requires_full_scan: true,
service_cohort: None,
#[cfg(test)]
@@ -90,7 +89,6 @@ pub(crate) struct ScannerCycleRequest {
pub(crate) scan_mode: HealScanMode,
pub(crate) scan_scope: ScannerBucketScanScope,
pub(crate) persisted_usage_baseline: Option<Bytes>,
pub(crate) observed_usage_candidate: Option<Bytes>,
/// Scheduled maintenance must visit clean buckets even with a valid dirty scope.
pub(crate) requires_full_scan: bool,
pub(crate) service_cohort: Option<Arc<StdMutex<ScannerServiceCohort>>>,
@@ -187,7 +185,6 @@ where
scan_mode,
scan_scope,
persisted_usage_baseline,
observed_usage_candidate,
requires_full_scan,
service_cohort,
#[cfg(test)]
@@ -300,8 +297,7 @@ where
ScannerBucketScopeResolution {
requested_scope: scan_scope,
baseline_proof: ScannerCacheBaselineProof {
authoritative_data: persisted_usage_baseline.as_ref(),
observed_candidate_data: observed_usage_candidate.as_ref(),
data: persisted_usage_baseline.as_ref(),
expected_sources: &expected_sources,
leader_epoch,
want_cycle,
@@ -338,21 +334,12 @@ where
dirty_usage_status,
activity_status,
);
let Some(candidate) = empty_namespace_usage_candidate(
&all_buckets,
&expected_sources,
&buckets_by_source,
ScannerSnapshotIdentity {
cycle: want_cycle,
leader_epoch,
plan_digest: scan_plan_digest,
coverage_digest: bucket_coverage_digest,
tier_registry_generation: Some(tier_registry_generation),
},
) else {
return Ok(ScannerCycleResult::new(ScannerCycleStatus::Incomplete, None).with_publication_epoch(publication_epoch));
let empty_usage = DataUsageInfo {
last_update: Some(SystemTime::now()),
scanner_cycle: Some(want_cycle),
usage_snapshot_complete: true,
..Default::default()
};
let (empty_usage, publication_expectation) = candidate.prepare(status);
let observational_snapshot_published = if should_publish_observational_snapshot(status) {
publish_observational_snapshot(&updates, empty_usage).await?
} else {
@@ -375,8 +362,7 @@ where
.with_activity_digest(activity_digest)
.with_observational_snapshot_published(observational_snapshot_published)
.with_remote_publication_lease_targets(remote_publication_lease_targets)
.with_remote_dirty_usage_acknowledgements(remote_dirty_usage_acknowledgements)
.with_publication_expectation(publication_expectation));
.with_remote_dirty_usage_acknowledgements(remote_dirty_usage_acknowledgements));
}
let total_results = expected_sources.len();
@@ -605,7 +591,7 @@ where
let (activity_status, remote_publication_lease_targets) =
scanner_cycle_activity_status(store, distributed, &activity_before).await;
let all_bucket_names = all_buckets.iter().map(|bucket| bucket.name.clone()).collect::<Vec<_>>();
let completed_usage = completed_usage_candidate(
let completed_usage = completed_data_usage_info(
&results,
&ScannerSnapshotScope {
sources: &expected_sources,
@@ -646,10 +632,7 @@ where
dirty_usage_status,
activity_status,
);
let mut publication_expectation = None;
let observational_snapshot_published = if let Some(candidate) = completed_usage {
let (data_usage_info, expectation) = candidate.prepare(cycle_status);
publication_expectation = expectation;
let observational_snapshot_published = if let Some((data_usage_info, _)) = completed_usage {
if should_publish_observational_snapshot(cycle_status) {
publish_observational_snapshot(&updates, data_usage_info).await?
} else {
@@ -687,6 +670,5 @@ where
.with_remote_dirty_usage_acknowledgements(remote_dirty_usage_acknowledgements)
.with_failed_dirty_usage(!failed_buckets.is_empty())
.with_pending_maintenance_work(pending_maintenance_work)
.with_required_cycle_floor(required_cycle_floor)
.with_publication_expectation(publication_expectation))
.with_required_cycle_floor(required_cycle_floor))
}
+8 -123
View File
@@ -397,7 +397,6 @@ async fn scoped_scan_production_entry_preserves_deep_and_full_maintenance_work()
scan_mode,
scan_scope: requested_scope,
persisted_usage_baseline: baseline,
observed_usage_candidate: None,
requires_full_scan,
resolved_scope_observer: Some(observer),
service_cohort: None,
@@ -454,7 +453,6 @@ async fn scoped_scan_same_cycle_maintenance_rewalks_after_root_delivery_failure(
.await
.expect("initial object should persist");
}
wait_for_namespace_commit_tails(&store).await;
let ctx = CancellationToken::new();
let budget = ScannerCycleBudget::new(&ctx, ScannerCycleBudgetConfig::default());
let (updates, receiver) = mpsc::channel(1);
@@ -472,7 +470,6 @@ async fn scoped_scan_same_cycle_maintenance_rewalks_after_root_delivery_failure(
scan_mode: HealScanMode::Normal,
scan_scope: ScannerBucketScanScope::default(),
persisted_usage_baseline: None,
observed_usage_candidate: None,
requires_full_scan: false,
resolved_scope_observer: None,
service_cohort: None,
@@ -504,7 +501,6 @@ async fn scoped_scan_same_cycle_maintenance_rewalks_after_root_delivery_failure(
.put_object("cold-bucket", "new", &mut reader, &ScannerObjectOptions::default())
.await
.expect("new cold object should persist");
wait_for_namespace_commit_tails(&store).await;
record_dirty_usage_bucket("hot-bucket");
if scan_mode == HealScanMode::Normal && !requires_full_scan {
record_dirty_usage_bucket("cold-bucket");
@@ -525,7 +521,6 @@ async fn scoped_scan_same_cycle_maintenance_rewalks_after_root_delivery_failure(
scan_mode,
scan_scope: ScannerBucketScanScope::default(),
persisted_usage_baseline: None,
observed_usage_candidate: None,
requires_full_scan,
resolved_scope_observer: None,
service_cohort: None,
@@ -996,8 +991,8 @@ fn dirty_usage_snapshot_clears_a_stably_absent_bucket_after_durable_save() {
assert!(dirty_usage_buckets().contains_key("temporarily-omitted"));
assert_eq!(dirty_usage_snapshot_status(&snapshot), DirtyUsageSnapshotStatus::Current);
let acknowledgements =
ScannerCycleResult::new(ScannerCycleStatus::Complete, Some(snapshot.buckets.as_ref().clone())).clear_verified_usage();
let acknowledgements = ScannerCycleResult::new(ScannerCycleStatus::Complete, Some(snapshot.buckets.as_ref().clone()))
.acknowledge_durable_usage();
assert!(acknowledgements.is_empty());
assert!(!dirty_usage_buckets().contains_key("temporarily-omitted"));
clear_dirty_usage_buckets_for_tests();
@@ -1103,7 +1098,7 @@ fn dirty_usage_is_acknowledged_only_after_durable_usage_confirmation() {
assert!(dirty_usage_buckets().contains_key("photos"));
let confirmed = ScannerCycleResult::new(ScannerCycleStatus::Complete, Some(snapshot.buckets.as_ref().clone()));
let acknowledgements = confirmed.clear_verified_usage();
let acknowledgements = confirmed.acknowledge_durable_usage();
assert!(acknowledgements.is_empty());
assert!(!dirty_usage_buckets().contains_key("photos"));
clear_dirty_usage_buckets_for_tests();
@@ -1336,8 +1331,7 @@ fn scoped_scan_requires_a_converged_complete_baseline_with_exact_set_provenance(
assert_eq!(
complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof {
authoritative_data: Some(&baseline),
observed_candidate_data: None,
data: Some(&baseline),
expected_sources: &expected_sources,
leader_epoch: 11,
want_cycle: 8,
@@ -1351,8 +1345,7 @@ fn scoped_scan_requires_a_converged_complete_baseline_with_exact_set_provenance(
let incomplete = bytes::Bytes::from(serde_json::to_vec(&incomplete).expect("test baseline should encode"));
assert_eq!(
complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof {
authoritative_data: Some(&incomplete),
observed_candidate_data: None,
data: Some(&incomplete),
expected_sources: &expected_sources,
leader_epoch: 11,
want_cycle: 8,
@@ -1366,8 +1359,7 @@ fn scoped_scan_requires_a_converged_complete_baseline_with_exact_set_provenance(
let wrong_provenance = bytes::Bytes::from(serde_json::to_vec(&wrong_provenance).expect("test baseline should encode"));
assert_eq!(
complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof {
authoritative_data: Some(&wrong_provenance),
observed_candidate_data: None,
data: Some(&wrong_provenance),
expected_sources: &expected_sources,
leader_epoch: 11,
want_cycle: 8,
@@ -1377,111 +1369,6 @@ fn scoped_scan_requires_a_converged_complete_baseline_with_exact_set_provenance(
);
}
#[test]
fn scoped_scan_accepts_only_a_complete_observation_tied_to_the_authoritative_baseline() {
let source = DataUsageCacheSource::new(1, 2);
let expected_sources = HashSet::from([source]);
let scan_plan_digest = DataUsageScanPlanDigest([9; 32]);
let authoritative = complete_usage_baseline(source, scan_plan_digest, 7, 11);
let authoritative_info =
serde_json::from_slice::<DataUsageInfo>(&authoritative).expect("authoritative baseline should decode");
let bootstrap_authoritative_info =
crate::scanner::scanner_usage_bootstrap_marker(SystemTime::UNIX_EPOCH + Duration::from_secs(9), Some(11));
let bootstrap_authoritative =
bytes::Bytes::from(serde_json::to_vec(&bootstrap_authoritative_info).expect("bootstrap baseline should encode"));
let mut observed_info = authoritative_info.clone();
observed_info.last_update = Some(SystemTime::UNIX_EPOCH + Duration::from_secs(11));
observed_info.scanner_cycle = Some(8);
observed_info.usage_snapshot_converged = Some(false);
observed_info.usage_snapshot_authoritative_baseline = Some(bootstrap_authoritative_info.snapshot_identity());
observed_info.usage_snapshot_set_states[0].scanner_cycle = Some(8);
let observed = bytes::Bytes::from(serde_json::to_vec(&observed_info).expect("observation should encode"));
macro_rules! proof {
($authoritative:expr, $candidate:expr) => {
ScannerCacheBaselineProof {
authoritative_data: Some($authoritative),
observed_candidate_data: $candidate,
expected_sources: &expected_sources,
leader_epoch: 11,
want_cycle: 9,
scan_plan_digest,
}
};
}
assert_eq!(
complete_scanner_cache_baseline_plan_digest(proof!(&bootstrap_authoritative, Some(&observed))),
Some(scan_plan_digest)
);
observed_info.usage_snapshot_authoritative_baseline = Some(DataUsageInfo::default().snapshot_identity());
let mismatched_baseline = bytes::Bytes::from(serde_json::to_vec(&observed_info).expect("observation should encode"));
assert_eq!(
complete_scanner_cache_baseline_plan_digest(proof!(&bootstrap_authoritative, Some(&mismatched_baseline))),
None
);
observed_info = serde_json::from_slice(&observed).expect("observation should decode");
observed_info.usage_snapshot_partial = true;
let partial = bytes::Bytes::from(serde_json::to_vec(&observed_info).expect("partial observation should encode"));
assert_eq!(
complete_scanner_cache_baseline_plan_digest(proof!(&bootstrap_authoritative, Some(&partial))),
None
);
observed_info = serde_json::from_slice(&observed).expect("observation should decode");
observed_info.usage_snapshot_converged = Some(true);
let converged = bytes::Bytes::from(serde_json::to_vec(&observed_info).expect("converged observation should encode"));
assert_eq!(
complete_scanner_cache_baseline_plan_digest(proof!(&bootstrap_authoritative, Some(&converged))),
None
);
let mut legacy_authoritative = authoritative_info.clone();
legacy_authoritative.usage_snapshot_converged = None;
let mut stale_info = serde_json::from_slice::<DataUsageInfo>(&observed).expect("observation should decode");
stale_info.scanner_cycle = Some(7);
stale_info.usage_snapshot_set_states[0].scanner_cycle = Some(7);
stale_info.usage_snapshot_authoritative_baseline = Some(legacy_authoritative.snapshot_identity());
let legacy_authoritative =
bytes::Bytes::from(serde_json::to_vec(&legacy_authoritative).expect("legacy baseline should encode"));
let stale = bytes::Bytes::from(serde_json::to_vec(&stale_info).expect("stale observation should encode"));
assert_eq!(
complete_scanner_cache_baseline_plan_digest(proof!(&legacy_authoritative, Some(&stale))),
None
);
let malformed = bytes::Bytes::from_static(b"not data usage json");
assert_eq!(
complete_scanner_cache_baseline_plan_digest(proof!(&bootstrap_authoritative, Some(&malformed))),
None
);
let mut nonconverged_authoritative = authoritative_info;
nonconverged_authoritative.usage_snapshot_converged = Some(false);
let mut observation_of_nonconverged_authoritative = nonconverged_authoritative.clone();
observation_of_nonconverged_authoritative.last_update = Some(SystemTime::UNIX_EPOCH + Duration::from_secs(12));
observation_of_nonconverged_authoritative.scanner_cycle = Some(8);
observation_of_nonconverged_authoritative.usage_snapshot_set_states[0].scanner_cycle = Some(8);
observation_of_nonconverged_authoritative.usage_snapshot_authoritative_baseline =
Some(nonconverged_authoritative.snapshot_identity());
let nonconverged_authoritative =
bytes::Bytes::from(serde_json::to_vec(&nonconverged_authoritative).expect("nonconverged baseline should encode"));
let observation_of_nonconverged_authoritative =
bytes::Bytes::from(serde_json::to_vec(&observation_of_nonconverged_authoritative).expect("observation should encode"));
assert_eq!(
complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof {
authoritative_data: Some(&nonconverged_authoritative),
observed_candidate_data: Some(&observation_of_nonconverged_authoritative),
expected_sources: &expected_sources,
leader_epoch: 11,
want_cycle: 9,
scan_plan_digest,
}),
None
);
}
#[test]
fn scoped_scan_selects_only_current_dirty_buckets_after_baseline_validation() {
let source = DataUsageCacheSource::new(1, 2);
@@ -1495,8 +1382,7 @@ fn scoped_scan_selects_only_current_dirty_buckets_after_baseline_validation() {
true,
&[bucket_info("photos")],
ScannerCacheBaselineProof {
authoritative_data: Some(&baseline),
observed_candidate_data: None,
data: Some(&baseline),
expected_sources: &expected_sources,
leader_epoch: 11,
want_cycle: 8,
@@ -1534,8 +1420,7 @@ fn scoped_scan_baseline_work_proof_requires_uniform_known_set_identity() {
let data = Bytes::from(serde_json::to_vec(&candidate).expect("candidate should encode"));
assert_eq!(
complete_scanner_cache_baseline_plan_digest(ScannerCacheBaselineProof {
authoritative_data: Some(&data),
observed_candidate_data: None,
data: Some(&data),
expected_sources: &sources,
leader_epoch: 11,
want_cycle: 8,
@@ -126,7 +126,6 @@ async fn run_entry(store: &Arc<ECStore>, cycle: u64, selected: Option<&str>, exp
scan_mode: HealScanMode::Normal,
scan_scope: ScannerBucketScanScope::default(),
persisted_usage_baseline: root_before.0.clone().map(Bytes::from),
observed_usage_candidate: None,
requires_full_scan: false,
service_cohort: None,
resolved_scope_observer: Some(observer),
@@ -61,7 +61,6 @@ async fn run_cohort_cycle(
scan_mode: HealScanMode::Normal,
scan_scope: ScannerBucketScanScope::default(),
persisted_usage_baseline: None,
observed_usage_candidate: None,
requires_full_scan: false,
service_cohort: Some(cohort),
resolved_scope_observer: None,
+49 -4
View File
@@ -63,6 +63,7 @@ const EVENT_ADMIN_REQUEST_STATE: &str = "admin_request_state";
const EVENT_ADMIN_REQUEST_REJECTED: &str = "admin_request_rejected";
const EVENT_ADMIN_REQUEST_FAILED: &str = "admin_request_failed";
const EVENT_ADMIN_RESPONSE_EMITTED: &str = "admin_response_emitted";
const POOL_ACTIVATION_FLEET_PROOF_REQUIRED: &str = "pool activation requires a live fleet capability proof";
fn admin_request_id(headers: &HeaderMap) -> Option<&str> {
headers
@@ -321,6 +322,17 @@ fn contextualize_admin_pool_api_error(
}
}
fn decommission_start_api_error(err: crate::storage_api::error::StorageError) -> ApiError {
if crate::storage_api::capacity::is_pool_activation_fleet_proof_error(&err) {
return ApiError {
code: S3ErrorCode::InternalError,
message: POOL_ACTIVATION_FLEET_PROOF_REQUIRED.to_string(),
source: Some(Box::new(err)),
};
}
ApiError::from(err)
}
fn decommission_admin_not_initialized_error_with_audit(operation: &str, audit: PoolAuditContext<'_>) -> S3Error {
error!(
event = EVENT_ADMIN_REQUEST_FAILED,
@@ -790,7 +802,24 @@ impl Operation for StartDecommission {
store
.decommission(ctx.clone(), pools_indices.clone())
.await
.map_err(ApiError::from)
.map_err(|err| {
error!(
event = EVENT_ADMIN_REQUEST_FAILED,
component = LOG_COMPONENT_ADMIN_API,
subsystem = LOG_SUBSYSTEM_POOL_ADMIN,
operation = "start_decommission",
action = "start_decommission",
result = "failed",
reason = "storage_decommission_failed",
request_id = %request_id,
actor = %actor,
remote_addr = %remote_addr,
pool_indices = ?pools_indices,
error = %err,
"admin request failed"
);
decommission_start_api_error(err)
})
.map_err(|err| contextualize_admin_pool_api_error(err, "start decommission", &pool_context))?;
}
}
@@ -1018,9 +1047,10 @@ impl Operation for ClearDecommission {
#[cfg(test)]
mod pools_handler_tests {
use super::{
AdminPoolStatus, Body, CancelDecommission, ClearDecommission, HeaderMap, ListPools, Method, Operation, Params,
PoolAuditContext, S3ErrorCode, S3Request, StartDecommission, StatusDecommission, StatusPool, Uri,
contextualize_admin_pool_api_error, decommission_admin_not_initialized_error_with_audit, decommission_peer_target,
AdminPoolStatus, Body, CancelDecommission, ClearDecommission, HeaderMap, ListPools, Method, Operation,
POOL_ACTIVATION_FLEET_PROOF_REQUIRED, Params, PoolAuditContext, S3ErrorCode, S3Request, StartDecommission,
StatusDecommission, StatusPool, Uri, contextualize_admin_pool_api_error,
decommission_admin_not_initialized_error_with_audit, decommission_peer_target, decommission_start_api_error,
has_duplicate_indices, parse_mutation_pool_query, parse_pool_idx_by_id, parse_status_pool_query,
pool_admin_missing_credentials_error, pool_admin_missing_credentials_error_with_request,
pool_admin_pool_index_error_with_audit, pool_admin_pool_not_found_error_with_audit,
@@ -1209,6 +1239,21 @@ mod pools_handler_tests {
);
}
#[test]
fn test_decommission_start_api_error_preserves_fleet_proof_retry_marker() {
let err = crate::storage_api::error::StorageError::other(POOL_ACTIVATION_FLEET_PROOF_REQUIRED);
let err = decommission_start_api_error(err);
assert_eq!(err.code, s3s::S3ErrorCode::InternalError);
assert_eq!(err.message, POOL_ACTIVATION_FLEET_PROOF_REQUIRED);
assert!(err.source.is_some());
let unrelated = decommission_start_api_error(crate::storage_api::error::StorageError::other("disk read failed"));
assert_eq!(unrelated.code, s3s::S3ErrorCode::InternalError);
assert_eq!(unrelated.message, "We encountered an internal error, please try again.");
}
#[test]
fn test_contextualize_admin_pool_api_error_preserves_source() {
let err = contextualize_admin_pool_api_error(
+1 -1
View File
@@ -417,7 +417,7 @@ pub(crate) mod ecstore_bucket {
pub(crate) mod ecstore_capacity {
pub(crate) use rustfs_ecstore::api::capacity::{
DecommissionUnresolvedEntry, PoolDecommissionInfo, PoolStatus, get_total_usable_capacity, get_total_usable_capacity_free,
is_reserved_or_invalid_bucket,
is_pool_activation_fleet_proof_error, is_reserved_or_invalid_bucket,
};
}
+2
View File
@@ -18,6 +18,8 @@
use rustfs_storage_api as storage_contracts;
pub(crate) mod capacity {
pub(crate) use crate::storage::storage_api::ecstore_capacity::is_pool_activation_fleet_proof_error;
pub(crate) mod service {
pub(crate) use crate::storage::storage_api::{all_local_disk, disk_drive_path, disk_endpoint};
}
@@ -1,8 +1,8 @@
{
"bytes": "eyJjaGFsbGVuZ2VJZCI6IjAxOGY3ZTZkLTlkNmEtN2Q5My04ZjY0LThiMjBiMzM4NDcxMiIsImNsdXN0ZXJOYW1lIjoib3JnYW5pemF0aW9ucy8wMThmN2U2ZC05ZDZhLTdkOTMtOGY2NC04YjIwYjMzODQ3MTMvY2x1c3RlcnMvMDE4ZjdlNmQtOWQ2YS03ZDkzLThmNjQtOGIyMGIzMzg0NzE0IiwiY29ubmVjdEtleUlkIjoiZGY4MmNlNjA4MGIwZWZkMDhkNGFlNjNhMjI0NGZmNmNiZWE5YTU0NzUyNjhjZGFkZTMzNzEyMmVjYzY3ZjJhNiIsImV4cGlyZXNBdCI6IjIwOTktMDEtMDhUMDA6MDA6MDBaIiwiZm9ybWF0VmVyc2lvbiI6InJ1c3Rmcy5jb25uZWN0Lm9mZmxpbmUuZW5yb2xsbWVudENoYWxsZW5nZS8xIiwiaXNzdWVkQXQiOiIyMDk5LTAxLTAxVDAwOjAwOjAwWiIsIm5vbmNlIjoicEtTa3BLU2twS1NrcEtTa3BLU2twS1NrcEtTa3BLU2twS1NrcEtTa3BLUSIsIm9yZ2FuaXphdGlvbk5hbWUiOiJvcmdhbml6YXRpb25zLzAxOGY3ZTZkLTlkNmEtN2Q5My04ZjY0LThiMjBiMzM4NDcxMyIsInByb3RvY29sVmVyc2lvbiI6InYxIiwidHJ1c3RDaGFpbiI6W3siYnl0ZXMiOiJleUptYjNKdFlYUldaWEp6YVc5dUlqb2ljblZ6ZEdaekxtTnZibTVsWTNRdWIyWm1iR2x1WlM1MGNuVnpkRXhwYm1zdk1TSXNJbkJ5YjNSdlkyOXNWbVZ5YzJsdmJpSTZJbll4SWl3aWMyVnlhV0ZzSWpvaVpUQTBZVEF3TURBd01EQXdNREF3TURBd01EQXdNREF3TURBd01EQXdNREVpTENKeWIyeGxJam9pYVc1MFpYSnRaV1JwWVhSbElpd2lhWE56ZFdWeVMyVjVTV1FpT2lKbU9EVTNZV1UwWldRMlpUTTBaVE16TnpZeFlXVmhNalZqWVdGbFpUTm1aVFUwWVRFMU9UWXdabUk1TW1SalpEWXpZVE0zTldGaU1USXhaR1ZpTW1FNUlpd2ljM1ZpYW1WamRFdGxlVWxrSWpvaU5XRXdOMkUzTWpOak56azFOamd6TUdGbE1tVTJZbVV3WldZM01XSXpPVGt4TVRZeVpqVmxZVFF5WkRKa05UUmlPRFk1WTJZeU1UWTNNelJrTnpoa015SXNJbk4xWW1wbFkzUlFkV0pzYVdOTFpYa2lPaUpDUmxWUVVuaEJSRGc1TFZoM09UbFJZWE5sV0RsdVNXWnpZVWczWlRRNWRtYzVTV3RUV1hCc2VVazBhMFV5UTFReGQwVjFWVXB3ZW1OV2VUbERkME5xZWtGZk1IUmpRV0pRWDI5YVlYSklOMDF1UVRKMVQxa2lMQ0p1YjNSQ1pXWnZjbVVpT2lJeU1EazRMVEV5TFRBeFZEQXdPakF3T2pBd1dpSXNJbTV2ZEVGbWRHVnlJam9pTWpBNU9TMHhNUzB6TUZRd01Eb3dNRG93TUZvaWZRPT0iLCJzaWduYXR1cmUiOnsiYWxnb3JpdGhtIjoiRVMyNTYiLCJrZXlJZCI6ImY4NTdhZTRlZDZlMzRlMzM3NjFhZWEyNWNhYWVlM2ZlNTRhMTU5NjBmYjkyZGNkNjNhMzc1YWIxMjFkZWIyYTkiLCJ2YWx1ZSI6IlhzbjVaRkI4OHU1bHg1TDU5UDJCWlNLMWNrTVc5ekpiNGFYUmNfdjB5cFIzd1pxc3BlTzJCNzZyV1V0RnhDRE9Vek9GdHFDZlFMRi1McWNsLWE5dFhRIn19LHsiYnl0ZXMiOiJleUptYjNKdFlYUldaWEp6YVc5dUlqb2ljblZ6ZEdaekxtTnZibTVsWTNRdWIyWm1iR2x1WlM1MGNuVnpkRXhwYm1zdk1TSXNJbkJ5YjNSdlkyOXNWbVZ5YzJsdmJpSTZJbll4SWl3aWMyVnlhV0ZzSWpvaVpUQTBZVEF3TURBd01EQXdNREF3TURBd01EQXdNREF3TURBd01EQXdNRElpTENKeWIyeGxJam9pYzJsbmJtbHVaeUlzSW1semMzVmxja3RsZVVsa0lqb2lOV0V3TjJFM01qTmpOemsxTmpnek1HRmxNbVUyWW1Vd1pXWTNNV0l6T1RreE1UWXlaalZsWVRReVpESmtOVFJpT0RZNVkyWXlNVFkzTXpSa056aGtNeUlzSW5OMVltcGxZM1JMWlhsSlpDSTZJbVJtT0RKalpUWXdPREJpTUdWbVpEQTRaRFJoWlRZellUSXlORFJtWmpaalltVmhPV0UxTkRjMU1qWTRZMlJoWkdVek16Y3hNakpsWTJNMk4yWXlZVFlpTENKemRXSnFaV04wVUhWaWJHbGpTMlY1SWpvaVFrWnJZWFF6U0hKMlVERjBia3hyU2xSU1FteExTek5TY0hBeFJYZHpTREpLWDBOS04wWnBOWGhvY21adU1EVnhkbmN3UlZoQmVIQlBhbmh2Y2xoNVdIbHVTeTFhVGpjd2IyMWZjekJ0VUdSdFMydHVaMUJCSWl3aWJtOTBRbVZtYjNKbElqb2lNakE1T0MweE1pMHhOVlF3TURvd01Eb3dNRm9pTENKdWIzUkJablJsY2lJNklqSXdPVGt0TURFdE1UVlVNREE2TURBNk1EQmFJbjA9Iiwic2lnbmF0dXJlIjp7ImFsZ29yaXRobSI6IkVTMjU2Iiwia2V5SWQiOiI1YTA3YTcyM2M3OTU2ODMwYWUyZTZiZTBlZjcxYjM5OTExNjJmNWVhNDJkMmQ1NGI4NjljZjIxNjczNGQ3OGQzIiwidmFsdWUiOiJSUGJHWEt4RVRRRDZKVklEVmcwNkRRaENCQ0NfYjJkay1oZTdVUjZOeXJ0bVlQRUEzb3RhM0hXd25YdUpnWG84TlVWbnF2Q2ZldkNQNGFFelVSbEdaQSJ9fV19",
"signature": {
"algorithm": "ES256",
"keyId": "df82ce6080b0efd08d4ae63a2244ff6cbea9a5475268cdade337122ecc67f2a6",
"value": "m83S9i09JTlxGnhAdnmawYxkM75u_v8F2-5ygK374fs-d0vWh6L_j_yMDS3B29CrViPVjo3FQKaGutw4MDI_LQ"
}
"bytes": "eyJmb3JtYXRWZXJzaW9uIjoicnVzdGZzLmNvbm5lY3Qub2ZmbGluZS5lbnJvbGxtZW50Q2hhbGxlbmdlLzEiLCJwcm90b2NvbFZlcnNpb24iOiJ2MSIsImNoYWxsZW5nZUlkIjoiMDE4ZjdlNmQtOWQ2YS03ZDkzLThmNjQtOGIyMGIzMzg0NzEyIiwib3JnYW5pemF0aW9uTmFtZSI6Im9yZ2FuaXphdGlvbnMvMDE4ZjdlNmQtOWQ2YS03ZDkzLThmNjQtOGIyMGIzMzg0NzEzIiwiY2x1c3Rlck5hbWUiOiJvcmdhbml6YXRpb25zLzAxOGY3ZTZkLTlkNmEtN2Q5My04ZjY0LThiMjBiMzM4NDcxMy9jbHVzdGVycy8wMThmN2U2ZC05ZDZhLTdkOTMtOGY2NC04YjIwYjMzODQ3MTQiLCJub25jZSI6InBLU2twS1NrcEtTa3BLU2twS1NrcEtTa3BLU2twS1NrcEtTa3BLU2twS1EiLCJpc3N1ZWRBdCI6IjIwOTktMDEtMDFUMDA6MDA6MDBaIiwiZXhwaXJlc0F0IjoiMjA5OS0wMS0wOFQwMDowMDowMFoiLCJjb25uZWN0S2V5SWQiOiIwNjA3NWY0ODRkY2U1YzA3NDVjYTc2NWE4ODk5MzA5YmQ3MTc2MTI5NDdhMzc3MGEwZDI3YjhjY2I4YmY5MDUyIiwidHJ1c3RDaGFpbiI6W3siYnl0ZXMiOiJleUptYjNKdFlYUldaWEp6YVc5dUlqb2ljblZ6ZEdaekxtTnZibTVsWTNRdWIyWm1iR2x1WlM1MGNuVnpkRXhwYm1zdk1TSXNJbkJ5YjNSdlkyOXNWbVZ5YzJsdmJpSTZJbll4SWl3aWMyVnlhV0ZzSWpvaVpUQTBZVEF3TURBd01EQXdNREF3TURBd01EQXdNREF3TURBd01EQXdNREVpTENKeWIyeGxJam9pYVc1MFpYSnRaV1JwWVhSbElpd2lhWE56ZFdWeVMyVjVTV1FpT2lJNVlUVmpNemcxWVdKa09UZGpOMlppT0RSak5UUXdZalZtTldVME1tUTNPVGt4T1RNd1lUUTBOMk01TURJMk56azVaVGhtWkRSak5ESXpNR0kxWXpOaklpd2ljM1ZpYW1WamRFdGxlVWxrSWpvaVlUazFOVGxpWkRSbE1EazBPV0ZqWXpFeU1tRm1abU0xTldVMlpHSTRaVGt4Wmpoa01EazNaVFUxWldNNE9UbGtZVFl4T1RGbE5EWTVZelE0T0RKaE1pSXNJbk4xWW1wbFkzUlFkV0pzYVdOTFpYa2lPaUpDVG13elZFeDFTazV0Tld4NFRGUmFNRGxUWVdwWVRHMUZSM0ZqZHpGaVQyZzROVWhEYmpWdVVUazJWVlZUYUZVNFZHeEhVVWsxYmxaSldWUnJRWFpLYjFWRVMwWmhVVlp0WVZGVk9WRkhaemsyVDJ0RU1EUWlMQ0p1YjNSQ1pXWnZjbVVpT2lJeU1ESXdMVEF4TFRBeFZEQXdPakF3T2pBd1dpSXNJbTV2ZEVGbWRHVnlJam9pTWpBNU9TMHdNUzB3TVZRd01Eb3dNRG93TUZvaWZRPT0iLCJzaWduYXR1cmUiOnsiYWxnb3JpdGhtIjoiRVMyNTYiLCJrZXlJZCI6IjlhNWMzODVhYmQ5N2M3ZmI4NGM1NDBiNWY1ZTQyZDc5OTE5MzBhNDQ3YzkwMjY3OTllOGZkNGM0MjMwYjVjM2MiLCJ2YWx1ZSI6IklUSThKVGFURjRoRHJ0bF9tSk0zYmZ0eUlhMlJaSjFNNVpPVk1xaEE2cVJSSzFNcEcxQlF0emRRVzB6bHo4eG1RZUJNWl84clViTUlFcnFBel8tSWZRIn19LHsiYnl0ZXMiOiJleUptYjNKdFlYUldaWEp6YVc5dUlqb2ljblZ6ZEdaekxtTnZibTVsWTNRdWIyWm1iR2x1WlM1MGNuVnpkRXhwYm1zdk1TSXNJbkJ5YjNSdlkyOXNWbVZ5YzJsdmJpSTZJbll4SWl3aWMyVnlhV0ZzSWpvaVpUQTBZVEF3TURBd01EQXdNREF3TURBd01EQXdNREF3TURBd01EQXdNRElpTENKeWIyeGxJam9pYzJsbmJtbHVaeUlzSW1semMzVmxja3RsZVVsa0lqb2lZVGsxTlRsaVpEUmxNRGswT1dGall6RXlNbUZtWm1NMU5XVTJaR0k0WlRreFpqaGtNRGszWlRVMVpXTTRPVGxrWVRZeE9URmxORFk1WXpRNE9ESmhNaUlzSW5OMVltcGxZM1JMWlhsSlpDSTZJakEyTURjMVpqUTROR1JqWlRWak1EYzBOV05oTnpZMVlUZzRPVGt6TURsaVpEY3hOell4TWprME4yRXpOemN3WVRCa01qZGlPR05qWWpoaVpqa3dOVElpTENKemRXSnFaV04wVUhWaWJHbGpTMlY1SWpvaVFrRk9lWFpITFdkR1YwdFRaM0pMWTNkUVpGUmhNRWN3TTBKdlVtRkNSbWhvZG10cWRFUjZlVk5tYm5weFYxcHROa2xLV0U1RWRVZEdlVXBFZVZWSmNYZzJSa1p5TlZSV2NsOXVhbUZCYzFWb2FEbE1VMlJ2SWl3aWJtOTBRbVZtYjNKbElqb2lNakF5TUMwd01TMHdNVlF3TURvd01Eb3dNRm9pTENKdWIzUkJablJsY2lJNklqSXdPVGt0TURFdE1ERlVNREE2TURBNk1EQmFJbjA9Iiwic2lnbmF0dXJlIjp7ImFsZ29yaXRobSI6IkVTMjU2Iiwia2V5SWQiOiJhOTU1OWJkNGUwOTQ5YWNjMTIyYWZmYzU1ZTZkYjhlOTFmOGQwOTdlNTVlYzg5OWRhNjE5MWU0NjljNDg4MmEyIiwidmFsdWUiOiI4MTVrVGNMY1hkZXZORFRia1dBck1KSnpORTVmaEdxbTNUb3hTM3Exdl9ZQkRVcmFieldmTFhjMEc5TjZiVXZFeEd3WEJUMVFpVzJvUDBCM3JKUS0ydyJ9fV19",
"signature": {
"algorithm": "ES256",
"keyId": "06075f484dce5c0745ca765a8899309bd717612947a3770a0d27b8ccb8bf9052",
"value": "8i2nvYhNsZbL1UPF82Otr4pcecYfKhdRV1Q7m3Geqcx-TdvYk486Q_cZHSQ7OaqnIhTZwPgMYFYKUuEhwAaYEQ"
}
}
+2 -2
View File
@@ -1,4 +1,4 @@
{
"keyId": "f857ae4ed6e34e33761aea25caaee3fe54a15960fb92dcd63a375ab121deb2a9",
"publicKey": "BG_wO5SSQc4drdQ1GeaWDgqFtBppoFwygQOqK84VlMoWPE91OlW_AdxT9sCwx-7ni0DG_30lqW4igrmJzvccFEo"
"keyId": "9a5c385abd97c7fb84c540b5f5e42d7991930a447c9026799e8fd4c4230b5c3c",
"publicKey": "BIPyiA1W2NDuy3ftHLtlY2tO2WNjkjGINZ_9QQvyyN9syzbMb91QYG2SN0AqZPylFsTL-loF4M1tVySZtXJpKhM"
}
@@ -1 +1 @@
[{"bytes":"eyJmb3JtYXRWZXJzaW9uIjoicnVzdGZzLmNvbm5lY3Qub2ZmbGluZS50cnVzdExpbmsvMSIsInByb3RvY29sVmVyc2lvbiI6InYxIiwic2VyaWFsIjoiZTA0YTAwMDAwMDAwMDAwMDAwMDAwMDAwMDAwMDAwMDEiLCJyb2xlIjoiaW50ZXJtZWRpYXRlIiwiaXNzdWVyS2V5SWQiOiJmODU3YWU0ZWQ2ZTM0ZTMzNzYxYWVhMjVjYWFlZTNmZTU0YTE1OTYwZmI5MmRjZDYzYTM3NWFiMTIxZGViMmE5Iiwic3ViamVjdEtleUlkIjoiNWEwN2E3MjNjNzk1NjgzMGFlMmU2YmUwZWY3MWIzOTkxMTYyZjVlYTQyZDJkNTRiODY5Y2YyMTY3MzRkNzhkMyIsInN1YmplY3RQdWJsaWNLZXkiOiJCRlVQUnhBRDg5LVh3OTlRYXNlWDluSWZzYUg3ZTQ5dmc5SWtTWXBseUk0a0UyQ1Qxd0V1VUpwemNWeTlDd0NqekFfMHRjQWJQX29aYXJIN01uQTJ1T1kiLCJub3RCZWZvcmUiOiIyMDk4LTEyLTAxVDAwOjAwOjAwWiIsIm5vdEFmdGVyIjoiMjA5OS0xMS0zMFQwMDowMDowMFoifQ==","signature":{"algorithm":"ES256","keyId":"f857ae4ed6e34e33761aea25caaee3fe54a15960fb92dcd63a375ab121deb2a9","value":"Xsn5ZFB88u5lx5L59P2BZSK1ckMW9zJb4aXRc_v0ypR3wZqspeO2B76rWUtFxCDOUzOFtqCfQLF-Lqcl-a9tXQ"}},{"bytes":"eyJmb3JtYXRWZXJzaW9uIjoicnVzdGZzLmNvbm5lY3Qub2ZmbGluZS50cnVzdExpbmsvMSIsInByb3RvY29sVmVyc2lvbiI6InYxIiwic2VyaWFsIjoiZTA0YTAwMDAwMDAwMDAwMDAwMDAwMDAwMDAwMDAwMDIiLCJyb2xlIjoic2lnbmluZyIsImlzc3VlcktleUlkIjoiNWEwN2E3MjNjNzk1NjgzMGFlMmU2YmUwZWY3MWIzOTkxMTYyZjVlYTQyZDJkNTRiODY5Y2YyMTY3MzRkNzhkMyIsInN1YmplY3RLZXlJZCI6ImRmODJjZTYwODBiMGVmZDA4ZDRhZTYzYTIyNDRmZjZjYmVhOWE1NDc1MjY4Y2RhZGUzMzcxMjJlY2M2N2YyYTYiLCJzdWJqZWN0UHVibGljS2V5IjoiQkZrYXQzSHJ2UDF0bkxrSlRSQmxLSzNScHAxRXdzSDJKX0NKN0ZpNXhocmZuMDVxdncwRVhBeHBPanhvclh5WHluSy1aTjcwb21fczBtUGRtS2tuZ1BBIiwibm90QmVmb3JlIjoiMjA5OC0xMi0xNVQwMDowMDowMFoiLCJub3RBZnRlciI6IjIwOTktMDEtMTVUMDA6MDA6MDBaIn0=","signature":{"algorithm":"ES256","keyId":"5a07a723c7956830ae2e6be0ef71b3991162f5ea42d2d54b869cf216734d78d3","value":"RPbGXKxETQD6JVIDVg06DQhCBCC_b2dk-he7UR6NyrtmYPEA3ota3HWwnXuJgXo8NUVnqvCfevCP4aEzURlGZA"}}]
[{"bytes":"eyJmb3JtYXRWZXJzaW9uIjoicnVzdGZzLmNvbm5lY3Qub2ZmbGluZS50cnVzdExpbmsvMSIsInByb3RvY29sVmVyc2lvbiI6InYxIiwic2VyaWFsIjoiZTA0YTAwMDAwMDAwMDAwMDAwMDAwMDAwMDAwMDAwMDEiLCJyb2xlIjoiaW50ZXJtZWRpYXRlIiwiaXNzdWVyS2V5SWQiOiI5YTVjMzg1YWJkOTdjN2ZiODRjNTQwYjVmNWU0MmQ3OTkxOTMwYTQ0N2M5MDI2Nzk5ZThmZDRjNDIzMGI1YzNjIiwic3ViamVjdEtleUlkIjoiYTk1NTliZDRlMDk0OWFjYzEyMmFmZmM1NWU2ZGI4ZTkxZjhkMDk3ZTU1ZWM4OTlkYTYxOTFlNDY5YzQ4ODJhMiIsInN1YmplY3RQdWJsaWNLZXkiOiJCTmwzVEx1Sk5tNWx4TFRaMDlTYWpYTG1FR3FjdzFiT2g4NUhDbjVuUTk2VVVTaFU4VGxHUUk1blZJWVRrQXZKb1VES0ZhUVZtYVFVOVFHZzk2T2tEMDQiLCJub3RCZWZvcmUiOiIyMDIwLTAxLTAxVDAwOjAwOjAwWiIsIm5vdEFmdGVyIjoiMjA5OS0wMS0wMVQwMDowMDowMFoifQ==","signature":{"algorithm":"ES256","keyId":"9a5c385abd97c7fb84c540b5f5e42d7991930a447c9026799e8fd4c4230b5c3c","value":"ITI8JTaTF4hDrtl_mJM3bftyIa2RZJ1M5ZOVMqhA6qRRK1MpG1BQtzdQW0zlz8xmQeBMZ_8rUbMIErqAz_-IfQ"}},{"bytes":"eyJmb3JtYXRWZXJzaW9uIjoicnVzdGZzLmNvbm5lY3Qub2ZmbGluZS50cnVzdExpbmsvMSIsInByb3RvY29sVmVyc2lvbiI6InYxIiwic2VyaWFsIjoiZTA0YTAwMDAwMDAwMDAwMDAwMDAwMDAwMDAwMDAwMDIiLCJyb2xlIjoic2lnbmluZyIsImlzc3VlcktleUlkIjoiYTk1NTliZDRlMDk0OWFjYzEyMmFmZmM1NWU2ZGI4ZTkxZjhkMDk3ZTU1ZWM4OTlkYTYxOTFlNDY5YzQ4ODJhMiIsInN1YmplY3RLZXlJZCI6IjA2MDc1ZjQ4NGRjZTVjMDc0NWNhNzY1YTg4OTkzMDliZDcxNzYxMjk0N2EzNzcwYTBkMjdiOGNjYjhiZjkwNTIiLCJzdWJqZWN0UHVibGljS2V5IjoiQkFOeXZHLWdGV0tTZ3JLY3dQZFRhMEcwM0JvUmFCRmhodmtqdER6eVNmbnpxV1ptNklKWE5EdUdGeUpEeVVJcXg2RkZyNVRWcl9uamFBc1VoaDlMU2RvIiwibm90QmVmb3JlIjoiMjAyMC0wMS0wMVQwMDowMDowMFoiLCJub3RBZnRlciI6IjIwOTktMDEtMDFUMDA6MDA6MDBaIn0=","signature":{"algorithm":"ES256","keyId":"a9559bd4e0949acc122affc55e6db8e91f8d097e55ec899da6191e469c4882a2","value":"815kTcLcXdevNDTbkWArMJJzNE5fhGqm3ToxS3q1v_YBDUrabzWfLXc0G9N6bUvExGwXBT1QiW2oP0B3rJQ-2w"}}]
+4 -12
View File
@@ -1,6 +1,6 @@
module rustfs.local/heal-outcome-compat
go 1.25.0
go 1.24.0
require github.com/minio/madmin-go/v3 v3.0.107-0.20250415152934-4b504b82db63
@@ -8,35 +8,27 @@ require (
github.com/cespare/xxhash/v2 v2.3.0 // indirect
github.com/dustin/go-humanize v1.0.1 // indirect
github.com/go-ini/ini v1.67.0 // indirect
github.com/go-ole/go-ole v1.3.0 // indirect
github.com/goccy/go-json v0.10.5 // indirect
github.com/golang-jwt/jwt/v4 v4.5.2 // indirect
github.com/golang/protobuf v1.5.4 // indirect
github.com/klauspost/cpuid/v2 v2.2.10 // indirect
github.com/lufia/plan9stats v0.0.0-20250317134145-8bc96cf8fc35 // indirect
github.com/matttproud/golang_protobuf_extensions v1.0.4 // indirect
github.com/minio/md5-simd v1.1.2 // indirect
github.com/minio/minio-go/v7 v7.0.90 // indirect
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
github.com/philhofer/fwd v1.1.3-0.20240916144458-20a13a1f6b7c // indirect
github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 // indirect
github.com/prometheus/client_model v0.6.2 // indirect
github.com/prometheus/common v0.63.0 // indirect
github.com/prometheus/procfs v0.16.0 // indirect
github.com/prometheus/prom2json v1.4.2 // indirect
github.com/prometheus/prometheus v0.303.0 // indirect
github.com/rs/xid v1.6.0 // indirect
github.com/safchain/ethtool v0.5.10 // indirect
github.com/secure-io/sio-go v0.3.1 // indirect
github.com/shirou/gopsutil/v3 v3.24.5 // indirect
github.com/shoenig/go-m1cpu v0.1.6 // indirect
github.com/tinylib/msgp v1.2.5 // indirect
github.com/tklauser/go-sysconf v0.3.15 // indirect
github.com/tklauser/numcpus v0.10.0 // indirect
github.com/yusufpapurcu/wmi v1.2.4 // indirect
golang.org/x/crypto v0.52.0 // indirect
golang.org/x/net v0.55.0 // indirect
golang.org/x/sync v0.13.0 // indirect
golang.org/x/sys v0.45.0 // indirect
golang.org/x/crypto v0.37.0 // indirect
golang.org/x/net v0.39.0 // indirect
golang.org/x/sys v0.32.0 // indirect
google.golang.org/protobuf v1.36.6 // indirect
)
+6 -39
View File
@@ -1,14 +1,9 @@
github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs=
github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs=
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM=
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY=
github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto=
github.com/go-ini/ini v1.67.0 h1:z6ZrTEZqSWOTyH2FlglNbNgARyHG8oLW9gMELqKr06A=
github.com/go-ini/ini v1.67.0/go.mod h1:ByCAeIL28uOIIG0E3PJtZPDL8WnHpFKFOtgjp+3Ies8=
github.com/go-ole/go-ole v1.2.6/go.mod h1:pprOEPIfldk/42T2oK7lQ4v4JSDwmV0As9GaiUsvbm0=
github.com/go-ole/go-ole v1.3.0 h1:Dt6ye7+vXGIKZ7Xtk4s6/xVdGDQynvom7xCFEdWr6uE=
github.com/go-ole/go-ole v1.3.0/go.mod h1:5LS6F96DhAwUc7C+1HLexzMXY1xGRSryjyPPKW6zv78=
github.com/goccy/go-json v0.10.5 h1:Fq85nIqj+gXn/S5ahsiTlK3TmC85qgirsdTP/+DeaC4=
github.com/goccy/go-json v0.10.5/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M=
github.com/golang-jwt/jwt/v4 v4.5.2 h1:YtQM7lnr8iZ+j5q71MGKkNw9Mn7AjHM68uc9g5fXeUI=
@@ -16,13 +11,7 @@ github.com/golang-jwt/jwt/v4 v4.5.2/go.mod h1:m21LjoU+eqJr34lmDMbreY2eSTRJ1cv77w
github.com/golang/protobuf v1.2.0/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U=
github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek=
github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps=
github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8=
github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU=
github.com/klauspost/cpuid/v2 v2.0.1/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg=
github.com/klauspost/cpuid/v2 v2.2.10 h1:tBs3QSyvjDyFTq3uoc/9xFpCuOsJQFNPiAhYdw2skhE=
github.com/klauspost/cpuid/v2 v2.2.10/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0=
github.com/lufia/plan9stats v0.0.0-20250317134145-8bc96cf8fc35 h1:PpXWgLPs+Fqr325bN2FD2ISlRRztXibcX6e8f5FR5Dc=
github.com/lufia/plan9stats v0.0.0-20250317134145-8bc96cf8fc35/go.mod h1:autxFIvghDt3jPTLoqZ9OZ7s9qTGNAWmYCjVFWPX/zg=
github.com/matttproud/golang_protobuf_extensions v1.0.4 h1:mmDVorXM7PCGKw94cs5zkfA9PSy5pEvNWRP0ET0TIVo=
github.com/matttproud/golang_protobuf_extensions v1.0.4/go.mod h1:BSXmuO+STAnVfrANrmjBb36TMTDstsz7MSK+HVaYKv4=
github.com/minio/madmin-go/v3 v3.0.107-0.20250415152934-4b504b82db63 h1:ktN/FrMuM9sjvjIbPZYRKeHEzBDOXQdpYUDiNO0CutE=
@@ -35,10 +24,6 @@ github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ=
github.com/philhofer/fwd v1.1.3-0.20240916144458-20a13a1f6b7c h1:dAMKvw0MlJT1GshSTtih8C2gDs04w8dReiOGXrGLNoY=
github.com/philhofer/fwd v1.1.3-0.20240916144458-20a13a1f6b7c/go.mod h1:RqIHx9QI14HlwKwm98g9Re5prTQ6LdeRQn+gXJFxsJM=
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U=
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 h1:o4JXh1EVt9k/+g42oCprj/FisM4qX9L3sZB3upGN2ZU=
github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55/go.mod h1:OmDBASR4679mdNQnz2pUhc2G8CO2JrUAVFDRBDP/hJE=
github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk=
github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE=
github.com/prometheus/common v0.63.0 h1:YR/EIY1o3mEFP/kZCD7iDMnLPlGyuU2Gb3HIcXnA98k=
@@ -51,47 +36,29 @@ github.com/prometheus/prometheus v0.303.0 h1:wsNNsbd4EycMCphYnTmNY9JASBVbp7NWwJn
github.com/prometheus/prometheus v0.303.0/go.mod h1:8PMRi+Fk1WzopMDeb0/6hbNs9nV6zgySkU/zds5Lu3o=
github.com/rs/xid v1.6.0 h1:fV591PaemRlL6JfRxGDEPl69wICngIQ3shQtzfy2gxU=
github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0=
github.com/safchain/ethtool v0.5.10 h1:Im294gZtuf4pSGJRAOGKaASNi3wMeFaGaWuSaomedpc=
github.com/safchain/ethtool v0.5.10/go.mod h1:w9jh2Lx7YBR4UwzLkzCmWl85UY0W2uZdd7/DckVE5+c=
github.com/secure-io/sio-go v0.3.1 h1:dNvY9awjabXTYGsTF1PiCySl9Ltofk9GA3VdWlo7rRc=
github.com/secure-io/sio-go v0.3.1/go.mod h1:+xbkjDzPjwh4Axd07pRKSNriS9SCiYksWnZqdnfpQxs=
github.com/shirou/gopsutil/v3 v3.24.5 h1:i0t8kL+kQTvpAYToeuiVk3TgDeKOFioZO3Ztz/iZ9pI=
github.com/shirou/gopsutil/v3 v3.24.5/go.mod h1:bsoOS1aStSs9ErQ1WWfxllSeS1K5D+U30r2NfcubMVk=
github.com/shoenig/go-m1cpu v0.1.6 h1:nxdKQNcEB6vzgA2E2bvzKIYRuNj7XNJ4S/aRSwKzFtM=
github.com/shoenig/go-m1cpu v0.1.6/go.mod h1:1JJMcUBvfNwpq05QDQVAnx3gUHr9IYF7GNg9SUEw2VQ=
github.com/shoenig/test v0.6.4 h1:kVTaSd7WLz5WZ2IaoM0RSzRsUD+m8wRR+5qvntpn4LU=
github.com/shoenig/test v0.6.4/go.mod h1:byHiCGXqrVaflBLAMq/srcZIHynQPQgeyvkvXnjqq0k=
github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA=
github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY=
github.com/tinylib/msgp v1.2.5 h1:WeQg1whrXRFiZusidTQqzETkRpGjFjcIhW6uqWH09po=
github.com/tinylib/msgp v1.2.5/go.mod h1:ykjzy2wzgrlvpDCRc4LA8UXy6D8bzMSuAF3WD57Gok0=
github.com/tklauser/go-sysconf v0.3.15 h1:VE89k0criAymJ/Os65CSn1IXaol+1wrsFHEB8Ol49K4=
github.com/tklauser/go-sysconf v0.3.15/go.mod h1:Dmjwr6tYFIseJw7a3dRLJfsHAMXZ3nEnL/aZY+0IuI4=
github.com/tklauser/numcpus v0.10.0 h1:18njr6LDBk1zuna922MgdjQuJFjrdppsZG60sHGfjso=
github.com/tklauser/numcpus v0.10.0/go.mod h1:BiTKazU708GQTYF4mB+cmlpT2Is1gLk7XVuEeem8LsQ=
github.com/yusufpapurcu/wmi v1.2.4 h1:zFUKzehAFReQwLys1b/iSMl+JQGSCSjtVqQn9bBrPo0=
github.com/yusufpapurcu/wmi v1.2.4/go.mod h1:SBZ9tNy3G9/m5Oi98Zks0QjeHVDvuK0qfxQmPyzfmi0=
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
golang.org/x/crypto v0.0.0-20200302210943-78000ba7a073/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto=
golang.org/x/crypto v0.52.0 h1:RMs7fP2rXdep0CftQlK8Uf+kibLm7qkCcradZWYz988=
golang.org/x/crypto v0.52.0/go.mod h1:1QgfPxDqh0T2M/elOJtp9RvuR95kVjir0e6/BvEmGbc=
golang.org/x/crypto v0.37.0 h1:kJNSjF/Xp7kU0iB2Z+9viTPMW4EqqsrywMXLJOOsXSE=
golang.org/x/crypto v0.37.0/go.mod h1:vg+k43peMZ0pUMhYmVAWysMK35e6ioLh3wB8ZCAfbVc=
golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
golang.org/x/net v0.55.0 h1:bcvxaJn3e1U6InsFWt1JUq1aSjnRxLzT2rtD2KfkDF8=
golang.org/x/net v0.55.0/go.mod h1:L5U2KuzuOe1lY7Z+aWVIKK6qEeJXnXV9yzGA+WCHJww=
golang.org/x/net v0.39.0 h1:ZCu7HMWDxpXpaiKdhzIfaltL9Lp31x/3fCP11bc6/fY=
golang.org/x/net v0.39.0/go.mod h1:X7NRbYVEA+ewNkCNyJ513WmMdQ3BineSwVtN2zD/d+E=
golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.13.0 h1:AauUjRAJ9OSnvULf/ARrrVywoJDy0YS2AwQ98I37610=
golang.org/x/sync v0.13.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA=
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20200302150141-5c8b2ff67527/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20201204225414-ed752295db88/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.1.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.29.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY=
golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/sys v0.32.0 h1:s77OFDvIQeibCmezSnk/q6iAfkdiQaJi4VzroCFrN20=
golang.org/x/sys v0.32.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k=
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
google.golang.org/protobuf v1.36.6 h1:z1NpPI8ku2WgiWnf+t9wTPsn6eP1L7ksHUlkfLvd9xY=
google.golang.org/protobuf v1.36.6/go.mod h1:jduwjTPXsFjZGTmRluh+L6NjiWu7pchiJ2/5YcXBHnY=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=