Compare commits

..

4 Commits

Author SHA1 Message Date
马登山 82fb0a8843 fix(scanner): correct observational usage arguments 2026-08-23 10:31:42 +08:00
马登山 307510749e style: format usage freshness changes 2026-08-23 10:28:47 +08:00
马登山 5496e14960 fix(ecstore): preserve quota baseline across restart 2026-08-23 10:20:39 +08:00
马登山 ec3b7a7dc6 fix(scanner): publish partial usage observations 2026-08-23 09:50:16 +08:00
14 changed files with 719 additions and 159 deletions
+108 -1
View File
@@ -236,6 +236,15 @@ pub struct DataUsageInfo {
/// without relying on synchronized clocks.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub usage_snapshot_authoritative_baseline: Option<DataUsageSnapshotIdentity>,
/// Per-set freshness for an observational aggregate. A set entry is
/// never sufficient to make the aggregate authoritative; it only records
/// which last-known-good generation contributed to the view.
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub usage_snapshot_set_states: Vec<DataUsageSnapshotSetState>,
/// An observational view may contain only the sets that completed this
/// cycle (or retained a compatible last-known-good cache).
#[serde(default)]
pub usage_snapshot_partial: bool,
/// Deprecated kept here for backward compatibility reasons
pub bucket_sizes: HashMap<String, u64>,
/// Per-disk snapshot information when available
@@ -252,6 +261,22 @@ pub struct DataUsageSnapshotIdentity {
pub scanner_epoch: Option<u64>,
}
#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
pub struct DataUsageSnapshotSetState {
pub pool_index: u64,
pub set_index: u64,
#[serde(default)]
pub scanner_cycle: Option<u64>,
#[serde(default)]
pub scanner_epoch: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub scan_plan_digest: Option<[u8; 32]>,
#[serde(default)]
pub complete: bool,
#[serde(default)]
pub tombstone: bool,
}
impl DataUsageInfo {
pub fn snapshot_identity(&self) -> DataUsageSnapshotIdentity {
DataUsageSnapshotIdentity {
@@ -291,7 +316,7 @@ pub fn data_usage_snapshot_is_newer(candidate: &DataUsageInfo, baseline: &DataUs
/// rollback delete/recreate fences the previous bucket incarnation too.
pub fn observed_data_usage_is_newer(observed: &DataUsageInfo, authoritative: &DataUsageInfo) -> bool {
observed.usage_snapshot_converged == Some(false)
&& observed.is_complete_bucket_usage_snapshot()
&& (observed.is_complete_bucket_usage_snapshot() || observed.is_valid_partial_snapshot())
&& observed.usage_snapshot_authoritative_baseline.as_ref() == Some(&authoritative.snapshot_identity())
&& data_usage_snapshot_is_newer(observed, authoritative)
}
@@ -1436,6 +1461,39 @@ impl DataUsageInfo {
&& u64::try_from(self.buckets_usage.len()).ok() == Some(self.buckets_count)
}
/// Validate provenance before an observational view can be selected for
/// admin display. Partial data is accepted only with unique set states,
/// a plan digest for every state, and at least one usable generation.
pub fn is_valid_partial_snapshot(&self) -> bool {
if !self.usage_snapshot_partial
|| self.usage_snapshot_converged != Some(false)
|| self.last_update.is_none()
|| self.scanner_cycle.is_none()
|| self.scanner_epoch.is_none()
|| self.usage_snapshot_set_states.is_empty()
|| u64::try_from(self.buckets_usage.len()).ok() != Some(self.buckets_count)
{
return false;
}
let mut previous = None;
let mut plan_digest = None;
let mut has_source = false;
for state in &self.usage_snapshot_set_states {
if state.scan_plan_digest.is_none()
|| plan_digest.is_some_and(|digest| Some(digest) != state.scan_plan_digest)
|| state.scanner_cycle.is_some() != state.scanner_epoch.is_some()
|| previous.is_some_and(|(pool, set)| (pool, set) >= (state.pool_index, state.set_index))
{
return false;
}
previous = Some((state.pool_index, state.set_index));
plan_digest = state.scan_plan_digest;
has_source |= state.scanner_cycle.is_some() && !state.tombstone;
}
has_source
}
/// Add object metadata to data usage statistics
pub fn add_object(&mut self, object_path: &str, meta_object: &rustfs_filemeta::MetaObject) {
// This method is kept for backward compatibility
@@ -2263,6 +2321,55 @@ mod tests {
assert!(!observed_data_usage_is_newer(&candidate(2, 9, Some(false), true), &authoritative));
assert!(!observed_data_usage_is_newer(&candidate(2, 11, Some(true), true), &authoritative));
assert!(!observed_data_usage_is_newer(&candidate(2, 11, Some(false), false), &authoritative));
let mut partial = candidate(2, 11, Some(false), false);
partial.usage_snapshot_partial = true;
partial.usage_snapshot_set_states = vec![DataUsageSnapshotSetState {
pool_index: 0,
set_index: 0,
scanner_cycle: Some(10),
scanner_epoch: Some(2),
scan_plan_digest: Some([1; 32]),
complete: false,
tombstone: false,
}];
assert!(observed_data_usage_is_newer(&partial, &authoritative));
}
#[test]
fn mixed_topology_snapshot_is_rejected() {
let mut partial = DataUsageInfo {
last_update: Some(SystemTime::UNIX_EPOCH + Duration::from_secs(2)),
scanner_cycle: Some(11),
scanner_epoch: Some(2),
buckets_count: 0,
usage_snapshot_converged: Some(false),
usage_snapshot_partial: true,
usage_snapshot_set_states: vec![
DataUsageSnapshotSetState {
pool_index: 0,
set_index: 0,
scanner_cycle: Some(11),
scanner_epoch: Some(2),
scan_plan_digest: Some([1; 32]),
complete: true,
tombstone: false,
},
DataUsageSnapshotSetState {
pool_index: 1,
set_index: 0,
scanner_cycle: Some(10),
scanner_epoch: Some(2),
scan_plan_digest: Some([2; 32]),
complete: false,
tombstone: false,
},
],
..Default::default()
};
assert!(!partial.is_valid_partial_snapshot());
partial.usage_snapshot_set_states[1].scan_plan_digest = Some([1; 32]);
assert!(partial.is_valid_partial_snapshot());
}
#[test]
+109 -3
View File
@@ -73,6 +73,16 @@ struct CachedBucketUsage {
// mutation. A strictly later generation is required before the mutation
// evidence can be discarded.
pending_scanner_position: Option<(u64, u64)>,
// Deletes are visible to admin immediately, but quota admission keeps
// them pending until a complete scanner generation reconciles the set.
// This marker intentionally remains process-local: the delete request
// updates this overlay before the scanner writes a durable snapshot. If
// the process restarts first, loading the persisted complete snapshot
// restores the pre-reconciliation (larger) baseline, which is
// conservative for quota admission. A persisted post-delete snapshot is
// necessarily a complete scanner reconciliation and therefore creates a
// fresh cache entry with no pending hold.
pending_negative_delta: u64,
}
type UsageMemoryCache = Arc<RwLock<HashMap<String, CachedBucketUsage>>>;
@@ -948,7 +958,12 @@ async fn load_observed_data_usage_snapshot(store: Arc<ECStore>) -> Option<DataUs
};
match parse_usage_snapshot(&data) {
Ok(info) if info.usage_snapshot_converged == Some(false) && info.is_complete_bucket_usage_snapshot() => Some(info),
Ok(info)
if info.usage_snapshot_converged == Some(false)
&& (info.is_complete_bucket_usage_snapshot() || info.is_valid_partial_snapshot()) =>
{
Some(info)
}
Ok(_) => {
error!(
event = "data_usage_snapshot_load_failed",
@@ -993,7 +1008,7 @@ async fn load_admin_data_usage_from_backend(store: Arc<ECStore>) -> Result<DataU
}
fn discard_incomplete_bucket_usage(data_usage_info: &mut DataUsageInfo) {
if !data_usage_info.is_complete_bucket_usage_snapshot() {
if !data_usage_info.is_complete_bucket_usage_snapshot() && !data_usage_info.usage_snapshot_partial {
data_usage_info.usage_snapshot_complete = false;
data_usage_info.buckets_usage.clear();
data_usage_info.bucket_sizes.clear();
@@ -1643,6 +1658,7 @@ fn cached_bucket_usage_from_backend(usage: BucketUsageInfo, updated_at: SystemTi
dirty: false,
stale_snapshot_pending: false,
pending_scanner_position: None,
pending_negative_delta: 0,
}
}
@@ -1656,6 +1672,7 @@ fn cached_bucket_usage_now(usage: BucketUsageInfo) -> CachedBucketUsage {
dirty: false,
stale_snapshot_pending: false,
pending_scanner_position: None,
pending_negative_delta: 0,
}
}
@@ -1808,6 +1825,7 @@ pub async fn record_bucket_object_delete_memory(bucket: &str, deleted_size: u64,
.or_insert_with(|| cached_bucket_usage_now(BucketUsageInfo::default()));
entry.usage.size = entry.usage.size.saturating_sub(deleted_size);
entry.pending_negative_delta = entry.pending_negative_delta.saturating_add(deleted_size);
if removed_current_object {
entry.usage.objects_count = entry.usage.objects_count.saturating_sub(1);
entry.usage.versions_count = entry.usage.versions_count.saturating_sub(1);
@@ -1863,7 +1881,7 @@ pub async fn get_bucket_usage_memory(bucket: &str) -> Option<u64> {
cache
.get(bucket)
.filter(|cached| cached.authoritative)
.map(|cached| cached.usage.size)
.map(|cached| cached.usage.size.saturating_add(cached.pending_negative_delta))
}
async fn update_usage_cache_if_needed() {
@@ -2943,6 +2961,45 @@ mod tests {
assert_eq!(selected.usage_snapshot_converged, Some(true));
}
#[test]
fn persisted_authoritative_stalls_but_memory_overlay_remains_visible() {
let authoritative = DataUsageInfo {
last_update: Some(SystemTime::UNIX_EPOCH),
scanner_epoch: Some(4),
scanner_cycle: Some(10),
usage_snapshot_complete: true,
..Default::default()
};
let mut partial = authoritative.clone();
partial.last_update = Some(SystemTime::UNIX_EPOCH + Duration::from_secs(1));
partial.scanner_cycle = Some(11);
partial.usage_snapshot_complete = false;
partial.usage_snapshot_partial = true;
partial.usage_snapshot_converged = Some(false);
partial.usage_snapshot_authoritative_baseline = Some(authoritative.snapshot_identity());
partial.usage_snapshot_set_states = vec![rustfs_data_usage::DataUsageSnapshotSetState {
pool_index: 0,
set_index: 0,
scanner_cycle: Some(10),
scanner_epoch: Some(4),
scan_plan_digest: Some([1; 32]),
complete: false,
tombstone: false,
}];
partial.buckets_usage.insert(
"bucket".to_string(),
BucketUsageInfo {
size: 100,
..Default::default()
},
);
partial.buckets_count = 1;
let (selected, _) = select_admin_data_usage_snapshot(authoritative, true, Some(partial));
assert!(selected.usage_snapshot_partial);
assert_eq!(selected.buckets_usage.get("bucket").map(|usage| usage.size), Some(100));
}
#[tokio::test]
async fn authoritative_save_cleanup_removes_observed_snapshot_best_effort() {
let store = UsageCasStore::default();
@@ -4665,6 +4722,55 @@ mod tests {
);
}
#[tokio::test]
#[serial]
async fn partial_usage_is_observational_not_authoritative_for_quota() {
clear_usage_memory_cache_for_test().await;
let mut partial = data_usage_info_for_test("bucket-a", 10, 100, SystemTime::now());
partial.usage_snapshot_complete = false;
partial.usage_snapshot_partial = true;
replace_bucket_usage_memory_from_info(&partial).await;
assert_eq!(get_bucket_usage_memory("bucket-a").await, None);
}
#[tokio::test]
#[serial]
async fn stale_quota_uses_complete_baseline_plus_positive_deltas() {
clear_usage_memory_cache_for_test().await;
let baseline = data_usage_info_for_test("bucket-a", 1, 100, SystemTime::now());
replace_bucket_usage_memory_from_info(&baseline).await;
record_bucket_object_write_memory("bucket-a", None, 25).await;
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(125));
}
#[tokio::test]
#[serial]
async fn negative_delta_waits_for_set_reconciliation() {
clear_usage_memory_cache_for_test().await;
let baseline = data_usage_info_for_test("bucket-a", 1, 100, SystemTime::UNIX_EPOCH + Duration::from_secs(100));
replace_bucket_usage_memory_from_info(&baseline).await;
record_bucket_object_delete_memory("bucket-a", 25, true).await;
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(100));
// Simulate a process restart: the request-path overlay is gone, but
// the persisted authoritative snapshot is still the pre-reconciliation
// baseline. Quota must remain conservative until a complete scanner
// result proves the delete.
clear_usage_memory_cache_for_test().await;
replace_bucket_usage_memory_from_info(&baseline).await;
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(100));
let reconciled = data_usage_info_for_test("bucket-a", 0, 75, SystemTime::UNIX_EPOCH + Duration::from_secs(101));
replace_bucket_usage_memory_from_info(&reconciled).await;
assert_eq!(get_bucket_usage_memory("bucket-a").await, Some(75));
}
#[tokio::test]
#[serial]
async fn memory_overlay_counts_versioned_overwrite_as_new_version() {
+21 -3
View File
@@ -28,8 +28,9 @@ use rustfs_common::heal_channel::HealScanMode;
use rustfs_config::ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS;
pub use rustfs_data_usage::{
AllTierStats, BucketTargetUsageInfo, BucketUsageInfo, DATA_USAGE_OBJECT_NAME, DATA_USAGE_OBSERVED_OBJECT_NAME,
DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageInfo, LEGACY_DATA_USAGE_OBJECT_NAME, PrefixUsageEntry,
PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeSummary, TierStats, hash_path, prefix_usage_in_cache,
DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageInfo, DataUsageSnapshotSetState, LEGACY_DATA_USAGE_OBJECT_NAME,
PrefixUsageEntry, PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeSummary, TierStats, hash_path,
prefix_usage_in_cache,
};
use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf};
use tokio::time::{Duration, Instant, sleep, timeout};
@@ -344,6 +345,18 @@ pub struct DataUsageCacheInfo {
pub scan_plan_digest: Option<DataUsageScanPlanDigest>,
#[serde(default)]
pub cache_key_format: u16,
/// Whether the entries retained while a set scan was incomplete come
/// from a prior complete set snapshot. This is observational input only.
#[serde(default)]
pub lkg_snapshot_complete: bool,
#[serde(default)]
pub lkg_next_cycle: Option<u64>,
#[serde(default)]
pub lkg_last_update: Option<SystemTime>,
#[serde(default)]
pub lkg_leader_epoch: Option<u64>,
#[serde(default)]
pub lkg_scan_plan_digest: Option<DataUsageScanPlanDigest>,
}
impl Serialize for DataUsageCacheInfo {
@@ -353,7 +366,7 @@ impl Serialize for DataUsageCacheInfo {
{
// Keep this metadata map-encoded so older readers can ignore fields
// appended by newer scanner versions during rolling upgrades.
let mut state = serializer.serialize_map(Some(16))?;
let mut state = serializer.serialize_map(Some(21))?;
state.serialize_entry("name", &self.name)?;
state.serialize_entry("next_cycle", &self.next_cycle)?;
state.serialize_entry("leader_epoch", &self.leader_epoch)?;
@@ -370,6 +383,11 @@ impl Serialize for DataUsageCacheInfo {
state.serialize_entry("snapshot_complete", &self.snapshot_complete)?;
state.serialize_entry("scan_plan_digest", &self.scan_plan_digest)?;
state.serialize_entry("cache_key_format", &self.cache_key_format)?;
state.serialize_entry("lkg_snapshot_complete", &self.lkg_snapshot_complete)?;
state.serialize_entry("lkg_next_cycle", &self.lkg_next_cycle)?;
state.serialize_entry("lkg_last_update", &self.lkg_last_update)?;
state.serialize_entry("lkg_leader_epoch", &self.lkg_leader_epoch)?;
state.serialize_entry("lkg_scan_plan_digest", &self.lkg_scan_plan_digest)?;
state.end()
}
}
+2 -3
View File
@@ -2274,9 +2274,8 @@ async fn final_data_usage_publication_defer_reason(
}
}
ScannerCycleStatus::Deferred(reason) => Some(reason),
// Incomplete cycles do not publish a usage snapshot. Keep the
// decision permissive so existing partial-cycle handling remains
// unchanged if a future scanner path emits a bookkeeping update.
// Incomplete cycles may publish a non-authoritative observational
// snapshot when at least one set has a usable current/LKG view.
ScannerCycleStatus::Incomplete => None,
}
}
+1 -1
View File
@@ -198,7 +198,7 @@ where
data_usage_info.usage_snapshot_authoritative_baseline = Some(authoritative.snapshot_identity());
}
if !data_usage_info.is_complete_bucket_usage_snapshot() {
if !data_usage_info.is_complete_bucket_usage_snapshot() && !data_usage_info.usage_snapshot_partial {
error!(
target: "rustfs::scanner",
event = EVENT_SCANNER_PERSIST_STATE,
+13 -2
View File
@@ -18,8 +18,8 @@ use crate::scanner_folder::{ScannerItem, scan_data_folder};
use crate::sleeper::SCANNER_SLEEPER;
use crate::{
DATA_USAGE_CACHE_NAME, DATA_USAGE_ROOT, DataUsageCache, DataUsageCacheInfo, DataUsageCachePrepareOutcome,
DataUsageCacheSource, DataUsageEntry, DataUsageEntryInfo, DataUsageInfo, DataUsageScanPlanDigest, ScannerError, SizeSummary,
TierStats,
DataUsageCacheSource, DataUsageEntry, DataUsageEntryInfo, DataUsageInfo, DataUsageScanPlanDigest, DataUsageSnapshotSetState,
ScannerError, SizeSummary, TierStats,
};
use futures::future::join_all;
use metrics::counter;
@@ -278,6 +278,17 @@ async fn publish_usage_snapshot(
Ok(true)
}
async fn publish_observational_snapshot(
updates: &mpsc::Sender<DataUsageInfo>,
mut data_usage_info: DataUsageInfo,
) -> Result<bool> {
data_usage_info.usage_snapshot_complete = false;
data_usage_info.usage_snapshot_partial = true;
data_usage_info.usage_snapshot_converged = Some(false);
send_data_usage_update(updates, data_usage_info).await?;
Ok(true)
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum ScannerCycleActivityStatus {
Unchanged,
+145 -2
View File
@@ -188,7 +188,7 @@ pub(super) fn completed_data_usage_info(
}
let mut total = DataUsageEntry::default();
let mut buckets_usage = HashMap::with_capacity(all_buckets.len());
let mut bucket_entries = HashMap::with_capacity(all_buckets.len());
for bucket in all_buckets {
let mut merged = DataUsageEntry::default();
for result in results {
@@ -200,10 +200,14 @@ pub(super) fn completed_data_usage_info(
if !total.checked_merge(&merged) {
return None;
}
buckets_usage.insert(bucket.clone(), checked_bucket_usage_info(&merged)?);
bucket_entries.insert(bucket.clone(), merged);
}
let merged_last_update = results.iter().filter_map(|result| result.info.last_update).max()?;
let buckets_usage = bucket_entries
.iter()
.map(|(bucket, entry)| Some((bucket.clone(), checked_bucket_usage_info(entry)?)))
.collect::<Option<HashMap<_, _>>>()?;
let bucket_sizes = buckets_usage
.iter()
.map(|(bucket, usage)| (bucket.clone(), usage.size))
@@ -225,6 +229,145 @@ pub(super) fn completed_data_usage_info(
Some((data_usage_info, merged_last_update))
}
/// Build a non-authoritative view from the set snapshots that completed this
/// cycle plus compatible per-set last-known-good caches. The caller must
/// persist this result only on the observational object; a missing set is
/// intentionally represented by an incomplete state and is never treated as
/// an empty set.
pub(super) fn observational_data_usage_info(
results: &[DataUsageCache],
expected_sources: &HashSet<DataUsageCacheSource>,
all_buckets: &[String],
expected_plan_digest: DataUsageScanPlanDigest,
scanner_cycle: u64,
leader_epoch: u64,
) -> Option<(DataUsageInfo, SystemTime)> {
let mut by_source = HashMap::with_capacity(results.len());
for result in results {
let source = result.info.source?;
if !expected_sources.contains(&source) || by_source.insert(source, result).is_some() {
return None;
}
}
let mut usable = Vec::new();
let mut set_states = Vec::with_capacity(expected_sources.len());
let mut sources = expected_sources.iter().copied().collect::<Vec<_>>();
sources.sort_by_key(|source| (source.pool_index, source.set_index));
for source in sources {
let result = by_source.get(&source).copied();
let current = result.filter(|result| {
result.info.snapshot_complete
&& result.info.next_cycle == scanner_cycle
&& result.info.leader_epoch == leader_epoch
&& result.info.scan_plan_digest == Some(expected_plan_digest)
});
let lkg = result.filter(|result| {
!result.info.snapshot_complete
&& result.info.lkg_snapshot_complete
&& result.info.lkg_scan_plan_digest == Some(expected_plan_digest)
&& result.info.lkg_leader_epoch.is_some_and(|epoch| {
epoch < leader_epoch
|| (epoch == leader_epoch && result.info.lkg_next_cycle.is_some_and(|cycle| cycle <= scanner_cycle))
})
});
let current_snapshot = current.is_some();
let selected = current.or(lkg);
if let Some(selected) = selected {
let (cycle, epoch, digest, last_update, complete) = if current_snapshot {
(
Some(selected.info.next_cycle),
Some(selected.info.leader_epoch),
selected.info.scan_plan_digest.map(|digest| digest.0),
selected.info.last_update,
true,
)
} else {
(
selected.info.lkg_next_cycle,
selected.info.lkg_leader_epoch,
selected.info.lkg_scan_plan_digest.map(|digest| digest.0),
selected.info.lkg_last_update,
false,
)
};
set_states.push(DataUsageSnapshotSetState {
pool_index: u64::try_from(source.pool_index).ok()?,
set_index: u64::try_from(source.set_index).ok()?,
scanner_cycle: cycle,
scanner_epoch: epoch,
scan_plan_digest: digest,
complete,
tombstone: false,
});
usable.push((selected, last_update));
} else {
set_states.push(DataUsageSnapshotSetState {
pool_index: u64::try_from(source.pool_index).ok()?,
set_index: u64::try_from(source.set_index).ok()?,
scanner_cycle: None,
scanner_epoch: None,
scan_plan_digest: Some(expected_plan_digest.0),
complete: false,
tombstone: false,
});
}
}
if usable.is_empty() {
return None;
}
let mut total = DataUsageEntry::default();
let mut bucket_entries = HashMap::with_capacity(all_buckets.len());
let mut merged_last_update = None;
for (result, last_update) in usable {
if let Some(update) = last_update {
merged_last_update = Some(merged_last_update.map_or(update, |current: SystemTime| current.max(update)));
}
for bucket in all_buckets {
let Some(entry) = result.checked_flatten(bucket) else {
continue;
};
let bucket_entry = bucket_entries.entry(bucket.clone()).or_insert_with(DataUsageEntry::default);
if !bucket_entry.checked_merge(&entry) {
return None;
}
if !total.checked_merge(&entry) {
return None;
}
}
}
let merged_last_update = merged_last_update?;
let buckets_usage = bucket_entries
.iter()
.map(|(bucket, entry)| Some((bucket.clone(), checked_bucket_usage_info(entry)?)))
.collect::<Option<HashMap<_, _>>>()?;
Some((
DataUsageInfo {
last_update: Some(merged_last_update),
scanner_cycle: Some(scanner_cycle),
scanner_epoch: Some(leader_epoch),
objects_total_count: u64::try_from(total.objects).ok()?,
versions_total_count: u64::try_from(total.versions).ok()?,
delete_markers_total_count: u64::try_from(total.delete_markers).ok()?,
objects_total_size: u64::try_from(total.size).ok()?,
tier_stats: total.all_tier_stats.filter(|tiers| !tiers.is_empty()),
buckets_count: u64::try_from(buckets_usage.len()).ok()?,
bucket_sizes: buckets_usage
.iter()
.map(|(bucket, usage)| (bucket.clone(), usage.size))
.collect(),
buckets_usage,
usage_snapshot_complete: false,
usage_snapshot_partial: true,
usage_snapshot_converged: Some(false),
usage_snapshot_set_states: set_states,
..Default::default()
},
merged_last_update,
))
}
pub(super) async fn send_cache_root_entry_info(
bucket_result_tx: &mpsc::Sender<DataUsageEntryInfo>,
cache: &DataUsageCache,
+89 -30
View File
@@ -40,6 +40,21 @@ impl ScannerIOCache for SetDisks {
let set_label = self.set_index.to_string();
let source = DataUsageCacheSource::new(self.pool_index, self.set_index);
let mut old_cache = DataUsageCache::default();
if let Err(e) = old_cache.load(self.clone(), DATA_USAGE_CACHE_NAME).await {
warn!(
target: "rustfs::scanner::io",
event = EVENT_SCANNER_CACHE_PERSIST_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_IO,
pool = self.pool_index,
set = self.set_index,
cache_name = DATA_USAGE_CACHE_NAME,
state = "old_cache_load_failed",
error = %e,
"Scanner old data usage cache load failed; rebuilding from bucket caches"
);
}
if buckets.is_empty() {
let now = SystemTime::now();
let mut cache = DataUsageCache {
@@ -80,6 +95,24 @@ impl ScannerIOCache for SetDisks {
"Scanner set state found no online disks"
);
reset_disk_bucket_scan_gauges(&pool_label, &set_label);
let lkg = old_cache.info.snapshot_complete.then(|| old_cache.clone());
let mut incomplete_scope = lkg.clone().unwrap_or_default();
incomplete_scope.info.name = DATA_USAGE_ROOT.to_string();
incomplete_scope.info.next_cycle = want_cycle;
incomplete_scope.info.last_update = None;
incomplete_scope.info.leader_epoch = leader_epoch;
incomplete_scope.info.source = Some(source);
incomplete_scope.info.snapshot_complete = false;
incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest);
incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT;
if let Some(lkg) = lkg {
incomplete_scope.info.lkg_snapshot_complete = true;
incomplete_scope.info.lkg_next_cycle = Some(lkg.info.next_cycle);
incomplete_scope.info.lkg_last_update = lkg.info.last_update;
incomplete_scope.info.lkg_leader_epoch = Some(lkg.info.leader_epoch);
incomplete_scope.info.lkg_scan_plan_digest = lkg.info.scan_plan_digest;
}
let _ = updates.send(incomplete_scope).await;
return Ok(());
}
// Preserve the original set topology across capability filtering. During
@@ -162,6 +195,24 @@ impl ScannerIOCache for SetDisks {
"Scanner set state found no usable namespace scanner disks"
);
reset_disk_bucket_scan_gauges(&pool_label, &set_label);
let lkg = old_cache.info.snapshot_complete.then(|| old_cache.clone());
let mut incomplete_scope = lkg.clone().unwrap_or_default();
incomplete_scope.info.name = DATA_USAGE_ROOT.to_string();
incomplete_scope.info.next_cycle = want_cycle;
incomplete_scope.info.last_update = None;
incomplete_scope.info.leader_epoch = leader_epoch;
incomplete_scope.info.source = Some(source);
incomplete_scope.info.snapshot_complete = false;
incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest);
incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT;
if let Some(lkg) = lkg {
incomplete_scope.info.lkg_snapshot_complete = true;
incomplete_scope.info.lkg_next_cycle = Some(lkg.info.next_cycle);
incomplete_scope.info.lkg_last_update = lkg.info.last_update;
incomplete_scope.info.lkg_leader_epoch = Some(lkg.info.leader_epoch);
incomplete_scope.info.lkg_scan_plan_digest = lkg.info.scan_plan_digest;
}
let _ = updates.send(incomplete_scope).await;
return Ok(());
}
let set_disk_inventory = Arc::new(scanner_set_disk_inventory(self.as_ref()).await);
@@ -203,22 +254,15 @@ impl ScannerIOCache for SetDisks {
record_disk_bucket_scans_active(0, &pool_label, &set_label);
let _reset_disk_bucket_scan_gauges = DiskBucketScanGaugeReset::new(pool_label.clone(), set_label.clone());
let mut old_cache = DataUsageCache::default();
if let Err(e) = old_cache.load(self.clone(), DATA_USAGE_CACHE_NAME).await {
warn!(
target: "rustfs::scanner::io",
event = EVENT_SCANNER_CACHE_PERSIST_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_IO,
pool = self.pool_index,
set = self.set_index,
cache_name = DATA_USAGE_CACHE_NAME,
state = "old_cache_load_failed",
error = %e,
"Scanner old data usage cache load failed; rebuilding from bucket caches"
);
}
match old_cache.prepare_for_scan(
let old_lkg = old_cache.info.snapshot_complete.then(|| {
(
old_cache.info.next_cycle,
old_cache.info.last_update,
old_cache.info.leader_epoch,
old_cache.info.scan_plan_digest,
)
});
let prepare_outcome = match old_cache.prepare_for_scan(
DATA_USAGE_ROOT,
want_cycle,
leader_epoch,
@@ -259,7 +303,16 @@ impl ScannerIOCache for SetDisks {
);
return Ok(());
}
DataUsageCachePrepareOutcome::Reused | DataUsageCachePrepareOutcome::Reset => {}
outcome => outcome,
};
if matches!(prepare_outcome, DataUsageCachePrepareOutcome::Reused)
&& let Some((cycle, last_update, epoch, digest)) = old_lkg
{
old_cache.info.lkg_snapshot_complete = true;
old_cache.info.lkg_next_cycle = Some(cycle);
old_cache.info.lkg_last_update = last_update;
old_cache.info.lkg_leader_epoch = Some(epoch);
old_cache.info.lkg_scan_plan_digest = digest;
}
let mut cache = DataUsageCache {
@@ -1099,23 +1152,29 @@ impl ScannerIOCache for SetDisks {
cache.info.next_cycle = want_cycle;
cache.info.last_update.get_or_insert_with(SystemTime::now);
cache.info.snapshot_complete = true;
cache.info.lkg_snapshot_complete = false;
cache.info.lkg_next_cycle = None;
cache.info.lkg_last_update = None;
cache.info.lkg_leader_epoch = None;
cache.info.lkg_scan_plan_digest = None;
cache.clone()
};
let _ = persist_and_publish_cache_snapshot(self.clone(), &updates, cache_snapshot, cache_cycle_floor.as_ref()).await;
} else {
let incomplete_scope = DataUsageCache {
info: DataUsageCacheInfo {
name: DATA_USAGE_ROOT.to_string(),
next_cycle: want_cycle,
leader_epoch,
source: Some(source),
snapshot_complete: false,
scan_plan_digest: Some(scan_plan_digest),
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
..Default::default()
},
cache: HashMap::new(),
};
let mut incomplete_scope = cache_mutex.lock().await.clone();
incomplete_scope.info.name = DATA_USAGE_ROOT.to_string();
incomplete_scope.info.next_cycle = want_cycle;
incomplete_scope.info.last_update = None;
incomplete_scope.info.leader_epoch = leader_epoch;
incomplete_scope.info.source = Some(source);
incomplete_scope.info.snapshot_complete = false;
incomplete_scope.info.scan_plan_digest = Some(scan_plan_digest);
incomplete_scope.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT;
incomplete_scope.info.lkg_snapshot_complete = old_cache.info.lkg_snapshot_complete;
incomplete_scope.info.lkg_next_cycle = old_cache.info.lkg_next_cycle;
incomplete_scope.info.lkg_last_update = old_cache.info.lkg_last_update;
incomplete_scope.info.lkg_leader_epoch = old_cache.info.lkg_leader_epoch;
incomplete_scope.info.lkg_scan_plan_digest = old_cache.info.lkg_scan_plan_digest;
if let Err(e) = updates.send(incomplete_scope).await {
error!(
target: "rustfs::scanner::io",
+33
View File
@@ -234,6 +234,7 @@ impl ScannerIOCycle for ECStore {
let active_set_scans_clone = active_set_scans.clone();
let (tx, mut rx) = mpsc::channel::<DataUsageCache>(1);
let failed_scope_tx = tx.clone();
// Spawn task to receive and store results
let receiver_fut = tokio::spawn(async move {
@@ -314,6 +315,21 @@ impl ScannerIOCycle for ECStore {
state = "set_scan_failed",
"Scanner set scan failed; continuing cycle"
);
let _ = failed_scope_tx
.send(DataUsageCache {
info: DataUsageCacheInfo {
name: DATA_USAGE_ROOT.to_string(),
next_cycle: want_cycle_clone,
leader_epoch,
source: Some(source),
snapshot_complete: false,
scan_plan_digest: Some(scan_plan_digest),
cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT,
..Default::default()
},
cache: HashMap::new(),
})
.await;
let mut first_err = first_err_mutex_clone.lock().await;
record_set_scan_failure(&mut first_err, e);
}
@@ -370,6 +386,19 @@ impl ScannerIOCycle for ECStore {
budget_elapsed,
ctx.is_cancelled(),
);
let observational_usage = completed_usage
.is_none()
.then(|| {
observational_data_usage_info(
&results,
&expected_sources,
&all_bucket_names,
scan_plan_digest,
want_cycle,
leader_epoch,
)
})
.flatten();
let structurally_complete_snapshot = result.is_ok() && completed_all_sets && completed_usage.is_some();
let cycle_status = classify_nsscanner_cycle(
structurally_complete_snapshot,
@@ -381,6 +410,10 @@ impl ScannerIOCycle for ECStore {
);
if let Some((data_usage_info, _)) = completed_usage {
publish_usage_snapshot(&updates, cycle_status, data_usage_info).await?;
} else if !ctx.is_cancelled()
&& let Some((data_usage_info, _)) = observational_usage
{
publish_observational_snapshot(&updates, data_usage_info).await?;
}
let dirty_usage_clear = should_clear_dirty_usage_snapshot(
result.is_ok(),
@@ -105,6 +105,160 @@ fn completed_data_usage_info_for_test(
completed_data_usage_info(results, &expected_sources, all_buckets, true, budget_elapsed, cancelled)
}
fn lkg_root_cache(bucket: &str, objects: usize, source: DataUsageCacheSource) -> DataUsageCache {
let mut cache = completed_root_cache(bucket, objects, 10, source);
cache.info.snapshot_complete = false;
cache.info.next_cycle = 8;
cache.info.leader_epoch = 3;
cache.info.lkg_snapshot_complete = true;
cache.info.lkg_next_cycle = Some(7);
cache.info.lkg_last_update = cache.info.last_update;
cache.info.lkg_leader_epoch = Some(3);
cache.info.lkg_scan_plan_digest = Some(TEST_PLAN_DIGEST);
cache
}
#[test]
fn partial_usage_is_observational_not_authoritative_for_quota() {
let all_buckets = vec!["bucket".to_string()];
let current_source = DataUsageCacheSource::new(0, 0);
let stalled_source = DataUsageCacheSource::new(1, 0);
let mut current = completed_root_cache("bucket", 2, 20, current_source);
current.info.next_cycle = 8;
current.info.leader_epoch = 3;
let stalled = lkg_root_cache("bucket", 1, stalled_source);
let expected = HashSet::from([current_source, stalled_source]);
assert!(
completed_data_usage_info(&[current.clone(), stalled.clone()], &expected, &all_buckets, true, false, false).is_none()
);
let (observed, _) = observational_data_usage_info(&[current, stalled], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
.expect("a completed set should produce an observational view");
assert!(observed.usage_snapshot_partial);
assert!(!observed.usage_snapshot_complete);
assert_eq!(observed.usage_snapshot_converged, Some(false));
assert_eq!(observed.usage_snapshot_set_states.len(), 2);
}
#[test]
fn lkg_scope_does_not_count_as_current_cycle_completion() {
let source = DataUsageCacheSource::new(0, 0);
let mut lkg = lkg_root_cache("bucket", 1, source);
lkg.info.last_update = None;
let expected = HashSet::from([source]);
assert!(!scanner_results_form_complete_snapshot(&[lkg], &expected));
}
#[test]
fn stale_quota_uses_complete_baseline_plus_positive_deltas() {
let all_buckets = vec!["bucket".to_string()];
let source = DataUsageCacheSource::new(0, 0);
let mut current = completed_root_cache("bucket", 3, 20, source);
current.info.next_cycle = 8;
current.info.leader_epoch = 3;
let expected = HashSet::from([source]);
let (observed, _) = observational_data_usage_info(&[current], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
.expect("complete set data is a valid observational baseline");
assert_eq!(observed.objects_total_size, 30);
assert_eq!(observed.usage_snapshot_set_states[0].complete, true);
}
#[test]
fn negative_delta_waits_for_set_reconciliation() {
let all_buckets = vec!["bucket".to_string()];
let source = DataUsageCacheSource::new(0, 0);
let mut stalled = lkg_root_cache("bucket", 4, source);
stalled.info.lkg_scan_plan_digest = Some(DataUsageScanPlanDigest([9; 32]));
let expected = HashSet::from([source]);
assert!(observational_data_usage_info(&[stalled], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3).is_none());
}
#[test]
fn set_membership_add_remove_uses_generation_and_tombstone() {
let state = DataUsageSnapshotSetState {
pool_index: 1,
set_index: 2,
scanner_cycle: Some(9),
scanner_epoch: Some(4),
scan_plan_digest: Some(TEST_PLAN_DIGEST.0),
complete: false,
tombstone: true,
};
let encoded = serde_json::to_vec(&state).expect("set state should serialize");
let decoded: DataUsageSnapshotSetState = serde_json::from_slice(&encoded).expect("set state should deserialize");
assert_eq!(decoded, state);
let snapshot = DataUsageInfo {
last_update: Some(SystemTime::UNIX_EPOCH + Duration::from_secs(10)),
scanner_cycle: Some(9),
scanner_epoch: Some(4),
buckets_count: 0,
usage_snapshot_converged: Some(false),
usage_snapshot_partial: true,
usage_snapshot_set_states: vec![
DataUsageSnapshotSetState {
pool_index: 0,
set_index: 0,
scanner_cycle: Some(9),
scanner_epoch: Some(4),
scan_plan_digest: Some(TEST_PLAN_DIGEST.0),
complete: true,
tombstone: false,
},
state,
],
..Default::default()
};
assert!(snapshot.is_valid_partial_snapshot());
}
#[test]
fn old_set_completion_cannot_overwrite_new_aggregate() {
let all_buckets = vec!["bucket".to_string()];
let source = DataUsageCacheSource::new(0, 0);
let mut old = completed_root_cache("bucket", 1, 20, source);
old.info.next_cycle = 7;
old.info.leader_epoch = 2;
let expected = HashSet::from([source]);
assert!(observational_data_usage_info(&[old], &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3).is_none());
}
#[test]
fn usage_aggregate_survives_restart_and_leader_failover() {
let all_buckets = vec!["bucket".to_string()];
let source = DataUsageCacheSource::new(0, 0);
let mut lkg = lkg_root_cache("bucket", 5, source);
lkg.info.lkg_leader_epoch = Some(4);
lkg.info.lkg_next_cycle = Some(9);
let expected = HashSet::from([source]);
let (observed, _) = observational_data_usage_info(&[lkg], &expected, &all_buckets, TEST_PLAN_DIGEST, 10, 5)
.expect("compatible LKG should survive a leader change");
assert_eq!(observed.usage_snapshot_set_states[0].scanner_epoch, Some(4));
assert_eq!(observed.objects_total_size, 50);
}
#[test]
fn usage_aggregate_cost_is_linear_in_set_count() {
let all_buckets = vec!["bucket".to_string()];
let mut results = Vec::new();
let mut expected = HashSet::new();
for index in 0..32 {
let source = DataUsageCacheSource::new(index, 0);
expected.insert(source);
let mut cache = completed_root_cache("bucket", 1, 20, source);
cache.info.next_cycle = 8;
cache.info.leader_epoch = 3;
results.push(cache);
}
let (observed, _) = observational_data_usage_info(&results, &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
.expect("all set snapshots should aggregate");
assert_eq!(observed.objects_total_count, 32);
let reversed = results.iter().rev().cloned().collect::<Vec<_>>();
let (reversed_observed, _) = observational_data_usage_info(&reversed, &expected, &all_buckets, TEST_PLAN_DIGEST, 8, 3)
.expect("reordered set snapshots should aggregate");
assert_eq!(observed.usage_snapshot_set_states, reversed_observed.usage_snapshot_set_states);
}
#[test]
fn completed_data_usage_info_publishes_tier_stats_across_sets() {
let all_buckets = vec!["bucket-a".to_string(), "bucket-b".to_string()];
+26 -11
View File
@@ -34,8 +34,8 @@ CONCURRENCY=8
DURATION="60s"
ROUNDS=3
COOLDOWN_SECS=20
DATASET_SETUP_DURATION="10s"
HEALTH_TIMEOUT_SECS=180
DATASET_OBJECTS_PER_WORKER=8
FAIL_PCT=10
WARN_PCT=5
ALLOW_REGRESSION=false
@@ -100,6 +100,9 @@ Benchmark:
--duration <dur> warp duration per cell (default 60s).
--rounds <n> rounds per cell; must be >= 3 (default 3).
--cooldown <n> cooldown seconds between rounds/sizes (default 20).
--dataset-setup-duration <dur>
isolated Warp PUT warm-up for get/mixed legs
(default 10s; not included in the measurement).
--concurrency <n> warp concurrency (default 8).
--warp-bin <path> warp binary (default warp).
@@ -170,6 +173,7 @@ while [[ $# -gt 0 ]]; do
--duration) DURATION="$2"; shift 2 ;;
--rounds) ROUNDS="$2"; shift 2 ;;
--cooldown) COOLDOWN_SECS="$2"; shift 2 ;;
--dataset-setup-duration) DATASET_SETUP_DURATION="$2"; shift 2 ;;
--health-timeout) HEALTH_TIMEOUT_SECS="$2"; shift 2 ;;
--fail-pct) FAIL_PCT="$2"; shift 2 ;;
--warn-pct) WARN_PCT="$2"; shift 2 ;;
@@ -362,18 +366,29 @@ measure() {
--duration "$DURATION" --rounds "$ROUNDS" --cooldown-secs "$COOLDOWN_SECS"
--out-dir "$cell"
)
if [[ "$mode" != "put" ]]; then
# Warp defaults to 2,500 setup objects per round. At 10 MiB that writes
# 25 GiB before every 12-second measurement, so the matrix cannot finish
# inside the workflow budget. Eight objects per worker keeps preparation
# bounded while retaining a multi-object working set for relative A/B.
args+=(--extra-args "--objects $((CONCURRENCY * DATASET_OBJECTS_PER_WORKER)) --noclear")
fi
[[ "$mode" == "put" ]] || args+=(--extra-args "--noclear")
[[ -n "$baseline_csv" ]] && args+=(--baseline-csv "$baseline_csv")
run "$ENHANCED_BENCH" "${args[@]}" >&2
echo "$cell"
}
prepare_dataset() {
local leg="$1" workload="$2" mode="$3" size="$4" sync_label="$5" bucket="$6"
[[ "$mode" != "put" ]] || return 0
local setup_cell="$OUT_DIR/$workload/$sync_label/$leg/dataset-setup"
local args=(
--tool warp --warp-bin "$WARP_BIN" --warp-mode put
--endpoint "$ADDRESS" --access-key "$ACCESS_KEY" --secret-key "$SECRET_KEY"
--region "$REGION" --bucket "$bucket" --sizes "$size" --concurrency "$CONCURRENCY"
--duration "$DATASET_SETUP_DURATION" --rounds 1 --cooldown-secs 0
--extra-args "--noclear"
--out-dir "$setup_cell"
)
log "preparing isolated dataset: $sync_label/$workload/$leg bucket=$bucket"
run "$ENHANCED_BENCH" "${args[@]}" >&2
}
write_schedule_header() {
echo "sync_label,drive_sync,workload,mode,size,leg,phase,binary,out_dir,bucket,dataset_setup" >"$OUT_DIR/abba_schedule.csv"
}
@@ -384,7 +399,7 @@ append_schedule() {
phase="$(phase_for_leg "$leg")"
bin="$(binary_for_leg "$leg")"
local dataset_setup="none"
[[ "$mode" == "put" ]] || dataset_setup="warp-native-bounded"
[[ "$mode" == "put" ]] || dataset_setup="warp-put"
echo "$sync_label,$drive_sync,$workload,$mode,$size,$leg,$phase,$bin,$OUT_DIR/$workload/$sync_label/$leg,$bucket,$dataset_setup" >>"$OUT_DIR/abba_schedule.csv"
}
@@ -466,8 +481,7 @@ dataset_namespace=$DATASET_NAMESPACE
local_run_data_root=$RUN_DATA_ROOT
bucket_isolation=per-leg
bucket_prefix=rustfs-abba-$DATASET_NAMESPACE
dataset_setup=get-and-mixed-via-bounded-warp-native
dataset_objects=$((CONCURRENCY * DATASET_OBJECTS_PER_WORKER))
dataset_setup=get-and-mixed-via-warp-put
endpoint=$ADDRESS
warp_version=$("$WARP_BIN" --version 2>/dev/null | head -n1 || echo unknown)
EOF
@@ -488,6 +502,7 @@ for ds_spec in "${DRIVE_SYNC_MATRIX[@]}"; do
log "=== $sync_label $workload leg $leg ($(phase_for_leg "$leg")) ==="
bucket="$(bucket_for_leg "$sync_label" "$workload" "$leg")"
bring_up "$leg" "$drive_sync" "$workload" "$mode" "$size" "$sync_label" "$bucket"
prepare_dataset "$leg" "$workload" "$mode" "$size" "$sync_label" "$bucket"
append_schedule "$sync_label" "$drive_sync" "$workload" "$mode" "$size" "$leg" "$bucket"
baseline_csv=""
+14 -29
View File
@@ -672,7 +672,7 @@ extract_report_line() {
local regex="$1"
local file="$2"
awk -v regex="$regex" '
/^(Report|Operation):/ {
/^Report:/ {
in_report = 1
next
}
@@ -703,9 +703,9 @@ normalize_duration_metric() {
extract_metrics() {
local log_file="$1"
local average_line request_line throughput reqps latency req_p90 req_p99 reqps_num
local average_line reqs_line throughput reqps latency req_p90 req_p99 reqps_num
average_line="$(extract_report_line '^[[:space:]]*[*][[:space:]]+Average:' "$log_file")"
request_line="$(extract_report_line '^[[:space:]]*[*][[:space:]]+(Reqs:[[:space:]]+)?Avg:' "$log_file")"
reqs_line="$(extract_report_line '^[[:space:]]*[*][[:space:]]+Reqs:' "$log_file")"
if [[ -n "$average_line" ]]; then
throughput="$(echo "$average_line" | sed -E 's/^.*Average:[[:space:]]*//; s/,[[:space:]]*.*$//')"
@@ -715,16 +715,20 @@ extract_metrics() {
reqps="$(extract_first '[0-9]+(\.[0-9]+)?[[:space:]]*(obj/s|req/s|ops/s|requests/s)' "$log_file")"
fi
if [[ -n "$request_line" ]]; then
latency="$(echo "$request_line" | rg -o 'Avg:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' | sed -E 's/^Avg:[[:space:]]+//')"
req_p90="$(echo "$request_line" | rg -o '90%:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' | sed -E 's/^90%:[[:space:]]+//')"
req_p99="$(echo "$request_line" | rg -o '99%:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' | sed -E 's/^99%:[[:space:]]+//')"
if [[ -n "$reqs_line" ]]; then
latency="$(echo "$reqs_line" | rg -o 'Avg:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' | sed -E 's/^Avg:[[:space:]]+//')"
req_p90="$(echo "$reqs_line" | rg -o '90%:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' | sed -E 's/^90%:[[:space:]]+//')"
req_p99="$(echo "$reqs_line" | rg -o '99%:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' | sed -E 's/^99%:[[:space:]]+//')"
else
latency="$(rg -o 'Reqs:[[:space:]]+Avg:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' "$log_file" | head -n1 | sed -E 's/^Reqs:[[:space:]]+Avg:[[:space:]]+//')"
req_p90="$(rg -o '90%:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' "$log_file" | head -n1 | sed -E 's/^90%:[[:space:]]+//')"
req_p99="$(rg -o '99%:[[:space:]]+[0-9]+(\.[0-9]+)?(ms|us|µs|s)' "$log_file" | head -n1 | sed -E 's/^99%:[[:space:]]+//')"
fi
if [[ -z "$latency" ]]; then
latency="$(extract_first '[0-9]+(\.[0-9]+)?[[:space:]]*(ms|us|µs|s)' "$log_file")"
fi
throughput="$(trim "${throughput:-N/A}")"
reqps="$(trim "${reqps:-N/A}")"
latency="$(trim "${latency:-N/A}")"
@@ -1124,8 +1128,6 @@ run_one_attempt() {
"--concurrent" "$CONCURRENCY"
"--duration" "$DURATION"
"--region" "$REGION"
"--no-color"
"--analyze.v"
)
if [[ "$INSECURE" == "true" ]]; then
cmd+=("--insecure")
@@ -1210,12 +1212,6 @@ run_one_attempt() {
req_p99_ms="$(to_ms "$req_p99_human")"
fi
if [[ "$DRY_RUN" != "true" && "$TOOL" == "warp" && "$status" == "ok" ]] \
&& rg -q '^[[:space:]]*(Total[[:space:]]+)?Errors:[[:space:]]+[1-9][0-9]*[.]?([[:space:]]|$)' "$log_file"; then
status="failed"
exit_code=1
fi
if [[ "$DRY_RUN" != "true" && "$status" == "ok" ]]; then
if [[ "$throughput_bps" == "N/A" && "$reqps" == "N/A" ]]; then
status="failed"
@@ -1322,21 +1318,10 @@ compare_baseline() {
dr="N/A"; dl="N/A"; dt="N/A"; dp90="N/A"; dp99="N/A"; ne="N/A"; be="N/A"; de="N/A"
if (br!="N/A" && n_req!="N/A" && br+0!=0) dr=sprintf("%.2f", ((n_req-br)/br)*100)
if (bl!="N/A" && n_lat!="N/A") {
if (bl+0!=0) dl=sprintf("%.2f", ((n_lat-bl)/bl)*100)
else if (n_lat+0==0) dl="0.00"
}
if (bl!="N/A" && n_lat!="N/A" && bl+0!=0) dl=sprintf("%.2f", ((n_lat-bl)/bl)*100)
if (bt!="N/A" && n_thr!="N/A" && bt+0!=0) dt=sprintf("%.2f", ((n_thr-bt)/bt)*100)
# Warp v1 rounds sub-millisecond latency to 0s. Two zero readings are
# the same below-resolution bucket; a nonzero candidate remains invalid.
if (bp90!="N/A" && n_p90!="N/A") {
if (bp90+0!=0) dp90=sprintf("%.2f", ((n_p90-bp90)/bp90)*100)
else if (n_p90+0==0) dp90="0.00"
}
if (bp99!="N/A" && n_p99!="N/A") {
if (bp99+0!=0) dp99=sprintf("%.2f", ((n_p99-bp99)/bp99)*100)
else if (n_p99+0==0) dp99="0.00"
}
if (bp90!="N/A" && n_p90!="N/A" && bp90+0!=0) dp90=sprintf("%.2f", ((n_p90-bp90)/bp90)*100)
if (bp99!="N/A" && n_p99!="N/A" && bp99+0!=0) dp99=sprintf("%.2f", ((n_p99-bp99)/bp99)*100)
if (n_ok!="N/A" && n_fail!="N/A" && n_ok+n_fail>0) ne=sprintf("%.2f", (n_fail/(n_ok+n_fail))*100)
if (bok!="N/A" && bfail!="N/A" && bok+bfail>0) be=sprintf("%.2f", (bfail/(bok+bfail))*100)
if (ne!="N/A" && be!="N/A") de=sprintf("%.2f", ne-be)
+2 -7
View File
@@ -58,13 +58,8 @@ rg -qx 'evidence_mode=dry-run' "$OUT_DIR/manifest.env"
rg -qx 'formal_evidence=false' "$OUT_DIR/manifest.env"
rg -qx 'performance_conclusion=not_measured_dry_run' "$OUT_DIR/manifest.env"
rg -qx 'bucket_isolation=per-leg' "$OUT_DIR/manifest.env"
rg -qx 'dataset_setup=get-and-mixed-via-bounded-warp-native' "$OUT_DIR/manifest.env"
rg -qx 'dataset_objects=64' "$OUT_DIR/manifest.env"
[[ "$(rg -c -- '--extra-args --objects\\ 64\\ --noclear' "$TRACE_FILE")" == "32" ]]
if rg -q -- 'dataset-setup' "$TRACE_FILE"; then
echo "unexpected redundant dataset setup command" >&2
exit 1
fi
rg -qx 'dataset_setup=get-and-mixed-via-warp-put' "$OUT_DIR/manifest.env"
[[ "$(rg -c -- '--extra-args --noclear' "$TRACE_FILE")" == "64" ]]
! rg -q -- 'rustfs-bench' "$TRACE_FILE"
if "$RUNNER" \
+2 -67
View File
@@ -74,28 +74,13 @@ FAKE_WARP="${TMP_DIR}/fake-warp"
cat >"$FAKE_WARP" <<'EOF'
#!/usr/bin/env bash
set -euo pipefail
[[ " $* " == *" --analyze.v "* ]]
[[ " $* " == *" --no-color "* ]]
if [[ "${FAKE_WARP_ZERO_LATENCY:-0}" == "1" ]]; then
cat <<'LOG'
Operation: GET. Concurrency: 8. Ran: 7s
Requests considered: 1000:
* Average: 160.00 MiB/s, 40960.00 obj/s
* Avg: 0s, 50%: 0s, 90%: 0s, 99%: 0s, Fastest: 0s, Slowest: 1ms, StdDev: 0s
LOG
exit 0
fi
cat <<'LOG'
- PUT Average: 161 Obj/s, 5.0MiB/s; Current 161 Obj/s, 5.0MiB/s.
Operation: GET. Concurrency: 64. Ran: 7s
Requests considered: 1000:
Report: GET. Concurrency: 64. Ran: 7s
* Average: 653.90 MiB/s, 20925.58 obj/s
* Avg: 3.5ms, 50%: 2.0ms, 90%: 3.6ms, 99%: 24.1ms, Fastest: 0.2ms, Slowest: 607.7ms, StdDev: 20.6ms
* Reqs: Avg: 3.5ms, 50%: 2.0ms, 90%: 3.6ms, 99%: 24.1ms, Fastest: 0.2ms, Slowest: 607.7ms, StdDev: 20.6ms
Throughput, split into 7 x 1s:
LOG
if [[ "${FAKE_WARP_ERRORS:-0}" == "1" ]]; then
echo 'Total Errors: 1.'
fi
EOF
chmod +x "$FAKE_WARP"
@@ -118,32 +103,6 @@ chmod +x "$FAKE_WARP"
rg -q '^32767B,warp,1,1,128,ok,0,[^,]+,[^,]+,653.90 MiB/s,685663846.400000,20925.58,3.5 ms,3.500000,[^,]+,3.6 ms,3.600000,24.1 ms,24.100000$' "${TMP_DIR}/fake-warp-run/round_results.csv"
cat >"${TMP_DIR}/warp-no-details.log" <<'EOF'
warp: Starting benchmark in 3s...
Operation: PUT. Concurrency: 8
* Average: 2.76 MiB/s, 707.03 obj/s
EOF
"$RUNNER" --extract-metrics-from-log "${TMP_DIR}/warp-no-details.log" >"${TMP_DIR}/warp-no-details.csv"
rg -qx '2.76 MiB/s,2894069.760000,707.03,N/A,N/A,N/A,N/A,N/A,N/A' "${TMP_DIR}/warp-no-details.csv"
if FAKE_WARP_ERRORS=1 "$RUNNER" \
--tool warp \
--endpoint http://127.0.0.1:9000 \
--access-key test-access \
--secret-key test-secret \
--sizes 32767B \
--rounds 1 \
--retry-per-round 1 \
--retry-sleep-secs 1 \
--cooldown-secs 0 \
--duration 1s \
--out-dir "${TMP_DIR}/fake-warp-errors" \
--warp-bin "$FAKE_WARP" >/dev/null 2>&1; then
echo "expected Warp request errors to fail the benchmark" >&2
exit 1
fi
rg -q ',failed,1,' "${TMP_DIR}/fake-warp-errors/round_results.csv"
"$RUNNER" \
--tool warp \
--endpoint http://127.0.0.1:9000 \
@@ -173,28 +132,4 @@ awk -F',' '
END { exit found ? 0 : 1 }
' "${TMP_DIR}/fake-warp-candidate/baseline_compare.csv"
for leg in baseline candidate; do
zero_args=(
--tool warp
--endpoint http://127.0.0.1:9000
--access-key test-access
--secret-key test-secret
--sizes 4KiB
--rounds 1
--retry-per-round 1
--cooldown-secs 0
--duration 1s
--out-dir "${TMP_DIR}/fake-warp-zero-${leg}"
--warp-bin "$FAKE_WARP"
)
if [[ "$leg" == "candidate" ]]; then
zero_args+=(--baseline-csv "${TMP_DIR}/fake-warp-zero-baseline/median_summary.csv")
fi
FAKE_WARP_ZERO_LATENCY=1 "$RUNNER" "${zero_args[@]}" >/dev/null 2>&1
done
"${SCRIPT_DIR}/hotpath_warp_ab_gate.sh" \
--compare-csv "${TMP_DIR}/fake-warp-zero-candidate/baseline_compare.csv" \
--require-tail-error >/dev/null
echo "object batch benchmark enhanced tests passed"