mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-28 00:58:59 +00:00
4d22ed4465
perf(capacity): remove per-PUT global lock and per-disk allocation from write dirty-scope Every successful write recorded its capacity dirty scope by allocating an endpoint/path String per online disk, deduplicating through a HashSet, entering the global dirty-scope Mutex, and — in the app response path — taking a global async RwLock to record the write frequency. Under small-object high concurrency this created a global serialization point and O(disks) allocation on the hot path (https://github.com/rustfs/backlog/issues/1315). This change makes the steady-state write path allocation-free and lock-free without altering capacity accounting semantics: - Memoize the per-set dirty scope. Each set resolves its disks' immutable endpoint/path identity lazily into a slot-indexed cache and reuses a shared `Arc<CapacityScope>`; steady-state writes clone the Arc under a read lock instead of rebuilding String/HashSet. The heal path keeps an ad-hoc scope builder because it passes disks in erasure-distribution order rather than physical-slot order. - Add a monotonic generation to the global dirty-scope registry, advanced only when a non-empty drain removes disks. A set upgrades the global registry mutex only on the first write of each generation and then skips it while the generation is unchanged; the observed generation is read under the registry lock so a concurrent drain forces a re-mark, preventing lost updates. The write commits its bytes before recording the scope, so any drain that could remove the mark is ordered after the commit and the following refresh reads the committed bytes. - Replace the write-frequency `RwLock<WriteRecord>` with lock-free atomics: per-second CAS buckets, an atomic last-write timestamp, and an atomic total counter. The frequency window and debounce semantics the refresh scheduler relies on are unchanged. Capacity marking remains a conservative superset of the disks actually written, so admin/scan totals are byte-for-byte identical: extra dirty marks only trigger a re-read of a disk whose usage is unchanged. White-box tests assert the memoized scope equals the previous ad-hoc construction, that the global registry is upgraded exactly once per generation and re-marked after a drain, and that the lock-free write record is exact under concurrent contention. Ref: https://github.com/rustfs/backlog/issues/1315
2406 lines
90 KiB
Rust
2406 lines
90 KiB
Rust
// Copyright 2024 RustFS Team
|
|
//
|
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
// you may not use this file except in compliance with the License.
|
|
// You may obtain a copy of the License at
|
|
//
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
//
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
// See the License for the specific language governing permissions and
|
|
// limitations under the License.
|
|
|
|
//! Hybrid Capacity Manager for efficient capacity statistics
|
|
|
|
use super::scan::refresh_capacity_with_scope;
|
|
use super::types::CapacityDiskRef;
|
|
use crate::capacity_scope::{CapacityScope, CapacityScopeDisk, drain_global_dirty_scopes, take_capacity_scope};
|
|
use futures::FutureExt;
|
|
use rustfs_config::{
|
|
DEFAULT_CAPACITY_ENABLE_DYNAMIC_TIMEOUT, DEFAULT_CAPACITY_FOLLOW_SYMLINKS, DEFAULT_CAPACITY_MAX_TIMEOUT_SECS,
|
|
DEFAULT_CAPACITY_METRICS_INTERVAL_SECS, DEFAULT_CAPACITY_MIN_TIMEOUT_SECS, DEFAULT_FAST_UPDATE_THRESHOLD_SECS,
|
|
DEFAULT_MAX_FILES_THRESHOLD, DEFAULT_SAMPLE_RATE, DEFAULT_SCHEDULED_UPDATE_INTERVAL_SECS, DEFAULT_STAT_TIMEOUT_SECS,
|
|
DEFAULT_WRITE_FREQUENCY_THRESHOLD, DEFAULT_WRITE_TRIGGER_DELAY_SECS, ENV_CAPACITY_ENABLE_DYNAMIC_TIMEOUT,
|
|
ENV_CAPACITY_FAST_UPDATE_THRESHOLD, ENV_CAPACITY_FOLLOW_SYMLINKS, ENV_CAPACITY_MAX_FILES_THRESHOLD, ENV_CAPACITY_MAX_TIMEOUT,
|
|
ENV_CAPACITY_METRICS_INTERVAL, ENV_CAPACITY_MIN_TIMEOUT, ENV_CAPACITY_SAMPLE_RATE, ENV_CAPACITY_SCHEDULED_INTERVAL,
|
|
ENV_CAPACITY_STAT_TIMEOUT, ENV_CAPACITY_WRITE_FREQUENCY_THRESHOLD, ENV_CAPACITY_WRITE_TRIGGER_DELAY,
|
|
};
|
|
use rustfs_io_metrics::capacity_metrics::{
|
|
record_capacity_current_bytes, record_capacity_degraded_reading, record_capacity_dirty_disk_count,
|
|
record_capacity_refresh_inflight, record_capacity_refresh_joiner, record_capacity_refresh_result,
|
|
record_capacity_update_completed, record_capacity_update_failed, record_capacity_write_operation,
|
|
};
|
|
use rustfs_utils::{get_env_bool, get_env_u64, get_env_usize};
|
|
use std::collections::{HashMap, HashSet};
|
|
use std::future::Future;
|
|
use std::panic::AssertUnwindSafe;
|
|
use std::sync::Arc;
|
|
use std::sync::atomic::{AtomicU64, Ordering};
|
|
use std::time::{Duration, Instant};
|
|
use tokio::sync::{Mutex, RwLock, watch};
|
|
use tracing::{debug, info, warn};
|
|
|
|
const LOG_COMPONENT_CAPACITY: &str = "capacity";
|
|
const LOG_SUBSYSTEM_REFRESH: &str = "refresh";
|
|
const LOG_SUBSYSTEM_RUNTIME: &str = "runtime";
|
|
const EVENT_CAPACITY_REFRESH_CACHE_UPDATED: &str = "capacity_refresh_cache_updated";
|
|
const EVENT_CAPACITY_REFRESH_WRITE_RECORDED: &str = "capacity_refresh_write_recorded";
|
|
const EVENT_CAPACITY_REFRESH_DEBOUNCE_STATE: &str = "capacity_refresh_debounce_state";
|
|
const EVENT_CAPACITY_REFRESH_PANIC: &str = "capacity_refresh_panic";
|
|
const EVENT_CAPACITY_CONFIG_CLAMPED: &str = "capacity_config_clamped";
|
|
const EVENT_CAPACITY_REFRESH_JOINER_TIMEOUT: &str = "capacity_refresh_joiner_timeout";
|
|
const EVENT_CAPACITY_REFRESH_CANCELLED: &str = "capacity_refresh_cancelled";
|
|
const EVENT_CAPACITY_REFRESH_RUNTIME_SUMMARY: &str = "capacity_refresh_runtime_summary";
|
|
const EVENT_CAPACITY_REFRESH_INTERVAL_CLAMPED: &str = "capacity_refresh_interval_clamped";
|
|
const EVENT_CAPACITY_REFRESH_SCHEDULED: &str = "capacity_refresh_scheduled";
|
|
const EVENT_CAPACITY_REFRESH_SKIPPED: &str = "capacity_refresh_skipped";
|
|
|
|
// ============================================================================
|
|
// Configuration Functions
|
|
// ============================================================================
|
|
|
|
/// Cached capacity configuration to avoid repeated environment variable reads
|
|
#[derive(Clone, Debug)]
|
|
struct CachedCapacityConfig {
|
|
/// Scheduled update interval
|
|
scheduled_update_interval: Duration,
|
|
/// Write trigger delay
|
|
write_trigger_delay: Duration,
|
|
/// Write frequency threshold
|
|
write_frequency_threshold: usize,
|
|
/// Fast update threshold
|
|
fast_update_threshold: Duration,
|
|
/// Max files threshold for sampling
|
|
max_files_threshold: usize,
|
|
/// Stat timeout
|
|
stat_timeout: Duration,
|
|
/// Sample rate
|
|
sample_rate: usize,
|
|
/// Metrics logging interval
|
|
metrics_interval: Duration,
|
|
/// Follow symlinks flag
|
|
follow_symlinks: bool,
|
|
/// Enable dynamic timeout flag
|
|
enable_dynamic_timeout: bool,
|
|
/// Min timeout
|
|
min_timeout: Duration,
|
|
/// Max timeout
|
|
max_timeout: Duration,
|
|
}
|
|
|
|
/// Fall back to the default when a capacity env var is set to a value that
|
|
/// would make every scan report 0 bytes or always fail (backlog#1019 S11).
|
|
fn env_u64_at_least(env: &'static str, default: u64, min: u64) -> u64 {
|
|
let value = get_env_u64(env, default);
|
|
if value < min {
|
|
warn!(
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_REFRESH,
|
|
event = EVENT_CAPACITY_CONFIG_CLAMPED,
|
|
env,
|
|
configured = value,
|
|
effective = default,
|
|
"capacity configuration value clamped to default"
|
|
);
|
|
default
|
|
} else {
|
|
value
|
|
}
|
|
}
|
|
|
|
impl CachedCapacityConfig {
|
|
/// Build configuration from environment variables
|
|
fn from_env() -> Self {
|
|
Self {
|
|
scheduled_update_interval: Duration::from_secs(get_env_u64(
|
|
ENV_CAPACITY_SCHEDULED_INTERVAL,
|
|
DEFAULT_SCHEDULED_UPDATE_INTERVAL_SECS,
|
|
)),
|
|
write_trigger_delay: Duration::from_secs(get_env_u64(
|
|
ENV_CAPACITY_WRITE_TRIGGER_DELAY,
|
|
DEFAULT_WRITE_TRIGGER_DELAY_SECS,
|
|
)),
|
|
write_frequency_threshold: get_env_usize(ENV_CAPACITY_WRITE_FREQUENCY_THRESHOLD, DEFAULT_WRITE_FREQUENCY_THRESHOLD),
|
|
fast_update_threshold: Duration::from_secs(get_env_u64(
|
|
ENV_CAPACITY_FAST_UPDATE_THRESHOLD,
|
|
DEFAULT_FAST_UPDATE_THRESHOLD_SECS,
|
|
)),
|
|
max_files_threshold: env_u64_at_least(ENV_CAPACITY_MAX_FILES_THRESHOLD, DEFAULT_MAX_FILES_THRESHOLD as u64, 1)
|
|
as usize,
|
|
stat_timeout: Duration::from_secs(env_u64_at_least(ENV_CAPACITY_STAT_TIMEOUT, DEFAULT_STAT_TIMEOUT_SECS, 1)),
|
|
sample_rate: get_env_usize(ENV_CAPACITY_SAMPLE_RATE, DEFAULT_SAMPLE_RATE),
|
|
metrics_interval: Duration::from_secs(get_env_u64(
|
|
ENV_CAPACITY_METRICS_INTERVAL,
|
|
DEFAULT_CAPACITY_METRICS_INTERVAL_SECS,
|
|
)),
|
|
follow_symlinks: get_env_bool(ENV_CAPACITY_FOLLOW_SYMLINKS, DEFAULT_CAPACITY_FOLLOW_SYMLINKS),
|
|
enable_dynamic_timeout: get_env_bool(ENV_CAPACITY_ENABLE_DYNAMIC_TIMEOUT, DEFAULT_CAPACITY_ENABLE_DYNAMIC_TIMEOUT),
|
|
min_timeout: Duration::from_secs(env_u64_at_least(ENV_CAPACITY_MIN_TIMEOUT, DEFAULT_CAPACITY_MIN_TIMEOUT_SECS, 1)),
|
|
max_timeout: Duration::from_secs(env_u64_at_least(ENV_CAPACITY_MAX_TIMEOUT, DEFAULT_CAPACITY_MAX_TIMEOUT_SECS, 1)),
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Get cached capacity configuration (reads environment variables once)
|
|
#[cfg(not(test))]
|
|
fn get_cached_config() -> &'static CachedCapacityConfig {
|
|
static CONFIG: std::sync::OnceLock<CachedCapacityConfig> = std::sync::OnceLock::new();
|
|
CONFIG.get_or_init(CachedCapacityConfig::from_env)
|
|
}
|
|
|
|
#[cfg(test)]
|
|
fn get_cached_config() -> CachedCapacityConfig {
|
|
// Don't cache in tests to allow temp_env::with_var to work
|
|
CachedCapacityConfig::from_env()
|
|
}
|
|
|
|
/// Get scheduled update interval from environment or default
|
|
#[cfg(not(test))]
|
|
pub fn get_scheduled_update_interval() -> Duration {
|
|
get_cached_config().scheduled_update_interval
|
|
}
|
|
|
|
/// Get scheduled update interval from environment or default (test mode)
|
|
#[cfg(test)]
|
|
pub fn get_scheduled_update_interval() -> Duration {
|
|
get_cached_config().scheduled_update_interval
|
|
}
|
|
|
|
/// Get write trigger delay from environment or default
|
|
#[cfg(not(test))]
|
|
pub fn get_write_trigger_delay() -> Duration {
|
|
get_cached_config().write_trigger_delay
|
|
}
|
|
|
|
/// Get write trigger delay from environment or default (test mode)
|
|
#[cfg(test)]
|
|
pub fn get_write_trigger_delay() -> Duration {
|
|
get_cached_config().write_trigger_delay
|
|
}
|
|
|
|
/// Get write frequency threshold from environment or default
|
|
#[cfg(not(test))]
|
|
pub fn get_write_frequency_threshold() -> usize {
|
|
get_cached_config().write_frequency_threshold
|
|
}
|
|
|
|
/// Get write frequency threshold from environment or default (test mode)
|
|
#[cfg(test)]
|
|
pub fn get_write_frequency_threshold() -> usize {
|
|
get_cached_config().write_frequency_threshold
|
|
}
|
|
|
|
/// Get fast update threshold from environment or default
|
|
#[cfg(not(test))]
|
|
pub fn get_fast_update_threshold() -> Duration {
|
|
get_cached_config().fast_update_threshold
|
|
}
|
|
|
|
/// Get fast update threshold from environment or default (test mode)
|
|
#[cfg(test)]
|
|
pub fn get_fast_update_threshold() -> Duration {
|
|
get_cached_config().fast_update_threshold
|
|
}
|
|
|
|
/// Get max files threshold from environment or default
|
|
#[cfg(not(test))]
|
|
pub fn get_max_files_threshold() -> usize {
|
|
get_cached_config().max_files_threshold
|
|
}
|
|
|
|
/// Get max files threshold from environment or default (test mode)
|
|
#[cfg(test)]
|
|
pub fn get_max_files_threshold() -> usize {
|
|
get_cached_config().max_files_threshold
|
|
}
|
|
|
|
/// Get stat timeout from environment or default
|
|
#[cfg(not(test))]
|
|
pub fn get_stat_timeout() -> Duration {
|
|
get_cached_config().stat_timeout
|
|
}
|
|
|
|
/// Get stat timeout from environment or default (test mode)
|
|
#[cfg(test)]
|
|
pub fn get_stat_timeout() -> Duration {
|
|
get_cached_config().stat_timeout
|
|
}
|
|
|
|
/// Get sample rate from environment or default
|
|
#[cfg(not(test))]
|
|
pub fn get_sample_rate() -> usize {
|
|
get_cached_config().sample_rate
|
|
}
|
|
|
|
/// Get sample rate from environment or default (test mode)
|
|
#[cfg(test)]
|
|
pub fn get_sample_rate() -> usize {
|
|
get_cached_config().sample_rate
|
|
}
|
|
|
|
/// Get capacity metrics logging interval from environment or default
|
|
#[cfg(not(test))]
|
|
pub fn get_metrics_interval() -> Duration {
|
|
get_cached_config().metrics_interval
|
|
}
|
|
|
|
/// Get capacity metrics logging interval from environment or default (test mode)
|
|
#[cfg(test)]
|
|
pub fn get_metrics_interval() -> Duration {
|
|
get_cached_config().metrics_interval
|
|
}
|
|
|
|
/// Get follow symlinks flag from environment or default
|
|
#[cfg(not(test))]
|
|
pub fn get_follow_symlinks() -> bool {
|
|
get_cached_config().follow_symlinks
|
|
}
|
|
|
|
/// Get follow symlinks flag from environment or default (test mode)
|
|
#[cfg(test)]
|
|
pub fn get_follow_symlinks() -> bool {
|
|
get_cached_config().follow_symlinks
|
|
}
|
|
|
|
/// Get enable dynamic timeout flag from environment or default
|
|
#[cfg(not(test))]
|
|
pub fn get_enable_dynamic_timeout() -> bool {
|
|
get_cached_config().enable_dynamic_timeout
|
|
}
|
|
|
|
/// Get enable dynamic timeout flag from environment or default (test mode)
|
|
#[cfg(test)]
|
|
pub fn get_enable_dynamic_timeout() -> bool {
|
|
get_cached_config().enable_dynamic_timeout
|
|
}
|
|
|
|
/// Get min timeout from environment or default
|
|
#[cfg(not(test))]
|
|
pub fn get_min_timeout() -> Duration {
|
|
get_cached_config().min_timeout
|
|
}
|
|
|
|
/// Get min timeout from environment or default (test mode)
|
|
#[cfg(test)]
|
|
pub fn get_min_timeout() -> Duration {
|
|
get_cached_config().min_timeout
|
|
}
|
|
|
|
/// Get max timeout from environment or default
|
|
#[cfg(not(test))]
|
|
pub fn get_max_timeout() -> Duration {
|
|
get_cached_config().max_timeout
|
|
}
|
|
|
|
/// Get max timeout from environment or default (test mode)
|
|
#[cfg(test)]
|
|
pub fn get_max_timeout() -> Duration {
|
|
get_cached_config().max_timeout
|
|
}
|
|
|
|
// ============================================================================
|
|
// Data Structures
|
|
// ============================================================================
|
|
|
|
/// Cached capacity data
|
|
#[derive(Clone, Debug)]
|
|
pub struct CachedCapacity {
|
|
/// Total used capacity in bytes
|
|
pub total_used: u64,
|
|
/// Last update time
|
|
pub last_update: Instant,
|
|
/// File count (optional)
|
|
pub file_count: usize,
|
|
/// Whether it's an estimated value
|
|
pub is_estimated: bool,
|
|
/// Whether the value comes from a refresh that had partial disk failures
|
|
/// (some disks kept their last-known values or were dropped entirely).
|
|
pub degraded: bool,
|
|
/// Data source
|
|
pub source: DataSource,
|
|
}
|
|
|
|
/// Structured capacity update payload.
|
|
#[derive(Clone, Debug)]
|
|
pub struct CapacityUpdate {
|
|
/// Total used capacity in bytes.
|
|
pub total_used: u64,
|
|
/// Number of files observed during scan.
|
|
pub file_count: usize,
|
|
/// Whether the value is estimated instead of exact.
|
|
pub is_estimated: bool,
|
|
/// Whether the refresh behind this update had partial disk failures.
|
|
/// `is_estimated` only reflects sampling; this flag is the only carrier of
|
|
/// the "some disks failed to scan" fact (backlog#1014).
|
|
pub degraded: bool,
|
|
/// Per-disk breakdown captured from a successful refresh.
|
|
pub per_disk: Vec<DiskCapacityUpdate>,
|
|
/// Expected disk count for a complete disk cache.
|
|
pub expected_disk_count: Option<usize>,
|
|
/// Whether this update should replace the current disk cache.
|
|
pub replaces_disk_cache: bool,
|
|
/// Dirty disks that can be cleared after the update is committed.
|
|
pub clear_dirty_disks: Vec<CapacityScopeDisk>,
|
|
/// When the scan behind this update started; dirty marks recorded after
|
|
/// this instant survive the commit (backlog#1020 S19). `None` clears
|
|
/// unconditionally.
|
|
pub scan_started_at: Option<Instant>,
|
|
}
|
|
|
|
impl CapacityUpdate {
|
|
/// Create an exact capacity update.
|
|
pub fn exact(total_used: u64, file_count: usize) -> Self {
|
|
Self {
|
|
total_used,
|
|
file_count,
|
|
is_estimated: false,
|
|
degraded: false,
|
|
per_disk: Vec::new(),
|
|
expected_disk_count: None,
|
|
replaces_disk_cache: false,
|
|
clear_dirty_disks: Vec::new(),
|
|
scan_started_at: None,
|
|
}
|
|
}
|
|
|
|
/// Create an estimated capacity update.
|
|
pub fn estimated(total_used: u64, file_count: usize) -> Self {
|
|
Self {
|
|
total_used,
|
|
file_count,
|
|
is_estimated: true,
|
|
degraded: false,
|
|
per_disk: Vec::new(),
|
|
expected_disk_count: None,
|
|
replaces_disk_cache: false,
|
|
clear_dirty_disks: Vec::new(),
|
|
scan_started_at: None,
|
|
}
|
|
}
|
|
|
|
/// Create a fallback capacity update.
|
|
pub fn fallback(total_used: u64) -> Self {
|
|
Self {
|
|
total_used,
|
|
file_count: 0,
|
|
is_estimated: true,
|
|
degraded: false,
|
|
per_disk: Vec::new(),
|
|
expected_disk_count: None,
|
|
replaces_disk_cache: false,
|
|
clear_dirty_disks: Vec::new(),
|
|
scan_started_at: None,
|
|
}
|
|
}
|
|
}
|
|
|
|
#[derive(Clone, Debug)]
|
|
pub struct DiskCapacityUpdate {
|
|
pub disk: CapacityScopeDisk,
|
|
pub used_bytes: u64,
|
|
pub file_count: usize,
|
|
pub is_estimated: bool,
|
|
}
|
|
|
|
#[derive(Clone, Debug)]
|
|
struct CachedDiskCapacity {
|
|
used_bytes: u64,
|
|
file_count: usize,
|
|
is_estimated: bool,
|
|
}
|
|
|
|
#[derive(Clone, Debug, PartialEq, Copy, Eq)]
|
|
pub enum DataSource {
|
|
/// Real-time statistics
|
|
RealTime,
|
|
/// Scheduled update
|
|
Scheduled,
|
|
/// Write triggered
|
|
WriteTriggered,
|
|
/// Fallback value
|
|
#[allow(dead_code)]
|
|
Fallback,
|
|
}
|
|
|
|
impl DataSource {
|
|
pub fn as_metric_label(self) -> &'static str {
|
|
match self {
|
|
Self::RealTime => "realtime",
|
|
Self::Scheduled => "scheduled",
|
|
Self::WriteTriggered => "write_triggered",
|
|
Self::Fallback => "fallback",
|
|
}
|
|
}
|
|
}
|
|
|
|
const WRITE_WINDOW_SECS: u64 = 60;
|
|
const WRITE_WINDOW_BUCKETS: usize = WRITE_WINDOW_SECS as usize;
|
|
|
|
/// Upper bound on how long a joiner waits for the in-flight refresh leader.
|
|
///
|
|
/// A healthy full refresh finishes well within this (each disk scan is capped
|
|
/// by the outer wall-clock budget); the bound only exists so admin queries
|
|
/// return a clear error instead of hanging until process restart if the
|
|
/// leader wedges in a way the guards don't cover (backlog#1017).
|
|
const REFRESH_JOINER_WAIT_TIMEOUT: Duration = Duration::from_secs(300);
|
|
|
|
/// Upper bound for background refresh/metrics intervals: 30 days. Far beyond
|
|
/// any sane configuration, but small enough that `Instant + Duration` can
|
|
/// never overflow.
|
|
const MAX_BACKGROUND_INTERVAL: Duration = Duration::from_secs(30 * 24 * 60 * 60);
|
|
|
|
/// A single second-granular write bucket packed into one `u64` so it can be
|
|
/// updated with a lock-free compare-and-swap. The high 32 bits hold the
|
|
/// monotonic-second key, the low 32 bits hold the write count for that second.
|
|
/// A per-second count never approaches `u32::MAX`, and the count saturates
|
|
/// rather than wrapping into the second key (backlog#1315).
|
|
#[derive(Default)]
|
|
struct AtomicWriteBucket(AtomicU64);
|
|
|
|
const WRITE_BUCKET_COUNT_MASK: u64 = 0xFFFF_FFFF;
|
|
|
|
impl AtomicWriteBucket {
|
|
#[inline]
|
|
fn unpack(packed: u64) -> (u64, usize) {
|
|
let second = packed >> 32;
|
|
let count = (packed & WRITE_BUCKET_COUNT_MASK) as usize;
|
|
(second, count)
|
|
}
|
|
|
|
#[inline]
|
|
fn pack(second: u64, count: usize) -> u64 {
|
|
let count = (count as u64).min(WRITE_BUCKET_COUNT_MASK);
|
|
(second << 32) | count
|
|
}
|
|
|
|
/// Record one write for `now_second`, resetting the bucket if it holds a
|
|
/// different (older or future) second. Lock-free CAS loop tolerating
|
|
/// concurrent writers landing on the same bucket.
|
|
fn record(&self, now_second: u64) {
|
|
loop {
|
|
let current = self.0.load(Ordering::Acquire);
|
|
let (second, count) = Self::unpack(current);
|
|
let next = if second == now_second {
|
|
Self::pack(now_second, count.saturating_add(1))
|
|
} else {
|
|
Self::pack(now_second, 1)
|
|
};
|
|
if self
|
|
.0
|
|
.compare_exchange_weak(current, next, Ordering::AcqRel, Ordering::Acquire)
|
|
.is_ok()
|
|
{
|
|
return;
|
|
}
|
|
}
|
|
}
|
|
|
|
fn snapshot(&self) -> (u64, usize) {
|
|
Self::unpack(self.0.load(Ordering::Acquire))
|
|
}
|
|
|
|
#[cfg(test)]
|
|
fn store(&self, second: u64, count: usize) {
|
|
self.0.store(Self::pack(second, count), Ordering::Release);
|
|
}
|
|
}
|
|
|
|
/// Lock-free write record for tracking write operations.
|
|
///
|
|
/// Previously guarded by an async `RwLock` taken on every successful write in
|
|
/// the PUT response path; the write side is now a set of relaxed/CAS atomic
|
|
/// updates so concurrent small-object PUTs no longer serialize on a single
|
|
/// writer lock (backlog#1315). Counters use saturating arithmetic and monotonic
|
|
/// second keys, preserving the exact debounce/frequency semantics the refresh
|
|
/// scheduler relies on.
|
|
pub struct WriteRecord {
|
|
/// Last write time, encoded as monotonic nanoseconds since `epoch`.
|
|
last_write_nanos: AtomicU64,
|
|
/// Whether any write has been recorded (`0` == never). Kept separate from
|
|
/// `last_write_nanos` so a genuine near-zero-nanos first write is still
|
|
/// distinguished from "no write yet".
|
|
has_write: AtomicU64,
|
|
/// Total write count (saturating). Observability only.
|
|
write_count: AtomicU64,
|
|
/// Fixed-size time buckets for the recent write window.
|
|
write_buckets: [AtomicWriteBucket; WRITE_WINDOW_BUCKETS],
|
|
/// Monotonic origin for bucket keys. Wall-clock keys made an NTP step
|
|
/// backwards mark recent buckets as "future" and silently suppress
|
|
/// write-triggered refreshes until the clock catches up (backlog#1022 S32).
|
|
epoch: Instant,
|
|
}
|
|
|
|
impl WriteRecord {
|
|
fn new() -> Self {
|
|
Self {
|
|
last_write_nanos: AtomicU64::new(0),
|
|
has_write: AtomicU64::new(0),
|
|
write_count: AtomicU64::new(0),
|
|
write_buckets: std::array::from_fn(|_| AtomicWriteBucket::default()),
|
|
epoch: Instant::now(),
|
|
}
|
|
}
|
|
|
|
fn monotonic_second(&self) -> u64 {
|
|
self.epoch.elapsed().as_secs()
|
|
}
|
|
|
|
/// Total number of writes recorded (saturating). Observability only.
|
|
fn total_write_count(&self) -> u64 {
|
|
self.write_count.load(Ordering::Relaxed)
|
|
}
|
|
|
|
/// Elapsed time since the last write, or `None` if no write has been
|
|
/// recorded yet.
|
|
fn time_since_last_write(&self) -> Option<Duration> {
|
|
if self.has_write.load(Ordering::Acquire) == 0 {
|
|
return None;
|
|
}
|
|
let last = self.last_write_nanos.load(Ordering::Acquire);
|
|
let now = self.epoch.elapsed().as_nanos() as u64;
|
|
Some(Duration::from_nanos(now.saturating_sub(last)))
|
|
}
|
|
|
|
fn recent_write_count(&self, now_second: u64) -> usize {
|
|
self.write_buckets
|
|
.iter()
|
|
.filter_map(|bucket| {
|
|
let (second, count) = bucket.snapshot();
|
|
if count > 0 && second <= now_second && now_second.saturating_sub(second) < WRITE_WINDOW_SECS {
|
|
Some(count)
|
|
} else {
|
|
None
|
|
}
|
|
})
|
|
.sum()
|
|
}
|
|
|
|
fn record_write(&self, _now: Instant) -> usize {
|
|
let now_second = self.monotonic_second();
|
|
let bucket_idx = (now_second % WRITE_WINDOW_BUCKETS as u64) as usize;
|
|
self.write_buckets[bucket_idx].record(now_second);
|
|
|
|
self.last_write_nanos
|
|
.store(self.epoch.elapsed().as_nanos() as u64, Ordering::Release);
|
|
self.has_write.store(1, Ordering::Release);
|
|
// Single atomic add (wraps only after 2^64 writes; effectively saturating
|
|
// for any real deployment). Observability only — the frequency window
|
|
// used for refresh decisions is tracked per-bucket above.
|
|
self.write_count.fetch_add(1, Ordering::Relaxed);
|
|
|
|
self.recent_write_count(now_second)
|
|
}
|
|
}
|
|
|
|
/// Hybrid strategy configuration
|
|
#[derive(Debug, Clone)]
|
|
#[allow(dead_code)]
|
|
pub struct HybridStrategyConfig {
|
|
/// Scheduled update interval
|
|
pub scheduled_update_interval: Duration,
|
|
/// Write trigger delay
|
|
pub write_trigger_delay: Duration,
|
|
/// Write frequency threshold (writes/minute)
|
|
pub write_frequency_threshold: usize,
|
|
/// Fast update threshold
|
|
pub fast_update_threshold: Duration,
|
|
/// Metrics logging interval
|
|
pub metrics_interval: Duration,
|
|
/// Enable smart update
|
|
pub enable_smart_update: bool,
|
|
/// Enable write trigger
|
|
pub enable_write_trigger: bool,
|
|
}
|
|
|
|
impl Default for HybridStrategyConfig {
|
|
fn default() -> Self {
|
|
Self {
|
|
scheduled_update_interval: get_scheduled_update_interval(),
|
|
write_trigger_delay: get_write_trigger_delay(),
|
|
write_frequency_threshold: get_write_frequency_threshold(),
|
|
fast_update_threshold: get_fast_update_threshold(),
|
|
metrics_interval: get_metrics_interval(),
|
|
enable_smart_update: true,
|
|
enable_write_trigger: true,
|
|
}
|
|
}
|
|
}
|
|
|
|
impl HybridStrategyConfig {
|
|
/// Create config from environment variables
|
|
pub fn from_env() -> Self {
|
|
Self::default()
|
|
}
|
|
}
|
|
|
|
// ============================================================================
|
|
// Hybrid Capacity Manager
|
|
// ============================================================================
|
|
|
|
struct RefreshState {
|
|
running: bool,
|
|
/// Sender for the current refresh cycle. Joiners subscribe to this before releasing the
|
|
/// mutex so they cannot miss the completion notification. A new channel is created at the
|
|
/// start of every refresh cycle so stale subscribers from previous cycles are not confused
|
|
/// by results that were already published.
|
|
result_tx: watch::Sender<Option<Result<CapacityUpdate, String>>>,
|
|
}
|
|
|
|
impl Default for RefreshState {
|
|
fn default() -> Self {
|
|
let (tx, _) = watch::channel(None);
|
|
Self {
|
|
running: false,
|
|
result_tx: tx,
|
|
}
|
|
}
|
|
}
|
|
|
|
fn reset_cancelled_refresh_state(state: &mut RefreshState) {
|
|
state.running = false;
|
|
record_capacity_refresh_inflight(0);
|
|
let _ = state
|
|
.result_tx
|
|
.send(Some(Err("capacity refresh leader was cancelled".to_string())));
|
|
}
|
|
|
|
/// Resets the singleflight leader state if the leading future is dropped before the
|
|
/// refresh cycle completes (e.g. the admin request that became leader is cancelled by
|
|
/// a client disconnect). Without this, `running` stays `true` forever: joiners block
|
|
/// indefinitely and no future refresh can start. `catch_unwind` covers panics but not
|
|
/// cancellation, so the reset must live in `Drop`.
|
|
struct RefreshLeaderGuard {
|
|
state: Option<Arc<Mutex<RefreshState>>>,
|
|
}
|
|
|
|
impl RefreshLeaderGuard {
|
|
fn disarm(&mut self) {
|
|
self.state = None;
|
|
}
|
|
}
|
|
|
|
impl Drop for RefreshLeaderGuard {
|
|
fn drop(&mut self) {
|
|
let Some(state) = self.state.take() else {
|
|
return;
|
|
};
|
|
warn!(
|
|
event = EVENT_CAPACITY_REFRESH_CANCELLED,
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_REFRESH,
|
|
result = "cancelled",
|
|
"capacity refresh leader dropped before completing; resetting refresh state"
|
|
);
|
|
if let Ok(mut guard) = state.try_lock() {
|
|
reset_cancelled_refresh_state(&mut guard);
|
|
return;
|
|
}
|
|
// The mutex is momentarily held by a joiner subscribing; finish the
|
|
// reset from a detached task since Drop cannot await.
|
|
match tokio::runtime::Handle::try_current() {
|
|
Ok(handle) => {
|
|
handle.spawn(async move {
|
|
reset_cancelled_refresh_state(&mut *state.lock().await);
|
|
});
|
|
}
|
|
Err(_) => {
|
|
// Dropped outside any tokio runtime (e.g. block_on teardown).
|
|
// Blocking is allowed here precisely because there is no
|
|
// executor to stall; skipping the reset would wedge the
|
|
// singleflight forever (backlog#1021 S31).
|
|
reset_cancelled_refresh_state(&mut state.blocking_lock());
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Hybrid capacity manager
|
|
pub struct HybridCapacityManager {
|
|
/// Capacity cache
|
|
cache: Arc<RwLock<Option<CachedCapacity>>>,
|
|
/// Write record. Lock-free atomics (backlog#1315): the PUT response path no
|
|
/// longer takes an async writer lock to record a write.
|
|
write_record: Arc<WriteRecord>,
|
|
/// Dirty disks recorded from write-side scope propagation, keyed to the
|
|
/// instant they were last marked so a commit only clears marks that
|
|
/// predate its scan (backlog#1020 S19).
|
|
dirty_disks: Arc<RwLock<HashMap<CapacityScopeDisk, Instant>>>,
|
|
/// Per-disk cache populated after a successful full refresh and updated by dirty subset refreshes.
|
|
disk_cache: Arc<RwLock<HashMap<CapacityScopeDisk, CachedDiskCapacity>>>,
|
|
/// Whether the per-disk cache currently covers all known disks.
|
|
disk_cache_complete: Arc<RwLock<bool>>,
|
|
/// Configuration
|
|
config: HybridStrategyConfig,
|
|
/// Shared singleflight refresh state
|
|
refresh_state: Arc<Mutex<RefreshState>>,
|
|
}
|
|
|
|
impl HybridCapacityManager {
|
|
async fn sync_global_dirty_scopes(&self) {
|
|
let scopes = drain_global_dirty_scopes();
|
|
if scopes.is_empty() {
|
|
return;
|
|
}
|
|
|
|
let now = Instant::now();
|
|
let mut dirty_disks = self.dirty_disks.write().await;
|
|
dirty_disks.extend(scopes.into_iter().map(|disk| (disk, now)));
|
|
record_capacity_dirty_disk_count(dirty_disks.len());
|
|
}
|
|
|
|
fn max_stale_age(&self) -> Duration {
|
|
self.config
|
|
.scheduled_update_interval
|
|
.max(self.config.fast_update_threshold.checked_mul(3).unwrap_or(Duration::MAX))
|
|
}
|
|
|
|
/// Create a new hybrid capacity manager
|
|
pub fn new(config: HybridStrategyConfig) -> Self {
|
|
Self {
|
|
cache: Arc::new(RwLock::new(None)),
|
|
write_record: Arc::new(WriteRecord::new()),
|
|
dirty_disks: Arc::new(RwLock::new(HashMap::new())),
|
|
disk_cache: Arc::new(RwLock::new(HashMap::new())),
|
|
disk_cache_complete: Arc::new(RwLock::new(false)),
|
|
config,
|
|
refresh_state: Arc::new(Mutex::new(RefreshState::default())),
|
|
}
|
|
}
|
|
|
|
/// Create with default config from environment
|
|
pub fn from_env() -> Self {
|
|
Self::new(HybridStrategyConfig::from_env())
|
|
}
|
|
|
|
/// Get capacity (core method)
|
|
pub async fn get_capacity(&self) -> Option<CachedCapacity> {
|
|
let cache = self.cache.read().await;
|
|
cache.clone()
|
|
}
|
|
|
|
/// Update capacity.
|
|
///
|
|
/// Returns the update with its cluster-level aggregates (`total_used`, `file_count`,
|
|
/// `is_estimated`) reconciled against the per-disk cache. A dirty-subset refresh carries
|
|
/// only the scanned subset's totals; callers that publish or return this value (e.g. the
|
|
/// `refresh_or_join` leader and joiners, admin used-capacity) must use the reconciled copy
|
|
/// so they never report the subset bytes as the whole-cluster total.
|
|
pub async fn update_capacity(&self, mut update: CapacityUpdate, source: DataSource) -> CapacityUpdate {
|
|
let start = Instant::now();
|
|
|
|
if !update.per_disk.is_empty() || update.degraded {
|
|
let mut disk_cache = self.disk_cache.write().await;
|
|
let mut disk_cache_complete = self.disk_cache_complete.write().await;
|
|
|
|
let recompute = if !update.per_disk.is_empty()
|
|
&& update.replaces_disk_cache
|
|
&& update.expected_disk_count == Some(update.per_disk.len())
|
|
{
|
|
disk_cache.clear();
|
|
for entry in &update.per_disk {
|
|
disk_cache.insert(
|
|
entry.disk.clone(),
|
|
CachedDiskCapacity {
|
|
used_bytes: entry.used_bytes,
|
|
file_count: entry.file_count,
|
|
is_estimated: entry.is_estimated,
|
|
},
|
|
);
|
|
}
|
|
*disk_cache_complete = true;
|
|
true
|
|
} else if *disk_cache_complete {
|
|
// Merging over a complete cache also covers a degraded
|
|
// (partial-failure) full refresh: successfully scanned disks
|
|
// are refreshed while failed disks keep their last-known
|
|
// values, so the published total does not dip to the
|
|
// surviving-subset sum and bounce back on the next refresh.
|
|
for entry in &update.per_disk {
|
|
disk_cache.insert(
|
|
entry.disk.clone(),
|
|
CachedDiskCapacity {
|
|
used_bytes: entry.used_bytes,
|
|
file_count: entry.file_count,
|
|
is_estimated: entry.is_estimated,
|
|
},
|
|
);
|
|
}
|
|
true
|
|
} else {
|
|
false
|
|
};
|
|
|
|
if recompute {
|
|
// Reconcile cluster-wide aggregates from the full per-disk cache so a
|
|
// dirty-subset refresh reports the merged cluster totals, not the subset sum.
|
|
update.total_used = disk_cache.values().map(|entry| entry.used_bytes).sum();
|
|
update.file_count = disk_cache.values().map(|entry| entry.file_count).sum();
|
|
update.is_estimated = disk_cache.values().any(|entry| entry.is_estimated);
|
|
}
|
|
}
|
|
|
|
let total_used = update.total_used;
|
|
|
|
let mut cache = self.cache.write().await;
|
|
*cache = Some(CachedCapacity {
|
|
total_used,
|
|
last_update: Instant::now(),
|
|
file_count: update.file_count,
|
|
is_estimated: update.is_estimated,
|
|
degraded: update.degraded,
|
|
source,
|
|
});
|
|
|
|
if !update.clear_dirty_disks.is_empty() {
|
|
let mut dirty_disks = self.dirty_disks.write().await;
|
|
for disk in &update.clear_dirty_disks {
|
|
// Only clear marks that predate the scan behind this update: a
|
|
// write landing while the walker was already past its prefix
|
|
// re-marks the disk, and that mark must survive the commit or
|
|
// the next dirty-subset refresh serves stale bytes as exact
|
|
// (backlog#1020 S19).
|
|
let clearable = match (update.scan_started_at, dirty_disks.get(disk)) {
|
|
(Some(scan_start), Some(marked_at)) => *marked_at <= scan_start,
|
|
_ => true,
|
|
};
|
|
if clearable {
|
|
dirty_disks.remove(disk);
|
|
}
|
|
}
|
|
record_capacity_dirty_disk_count(dirty_disks.len());
|
|
}
|
|
|
|
debug!(
|
|
event = EVENT_CAPACITY_REFRESH_CACHE_UPDATED,
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_REFRESH,
|
|
result = "updated",
|
|
total_used,
|
|
file_count = update.file_count,
|
|
estimated = update.is_estimated,
|
|
degraded = update.degraded,
|
|
source = source.as_metric_label(),
|
|
elapsed_ms = start.elapsed().as_millis() as u64,
|
|
"capacity refresh cache updated"
|
|
);
|
|
record_capacity_current_bytes(total_used);
|
|
record_capacity_update_completed(source.as_metric_label(), start.elapsed(), total_used, update.is_estimated);
|
|
if update.degraded {
|
|
record_capacity_degraded_reading(source.as_metric_label());
|
|
}
|
|
update
|
|
}
|
|
|
|
/// Record write operation
|
|
pub async fn record_write_operation(&self) {
|
|
let now = Instant::now();
|
|
let recent_write_count = self.write_record.record_write(now);
|
|
|
|
record_capacity_write_operation(recent_write_count);
|
|
debug!(
|
|
event = EVENT_CAPACITY_REFRESH_WRITE_RECORDED,
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_REFRESH,
|
|
state = "recorded",
|
|
total_writes = self.write_record.total_write_count(),
|
|
recent_writes = recent_write_count,
|
|
"capacity refresh write recorded"
|
|
);
|
|
}
|
|
|
|
/// Record write scope propagated from the storage layer.
|
|
pub async fn mark_dirty_scope(&self, scope: &CapacityScope) {
|
|
if scope.disks.is_empty() {
|
|
return;
|
|
}
|
|
|
|
let now = Instant::now();
|
|
let mut dirty_disks = self.dirty_disks.write().await;
|
|
dirty_disks.extend(scope.disks.iter().cloned().map(|disk| (disk, now)));
|
|
record_capacity_dirty_disk_count(dirty_disks.len());
|
|
}
|
|
|
|
/// Record a write operation and consume any propagated disk scope bound to the token.
|
|
pub async fn record_write_operation_with_scope_token(&self, scope_token: Option<uuid::Uuid>) {
|
|
if let Some(token) = scope_token
|
|
&& let Some(scope) = take_capacity_scope(token)
|
|
{
|
|
self.mark_dirty_scope(&scope).await;
|
|
}
|
|
|
|
self.record_write_operation().await;
|
|
}
|
|
|
|
/// Check if fast update is needed
|
|
pub async fn needs_fast_update(&self) -> bool {
|
|
if !self.config.enable_smart_update {
|
|
return false;
|
|
}
|
|
|
|
let cache = self.cache.read().await;
|
|
if let Some(cached) = cache.as_ref() {
|
|
let cache_age = cached.last_update.elapsed();
|
|
|
|
// Cache is fresh, no need to update
|
|
if cache_age < self.config.fast_update_threshold {
|
|
return false;
|
|
}
|
|
|
|
if !self.config.enable_write_trigger {
|
|
return false;
|
|
}
|
|
|
|
let write_record = &self.write_record;
|
|
let write_frequency = write_record.recent_write_count(write_record.monotonic_second());
|
|
if write_frequency <= self.config.write_frequency_threshold {
|
|
return false;
|
|
}
|
|
|
|
if let Some(time_since_write) = write_record.time_since_last_write() {
|
|
if time_since_write < self.config.write_trigger_delay {
|
|
debug!(
|
|
event = EVENT_CAPACITY_REFRESH_DEBOUNCE_STATE,
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_REFRESH,
|
|
state = "debounced",
|
|
time_since_write_ms = time_since_write.as_millis() as u64,
|
|
trigger_delay_ms = self.config.write_trigger_delay.as_millis() as u64,
|
|
writes_per_minute = write_frequency,
|
|
"capacity refresh debounce state changed"
|
|
);
|
|
return false;
|
|
}
|
|
|
|
debug!(
|
|
event = EVENT_CAPACITY_REFRESH_DEBOUNCE_STATE,
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_REFRESH,
|
|
state = "eligible",
|
|
time_since_write_ms = time_since_write.as_millis() as u64,
|
|
trigger_delay_ms = self.config.write_trigger_delay.as_millis() as u64,
|
|
writes_per_minute = write_frequency,
|
|
"capacity refresh debounce state changed"
|
|
);
|
|
return true;
|
|
}
|
|
}
|
|
|
|
false
|
|
}
|
|
|
|
/// Get cache age
|
|
#[allow(dead_code)]
|
|
pub async fn get_cache_age(&self) -> Option<Duration> {
|
|
let cache = self.cache.read().await;
|
|
cache.as_ref().map(|c| c.last_update.elapsed())
|
|
}
|
|
|
|
/// Get write frequency (writes/minute)
|
|
#[allow(dead_code)]
|
|
pub async fn get_write_frequency(&self) -> usize {
|
|
let record = &self.write_record;
|
|
record.recent_write_count(record.monotonic_second())
|
|
}
|
|
|
|
/// Snapshot the currently dirty disks recorded from write-side scope propagation.
|
|
pub async fn get_dirty_disks(&self) -> Vec<CapacityScopeDisk> {
|
|
self.sync_global_dirty_scopes().await;
|
|
let dirty_disks = self.dirty_disks.read().await;
|
|
dirty_disks.keys().cloned().collect()
|
|
}
|
|
|
|
/// Drop dirty marks for disks outside the current local topology.
|
|
///
|
|
/// The write side records every disk of an EC set — including remote
|
|
/// peers — while scans and dirty clearing only ever cover local disks, so
|
|
/// remote or removed disks would otherwise stay marked forever and keep
|
|
/// the dirty-disk gauge permanently non-zero (backlog#1020 S30).
|
|
pub async fn retain_dirty_disks_within(&self, local: &HashSet<CapacityScopeDisk>) {
|
|
let mut dirty_disks = self.dirty_disks.write().await;
|
|
let before = dirty_disks.len();
|
|
dirty_disks.retain(|disk, _| local.contains(disk));
|
|
if dirty_disks.len() != before {
|
|
record_capacity_dirty_disk_count(dirty_disks.len());
|
|
}
|
|
}
|
|
|
|
/// Returns true if the manager has a complete per-disk cache and can safely refresh only dirty disks.
|
|
pub async fn can_refresh_dirty_subset(&self) -> bool {
|
|
*self.disk_cache_complete.read().await
|
|
}
|
|
|
|
/// Run a singleflight refresh. Callers either join an existing in-flight refresh or become the leader.
|
|
///
|
|
/// Joiners subscribe to the watch channel *before* releasing the mutex, which guarantees
|
|
/// they cannot miss the completion notification even if the leader finishes very quickly.
|
|
pub async fn refresh_or_join<F, Fut>(&self, source: DataSource, refresh_fn: F) -> Result<CapacityUpdate, String>
|
|
where
|
|
F: FnOnce() -> Fut,
|
|
Fut: Future<Output = Result<CapacityUpdate, String>>,
|
|
{
|
|
let maybe_rx = {
|
|
let mut state = self.refresh_state.lock().await;
|
|
if state.running {
|
|
// Subscribe while holding the lock so the send that completes the current
|
|
// refresh cycle cannot happen before we are subscribed.
|
|
record_capacity_refresh_joiner(source.as_metric_label());
|
|
Some(state.result_tx.subscribe())
|
|
} else {
|
|
// Become the leader. Create a fresh channel so that joiners from a previous
|
|
// cycle cannot observe the result that was published for the new cycle.
|
|
let (tx, _) = watch::channel(None);
|
|
state.result_tx = tx;
|
|
state.running = true;
|
|
record_capacity_refresh_inflight(1);
|
|
None
|
|
}
|
|
};
|
|
|
|
if let Some(mut result_rx) = maybe_rx {
|
|
// Wait until the leader publishes Some(result). Because we subscribed before
|
|
// releasing the mutex, we cannot miss the notification. The wait is bounded so
|
|
// a wedged leader degrades admin queries into a clear error, never a hang.
|
|
match tokio::time::timeout(REFRESH_JOINER_WAIT_TIMEOUT, result_rx.wait_for(|v| v.is_some())).await {
|
|
Err(_) => {
|
|
warn!(
|
|
event = EVENT_CAPACITY_REFRESH_JOINER_TIMEOUT,
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_REFRESH,
|
|
result = "timeout",
|
|
source = source.as_metric_label(),
|
|
waited_ms = REFRESH_JOINER_WAIT_TIMEOUT.as_millis() as u64,
|
|
"capacity refresh joiner timed out waiting for the leader"
|
|
);
|
|
return Err("timed out waiting for the in-flight capacity refresh to publish a result".to_string());
|
|
}
|
|
// The leader's sender was dropped (e.g. due to a panic) without publishing
|
|
// a result. Surface a clear error rather than silently returning the default.
|
|
Ok(Err(_)) => return Err("capacity refresh leader exited without publishing a result".to_string()),
|
|
Ok(Ok(_)) => {}
|
|
}
|
|
return result_rx
|
|
.borrow()
|
|
.as_ref()
|
|
.cloned()
|
|
.unwrap_or_else(|| Err("capacity refresh completed without a result".to_string()));
|
|
}
|
|
|
|
// From here on this future is the leader; if it is dropped at any await point
|
|
// below (request cancellation), the guard resets the singleflight state so
|
|
// joiners unblock and later refreshes are not wedged behind `running = true`.
|
|
let mut leader_guard = RefreshLeaderGuard {
|
|
state: Some(self.refresh_state.clone()),
|
|
};
|
|
|
|
let refresh_start = Instant::now();
|
|
// `refresh_fn` is invoked inside the wrapped future so a panic while
|
|
// *constructing* the future is caught too, not only panics while
|
|
// polling it (backlog#1021 S20).
|
|
let result = AssertUnwindSafe(async move { refresh_fn().await })
|
|
.catch_unwind()
|
|
.await
|
|
.unwrap_or_else(|err| {
|
|
warn!(
|
|
event = EVENT_CAPACITY_REFRESH_PANIC,
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_REFRESH,
|
|
result = "panic",
|
|
source = source.as_metric_label(),
|
|
error = ?err,
|
|
"capacity refresh panicked"
|
|
);
|
|
Err("capacity refresh panicked".to_string())
|
|
});
|
|
// Commit the update and rebind `result` to the reconciled copy so both the value
|
|
// returned to this caller and the one published to joiners carry the cluster totals
|
|
// rather than a dirty-subset's partial bytes.
|
|
let result = match result {
|
|
Ok(update) => Ok(self.update_capacity(update, source).await),
|
|
Err(err) => Err(err),
|
|
};
|
|
let refresh_duration = refresh_start.elapsed();
|
|
if result.is_err() {
|
|
record_capacity_update_failed(source.as_metric_label());
|
|
}
|
|
record_capacity_refresh_result(
|
|
source.as_metric_label(),
|
|
if result.is_ok() { "success" } else { "error" },
|
|
refresh_duration,
|
|
);
|
|
|
|
{
|
|
let mut state = self.refresh_state.lock().await;
|
|
leader_guard.disarm();
|
|
state.running = false;
|
|
record_capacity_refresh_inflight(0);
|
|
let _ = state.result_tx.send(Some(result.clone()));
|
|
}
|
|
|
|
result
|
|
}
|
|
|
|
/// Start a background refresh if one is not already in flight.
|
|
pub async fn spawn_refresh_if_needed<F, Fut>(self: Arc<Self>, source: DataSource, refresh_fn: F) -> bool
|
|
where
|
|
F: FnOnce() -> Fut + Send + 'static,
|
|
Fut: Future<Output = Result<CapacityUpdate, String>> + Send + 'static,
|
|
{
|
|
let should_spawn = {
|
|
let mut state = self.refresh_state.lock().await;
|
|
if state.running {
|
|
false
|
|
} else {
|
|
let (tx, _) = watch::channel(None);
|
|
state.result_tx = tx;
|
|
state.running = true;
|
|
record_capacity_refresh_inflight(1);
|
|
true
|
|
}
|
|
};
|
|
|
|
if !should_spawn {
|
|
return false;
|
|
}
|
|
|
|
tokio::spawn(async move {
|
|
// Symmetric with `refresh_or_join`: `catch_unwind` only covers the
|
|
// refresh future itself, so a panic in the commit/metrics code
|
|
// below would otherwise kill this task before the reset block,
|
|
// leaving `running` true forever — scheduled refreshes silently
|
|
// stop and every joiner hangs (backlog#1021 S20).
|
|
let mut leader_guard = RefreshLeaderGuard {
|
|
state: Some(self.refresh_state.clone()),
|
|
};
|
|
|
|
let refresh_start = Instant::now();
|
|
let result = AssertUnwindSafe(async move { refresh_fn().await })
|
|
.catch_unwind()
|
|
.await
|
|
.unwrap_or_else(|err| {
|
|
warn!(
|
|
event = EVENT_CAPACITY_REFRESH_PANIC,
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_REFRESH,
|
|
result = "panic",
|
|
source = source.as_metric_label(),
|
|
error = ?err,
|
|
"capacity refresh panicked"
|
|
);
|
|
Err("capacity refresh panicked".to_string())
|
|
});
|
|
// Publish the reconciled update to joiners so a dirty-subset refresh does not
|
|
// broadcast a subset's partial bytes as the cluster total.
|
|
let result = match result {
|
|
Ok(update) => Ok(self.update_capacity(update, source).await),
|
|
Err(err) => Err(err),
|
|
};
|
|
let refresh_duration = refresh_start.elapsed();
|
|
if result.is_err() {
|
|
record_capacity_update_failed(source.as_metric_label());
|
|
}
|
|
record_capacity_refresh_result(
|
|
source.as_metric_label(),
|
|
if result.is_ok() { "success" } else { "error" },
|
|
refresh_duration,
|
|
);
|
|
|
|
let mut state = self.refresh_state.lock().await;
|
|
leader_guard.disarm();
|
|
state.running = false;
|
|
record_capacity_refresh_inflight(0);
|
|
let _ = state.result_tx.send(Some(result));
|
|
});
|
|
|
|
true
|
|
}
|
|
|
|
/// Get config
|
|
pub fn get_config(&self) -> &HybridStrategyConfig {
|
|
&self.config
|
|
}
|
|
|
|
/// Check if the cache is too stale to keep serving without a foreground refresh.
|
|
pub fn should_block_on_refresh(&self, cache_age: Duration) -> bool {
|
|
cache_age >= self.max_stale_age()
|
|
}
|
|
|
|
/// Return whether a refresh is currently in flight.
|
|
pub async fn refresh_in_progress(&self) -> bool {
|
|
self.refresh_state.lock().await.running
|
|
}
|
|
|
|
/// Log capacity runtime summary for observability.
|
|
async fn log_runtime_summary(&self) {
|
|
let cached = self.get_capacity().await;
|
|
let recent_write_frequency = self.get_write_frequency().await;
|
|
let dirty_disks = self.get_dirty_disks().await;
|
|
let refresh_running = self.refresh_in_progress().await;
|
|
|
|
if let Some(cached) = cached {
|
|
info!(
|
|
event = EVENT_CAPACITY_REFRESH_RUNTIME_SUMMARY,
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_RUNTIME,
|
|
state = "cache_present",
|
|
total_used = cached.total_used,
|
|
file_count = cached.file_count,
|
|
estimated = cached.is_estimated,
|
|
source = cached.source.as_metric_label(),
|
|
cache_age_secs = cached.last_update.elapsed().as_secs(),
|
|
writes_per_minute = recent_write_frequency,
|
|
dirty_disk_count = dirty_disks.len(),
|
|
refresh_inflight = refresh_running,
|
|
"capacity refresh runtime summary"
|
|
);
|
|
} else {
|
|
info!(
|
|
event = EVENT_CAPACITY_REFRESH_RUNTIME_SUMMARY,
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_RUNTIME,
|
|
state = "cache_empty",
|
|
writes_per_minute = recent_write_frequency,
|
|
dirty_disk_count = dirty_disks.len(),
|
|
refresh_inflight = refresh_running,
|
|
"capacity refresh runtime summary"
|
|
);
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Global capacity manager instance
|
|
static GLOBAL_CAPACITY_MANAGER: std::sync::OnceLock<Arc<HybridCapacityManager>> = std::sync::OnceLock::new();
|
|
|
|
/// Get or initialize the global capacity manager
|
|
pub fn get_capacity_manager() -> Arc<HybridCapacityManager> {
|
|
GLOBAL_CAPACITY_MANAGER
|
|
.get_or_init(|| Arc::new(HybridCapacityManager::from_env()))
|
|
.clone()
|
|
}
|
|
|
|
/// Create an isolated capacity manager instance for testing
|
|
///
|
|
/// This factory function allows tests to create independent instances
|
|
/// without affecting the global singleton, avoiding test pollution.
|
|
///
|
|
/// # Example
|
|
/// ```ignore
|
|
/// let manager = create_isolated_manager(HybridStrategyConfig::default());
|
|
/// manager
|
|
/// .update_capacity(CapacityUpdate::exact(1000, 0), DataSource::RealTime)
|
|
/// .await;
|
|
/// ```
|
|
#[allow(dead_code)]
|
|
pub fn create_isolated_manager(config: HybridStrategyConfig) -> Arc<HybridCapacityManager> {
|
|
Arc::new(HybridCapacityManager::new(config))
|
|
}
|
|
|
|
/// Start background update task
|
|
/// Clamp a background interval into [1s, 30 days].
|
|
///
|
|
/// Zero panics tokio::time::interval; absurd values (e.g. u64::MAX seconds)
|
|
/// overflow `Instant + Duration` and panic the spawned task, silently killing
|
|
/// scheduled refreshes and runtime metrics (backlog#1022 S33).
|
|
fn clamp_background_interval(value: Duration, env_var: &'static str) -> Duration {
|
|
let clamped = value.clamp(Duration::from_secs(1), MAX_BACKGROUND_INTERVAL);
|
|
if clamped != value {
|
|
warn!(
|
|
event = EVENT_CAPACITY_REFRESH_INTERVAL_CLAMPED,
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_RUNTIME,
|
|
result = "clamped",
|
|
env_var,
|
|
configured_secs = value.as_secs(),
|
|
effective_secs = clamped.as_secs(),
|
|
"capacity refresh interval clamped"
|
|
);
|
|
}
|
|
clamped
|
|
}
|
|
|
|
pub async fn start_background_task(disks: Vec<CapacityDiskRef>) {
|
|
let manager = get_capacity_manager();
|
|
let manager_for_refresh = manager.clone();
|
|
let manager_for_metrics = manager.clone();
|
|
let mut refresh_interval = manager.get_config().scheduled_update_interval;
|
|
let mut metrics_interval = manager.get_config().metrics_interval;
|
|
|
|
refresh_interval = clamp_background_interval(refresh_interval, ENV_CAPACITY_SCHEDULED_INTERVAL);
|
|
metrics_interval = clamp_background_interval(metrics_interval, ENV_CAPACITY_METRICS_INTERVAL);
|
|
|
|
tokio::spawn(async move {
|
|
let mut timer = tokio::time::interval_at(tokio::time::Instant::now() + refresh_interval, refresh_interval);
|
|
|
|
loop {
|
|
timer.tick().await;
|
|
|
|
let start = Instant::now();
|
|
let manager = manager_for_refresh.clone();
|
|
let disks = disks.clone();
|
|
let disk_count = disks.len();
|
|
let started = manager
|
|
.clone()
|
|
.spawn_refresh_if_needed(
|
|
DataSource::Scheduled,
|
|
move || async move { refresh_capacity_with_scope(disks, false).await },
|
|
)
|
|
.await;
|
|
|
|
if started {
|
|
debug!(
|
|
event = EVENT_CAPACITY_REFRESH_SCHEDULED,
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_RUNTIME,
|
|
state = "started",
|
|
source = DataSource::Scheduled.as_metric_label(),
|
|
disk_count,
|
|
enqueue_latency_ms = start.elapsed().as_millis() as u64,
|
|
"capacity refresh scheduled"
|
|
);
|
|
} else {
|
|
debug!(
|
|
event = EVENT_CAPACITY_REFRESH_SKIPPED,
|
|
component = LOG_COMPONENT_CAPACITY,
|
|
subsystem = LOG_SUBSYSTEM_RUNTIME,
|
|
state = "inflight",
|
|
source = DataSource::Scheduled.as_metric_label(),
|
|
disk_count,
|
|
"capacity refresh skipped"
|
|
);
|
|
}
|
|
}
|
|
});
|
|
|
|
tokio::spawn(async move {
|
|
let mut timer = tokio::time::interval_at(tokio::time::Instant::now() + metrics_interval, metrics_interval);
|
|
loop {
|
|
timer.tick().await;
|
|
manager_for_metrics.log_runtime_summary().await;
|
|
}
|
|
});
|
|
}
|
|
|
|
// ============================================================================
|
|
// Tests
|
|
// ============================================================================
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use crate::capacity_scope::{CapacityScope, CapacityScopeDisk, record_capacity_scope, record_global_dirty_scope};
|
|
use rustfs_config::{
|
|
ENV_CAPACITY_FAST_UPDATE_THRESHOLD, ENV_CAPACITY_MAX_FILES_THRESHOLD, ENV_CAPACITY_METRICS_INTERVAL,
|
|
ENV_CAPACITY_SAMPLE_RATE, ENV_CAPACITY_STAT_TIMEOUT, ENV_CAPACITY_WRITE_FREQUENCY_THRESHOLD,
|
|
ENV_CAPACITY_WRITE_TRIGGER_DELAY,
|
|
};
|
|
use serial_test::serial;
|
|
use std::sync::Arc;
|
|
use std::sync::atomic::{AtomicUsize, Ordering};
|
|
|
|
type ConfigGetterCase = (&'static str, fn() -> u64, u64, &'static str, u64);
|
|
|
|
/// Table of env-configurable getters: (env var, getter normalized to u64,
|
|
/// expected default, override string, expected override value).
|
|
/// Durations are normalized to whole seconds.
|
|
fn config_getter_cases() -> Vec<ConfigGetterCase> {
|
|
vec![
|
|
(
|
|
ENV_CAPACITY_SCHEDULED_INTERVAL,
|
|
|| get_scheduled_update_interval().as_secs(),
|
|
120,
|
|
"600",
|
|
600,
|
|
),
|
|
(ENV_CAPACITY_WRITE_TRIGGER_DELAY, || get_write_trigger_delay().as_secs(), 5, "20", 20),
|
|
(
|
|
ENV_CAPACITY_WRITE_FREQUENCY_THRESHOLD,
|
|
|| get_write_frequency_threshold() as u64,
|
|
5,
|
|
"20",
|
|
20,
|
|
),
|
|
(
|
|
ENV_CAPACITY_FAST_UPDATE_THRESHOLD,
|
|
|| get_fast_update_threshold().as_secs(),
|
|
30,
|
|
"120",
|
|
120,
|
|
),
|
|
(
|
|
ENV_CAPACITY_MAX_FILES_THRESHOLD,
|
|
|| get_max_files_threshold() as u64,
|
|
200_000,
|
|
"2000000",
|
|
2_000_000,
|
|
),
|
|
(ENV_CAPACITY_STAT_TIMEOUT, || get_stat_timeout().as_secs(), 3, "10", 10),
|
|
(ENV_CAPACITY_SAMPLE_RATE, || get_sample_rate() as u64, 200, "500", 500),
|
|
(ENV_CAPACITY_METRICS_INTERVAL, || get_metrics_interval().as_secs(), 600, "90", 90),
|
|
]
|
|
}
|
|
|
|
#[test]
|
|
#[serial]
|
|
fn test_config_getter_defaults() {
|
|
for (env_var, getter, default, _, _) in config_getter_cases() {
|
|
temp_env::with_var(env_var, None::<&str>, || {
|
|
assert_eq!(getter(), default, "{env_var}: unexpected default value");
|
|
});
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
#[serial]
|
|
fn test_config_getter_env_overrides() {
|
|
for (env_var, getter, _, override_value, expected) in config_getter_cases() {
|
|
temp_env::with_var(env_var, Some(override_value), || {
|
|
assert_eq!(getter(), expected, "{env_var}: override not applied");
|
|
});
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
#[serial]
|
|
fn test_zero_env_values_clamp_to_defaults() {
|
|
// A zero threshold makes small disks report 0 bytes; a zero timeout
|
|
// (with dynamic timeout off) makes every scan fail. Both must fall
|
|
// back to the defaults instead of being accepted (backlog#1019 S11).
|
|
type ZeroCase = (&'static str, fn() -> u64, u64);
|
|
let zero_cases: Vec<ZeroCase> = vec![
|
|
(ENV_CAPACITY_MAX_FILES_THRESHOLD, || get_max_files_threshold() as u64, 200_000),
|
|
(ENV_CAPACITY_STAT_TIMEOUT, || get_stat_timeout().as_secs(), 3),
|
|
(ENV_CAPACITY_MIN_TIMEOUT, || get_min_timeout().as_secs(), 2),
|
|
(ENV_CAPACITY_MAX_TIMEOUT, || get_max_timeout().as_secs(), 15),
|
|
];
|
|
for (env_var, getter, default) in zero_cases {
|
|
temp_env::with_var(env_var, Some("0"), || {
|
|
assert_eq!(getter(), default, "{env_var}: zero must clamp to default");
|
|
});
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_update_capacity_preserves_retrieval_metadata() {
|
|
let manager = HybridCapacityManager::from_env();
|
|
|
|
manager
|
|
.update_capacity(CapacityUpdate::exact(1000, 10), DataSource::RealTime)
|
|
.await;
|
|
|
|
let cached = manager.get_capacity().await.unwrap();
|
|
assert_eq!(cached.total_used, 1000);
|
|
assert_eq!(cached.file_count, 10);
|
|
assert_eq!(cached.source, DataSource::RealTime);
|
|
assert!(!cached.is_estimated);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_record_write_operation() {
|
|
let manager = HybridCapacityManager::from_env();
|
|
|
|
manager.record_write_operation().await;
|
|
|
|
let frequency = manager.get_write_frequency().await;
|
|
assert_eq!(frequency, 1);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_write_frequency_window() {
|
|
let manager = HybridCapacityManager::from_env();
|
|
|
|
for _ in 0..20 {
|
|
manager.record_write_operation().await;
|
|
}
|
|
|
|
assert_eq!(manager.get_write_frequency().await, 20);
|
|
}
|
|
|
|
#[test]
|
|
fn test_clamp_background_interval_bounds() {
|
|
// Zero panics tokio interval; u64::MAX seconds overflows Instant +
|
|
// Duration and used to panic the spawned background task.
|
|
assert_eq!(
|
|
clamp_background_interval(Duration::ZERO, ENV_CAPACITY_SCHEDULED_INTERVAL),
|
|
Duration::from_secs(1)
|
|
);
|
|
assert_eq!(
|
|
clamp_background_interval(Duration::from_secs(u64::MAX), ENV_CAPACITY_SCHEDULED_INTERVAL),
|
|
MAX_BACKGROUND_INTERVAL
|
|
);
|
|
assert_eq!(
|
|
clamp_background_interval(Duration::from_secs(120), ENV_CAPACITY_SCHEDULED_INTERVAL),
|
|
Duration::from_secs(120)
|
|
);
|
|
// The cap itself must be safely addable to a tokio Instant.
|
|
let _ = tokio::time::Instant::now() + MAX_BACKGROUND_INTERVAL;
|
|
}
|
|
|
|
#[test]
|
|
#[serial]
|
|
fn test_recent_write_count_ignores_future_buckets() {
|
|
let record = WriteRecord::new();
|
|
record.write_buckets[0].store(120, 3);
|
|
record.write_buckets[1].store(90, 2);
|
|
|
|
assert_eq!(
|
|
record.recent_write_count(100),
|
|
2,
|
|
"buckets from future seconds should not inflate recent write frequency"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_needs_fast_update() {
|
|
let manager = HybridCapacityManager::from_env();
|
|
|
|
// No cache, should not need update
|
|
assert!(!manager.needs_fast_update().await);
|
|
|
|
// Update cache
|
|
manager
|
|
.update_capacity(CapacityUpdate::exact(1000, 0), DataSource::RealTime)
|
|
.await;
|
|
|
|
// Fresh cache, should not need update
|
|
assert!(!manager.needs_fast_update().await);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_cache_age_tracking() {
|
|
let manager = HybridCapacityManager::from_env();
|
|
|
|
assert!(manager.get_cache_age().await.is_none());
|
|
|
|
manager
|
|
.update_capacity(CapacityUpdate::exact(1000, 1), DataSource::RealTime)
|
|
.await;
|
|
|
|
let age = manager.get_cache_age().await.unwrap();
|
|
assert!(age < Duration::from_secs(1));
|
|
|
|
tokio::time::sleep(Duration::from_millis(100)).await;
|
|
|
|
let age = manager.get_cache_age().await.unwrap();
|
|
assert!(age >= Duration::from_millis(100));
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_data_source_tracking() {
|
|
let manager = HybridCapacityManager::from_env();
|
|
|
|
for source in [
|
|
DataSource::RealTime,
|
|
DataSource::Scheduled,
|
|
DataSource::WriteTriggered,
|
|
DataSource::Fallback,
|
|
] {
|
|
manager.update_capacity(CapacityUpdate::exact(1000, 1), source).await;
|
|
assert_eq!(manager.get_capacity().await.unwrap().source, source);
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_needs_fast_update_waits_for_write_trigger_delay() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig {
|
|
scheduled_update_interval: Duration::from_secs(60),
|
|
write_trigger_delay: Duration::from_millis(50),
|
|
write_frequency_threshold: 1,
|
|
fast_update_threshold: Duration::from_millis(10),
|
|
metrics_interval: Duration::from_secs(600),
|
|
enable_smart_update: true,
|
|
enable_write_trigger: true,
|
|
});
|
|
|
|
manager
|
|
.update_capacity(CapacityUpdate::exact(1000, 0), DataSource::RealTime)
|
|
.await;
|
|
tokio::time::sleep(Duration::from_millis(15)).await;
|
|
|
|
manager.record_write_operation().await;
|
|
manager.record_write_operation().await;
|
|
tokio::time::sleep(Duration::from_millis(5)).await;
|
|
|
|
assert!(
|
|
!manager.needs_fast_update().await,
|
|
"write-triggered refresh should wait for debounce delay after a qualifying burst"
|
|
);
|
|
|
|
tokio::time::sleep(Duration::from_millis(60)).await;
|
|
assert!(manager.needs_fast_update().await);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_needs_fast_update_respects_enable_write_trigger() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig {
|
|
scheduled_update_interval: Duration::from_secs(60),
|
|
write_trigger_delay: Duration::from_secs(60),
|
|
write_frequency_threshold: 1,
|
|
fast_update_threshold: Duration::from_millis(10),
|
|
metrics_interval: Duration::from_secs(600),
|
|
enable_smart_update: true,
|
|
enable_write_trigger: false,
|
|
});
|
|
|
|
manager
|
|
.update_capacity(CapacityUpdate::exact(1000, 0), DataSource::RealTime)
|
|
.await;
|
|
tokio::time::sleep(Duration::from_millis(15)).await;
|
|
|
|
manager.record_write_operation().await;
|
|
manager.record_write_operation().await;
|
|
|
|
assert!(
|
|
!manager.needs_fast_update().await,
|
|
"write-triggered refresh should be disabled when enable_write_trigger is false"
|
|
);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_concurrent_access() {
|
|
let manager = Arc::new(HybridCapacityManager::from_env());
|
|
let mut handles = Vec::new();
|
|
|
|
for i in 0..10 {
|
|
let mgr = manager.clone();
|
|
handles.push(tokio::spawn(async move {
|
|
mgr.update_capacity(CapacityUpdate::exact(i as u64 * 100, i), DataSource::RealTime)
|
|
.await;
|
|
mgr.record_write_operation().await;
|
|
}));
|
|
}
|
|
|
|
for handle in handles {
|
|
handle.await.unwrap();
|
|
}
|
|
|
|
assert!(manager.get_capacity().await.is_some());
|
|
assert_eq!(manager.get_write_frequency().await, 10);
|
|
}
|
|
|
|
// backlog#1315: the lock-free write record must not drop concurrent writes.
|
|
// The previous async RwLock serialized every write; the CAS bucket must be
|
|
// exact under heavy same-second contention or the frequency window (and the
|
|
// write-trigger decision) would undercount.
|
|
#[tokio::test(flavor = "multi_thread", worker_threads = 8)]
|
|
#[serial]
|
|
async fn test_record_write_operation_lock_free_is_exact_under_contention() {
|
|
let manager = Arc::new(HybridCapacityManager::from_env());
|
|
let mut handles = Vec::new();
|
|
const WRITERS: usize = 16;
|
|
const PER_WRITER: usize = 64;
|
|
|
|
for _ in 0..WRITERS {
|
|
let mgr = manager.clone();
|
|
handles.push(tokio::spawn(async move {
|
|
for _ in 0..PER_WRITER {
|
|
mgr.record_write_operation().await;
|
|
}
|
|
}));
|
|
}
|
|
for handle in handles {
|
|
handle.await.unwrap();
|
|
}
|
|
|
|
// All writes land in the same monotonic second (the test runs well under
|
|
// one second), so the recent-window frequency must equal the exact total.
|
|
assert_eq!(manager.get_write_frequency().await, WRITERS * PER_WRITER);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_performance_overhead() {
|
|
let manager = Arc::new(HybridCapacityManager::from_env());
|
|
let start = Instant::now();
|
|
|
|
for i in 0..1000 {
|
|
manager
|
|
.update_capacity(CapacityUpdate::exact(i as u64, i), DataSource::RealTime)
|
|
.await;
|
|
manager.record_write_operation().await;
|
|
let _ = manager.get_capacity().await;
|
|
}
|
|
|
|
assert!(start.elapsed() < Duration::from_secs(1));
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_refresh_or_join_singleflight() {
|
|
let manager = Arc::new(HybridCapacityManager::from_env());
|
|
let calls = Arc::new(AtomicUsize::new(0));
|
|
|
|
let mgr1 = manager.clone();
|
|
let calls1 = calls.clone();
|
|
let first = tokio::spawn(async move {
|
|
mgr1.refresh_or_join(DataSource::Scheduled, move || async move {
|
|
calls1.fetch_add(1, Ordering::SeqCst);
|
|
tokio::time::sleep(Duration::from_millis(50)).await;
|
|
Ok(CapacityUpdate::exact(2048, 8))
|
|
})
|
|
.await
|
|
});
|
|
|
|
tokio::time::sleep(Duration::from_millis(10)).await;
|
|
|
|
let mgr2 = manager.clone();
|
|
let calls2 = calls.clone();
|
|
let second = tokio::spawn(async move {
|
|
mgr2.refresh_or_join(DataSource::WriteTriggered, move || async move {
|
|
calls2.fetch_add(1, Ordering::SeqCst);
|
|
Ok(CapacityUpdate::exact(4096, 16))
|
|
})
|
|
.await
|
|
});
|
|
|
|
let first = first.await.unwrap().unwrap();
|
|
let second = second.await.unwrap().unwrap();
|
|
|
|
assert_eq!(calls.load(Ordering::SeqCst), 1);
|
|
assert_eq!(first.total_used, 2048);
|
|
assert_eq!(second.total_used, 2048);
|
|
let cached = manager.get_capacity().await.unwrap();
|
|
assert_eq!(cached.total_used, 2048);
|
|
assert_eq!(cached.file_count, 8);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_refresh_or_join_recovers_after_leader_cancellation() {
|
|
let manager = Arc::new(HybridCapacityManager::from_env());
|
|
|
|
// Become the leader with a refresh that never completes, then drop the
|
|
// future mid-flight to simulate a cancelled admin request.
|
|
let mgr = manager.clone();
|
|
let mut leader = Box::pin(mgr.refresh_or_join(DataSource::Scheduled, || async {
|
|
futures::future::pending::<Result<CapacityUpdate, String>>().await
|
|
}));
|
|
assert!(futures::poll!(leader.as_mut()).is_pending());
|
|
drop(leader);
|
|
// Let a possibly-spawned reset task run.
|
|
tokio::task::yield_now().await;
|
|
|
|
// A joiner that subscribed to the cancelled cycle must unblock with an error
|
|
// (not hang), and a subsequent refresh must be able to become the new leader.
|
|
let refreshed = tokio::time::timeout(
|
|
Duration::from_secs(1),
|
|
manager.refresh_or_join(DataSource::WriteTriggered, || async { Ok(CapacityUpdate::exact(1024, 4)) }),
|
|
)
|
|
.await
|
|
.expect("refresh after cancelled leader must not hang")
|
|
.expect("new leader refresh should succeed");
|
|
assert_eq!(refreshed.total_used, 1024);
|
|
assert!(!manager.refresh_in_progress().await);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_refresh_or_join_cancelled_leader_unblocks_joiner() {
|
|
let manager = Arc::new(HybridCapacityManager::from_env());
|
|
|
|
let mgr = manager.clone();
|
|
let mut leader = Box::pin(mgr.refresh_or_join(DataSource::Scheduled, || async {
|
|
futures::future::pending::<Result<CapacityUpdate, String>>().await
|
|
}));
|
|
assert!(futures::poll!(leader.as_mut()).is_pending());
|
|
|
|
// Subscribe a joiner while the leader is still alive.
|
|
let mgr2 = manager.clone();
|
|
let joiner = tokio::spawn(async move {
|
|
mgr2.refresh_or_join(DataSource::WriteTriggered, || async { Ok(CapacityUpdate::exact(2048, 8)) })
|
|
.await
|
|
});
|
|
tokio::time::sleep(Duration::from_millis(20)).await;
|
|
|
|
drop(leader);
|
|
|
|
let joined = tokio::time::timeout(Duration::from_secs(1), joiner)
|
|
.await
|
|
.expect("joiner must unblock after leader cancellation")
|
|
.expect("joiner task must not panic");
|
|
assert!(joined.is_err(), "joiner should observe the cancellation error, got {joined:?}");
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_spawn_refresh_if_needed_deduplicates_background_refresh() {
|
|
let manager = Arc::new(HybridCapacityManager::from_env());
|
|
let calls = Arc::new(AtomicUsize::new(0));
|
|
|
|
let first_manager = manager.clone();
|
|
let first_calls = calls.clone();
|
|
let started = first_manager
|
|
.clone()
|
|
.spawn_refresh_if_needed(DataSource::Scheduled, move || async move {
|
|
first_calls.fetch_add(1, Ordering::SeqCst);
|
|
tokio::time::sleep(Duration::from_millis(50)).await;
|
|
Ok(CapacityUpdate::estimated(8192, 32))
|
|
})
|
|
.await;
|
|
assert!(started);
|
|
|
|
let second_manager = manager.clone();
|
|
let second_calls = calls.clone();
|
|
let started = second_manager
|
|
.clone()
|
|
.spawn_refresh_if_needed(DataSource::Scheduled, move || async move {
|
|
second_calls.fetch_add(1, Ordering::SeqCst);
|
|
Ok(CapacityUpdate::exact(1, 1))
|
|
})
|
|
.await;
|
|
assert!(!started);
|
|
|
|
tokio::time::sleep(Duration::from_millis(100)).await;
|
|
|
|
assert_eq!(calls.load(Ordering::SeqCst), 1);
|
|
assert!(!manager.refresh_in_progress().await);
|
|
let cached = manager.get_capacity().await.unwrap();
|
|
assert_eq!(cached.total_used, 8192);
|
|
assert!(cached.is_estimated);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_record_write_operation_with_scope_token_marks_dirty_disks() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig::default());
|
|
let token = uuid::Uuid::new_v4();
|
|
record_capacity_scope(
|
|
token,
|
|
CapacityScope {
|
|
disks: vec![CapacityScopeDisk {
|
|
endpoint: "node-a".to_string(),
|
|
drive_path: "/tmp/disk-a".to_string(),
|
|
}],
|
|
},
|
|
);
|
|
|
|
manager.record_write_operation_with_scope_token(Some(token)).await;
|
|
|
|
let dirty_disks = manager.get_dirty_disks().await;
|
|
assert_eq!(dirty_disks.len(), 1);
|
|
assert_eq!(dirty_disks[0].endpoint, "node-a");
|
|
assert_eq!(dirty_disks[0].drive_path, "/tmp/disk-a");
|
|
assert_eq!(manager.get_write_frequency().await, 1);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_get_dirty_disks_drains_global_dirty_scope_registry() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig::default());
|
|
record_global_dirty_scope(CapacityScope {
|
|
disks: vec![CapacityScopeDisk {
|
|
endpoint: "node-bg".to_string(),
|
|
drive_path: "/tmp/disk-bg".to_string(),
|
|
}],
|
|
});
|
|
|
|
let dirty_disks = manager.get_dirty_disks().await;
|
|
assert_eq!(dirty_disks.len(), 1);
|
|
assert_eq!(dirty_disks[0].endpoint, "node-bg");
|
|
assert_eq!(dirty_disks[0].drive_path, "/tmp/disk-bg");
|
|
|
|
let second_read = manager.get_dirty_disks().await;
|
|
assert_eq!(second_read.len(), 1);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_update_capacity_recomputes_total_from_disk_cache_for_subset_refresh() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig::default());
|
|
|
|
manager
|
|
.update_capacity(
|
|
CapacityUpdate {
|
|
total_used: 300,
|
|
file_count: 3,
|
|
is_estimated: false,
|
|
degraded: false,
|
|
per_disk: vec![
|
|
DiskCapacityUpdate {
|
|
disk: CapacityScopeDisk {
|
|
endpoint: "node-a".to_string(),
|
|
drive_path: "/tmp/disk-a".to_string(),
|
|
},
|
|
used_bytes: 100,
|
|
file_count: 1,
|
|
is_estimated: false,
|
|
},
|
|
DiskCapacityUpdate {
|
|
disk: CapacityScopeDisk {
|
|
endpoint: "node-b".to_string(),
|
|
drive_path: "/tmp/disk-b".to_string(),
|
|
},
|
|
used_bytes: 200,
|
|
file_count: 2,
|
|
is_estimated: false,
|
|
},
|
|
],
|
|
expected_disk_count: Some(2),
|
|
replaces_disk_cache: true,
|
|
clear_dirty_disks: Vec::new(),
|
|
scan_started_at: None,
|
|
},
|
|
DataSource::RealTime,
|
|
)
|
|
.await;
|
|
|
|
let returned = manager
|
|
.update_capacity(
|
|
CapacityUpdate {
|
|
total_used: 150,
|
|
file_count: 1,
|
|
is_estimated: true,
|
|
degraded: false,
|
|
per_disk: vec![DiskCapacityUpdate {
|
|
disk: CapacityScopeDisk {
|
|
endpoint: "node-a".to_string(),
|
|
drive_path: "/tmp/disk-a".to_string(),
|
|
},
|
|
used_bytes: 150,
|
|
file_count: 1,
|
|
is_estimated: true,
|
|
}],
|
|
expected_disk_count: Some(1),
|
|
replaces_disk_cache: false,
|
|
clear_dirty_disks: Vec::new(),
|
|
scan_started_at: None,
|
|
},
|
|
DataSource::WriteTriggered,
|
|
)
|
|
.await;
|
|
|
|
// The returned update must carry the reconciled cluster totals (node-a 150 + node-b 200),
|
|
// not the dirty-subset's own bytes (150). file_count merges the full cache (1 + 2), and
|
|
// is_estimated is true because node-a is now estimated.
|
|
assert_eq!(returned.total_used, 350);
|
|
assert_eq!(returned.file_count, 3);
|
|
assert!(returned.is_estimated);
|
|
|
|
let cached = manager.get_capacity().await.unwrap();
|
|
assert_eq!(cached.total_used, 350);
|
|
assert_eq!(cached.file_count, 3);
|
|
assert!(cached.is_estimated);
|
|
}
|
|
|
|
fn disk_entry(endpoint: &str, drive_path: &str, used_bytes: u64, file_count: usize) -> DiskCapacityUpdate {
|
|
DiskCapacityUpdate {
|
|
disk: CapacityScopeDisk {
|
|
endpoint: endpoint.to_string(),
|
|
drive_path: drive_path.to_string(),
|
|
},
|
|
used_bytes,
|
|
file_count,
|
|
is_estimated: false,
|
|
}
|
|
}
|
|
|
|
fn full_two_disk_update() -> CapacityUpdate {
|
|
CapacityUpdate {
|
|
total_used: 200,
|
|
file_count: 2,
|
|
is_estimated: false,
|
|
degraded: false,
|
|
per_disk: vec![
|
|
disk_entry("node-a", "/tmp/disk-a", 100, 1),
|
|
disk_entry("node-b", "/tmp/disk-b", 100, 1),
|
|
],
|
|
expected_disk_count: Some(2),
|
|
replaces_disk_cache: true,
|
|
clear_dirty_disks: Vec::new(),
|
|
scan_started_at: None,
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_update_capacity_degraded_full_refresh_merges_cache_and_does_not_oscillate() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig::default());
|
|
|
|
// Round 1: healthy full refresh seeds a complete cache (A=100, B=100).
|
|
manager.update_capacity(full_two_disk_update(), DataSource::RealTime).await;
|
|
|
|
// Round 2: disk B fails mid-scan; the degraded update carries only the
|
|
// surviving subset's totals and per-disk entries.
|
|
let degraded = manager
|
|
.update_capacity(
|
|
CapacityUpdate {
|
|
total_used: 100,
|
|
file_count: 1,
|
|
is_estimated: false,
|
|
degraded: true,
|
|
per_disk: vec![disk_entry("node-a", "/tmp/disk-a", 100, 1)],
|
|
expected_disk_count: None,
|
|
replaces_disk_cache: false,
|
|
clear_dirty_disks: Vec::new(),
|
|
scan_started_at: None,
|
|
},
|
|
DataSource::Scheduled,
|
|
)
|
|
.await;
|
|
|
|
// Disk B keeps its last-known 100 bytes: the published total must not
|
|
// dip to the surviving-subset sum.
|
|
assert_eq!(degraded.total_used, 200);
|
|
assert!(degraded.degraded);
|
|
let cached = manager.get_capacity().await.unwrap();
|
|
assert_eq!(cached.total_used, 200);
|
|
assert!(cached.degraded);
|
|
|
|
// Round 3: disk B recovers; the healthy full refresh clears the flag
|
|
// and the reported total never oscillated (200 → 200 → 200).
|
|
let recovered = manager.update_capacity(full_two_disk_update(), DataSource::Scheduled).await;
|
|
assert_eq!(recovered.total_used, 200);
|
|
assert!(!recovered.degraded);
|
|
let cached = manager.get_capacity().await.unwrap();
|
|
assert_eq!(cached.total_used, 200);
|
|
assert!(!cached.degraded);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_update_capacity_degraded_with_empty_per_disk_serves_merged_cache() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig::default());
|
|
manager.update_capacity(full_two_disk_update(), DataSource::RealTime).await;
|
|
|
|
// Every surviving disk had intra-disk errors, so no per-disk entry is
|
|
// trustworthy; the complete cache must still back the published total.
|
|
let degraded = manager
|
|
.update_capacity(
|
|
CapacityUpdate {
|
|
total_used: 40,
|
|
file_count: 1,
|
|
is_estimated: false,
|
|
degraded: true,
|
|
per_disk: Vec::new(),
|
|
expected_disk_count: None,
|
|
replaces_disk_cache: false,
|
|
clear_dirty_disks: Vec::new(),
|
|
scan_started_at: None,
|
|
},
|
|
DataSource::Scheduled,
|
|
)
|
|
.await;
|
|
|
|
assert_eq!(degraded.total_used, 200);
|
|
assert!(degraded.degraded);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_update_capacity_degraded_without_complete_cache_keeps_partial_sum() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig::default());
|
|
|
|
// No complete cache exists: nothing to merge, the partial sum is the
|
|
// best available value, but it must be visibly marked degraded and the
|
|
// partial per-disk data must not mark the cache complete (#805).
|
|
let degraded = manager
|
|
.update_capacity(
|
|
CapacityUpdate {
|
|
total_used: 100,
|
|
file_count: 1,
|
|
is_estimated: false,
|
|
degraded: true,
|
|
per_disk: vec![disk_entry("node-a", "/tmp/disk-a", 100, 1)],
|
|
expected_disk_count: None,
|
|
replaces_disk_cache: false,
|
|
clear_dirty_disks: Vec::new(),
|
|
scan_started_at: None,
|
|
},
|
|
DataSource::Scheduled,
|
|
)
|
|
.await;
|
|
|
|
assert_eq!(degraded.total_used, 100);
|
|
assert!(degraded.degraded);
|
|
let cached = manager.get_capacity().await.unwrap();
|
|
assert_eq!(cached.total_used, 100);
|
|
assert!(cached.degraded);
|
|
assert!(!manager.can_refresh_dirty_subset().await);
|
|
}
|
|
|
|
fn scope_disk(endpoint: &str, drive_path: &str) -> CapacityScopeDisk {
|
|
CapacityScopeDisk {
|
|
endpoint: endpoint.to_string(),
|
|
drive_path: drive_path.to_string(),
|
|
}
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_commit_keeps_dirty_marks_recorded_after_scan_start() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig::default());
|
|
let disk = scope_disk("node-a", "/tmp/disk-a");
|
|
|
|
// Dirty before the scan starts: this mark is what the scan observed.
|
|
manager
|
|
.mark_dirty_scope(&CapacityScope {
|
|
disks: vec![disk.clone()],
|
|
})
|
|
.await;
|
|
let scan_started_at = Instant::now();
|
|
|
|
// A write lands while the scan is in flight and re-marks the disk;
|
|
// the commit must not clear this newer mark.
|
|
tokio::time::sleep(Duration::from_millis(5)).await;
|
|
manager
|
|
.mark_dirty_scope(&CapacityScope {
|
|
disks: vec![disk.clone()],
|
|
})
|
|
.await;
|
|
|
|
let mut update = CapacityUpdate::exact(100, 1);
|
|
update.clear_dirty_disks = vec![disk.clone()];
|
|
update.scan_started_at = Some(scan_started_at);
|
|
manager.update_capacity(update, DataSource::Scheduled).await;
|
|
|
|
assert_eq!(manager.get_dirty_disks().await, vec![disk]);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_commit_clears_dirty_marks_recorded_before_scan_start() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig::default());
|
|
let disk = scope_disk("node-a", "/tmp/disk-a");
|
|
|
|
manager
|
|
.mark_dirty_scope(&CapacityScope {
|
|
disks: vec![disk.clone()],
|
|
})
|
|
.await;
|
|
tokio::time::sleep(Duration::from_millis(5)).await;
|
|
|
|
let mut update = CapacityUpdate::exact(100, 1);
|
|
update.clear_dirty_disks = vec![disk.clone()];
|
|
update.scan_started_at = Some(Instant::now());
|
|
manager.update_capacity(update, DataSource::Scheduled).await;
|
|
|
|
assert!(manager.get_dirty_disks().await.is_empty());
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_retain_dirty_disks_within_drops_ghost_entries() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig::default());
|
|
let local = scope_disk("node-a", "/tmp/disk-a");
|
|
let remote = scope_disk("http://peer:9000", "/data/disk-1");
|
|
|
|
manager
|
|
.mark_dirty_scope(&CapacityScope {
|
|
disks: vec![local.clone(), remote],
|
|
})
|
|
.await;
|
|
|
|
let local_set: HashSet<CapacityScopeDisk> = [local.clone()].into_iter().collect();
|
|
manager.retain_dirty_disks_within(&local_set).await;
|
|
|
|
assert_eq!(manager.get_dirty_disks().await, vec![local]);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_spawn_refresh_recovers_from_construction_panic() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig::default());
|
|
|
|
// A refresh_fn that panics while *constructing* the future (before any
|
|
// poll). This used to escape catch_unwind and kill the spawned task
|
|
// before the reset block, wedging the singleflight forever.
|
|
let spawned = manager
|
|
.clone()
|
|
.spawn_refresh_if_needed(DataSource::Scheduled, || {
|
|
#[allow(unreachable_code)]
|
|
{
|
|
panic!("refresh_fn construction panic");
|
|
std::future::ready(Err::<CapacityUpdate, String>("unreachable".into()))
|
|
}
|
|
})
|
|
.await;
|
|
assert!(spawned);
|
|
|
|
for _ in 0..200 {
|
|
if !manager.refresh_in_progress().await {
|
|
break;
|
|
}
|
|
tokio::time::sleep(Duration::from_millis(5)).await;
|
|
}
|
|
assert!(!manager.refresh_in_progress().await, "singleflight must reset after the panic");
|
|
|
|
// The next scheduled refresh must be able to start (and succeed).
|
|
let spawned_again = manager
|
|
.clone()
|
|
.spawn_refresh_if_needed(DataSource::Scheduled, || async { Ok(CapacityUpdate::exact(1, 1)) })
|
|
.await;
|
|
assert!(spawned_again);
|
|
}
|
|
|
|
#[test]
|
|
fn test_leader_guard_drop_without_runtime_resets_state() {
|
|
let (tx, _) = watch::channel(None);
|
|
let state = Arc::new(Mutex::new(RefreshState {
|
|
running: true,
|
|
result_tx: tx,
|
|
}));
|
|
|
|
// Hold the mutex on another thread so the guard's try_lock fails and
|
|
// it must take the no-runtime fallback path.
|
|
let locked = state.clone();
|
|
let barrier = Arc::new(std::sync::Barrier::new(2));
|
|
let holder_barrier = barrier.clone();
|
|
let holder = std::thread::spawn(move || {
|
|
let guard = locked.blocking_lock();
|
|
holder_barrier.wait();
|
|
std::thread::sleep(Duration::from_millis(100));
|
|
drop(guard);
|
|
});
|
|
barrier.wait();
|
|
|
|
// No tokio runtime on this thread: the drop must block until the lock
|
|
// frees and reset the state, not silently skip the reset.
|
|
drop(RefreshLeaderGuard {
|
|
state: Some(state.clone()),
|
|
});
|
|
|
|
holder.join().unwrap();
|
|
assert!(!state.try_lock().unwrap().running);
|
|
}
|
|
|
|
#[tokio::test(start_paused = true)]
|
|
#[serial]
|
|
async fn test_refresh_or_join_joiner_times_out_when_leader_wedges() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig::default());
|
|
|
|
// A leader that never publishes (models a refresh wedged beyond what
|
|
// the drop/panic guards cover).
|
|
let leader_manager = manager.clone();
|
|
let leader = tokio::spawn(async move {
|
|
leader_manager
|
|
.refresh_or_join(DataSource::Scheduled, || async {
|
|
futures::future::pending::<Result<CapacityUpdate, String>>().await
|
|
})
|
|
.await
|
|
});
|
|
// Let the leader claim the singleflight slot before joining.
|
|
tokio::task::yield_now().await;
|
|
|
|
// Paused time auto-advances past REFRESH_JOINER_WAIT_TIMEOUT: the
|
|
// joiner must surface a clear error instead of hanging forever.
|
|
let err = manager
|
|
.refresh_or_join(DataSource::RealTime, || async { Ok(CapacityUpdate::exact(1, 0)) })
|
|
.await
|
|
.unwrap_err();
|
|
assert!(err.contains("timed out"), "unexpected joiner error: {err}");
|
|
leader.abort();
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_refresh_or_join_returns_cluster_total_for_dirty_subset() {
|
|
let manager = create_isolated_manager(HybridStrategyConfig::default());
|
|
|
|
// Seed a complete two-disk cache: cluster total 300, 3 files, exact.
|
|
manager
|
|
.update_capacity(
|
|
CapacityUpdate {
|
|
total_used: 300,
|
|
file_count: 3,
|
|
is_estimated: false,
|
|
degraded: false,
|
|
per_disk: vec![
|
|
DiskCapacityUpdate {
|
|
disk: CapacityScopeDisk {
|
|
endpoint: "node-a".to_string(),
|
|
drive_path: "/tmp/disk-a".to_string(),
|
|
},
|
|
used_bytes: 100,
|
|
file_count: 1,
|
|
is_estimated: false,
|
|
},
|
|
DiskCapacityUpdate {
|
|
disk: CapacityScopeDisk {
|
|
endpoint: "node-b".to_string(),
|
|
drive_path: "/tmp/disk-b".to_string(),
|
|
},
|
|
used_bytes: 200,
|
|
file_count: 2,
|
|
is_estimated: false,
|
|
},
|
|
],
|
|
expected_disk_count: Some(2),
|
|
replaces_disk_cache: true,
|
|
clear_dirty_disks: Vec::new(),
|
|
scan_started_at: None,
|
|
},
|
|
DataSource::RealTime,
|
|
)
|
|
.await;
|
|
|
|
// A dirty-subset refresh re-scans only node-a; its raw update carries subset-only totals.
|
|
let subset_update = CapacityUpdate {
|
|
total_used: 150,
|
|
file_count: 1,
|
|
is_estimated: true,
|
|
degraded: false,
|
|
per_disk: vec![DiskCapacityUpdate {
|
|
disk: CapacityScopeDisk {
|
|
endpoint: "node-a".to_string(),
|
|
drive_path: "/tmp/disk-a".to_string(),
|
|
},
|
|
used_bytes: 150,
|
|
file_count: 1,
|
|
is_estimated: true,
|
|
}],
|
|
expected_disk_count: Some(1),
|
|
replaces_disk_cache: false,
|
|
clear_dirty_disks: Vec::new(),
|
|
scan_started_at: None,
|
|
};
|
|
|
|
let returned = manager
|
|
.refresh_or_join(DataSource::WriteTriggered, || async { Ok(subset_update) })
|
|
.await
|
|
.unwrap();
|
|
|
|
// The leader must return the merged cluster total (150 + 200), not the subset sum (150).
|
|
assert_eq!(returned.total_used, 350);
|
|
assert_eq!(returned.file_count, 3);
|
|
assert!(returned.is_estimated);
|
|
|
|
// The cache the joiners read from must agree with the returned value.
|
|
let cached = manager.get_capacity().await.unwrap();
|
|
assert_eq!(cached.total_used, returned.total_used);
|
|
assert_eq!(cached.file_count, returned.file_count);
|
|
assert_eq!(cached.is_estimated, returned.is_estimated);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_config_from_env() {
|
|
let config = HybridStrategyConfig::from_env();
|
|
|
|
// Check default values
|
|
assert_eq!(config.scheduled_update_interval, Duration::from_secs(120));
|
|
assert_eq!(config.write_trigger_delay, Duration::from_secs(5));
|
|
assert_eq!(config.write_frequency_threshold, 5);
|
|
assert_eq!(config.fast_update_threshold, Duration::from_secs(30));
|
|
assert!(config.enable_smart_update);
|
|
assert!(config.enable_write_trigger);
|
|
}
|
|
|
|
#[tokio::test]
|
|
#[serial]
|
|
async fn test_config_from_env_with_override() {
|
|
temp_env::with_var(ENV_CAPACITY_SCHEDULED_INTERVAL, Some("600"), || {
|
|
let config = HybridStrategyConfig::from_env();
|
|
assert_eq!(config.scheduled_update_interval, Duration::from_secs(600));
|
|
});
|
|
}
|
|
}
|