feat(ecstore): add on-demand migration runtime OnDemandMigrationSys (#7074)

* feat(ecstore): add on-demand migration bucket config model

Introduce OnDemandMigrationConfig (deny_unknown_fields, version 1) with typed validation, credential redaction, a secret-free Debug impl, and the OnceLock publish hook the runtime registers into. Exported through the api facade.

* feat(ecstore): persist on-demand migration config in bucket metadata

Store the config as a RustFS extension entry (on-demand-migration.json) with its update time in .metadata.bin, add the typed BucketMetadataSys accessor, and publish the config through the hook on every cache-install path alongside the durability sync.

* refactor(ecstore): extract shared remote S3 client builder

Move the aws_sdk_s3 client construction out of bucket_target_sys into
bucket/remote_s3_client.rs: endpoint assembly, credential provider,
path-style selection, custom CA / skip-TLS transports and the outbound
SSRF gate now build from a neutral RemoteS3EndpointSpec so replication
targets and the upcoming on-demand migration source client share one
policy. Replication builds its client through From<&BucketTarget>; the
gate keeps its relaxed semantics (private allowed, loopback only behind
RUSTFS_REPLICATION_ALLOW_LOOPBACK_TARGET) verbatim. The builder also
gains optional connect/read timeouts and a User-Agent suffix
interceptor, both unset for replication.

Refs rustfs/backlog#2149

* feat(ecstore): add on-demand migration SourceClient

Add bucket/on_demand_migration/source_client.rs on top of the shared
remote S3 builder: HEAD, ranged streaming GET, ListObjectsV2 with
source-prefix mapping, GetObjectTagging and an admin probe. Every request
carries the x-rustfs-/x-minio-source-proxy-request anti-loop markers and
a RustFS-OnDemandMigration/<version> User-Agent suffix; SSE-C source
objects are rejected as unsupported. SourceError classifies SDK failures
(not found, access denied, throttled, timeout, connect, server error)
with retryability and a stable metrics label. Debug output redacts
credentials.

Refs rustfs/backlog#2149

* docs(operations): point outbound policy at shared remote S3 client builder

* chore: integrate ODM-01 and ODM-02 as B1 base (fix facade merge)

* feat(ecstore): add on-demand migration runtime OnDemandMigrationSys

Per-node runtime for On-Demand Migration (rustfs/backlog#2152): turns each
bucket's persisted config into a live SourceClient guarded by a three-state
circuit breaker, a TTL negative cache, per-key singleflight, a pull
concurrency semaphore and lock-free counters with a serializable snapshot.

- sys.rs: OnceLock singleton; `apply` installs/rebuilds/removes bucket state
  (config compared by value, counters preserved across rebuilds, old
  cancellation token fired); `publish` is the metadata publish-hook entry
  (sync removal, spawned install, generation-ordered so a slow older install
  cannot overwrite a newer one); `resolve(bucket, key)` judges module switch,
  bucket state, prefix filter, client availability, negative cache, breaker.
- breaker.rs: Closed/Open/HalfOpen with fixed constants (5 failures / 30 s
  window / 30 s open / 1 probe); NotFound resets, AccessDenied is neutral.
- negative_cache.rs: moka sync cache keyed by local key, ttl=0 disables.
- stats.rs: requests_total{op,outcome}, pulled_bytes/objects, pull_failures,
  inflight/queue gauges, log-bucket latency histogram, last_source_error;
  snake_case snapshot pinned by a golden JSON test.
- Anonymous sources surface as a typed `OdmStateError::AnonymousUnsupported`
  until the shared client builder gains an anonymous mode.
- rustfs: `RUSTFS_ON_DEMAND_MIGRATION_ENABLED` module switch (default false)
  published to module_switches and injected into ecstore before bucket
  metadata loads; hook registered at the same point.
This commit is contained in:
Zhengchao An
2026-09-03 00:51:03 +08:00
committed by GitHub
parent 183b5c9ede
commit a23d4b05a3
11 changed files with 2368 additions and 8 deletions
+2 -2
View File
@@ -166,7 +166,7 @@ uuid = { workspace = true, features = ["v4", "fast-rng", "serde", "macro-diagnos
reed-solomon-erasure = { workspace = true, features = ["simd-accel"] }
reed-solomon-simd = { workspace = true }
lazy_static.workspace = true
moka = { workspace = true, features = ["future"] }
moka = { workspace = true, features = ["future", "sync"] }
rustfs-lock.workspace = true
rustfs-io-metrics.workspace = true
regex = { workspace = true }
@@ -185,7 +185,7 @@ hyper-rustls = { workspace = true, default-features = false, features = ["native
hostname.workspace = true
rustls = { workspace = true, default-features = false, features = ["aws-lc-rs", "logging", "tls12", "prefer-post-quantum", "std"] }
rustls-pki-types.workspace = true
tokio = { workspace = true, features = ["io-util", "sync", "signal", "fs", "rt-multi-thread"] }
tokio = { workspace = true, features = ["io-util", "sync", "signal", "fs", "rt-multi-thread", "time"] }
tonic = { workspace = true, features = ["gzip", "deflate"] }
xxhash-rust = { workspace = true, features = ["xxh64", "xxh3"] }
tower = { workspace = true, features = ["timeout"] }
+8
View File
@@ -146,6 +146,14 @@ pub mod bucket {
}
pub mod on_demand_migration {
pub use crate::bucket::on_demand_migration::{
ApplyOutcome, BREAKER_FAILURE_THRESHOLD, BREAKER_FAILURE_WINDOW, BREAKER_HALF_OPEN_MAX_PROBES, BREAKER_OPEN_DURATION,
Breaker, BreakerState, BreakerTransition, BreakerVerdict, BucketOdmState, GLOBAL_ON_DEMAND_MIGRATION_SYS, GaugeGuard,
LastSourceError, LatencyBucketSnapshot, NEGATIVE_CACHE_MAX_ENTRIES, NegativeCache, OdmBucketSnapshot, OdmLookup,
OdmOp, OdmOutcome, OdmStateError, OdmStats, OdmStatsSnapshot, OnDemandMigrationSys, PullError, PullFailureReason,
PullFollower, PullLeader, PullOutcome, PullPath, PullResult, PullSlot, SOURCE_LATENCY_BUCKET_BOUNDS_MS,
SourceLatencySnapshot, source_client_spec,
};
pub use crate::bucket::on_demand_migration::{
ConfigPublishHook, FilterConfig, HeadPolicy, ON_DEMAND_MIGRATION_CONFIG_HOOK, ON_DEMAND_MIGRATION_CONFIG_VERSION,
OnDemandMigrationConfig, OnDemandMigrationConfigError, PathStyle, PolicyConfig, Provider, RangeGetPolicy,
@@ -0,0 +1,362 @@
// 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.
//! Per-bucket three-state circuit breaker protecting an on-demand migration
//! source (rustfs/backlog#2152).
//!
//! `Closed` lets every request through and counts consecutive failures
//! inside a sliding window; reaching the threshold opens the breaker. `Open`
//! rejects everything until the open duration elapses, then moves to
//! `HalfOpen`, which admits a single probe: success closes the breaker,
//! failure re-opens it. Timing uses `tokio::time::Instant` so tests can drive
//! it with `tokio::time::pause`.
//!
//! Only transport-level failures count (`Throttled`, `Timeout`, `Connect`,
//! `ServerError`). `NotFound` is a healthy answer and resets the failure
//! streak; `AccessDenied`, `Unsupported` and `Other` are configuration or
//! object problems that neither open nor close the breaker.
use super::source_client::SourceError;
use parking_lot::Mutex;
use serde::{Deserialize, Serialize};
use std::time::Duration;
use tokio::time::Instant;
/// Consecutive counted failures that open the breaker.
pub const BREAKER_FAILURE_THRESHOLD: u32 = 5;
/// Failures further apart than this do not accumulate.
pub const BREAKER_FAILURE_WINDOW: Duration = Duration::from_secs(30);
/// How long an open breaker rejects before admitting a probe.
pub const BREAKER_OPEN_DURATION: Duration = Duration::from_secs(30);
/// Probes admitted while half-open.
pub const BREAKER_HALF_OPEN_MAX_PROBES: u32 = 1;
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum BreakerState {
Closed,
Open,
HalfOpen,
}
impl BreakerState {
pub fn as_str(self) -> &'static str {
match self {
BreakerState::Closed => "closed",
BreakerState::Open => "open",
BreakerState::HalfOpen => "half_open",
}
}
}
/// A state change the caller may want to log.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct BreakerTransition {
pub from: BreakerState,
pub to: BreakerState,
}
/// How a source result is scored by the breaker.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum BreakerVerdict {
/// Resets the failure streak; closes a half-open breaker.
Success,
/// Counts toward the threshold; re-opens a half-open breaker.
Failure,
/// Leaves the breaker untouched.
Neutral,
}
impl BreakerVerdict {
/// `None` is a successful source call.
pub fn for_result(error: Option<&SourceError>) -> Self {
match error {
None | Some(SourceError::NotFound) => BreakerVerdict::Success,
Some(SourceError::Throttled | SourceError::Timeout | SourceError::Connect(_) | SourceError::ServerError(_)) => {
BreakerVerdict::Failure
}
Some(SourceError::AccessDenied | SourceError::Unsupported(_) | SourceError::Other(_)) => BreakerVerdict::Neutral,
}
}
}
#[derive(Debug)]
struct Inner {
state: BreakerState,
consecutive_failures: u32,
last_failure_at: Option<Instant>,
opened_at: Option<Instant>,
half_open_probes: u32,
}
#[derive(Debug)]
pub struct Breaker {
inner: Mutex<Inner>,
}
impl Default for Breaker {
fn default() -> Self {
Self::new()
}
}
impl Breaker {
pub fn new() -> Self {
Self {
inner: Mutex::new(Inner {
state: BreakerState::Closed,
consecutive_failures: 0,
last_failure_at: None,
opened_at: None,
half_open_probes: 0,
}),
}
}
/// Current state after applying the open-duration timeout.
pub fn state(&self) -> BreakerState {
let mut inner = self.inner.lock();
Self::advance(&mut inner, Instant::now());
inner.state
}
/// Whether a request may reach the source right now. Consumes the
/// half-open probe budget when it grants one.
pub fn allow_request(&self) -> bool {
let mut inner = self.inner.lock();
Self::advance(&mut inner, Instant::now());
match inner.state {
BreakerState::Closed => true,
BreakerState::Open => false,
BreakerState::HalfOpen => {
if inner.half_open_probes < BREAKER_HALF_OPEN_MAX_PROBES {
inner.half_open_probes += 1;
true
} else {
false
}
}
}
}
/// Scores a source result; returns the transition it caused, if any.
pub fn record(&self, verdict: BreakerVerdict) -> Option<BreakerTransition> {
match verdict {
BreakerVerdict::Success => self.record_success(),
BreakerVerdict::Failure => self.record_failure(),
BreakerVerdict::Neutral => None,
}
}
pub fn record_success(&self) -> Option<BreakerTransition> {
let mut inner = self.inner.lock();
let now = Instant::now();
Self::advance(&mut inner, now);
inner.consecutive_failures = 0;
inner.last_failure_at = None;
match inner.state {
BreakerState::Closed => None,
// A success while open can only come from a request admitted
// before the breaker opened; it says nothing about recovery.
BreakerState::Open => None,
BreakerState::HalfOpen => Some(Self::transition(&mut inner, BreakerState::Closed, now)),
}
}
pub fn record_failure(&self) -> Option<BreakerTransition> {
let mut inner = self.inner.lock();
let now = Instant::now();
Self::advance(&mut inner, now);
match inner.state {
BreakerState::Closed => {
let within_window = inner
.last_failure_at
.is_some_and(|last| now.saturating_duration_since(last) <= BREAKER_FAILURE_WINDOW);
inner.consecutive_failures = if within_window { inner.consecutive_failures + 1 } else { 1 };
inner.last_failure_at = Some(now);
if inner.consecutive_failures >= BREAKER_FAILURE_THRESHOLD {
Some(Self::transition(&mut inner, BreakerState::Open, now))
} else {
None
}
}
BreakerState::Open => None,
BreakerState::HalfOpen => Some(Self::transition(&mut inner, BreakerState::Open, now)),
}
}
fn advance(inner: &mut Inner, now: Instant) {
if inner.state == BreakerState::Open
&& inner
.opened_at
.is_some_and(|opened| now.saturating_duration_since(opened) >= BREAKER_OPEN_DURATION)
{
Self::transition(inner, BreakerState::HalfOpen, now);
}
}
fn transition(inner: &mut Inner, to: BreakerState, now: Instant) -> BreakerTransition {
let from = inner.state;
inner.state = to;
match to {
BreakerState::Open => {
inner.opened_at = Some(now);
inner.half_open_probes = 0;
}
BreakerState::HalfOpen => {
inner.half_open_probes = 0;
}
BreakerState::Closed => {
inner.opened_at = None;
inner.half_open_probes = 0;
inner.consecutive_failures = 0;
inner.last_failure_at = None;
}
}
BreakerTransition { from, to }
}
}
#[cfg(test)]
mod tests {
use super::*;
fn server_error() -> SourceError {
SourceError::ServerError(503)
}
#[tokio::test(start_paused = true)]
async fn five_failures_open_then_half_open_after_timeout() {
let breaker = Breaker::new();
for i in 0..BREAKER_FAILURE_THRESHOLD - 1 {
assert_eq!(breaker.record(BreakerVerdict::for_result(Some(&server_error()))), None, "failure {i}");
assert_eq!(breaker.state(), BreakerState::Closed);
}
assert_eq!(
breaker.record(BreakerVerdict::for_result(Some(&server_error()))),
Some(BreakerTransition {
from: BreakerState::Closed,
to: BreakerState::Open
})
);
assert_eq!(breaker.state(), BreakerState::Open);
assert!(!breaker.allow_request());
tokio::time::advance(BREAKER_OPEN_DURATION - Duration::from_secs(1)).await;
assert!(!breaker.allow_request());
assert_eq!(breaker.state(), BreakerState::Open);
tokio::time::advance(Duration::from_secs(1)).await;
assert_eq!(breaker.state(), BreakerState::HalfOpen);
assert!(breaker.allow_request(), "one probe is admitted");
assert!(!breaker.allow_request(), "second probe is rejected");
}
#[tokio::test(start_paused = true)]
async fn half_open_probe_success_closes_and_failure_reopens() {
let breaker = Breaker::new();
for _ in 0..BREAKER_FAILURE_THRESHOLD {
breaker.record_failure();
}
tokio::time::advance(BREAKER_OPEN_DURATION).await;
assert!(breaker.allow_request());
assert_eq!(
breaker.record_failure(),
Some(BreakerTransition {
from: BreakerState::HalfOpen,
to: BreakerState::Open
})
);
assert!(!breaker.allow_request());
tokio::time::advance(BREAKER_OPEN_DURATION).await;
assert!(breaker.allow_request());
assert_eq!(
breaker.record_success(),
Some(BreakerTransition {
from: BreakerState::HalfOpen,
to: BreakerState::Closed
})
);
assert_eq!(breaker.state(), BreakerState::Closed);
assert!(breaker.allow_request());
// The streak restarts from zero after closing.
for _ in 0..BREAKER_FAILURE_THRESHOLD - 1 {
assert_eq!(breaker.record_failure(), None);
}
assert_eq!(breaker.state(), BreakerState::Closed);
}
#[tokio::test(start_paused = true)]
async fn failures_outside_window_do_not_accumulate() {
let breaker = Breaker::new();
for _ in 0..BREAKER_FAILURE_THRESHOLD - 1 {
breaker.record_failure();
}
tokio::time::advance(BREAKER_FAILURE_WINDOW + Duration::from_secs(1)).await;
assert_eq!(breaker.record_failure(), None, "stale streak restarts at one");
assert_eq!(breaker.state(), BreakerState::Closed);
}
#[test]
fn not_found_and_access_denied_do_not_count() {
let breaker = Breaker::new();
for _ in 0..BREAKER_FAILURE_THRESHOLD - 1 {
breaker.record(BreakerVerdict::for_result(Some(&server_error())));
}
assert_eq!(breaker.record(BreakerVerdict::for_result(Some(&SourceError::AccessDenied))), None);
assert_eq!(breaker.state(), BreakerState::Closed);
// AccessDenied is neutral: the streak is still one short of opening.
assert_eq!(
breaker.record(BreakerVerdict::for_result(Some(&SourceError::Unsupported("sse-c".into())))),
None
);
assert_eq!(breaker.record(BreakerVerdict::for_result(Some(&SourceError::Other("x".into())))), None);
// NotFound is a healthy answer and resets the streak entirely.
assert_eq!(breaker.record(BreakerVerdict::for_result(Some(&SourceError::NotFound))), None);
for _ in 0..BREAKER_FAILURE_THRESHOLD - 1 {
assert_eq!(breaker.record(BreakerVerdict::for_result(Some(&SourceError::Timeout))), None);
}
assert_eq!(breaker.state(), BreakerState::Closed);
}
#[test]
fn verdicts_cover_every_source_error_class() {
assert_eq!(BreakerVerdict::for_result(None), BreakerVerdict::Success);
assert_eq!(BreakerVerdict::for_result(Some(&SourceError::NotFound)), BreakerVerdict::Success);
for failure in [
SourceError::Throttled,
SourceError::Timeout,
SourceError::Connect("refused".into()),
SourceError::ServerError(500),
] {
assert_eq!(BreakerVerdict::for_result(Some(&failure)), BreakerVerdict::Failure, "{failure:?}");
}
for neutral in [
SourceError::AccessDenied,
SourceError::Unsupported("sse-c".into()),
SourceError::Other("x".into()),
] {
assert_eq!(BreakerVerdict::for_result(Some(&neutral)), BreakerVerdict::Neutral, "{neutral:?}");
}
}
#[test]
fn state_labels_are_stable() {
assert_eq!(BreakerState::Closed.as_str(), "closed");
assert_eq!(BreakerState::Open.as_str(), "open");
assert_eq!(BreakerState::HalfOpen.as_str(), "half_open");
assert_eq!(serde_json::to_string(&BreakerState::HalfOpen).unwrap(), "\"half_open\"");
}
}
@@ -15,14 +15,33 @@
//! On-Demand Migration (ODM): a bucket can name an external S3-compatible
//! source bucket; GET misses are served from that source and backfilled
//! locally. This module owns the bucket-level configuration model
//! (`on-demand-migration.json` in the bucket metadata file); the runtime is
//! layered on top of it by later tasks (rustfs/backlog#2147).
//! (`on-demand-migration.json` in the bucket metadata file), the source
//! client, and the per-node runtime (`sys`) that turns configs into live
//! clients guarded by a breaker, a negative cache, singleflight and a pull
//! concurrency limit (rustfs/backlog#2147).
pub mod breaker;
pub mod config;
pub mod negative_cache;
pub mod source_client;
pub mod stats;
pub mod sys;
pub use breaker::{
BREAKER_FAILURE_THRESHOLD, BREAKER_FAILURE_WINDOW, BREAKER_HALF_OPEN_MAX_PROBES, BREAKER_OPEN_DURATION, Breaker,
BreakerState, BreakerTransition, BreakerVerdict,
};
pub use config::{
ConfigPublishHook, FilterConfig, HeadPolicy, ON_DEMAND_MIGRATION_CONFIG_HOOK, ON_DEMAND_MIGRATION_CONFIG_VERSION,
OnDemandMigrationConfig, OnDemandMigrationConfigError, PathStyle, PolicyConfig, Provider, RangeGetPolicy, SourceConfig,
SourceCredentials, SourceErrorPolicy, SourceTimeout, TlsConfig, ValidationContext,
};
pub use negative_cache::{NEGATIVE_CACHE_MAX_ENTRIES, NegativeCache};
pub use stats::{
GaugeGuard, LastSourceError, LatencyBucketSnapshot, OdmOp, OdmOutcome, OdmStats, OdmStatsSnapshot, PullFailureReason,
PullPath, SOURCE_LATENCY_BUCKET_BOUNDS_MS, SourceLatencySnapshot,
};
pub use sys::{
ApplyOutcome, BucketOdmState, GLOBAL_ON_DEMAND_MIGRATION_SYS, OdmBucketSnapshot, OdmLookup, OdmStateError,
OnDemandMigrationSys, PullError, PullFollower, PullLeader, PullOutcome, PullResult, PullSlot, source_client_spec,
};
@@ -0,0 +1,130 @@
// 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.
//! Per-bucket cache of keys the source answered 404 for
//! (rustfs/backlog#2152). A hit short-circuits the source lookup for
//! `policy.negative_cache_ttl_secs`; a TTL of zero disables the cache.
//!
//! Entries are never invalidated on a local PUT: once the object exists
//! locally the handler never consults ODM for it, so a stale negative entry
//! is harmless.
use std::time::Duration;
/// Upper bound on remembered keys per bucket; LRU eviction beyond it.
pub const NEGATIVE_CACHE_MAX_ENTRIES: u64 = 100_000;
#[derive(Debug)]
pub struct NegativeCache {
cache: Option<moka::sync::Cache<String, ()>>,
ttl: Duration,
}
impl NegativeCache {
/// `ttl == 0` builds a disabled cache that never records anything.
pub fn new(ttl: Duration) -> Self {
Self::with_capacity(ttl, NEGATIVE_CACHE_MAX_ENTRIES)
}
pub fn with_capacity(ttl: Duration, max_entries: u64) -> Self {
let cache = (!ttl.is_zero()).then(|| {
moka::sync::Cache::builder()
.max_capacity(max_entries)
.time_to_live(ttl)
.build()
});
Self { cache, ttl }
}
pub fn is_enabled(&self) -> bool {
self.cache.is_some()
}
pub fn ttl(&self) -> Duration {
self.ttl
}
/// Whether `key` is currently remembered as absent on the source.
pub fn contains(&self, key: &str) -> bool {
self.cache.as_ref().is_some_and(|cache| cache.get(key).is_some())
}
/// Remembers `key` as absent; no-op when disabled.
pub fn insert(&self, key: &str) {
if let Some(cache) = &self.cache {
cache.insert(key.to_string(), ());
}
}
/// Forgets `key` (e.g. after an admin-triggered backfill found it).
pub fn remove(&self, key: &str) {
if let Some(cache) = &self.cache {
cache.invalidate(key);
}
}
/// Approximate live entry count, for status snapshots only.
pub fn len(&self) -> u64 {
self.cache.as_ref().map_or(0, |cache| {
cache.run_pending_tasks();
cache.entry_count()
})
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn entry_expires_after_ttl() {
let cache = NegativeCache::new(Duration::from_millis(80));
assert!(cache.is_enabled());
cache.insert("a/x");
assert!(cache.contains("a/x"));
assert!(!cache.contains("a/y"));
std::thread::sleep(Duration::from_millis(160));
assert!(!cache.contains("a/x"), "entry must expire after the TTL");
}
#[test]
fn zero_ttl_disables_the_cache() {
let cache = NegativeCache::new(Duration::ZERO);
assert!(!cache.is_enabled());
cache.insert("a/x");
assert!(!cache.contains("a/x"));
assert!(cache.is_empty());
}
#[test]
fn remove_forgets_a_key() {
let cache = NegativeCache::new(Duration::from_secs(30));
cache.insert("a/x");
cache.remove("a/x");
assert!(!cache.contains("a/x"));
}
#[test]
fn capacity_bounds_entries() {
let cache = NegativeCache::with_capacity(Duration::from_secs(30), 4);
for i in 0..64 {
cache.insert(&format!("k{i}"));
}
assert!(cache.len() <= 4, "len {} exceeds capacity", cache.len());
}
}
@@ -0,0 +1,527 @@
// 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.
//! Per-bucket on-demand migration counters (rustfs/backlog#2152).
//!
//! `OdmStats` is lock-free and survives config rebuilds; `snapshot()` turns
//! it into the serializable `OdmStatsSnapshot` that the metrics collector
//! and the admin status route (ODM-10/14/15) consume. Field names and label
//! values are a wire contract: the golden JSON test below pins them.
use super::breaker::BreakerState;
use super::source_client::SourceError;
use parking_lot::Mutex;
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
use time::OffsetDateTime;
/// Request operations that can enter ODM.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum OdmOp {
Get,
Head,
}
impl OdmOp {
pub const ALL: [OdmOp; 2] = [OdmOp::Get, OdmOp::Head];
pub fn as_str(self) -> &'static str {
match self {
OdmOp::Get => "get",
OdmOp::Head => "head",
}
}
}
/// How a request that entered ODM ended. `local_hit` is deliberately absent:
/// requests served locally never reach the runtime.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum OdmOutcome {
SourceHit,
SourceMiss,
SourceError,
BreakerOpen,
NegativeCached,
Filtered,
Unsupported,
}
impl OdmOutcome {
pub const ALL: [OdmOutcome; 7] = [
OdmOutcome::SourceHit,
OdmOutcome::SourceMiss,
OdmOutcome::SourceError,
OdmOutcome::BreakerOpen,
OdmOutcome::NegativeCached,
OdmOutcome::Filtered,
OdmOutcome::Unsupported,
];
pub fn as_str(self) -> &'static str {
match self {
OdmOutcome::SourceHit => "source_hit",
OdmOutcome::SourceMiss => "source_miss",
OdmOutcome::SourceError => "source_error",
OdmOutcome::BreakerOpen => "breaker_open",
OdmOutcome::NegativeCached => "negative_cached",
OdmOutcome::Filtered => "filtered",
OdmOutcome::Unsupported => "unsupported",
}
}
}
/// Which pipeline stored a pulled object locally.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum PullPath {
/// Streamed to the client and written locally in one pass.
Inline,
/// Pulled by a background task after a partial/large read.
Background,
/// Pulled by the backfill job.
Backfill,
}
impl PullPath {
pub const ALL: [PullPath; 3] = [PullPath::Inline, PullPath::Background, PullPath::Backfill];
pub fn as_str(self) -> &'static str {
match self {
PullPath::Inline => "inline",
PullPath::Background => "background",
PullPath::Backfill => "backfill",
}
}
}
/// Why a pull did not produce a local object.
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum PullFailureReason {
SourceNotFound,
SourceAccessDenied,
SourceThrottled,
SourceTimeout,
SourceConnect,
SourceServerError,
SourceUnsupported,
SourceOther,
/// Source bytes did not match the ETag advertised by HEAD/GET.
EtagMismatch,
/// The local write (internal PUT) failed.
LocalWrite,
/// The bucket state was removed or the process is shutting down.
Canceled,
/// The background pull queue was full.
QueueFull,
}
impl PullFailureReason {
pub const ALL: [PullFailureReason; 12] = [
PullFailureReason::SourceNotFound,
PullFailureReason::SourceAccessDenied,
PullFailureReason::SourceThrottled,
PullFailureReason::SourceTimeout,
PullFailureReason::SourceConnect,
PullFailureReason::SourceServerError,
PullFailureReason::SourceUnsupported,
PullFailureReason::SourceOther,
PullFailureReason::EtagMismatch,
PullFailureReason::LocalWrite,
PullFailureReason::Canceled,
PullFailureReason::QueueFull,
];
pub fn as_str(self) -> &'static str {
match self {
PullFailureReason::SourceNotFound => "source_not_found",
PullFailureReason::SourceAccessDenied => "source_access_denied",
PullFailureReason::SourceThrottled => "source_throttled",
PullFailureReason::SourceTimeout => "source_timeout",
PullFailureReason::SourceConnect => "source_connect",
PullFailureReason::SourceServerError => "source_server_error",
PullFailureReason::SourceUnsupported => "source_unsupported",
PullFailureReason::SourceOther => "source_other",
PullFailureReason::EtagMismatch => "etag_mismatch",
PullFailureReason::LocalWrite => "local_write",
PullFailureReason::Canceled => "canceled",
PullFailureReason::QueueFull => "queue_full",
}
}
}
impl From<&SourceError> for PullFailureReason {
fn from(err: &SourceError) -> Self {
match err {
SourceError::NotFound => PullFailureReason::SourceNotFound,
SourceError::AccessDenied => PullFailureReason::SourceAccessDenied,
SourceError::Throttled => PullFailureReason::SourceThrottled,
SourceError::Timeout => PullFailureReason::SourceTimeout,
SourceError::Connect(_) => PullFailureReason::SourceConnect,
SourceError::ServerError(_) => PullFailureReason::SourceServerError,
SourceError::Unsupported(_) => PullFailureReason::SourceUnsupported,
SourceError::Other(_) => PullFailureReason::SourceOther,
}
}
}
/// Upper bounds (milliseconds) of the source latency histogram buckets; the
/// implicit last bucket is unbounded. Roughly logarithmic from 5 ms to 60 s.
pub const SOURCE_LATENCY_BUCKET_BOUNDS_MS: [u64; 14] = [
5, 10, 20, 50, 100, 200, 500, 1_000, 2_000, 5_000, 10_000, 20_000, 30_000, 60_000,
];
#[derive(Debug, Default)]
struct LatencyHistogram {
/// One counter per bound plus one for the overflow bucket.
buckets: [AtomicU64; SOURCE_LATENCY_BUCKET_BOUNDS_MS.len() + 1],
count: AtomicU64,
sum_ms: AtomicU64,
}
impl LatencyHistogram {
fn observe(&self, latency: Duration) {
let ms = u64::try_from(latency.as_millis()).unwrap_or(u64::MAX);
let index = SOURCE_LATENCY_BUCKET_BOUNDS_MS
.iter()
.position(|bound| ms <= *bound)
.unwrap_or(SOURCE_LATENCY_BUCKET_BOUNDS_MS.len());
self.buckets[index].fetch_add(1, Ordering::Relaxed);
self.count.fetch_add(1, Ordering::Relaxed);
self.sum_ms.fetch_add(ms, Ordering::Relaxed);
}
fn snapshot(&self) -> SourceLatencySnapshot {
let mut cumulative = 0;
let buckets = SOURCE_LATENCY_BUCKET_BOUNDS_MS
.iter()
.zip(self.buckets.iter())
.map(|(bound, counter)| {
cumulative += counter.load(Ordering::Relaxed);
LatencyBucketSnapshot {
le_ms: *bound,
count: cumulative,
}
})
.collect();
SourceLatencySnapshot {
buckets,
count: self.count.load(Ordering::Relaxed),
sum_ms: self.sum_ms.load(Ordering::Relaxed),
}
}
}
/// The most recent source failure, kept for operators: class only, never the
/// key or the message (which may echo attacker-controlled input).
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct LastSourceError {
pub class: String,
#[serde(with = "time::serde::rfc3339")]
pub at: OffsetDateTime,
}
#[derive(Debug, Default)]
pub struct OdmStats {
requests_total: [[AtomicU64; OdmOutcome::ALL.len()]; OdmOp::ALL.len()],
pulled_bytes_total: AtomicU64,
pulled_objects_total: [AtomicU64; PullPath::ALL.len()],
pull_failures_total: [AtomicU64; PullFailureReason::ALL.len()],
inflight_pulls: AtomicU64,
queue_depth: AtomicU64,
source_latency: LatencyHistogram,
last_source_error: Mutex<Option<LastSourceError>>,
}
impl OdmStats {
pub fn new() -> Self {
Self::default()
}
pub fn record_request(&self, op: OdmOp, outcome: OdmOutcome) {
self.requests_total[op as usize][outcome as usize].fetch_add(1, Ordering::Relaxed);
}
pub fn record_pulled_bytes(&self, bytes: u64) {
self.pulled_bytes_total.fetch_add(bytes, Ordering::Relaxed);
}
pub fn record_pulled_object(&self, path: PullPath) {
self.pulled_objects_total[path as usize].fetch_add(1, Ordering::Relaxed);
}
pub fn record_pull_failure(&self, reason: PullFailureReason) {
self.pull_failures_total[reason as usize].fetch_add(1, Ordering::Relaxed);
}
pub fn record_source_latency(&self, latency: Duration) {
self.source_latency.observe(latency);
}
pub fn record_source_error(&self, err: &SourceError) {
self.record_source_error_at(err, OffsetDateTime::now_utc());
}
pub fn record_source_error_at(&self, err: &SourceError, at: OffsetDateTime) {
*self.last_source_error.lock() = Some(LastSourceError {
class: err.class_label().to_string(),
at,
});
}
pub fn last_source_error(&self) -> Option<LastSourceError> {
self.last_source_error.lock().clone()
}
pub fn inflight_pulls(&self) -> u64 {
self.inflight_pulls.load(Ordering::Relaxed)
}
pub fn queue_depth(&self) -> u64 {
self.queue_depth.load(Ordering::Relaxed)
}
/// RAII increment of `inflight_pulls`.
pub fn inflight_guard(self: &Arc<Self>) -> GaugeGuard {
GaugeGuard::new(Arc::clone(self), OdmGauge::InflightPulls)
}
/// RAII increment of `queue_depth`.
pub fn queue_guard(self: &Arc<Self>) -> GaugeGuard {
GaugeGuard::new(Arc::clone(self), OdmGauge::QueueDepth)
}
fn gauge(&self, gauge: OdmGauge) -> &AtomicU64 {
match gauge {
OdmGauge::InflightPulls => &self.inflight_pulls,
OdmGauge::QueueDepth => &self.queue_depth,
}
}
/// Read-only, side-effect-free copy of every counter. The breaker lives
/// next to the stats in the bucket state; its state is passed in so the
/// snapshot stays a single document.
pub fn snapshot(&self, breaker_state: BreakerState) -> OdmStatsSnapshot {
let mut requests_total = BTreeMap::new();
for op in OdmOp::ALL {
let mut by_outcome = BTreeMap::new();
for outcome in OdmOutcome::ALL {
by_outcome.insert(
outcome.as_str().to_string(),
self.requests_total[op as usize][outcome as usize].load(Ordering::Relaxed),
);
}
requests_total.insert(op.as_str().to_string(), by_outcome);
}
let pulled_objects_total = PullPath::ALL
.iter()
.map(|path| {
(
path.as_str().to_string(),
self.pulled_objects_total[*path as usize].load(Ordering::Relaxed),
)
})
.collect();
let pull_failures_total = PullFailureReason::ALL
.iter()
.map(|reason| {
(
reason.as_str().to_string(),
self.pull_failures_total[*reason as usize].load(Ordering::Relaxed),
)
})
.collect();
OdmStatsSnapshot {
requests_total,
pulled_bytes_total: self.pulled_bytes_total.load(Ordering::Relaxed),
pulled_objects_total,
pull_failures_total,
inflight_pulls: self.inflight_pulls(),
queue_depth: self.queue_depth(),
source_latency: self.source_latency.snapshot(),
last_source_error: self.last_source_error(),
breaker_state,
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum OdmGauge {
InflightPulls,
QueueDepth,
}
/// Increments a gauge on creation and decrements it on drop. Owns its
/// `OdmStats` so it can live inside the pull slot handed to callers.
#[derive(Debug)]
pub struct GaugeGuard {
stats: Arc<OdmStats>,
gauge: OdmGauge,
}
impl GaugeGuard {
fn new(stats: Arc<OdmStats>, gauge: OdmGauge) -> Self {
stats.gauge(gauge).fetch_add(1, Ordering::Relaxed);
Self { stats, gauge }
}
}
impl Drop for GaugeGuard {
fn drop(&mut self) {
self.stats.gauge(self.gauge).fetch_sub(1, Ordering::Relaxed);
}
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct LatencyBucketSnapshot {
/// Upper bound of the bucket in milliseconds.
pub le_ms: u64,
/// Cumulative observations at or below `le_ms`.
pub count: u64,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct SourceLatencySnapshot {
pub buckets: Vec<LatencyBucketSnapshot>,
/// Total observations, including those above the last bound.
pub count: u64,
pub sum_ms: u64,
}
/// Serializable copy of [`OdmStats`]. Every key is snake_case and every
/// label set is fixed, so consumers can rely on the document shape.
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct OdmStatsSnapshot {
/// `op -> outcome -> count`.
pub requests_total: BTreeMap<String, BTreeMap<String, u64>>,
pub pulled_bytes_total: u64,
/// `path -> count`.
pub pulled_objects_total: BTreeMap<String, u64>,
/// `reason -> count`.
pub pull_failures_total: BTreeMap<String, u64>,
pub inflight_pulls: u64,
pub queue_depth: u64,
pub source_latency: SourceLatencySnapshot,
pub last_source_error: Option<LastSourceError>,
pub breaker_state: BreakerState,
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
use time::macros::datetime;
#[test]
fn snapshot_matches_golden_json() {
let stats = Arc::new(OdmStats::new());
stats.record_request(OdmOp::Get, OdmOutcome::SourceHit);
stats.record_request(OdmOp::Get, OdmOutcome::SourceHit);
stats.record_request(OdmOp::Head, OdmOutcome::NegativeCached);
stats.record_pulled_bytes(4096);
stats.record_pulled_object(PullPath::Inline);
stats.record_pull_failure(PullFailureReason::from(&SourceError::Timeout));
stats.record_source_latency(Duration::from_millis(3));
stats.record_source_latency(Duration::from_millis(750));
stats.record_source_latency(Duration::from_secs(90));
stats.record_source_error_at(&SourceError::ServerError(502), datetime!(2026-09-02 10:00:00 UTC));
let _inflight = stats.inflight_guard();
let _queued = stats.queue_guard();
let snapshot = stats.snapshot(BreakerState::HalfOpen);
let actual = serde_json::to_value(&snapshot).unwrap();
let expected = json!({
"requests_total": {
"get": {
"breaker_open": 0, "filtered": 0, "negative_cached": 0, "source_error": 0,
"source_hit": 2, "source_miss": 0, "unsupported": 0
},
"head": {
"breaker_open": 0, "filtered": 0, "negative_cached": 1, "source_error": 0,
"source_hit": 0, "source_miss": 0, "unsupported": 0
}
},
"pulled_bytes_total": 4096,
"pulled_objects_total": { "backfill": 0, "background": 0, "inline": 1 },
"pull_failures_total": {
"canceled": 0, "etag_mismatch": 0, "local_write": 0, "queue_full": 0,
"source_access_denied": 0, "source_connect": 0, "source_not_found": 0, "source_other": 0,
"source_server_error": 0, "source_throttled": 0, "source_timeout": 1, "source_unsupported": 0
},
"inflight_pulls": 1,
"queue_depth": 1,
"source_latency": {
"buckets": [
{ "le_ms": 5, "count": 1 }, { "le_ms": 10, "count": 1 }, { "le_ms": 20, "count": 1 },
{ "le_ms": 50, "count": 1 }, { "le_ms": 100, "count": 1 }, { "le_ms": 200, "count": 1 },
{ "le_ms": 500, "count": 1 }, { "le_ms": 1000, "count": 2 }, { "le_ms": 2000, "count": 2 },
{ "le_ms": 5000, "count": 2 }, { "le_ms": 10000, "count": 2 }, { "le_ms": 20000, "count": 2 },
{ "le_ms": 30000, "count": 2 }, { "le_ms": 60000, "count": 2 }
],
"count": 3,
"sum_ms": 90753
},
"last_source_error": { "class": "server_error", "at": "2026-09-02T10:00:00Z" },
"breaker_state": "half_open"
});
assert_eq!(actual, expected);
let round_trip: OdmStatsSnapshot = serde_json::from_value(actual).unwrap();
assert_eq!(round_trip, snapshot);
}
#[test]
fn gauges_return_to_zero_when_guards_drop() {
let stats = Arc::new(OdmStats::new());
{
let _a = stats.inflight_guard();
let _b = stats.inflight_guard();
let _c = stats.queue_guard();
assert_eq!(stats.inflight_pulls(), 2);
assert_eq!(stats.queue_depth(), 1);
}
assert_eq!(stats.inflight_pulls(), 0);
assert_eq!(stats.queue_depth(), 0);
}
#[test]
fn pull_failure_reason_covers_every_source_error_class() {
let cases = [
(SourceError::NotFound, PullFailureReason::SourceNotFound),
(SourceError::AccessDenied, PullFailureReason::SourceAccessDenied),
(SourceError::Throttled, PullFailureReason::SourceThrottled),
(SourceError::Timeout, PullFailureReason::SourceTimeout),
(SourceError::Connect("x".into()), PullFailureReason::SourceConnect),
(SourceError::ServerError(500), PullFailureReason::SourceServerError),
(SourceError::Unsupported("x".into()), PullFailureReason::SourceUnsupported),
(SourceError::Other("x".into()), PullFailureReason::SourceOther),
];
for (err, reason) in cases {
assert_eq!(PullFailureReason::from(&err), reason, "{err:?}");
assert_eq!(serde_json::to_string(&reason).unwrap(), format!("\"{}\"", reason.as_str()));
}
}
#[test]
fn label_lists_are_exhaustive_and_unique() {
let outcomes: std::collections::BTreeSet<_> = OdmOutcome::ALL.iter().map(|o| o.as_str()).collect();
assert_eq!(outcomes.len(), OdmOutcome::ALL.len());
let reasons: std::collections::BTreeSet<_> = PullFailureReason::ALL.iter().map(|r| r.as_str()).collect();
assert_eq!(reasons.len(), PullFailureReason::ALL.len());
let paths: std::collections::BTreeSet<_> = PullPath::ALL.iter().map(|p| p.as_str()).collect();
assert_eq!(paths.len(), PullPath::ALL.len());
}
}
File diff suppressed because it is too large Load Diff