use crate::backend::interface::{BackendKind, CacheSetItem};
use crate::backend::{CacheBackend, CacheConnector, CacheReader, CacheWriter};
use crate::error::{OxCacheError, OxCacheResult};
use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DegradationState {
Active,
Degraded,
HalfOpen,
}
impl DegradationState {
pub fn as_str(&self) -> &'static str {
match self {
DegradationState::Active => "active",
DegradationState::Degraded => "degraded",
DegradationState::HalfOpen => "half_open",
}
}
}
struct ControllerInner {
state: DegradationState,
consecutive_failures: u32,
degraded_since: Option<Instant>,
}
pub struct DegradationController {
failure_threshold: u32,
recovery_timeout: Duration,
inner: Mutex<ControllerInner>,
state_listeners: Vec<Box<dyn Fn(DegradationState) + Send + Sync>>,
}
impl DegradationController {
pub fn new(failure_threshold: u32, recovery_timeout: Duration) -> Self {
Self {
failure_threshold: failure_threshold.max(1),
recovery_timeout,
inner: Mutex::new(ControllerInner {
state: DegradationState::Active,
consecutive_failures: 0,
degraded_since: None,
}),
state_listeners: Vec::new(),
}
}
pub fn on_state_change(
mut self,
listener: impl Fn(DegradationState) + Send + Sync + 'static,
) -> Self {
self.state_listeners.push(Box::new(listener));
self
}
fn transition(&self, inner: &mut ControllerInner, next: DegradationState) {
if inner.state != next {
inner.state = next;
if next == DegradationState::Degraded {
inner.degraded_since = Some(Instant::now());
#[cfg(feature = "metrics")]
crate::infra::GLOBAL_UNIFIED_METRICS.record_l2_degraded();
}
if next == DegradationState::Active {
inner.consecutive_failures = 0;
inner.degraded_since = None;
}
for listener in &self.state_listeners {
listener(next);
}
}
}
pub fn state(&self) -> DegradationState {
self.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.state
}
pub fn is_l1_only(&self) -> bool {
self.state() == DegradationState::Degraded
}
pub fn record_failure(&self) -> DegradationState {
let mut inner = self
.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
match inner.state {
DegradationState::Active | DegradationState::HalfOpen => {
inner.consecutive_failures = inner.consecutive_failures.saturating_add(1);
if inner.consecutive_failures >= self.failure_threshold {
self.transition(&mut inner, DegradationState::Degraded);
}
}
DegradationState::Degraded => {}
}
inner.state
}
pub fn record_success(&self) -> DegradationState {
let mut inner = self
.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
match inner.state {
DegradationState::Active => {
inner.consecutive_failures = 0;
}
DegradationState::HalfOpen => {
self.transition(&mut inner, DegradationState::Active);
}
DegradationState::Degraded => {}
}
inner.state
}
pub fn allow_l2(&self) -> bool {
let mut inner = self
.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
match inner.state {
DegradationState::Active => true,
DegradationState::HalfOpen => true,
DegradationState::Degraded => {
let elapsed = inner
.degraded_since
.map(|t| t.elapsed())
.unwrap_or(self.recovery_timeout);
if elapsed >= self.recovery_timeout {
self.transition(&mut inner, DegradationState::HalfOpen);
true
} else {
false
}
}
}
}
pub fn snapshot(&self) -> DegradationSnapshot {
let inner = self
.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
DegradationSnapshot {
state: inner.state,
consecutive_failures: inner.consecutive_failures,
failure_threshold: self.failure_threshold,
degraded_elapsed: if inner.state == DegradationState::Degraded {
inner.degraded_since.map(|t| t.elapsed())
} else {
None
},
recovery_timeout: self.recovery_timeout,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct DegradationSnapshot {
pub state: DegradationState,
pub consecutive_failures: u32,
pub failure_threshold: u32,
pub degraded_elapsed: Option<Duration>,
pub recovery_timeout: Duration,
}
impl DegradationSnapshot {
pub fn is_l1_only(&self) -> bool {
self.state == DegradationState::Degraded
}
pub fn l2_probe_due(&self) -> bool {
match self.state {
DegradationState::Active | DegradationState::HalfOpen => true,
DegradationState::Degraded => {
self.degraded_elapsed.unwrap_or_default() >= self.recovery_timeout
}
}
}
}
pub struct DegradableBackend {
inner: Arc<dyn CacheBackend>,
controller: Arc<DegradationController>,
}
impl DegradableBackend {
pub fn new(inner: Arc<dyn CacheBackend>, controller: Arc<DegradationController>) -> Self {
Self { inner, controller }
}
pub fn controller(&self) -> &Arc<DegradationController> {
&self.controller
}
fn degraded_err() -> OxCacheError {
OxCacheError::Degraded(
"L2 degraded: serving L1-only until half-open probe succeeds".to_string(),
)
}
}
#[cfg(feature = "telemetry")]
pub struct DegradationTracing;
#[cfg(feature = "telemetry")]
impl DegradationTracing {
pub fn listener() -> impl Fn(DegradationState) + Send + Sync + 'static {
use crate::i18n::messages::{
MSG_LOG_DEGRADATION_ENTERED, MSG_LOG_DEGRADATION_HALF_OPEN,
MSG_LOG_DEGRADATION_RECOVERED, t,
};
|state| match state {
DegradationState::Degraded => tracing::warn!(
target = "oxcache::degradation",
state = state.as_str(),
"{}",
t(MSG_LOG_DEGRADATION_ENTERED, &[])
),
DegradationState::HalfOpen => tracing::info!(
target = "oxcache::degradation",
state = state.as_str(),
"{}",
t(MSG_LOG_DEGRADATION_HALF_OPEN, &[])
),
DegradationState::Active => tracing::info!(
target = "oxcache::degradation",
state = state.as_str(),
"{}",
t(MSG_LOG_DEGRADATION_RECOVERED, &[])
),
}
}
}
#[async_trait::async_trait]
impl CacheReader for DegradableBackend {
async fn get(&self, key: &str) -> OxCacheResult<Option<Vec<u8>>> {
if !self.controller.allow_l2() {
return Err(Self::degraded_err());
}
match self.inner.get(key).await {
Ok(v) => {
self.controller.record_success();
Ok(v)
}
Err(e) => {
self.controller.record_failure();
Err(e)
}
}
}
async fn exists(&self, key: &str) -> OxCacheResult<bool> {
if !self.controller.allow_l2() {
return Err(Self::degraded_err());
}
match self.inner.exists(key).await {
Ok(v) => {
self.controller.record_success();
Ok(v)
}
Err(e) => {
self.controller.record_failure();
Err(e)
}
}
}
async fn ttl(&self, key: &str) -> OxCacheResult<Option<Duration>> {
self.inner.ttl(key).await
}
async fn len(&self) -> OxCacheResult<u64> {
self.inner.len().await
}
async fn capacity(&self) -> OxCacheResult<u64> {
self.inner.capacity().await
}
async fn stats(&self) -> OxCacheResult<HashMap<String, String>> {
self.inner.stats().await
}
async fn keys(&self, pattern: &str) -> OxCacheResult<Vec<String>> {
self.inner.keys(pattern).await
}
}
#[async_trait::async_trait]
impl CacheWriter for DegradableBackend {
async fn set(
&self,
key: Arc<str>,
value: Arc<Vec<u8>>,
ttl: Option<Duration>,
) -> OxCacheResult<()> {
if !self.controller.allow_l2() {
return Err(Self::degraded_err());
}
match self.inner.set(key, value, ttl).await {
Ok(v) => {
self.controller.record_success();
Ok(v)
}
Err(e) => {
self.controller.record_failure();
Err(e)
}
}
}
async fn delete(&self, key: &str) -> OxCacheResult<()> {
if !self.controller.allow_l2() {
return Err(Self::degraded_err());
}
match self.inner.delete(key).await {
Ok(v) => {
self.controller.record_success();
Ok(v)
}
Err(e) => {
self.controller.record_failure();
Err(e)
}
}
}
async fn clear(&self) -> OxCacheResult<()> {
if !self.controller.allow_l2() {
return Err(Self::degraded_err());
}
match self.inner.clear().await {
Ok(v) => {
self.controller.record_success();
Ok(v)
}
Err(e) => {
self.controller.record_failure();
Err(e)
}
}
}
async fn expire(&self, key: &str, ttl: Duration) -> OxCacheResult<bool> {
self.inner.expire(key, ttl).await
}
async fn set_many(&self, items: &[CacheSetItem]) -> OxCacheResult<()> {
if !self.controller.allow_l2() {
return Err(Self::degraded_err());
}
match self.inner.set_many(items).await {
Ok(v) => {
self.controller.record_success();
Ok(v)
}
Err(e) => {
self.controller.record_failure();
Err(e)
}
}
}
async fn delete_many(&self, keys: &[String]) -> OxCacheResult<()> {
if !self.controller.allow_l2() {
return Err(Self::degraded_err());
}
match self.inner.delete_many(keys).await {
Ok(v) => {
self.controller.record_success();
Ok(v)
}
Err(e) => {
self.controller.record_failure();
Err(e)
}
}
}
}
#[async_trait::async_trait]
impl CacheConnector for DegradableBackend {
async fn health_check(&self) -> OxCacheResult<()> {
self.inner.health_check().await
}
async fn shutdown(&self) {
self.inner.shutdown().await;
}
fn backend_kind(&self) -> BackendKind {
self.inner.backend_kind()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::backend::MockBackend;
fn failing_backend() -> Arc<dyn CacheBackend> {
Arc::new(
MockBackend::new("mock", 50, false)
.with_fail_get()
.with_fail_set(),
)
}
fn healthy_backend() -> Arc<dyn CacheBackend> {
Arc::new(MockBackend::new("mock", 50, false))
}
#[test]
fn threshold_triggers_degradation() {
let controller = Arc::new(DegradationController::new(3, Duration::from_secs(60)));
assert_eq!(controller.state(), DegradationState::Active);
assert_eq!(controller.record_failure(), DegradationState::Active);
assert_eq!(controller.record_failure(), DegradationState::Active);
assert_eq!(
controller.record_failure(),
DegradationState::Degraded,
"连续失败达阈值应降级"
);
assert!(controller.is_l1_only());
}
#[test]
fn degraded_blocks_l2_until_timeout() {
let controller = Arc::new(DegradationController::new(1, Duration::from_millis(80)));
controller.record_failure();
assert!(!controller.allow_l2(), "降级窗口内应拦截 L2");
std::thread::sleep(Duration::from_millis(100));
assert!(controller.allow_l2(), "超时后应放行半开探测");
assert_eq!(controller.state(), DegradationState::HalfOpen);
}
#[test]
fn half_open_success_recovers() {
let controller = Arc::new(DegradationController::new(1, Duration::from_millis(50)));
controller.record_failure();
std::thread::sleep(Duration::from_millis(60));
assert!(controller.allow_l2());
controller.record_success();
assert_eq!(
controller.state(),
DegradationState::Active,
"探测成功应恢复"
);
assert!(controller.allow_l2());
}
#[test]
fn half_open_failure_re_degrades() {
let controller = Arc::new(DegradationController::new(1, Duration::from_millis(50)));
controller.record_failure();
std::thread::sleep(Duration::from_millis(60));
assert!(controller.allow_l2());
assert_eq!(controller.state(), DegradationState::HalfOpen);
assert_eq!(
controller.record_failure(),
DegradationState::Degraded,
"探测失败应重新降级"
);
}
#[test]
fn success_in_active_resets_failure_count() {
let controller = Arc::new(DegradationController::new(3, Duration::from_secs(60)));
controller.record_failure();
controller.record_failure();
controller.record_success();
assert_eq!(controller.record_failure(), DegradationState::Active);
assert_eq!(controller.record_failure(), DegradationState::Active);
assert_eq!(controller.record_failure(), DegradationState::Degraded);
}
#[test]
fn state_change_listeners_notified() {
let states = Arc::new(Mutex::new(Vec::new()));
let sink = states.clone();
let controller = Arc::new(
DegradationController::new(1, Duration::from_millis(40))
.on_state_change(move |state| sink.lock().unwrap().push(state)),
);
controller.record_failure(); std::thread::sleep(Duration::from_millis(50));
controller.allow_l2(); controller.record_success();
let observed = states.lock().unwrap().clone();
assert_eq!(
observed,
vec![
DegradationState::Degraded,
DegradationState::HalfOpen,
DegradationState::Active
]
);
}
#[tokio::test]
async fn degradable_backend_short_circuits_when_degraded() {
let controller = Arc::new(DegradationController::new(2, Duration::from_millis(50)));
let backend = DegradableBackend::new(failing_backend(), controller.clone());
let _ = backend
.set(Arc::from("k"), Arc::new(b"v".to_vec()), None)
.await;
let _ = backend
.set(Arc::from("k"), Arc::new(b"v".to_vec()), None)
.await;
assert_eq!(controller.state(), DegradationState::Degraded);
let err = backend.get("k").await.expect_err("降级期应拒绝 L2");
assert!(
matches!(err, OxCacheError::Degraded(_)),
"应返回 Degraded 错误供上层回落 L1,got {err:?}"
);
std::thread::sleep(Duration::from_millis(60));
let recovered = DegradableBackend::new(healthy_backend(), controller.clone());
assert!(recovered.get("k").await.unwrap().is_none());
assert_eq!(controller.state(), DegradationState::Active);
}
#[tokio::test]
async fn healthy_backend_stays_active() {
let controller = Arc::new(DegradationController::new(2, Duration::from_secs(60)));
let backend = DegradableBackend::new(healthy_backend(), controller.clone());
backend
.set(Arc::from("k"), Arc::new(b"v".to_vec()), None)
.await
.unwrap();
assert_eq!(backend.get("k").await.unwrap(), Some(b"v".to_vec()));
assert_eq!(controller.state(), DegradationState::Active);
}
#[test]
fn snapshot_reflects_lifecycle_without_mutating_state() {
let controller = Arc::new(DegradationController::new(2, Duration::from_millis(60)));
let snap = controller.snapshot();
assert_eq!(snap.state, DegradationState::Active);
assert_eq!(snap.consecutive_failures, 0);
assert_eq!(snap.failure_threshold, 2);
assert_eq!(snap.degraded_elapsed, None);
assert!(!snap.is_l1_only());
assert!(snap.l2_probe_due(), "Active 应预测放行 L2");
assert_eq!(
controller.state(),
DegradationState::Active,
"snapshot 不得触发状态迁移"
);
controller.record_failure();
controller.record_failure();
let snap = controller.snapshot();
assert_eq!(snap.state, DegradationState::Degraded);
assert_eq!(snap.consecutive_failures, 2);
assert!(snap.degraded_elapsed.is_some());
assert!(snap.is_l1_only());
assert!(!snap.l2_probe_due(), "降级窗口内应预测拦截");
assert_eq!(
controller.state(),
DegradationState::Degraded,
"snapshot 不得把 Degraded 推进到 HalfOpen"
);
assert!(!controller.allow_l2(), "降级窗口内 allow_l2 应拦截");
assert_eq!(controller.state(), DegradationState::Degraded);
std::thread::sleep(Duration::from_millis(80));
let snap = controller.snapshot();
assert!(snap.l2_probe_due(), "超时后应预测放行,但不产生迁移");
assert_eq!(controller.state(), DegradationState::Degraded);
assert!(controller.allow_l2());
assert_eq!(controller.state(), DegradationState::HalfOpen);
let snap = controller.snapshot();
assert_eq!(snap.state, DegradationState::HalfOpen);
assert_eq!(
snap.degraded_elapsed, None,
"degraded_elapsed 仅在 Degraded 状态上报"
);
controller.record_success();
let snap = controller.snapshot();
assert_eq!(snap.state, DegradationState::Active);
assert_eq!(snap.consecutive_failures, 0);
assert_eq!(snap.degraded_elapsed, None);
}
struct InnerCircuitBackend {
reset: Duration,
open_until: Mutex<Option<Instant>>,
first_failure_pending: std::sync::atomic::AtomicBool,
admitted: std::sync::atomic::AtomicUsize,
}
impl InnerCircuitBackend {
fn new(reset: Duration) -> Self {
Self {
reset,
open_until: Mutex::new(None),
first_failure_pending: std::sync::atomic::AtomicBool::new(true),
admitted: std::sync::atomic::AtomicUsize::new(0),
}
}
fn admitted(&self) -> usize {
self.admitted.load(std::sync::atomic::Ordering::SeqCst)
}
}
#[async_trait::async_trait]
impl crate::backend::CacheReader for InnerCircuitBackend {
async fn get(&self, _key: &str) -> OxCacheResult<Option<Vec<u8>>> {
use std::sync::atomic::Ordering;
let mut open_until = self.open_until.lock().unwrap();
if let Some(t) = *open_until {
if Instant::now() < t {
return Err(OxCacheError::Degraded("inner breaker is open".to_string()));
}
*open_until = None;
self.admitted.fetch_add(1, Ordering::SeqCst);
return Ok(Some(b"probe-ok".to_vec()));
}
if self.first_failure_pending.swap(false, Ordering::SeqCst) {
*open_until = Some(Instant::now() + self.reset);
return Err(OxCacheError::Operation("inner real failure".to_string()));
}
self.admitted.fetch_add(1, Ordering::SeqCst);
Ok(Some(b"probe-ok".to_vec()))
}
async fn exists(&self, _key: &str) -> OxCacheResult<bool> {
Ok(false)
}
async fn ttl(&self, _key: &str) -> OxCacheResult<Option<Duration>> {
Ok(None)
}
async fn len(&self) -> OxCacheResult<u64> {
Ok(0)
}
async fn capacity(&self) -> OxCacheResult<u64> {
Ok(0)
}
async fn stats(&self) -> OxCacheResult<HashMap<String, String>> {
Ok(HashMap::new())
}
async fn keys(&self, _pattern: &str) -> OxCacheResult<Vec<String>> {
Ok(Vec::new())
}
}
#[async_trait::async_trait]
impl crate::backend::CacheWriter for InnerCircuitBackend {
async fn set(
&self,
_key: Arc<str>,
_value: Arc<Vec<u8>>,
_ttl: Option<Duration>,
) -> OxCacheResult<()> {
Ok(())
}
async fn delete(&self, _key: &str) -> OxCacheResult<()> {
Ok(())
}
async fn clear(&self) -> OxCacheResult<()> {
Ok(())
}
async fn expire(&self, _key: &str, _ttl: Duration) -> OxCacheResult<bool> {
Ok(false)
}
async fn set_many(&self, _items: &[CacheSetItem]) -> OxCacheResult<()> {
Ok(())
}
async fn delete_many(&self, _keys: &[String]) -> OxCacheResult<()> {
Ok(())
}
}
#[async_trait::async_trait]
impl crate::backend::CacheConnector for InnerCircuitBackend {
async fn health_check(&self) -> OxCacheResult<()> {
Ok(())
}
async fn shutdown(&self) {}
fn backend_kind(&self) -> BackendKind {
BackendKind::Mock
}
}
#[tokio::test]
async fn outer_half_open_probe_rejected_by_inner_open_is_false_recovery() {
let outer = Arc::new(DegradationController::new(1, Duration::from_millis(80)));
let inner = Arc::new(InnerCircuitBackend::new(Duration::from_millis(400)));
let inner_dyn: Arc<dyn CacheBackend> = inner.clone();
let backend = DegradableBackend::new(inner_dyn, outer.clone());
let err = backend.get("k").await.expect_err("首次访问应失败");
assert!(
matches!(err, OxCacheError::Operation(_)),
"首次应为真实故障而非熔断拒绝, got {err:?}"
);
assert_eq!(outer.state(), DegradationState::Degraded);
assert_eq!(inner.admitted(), 0);
let err = backend.get("k").await.expect_err("降级窗口内应短路");
assert!(matches!(err, OxCacheError::Degraded(_)));
assert_eq!(inner.admitted(), 0);
std::thread::sleep(Duration::from_millis(120));
let snap = outer.snapshot();
assert!(snap.l2_probe_due(), "外层窗口已到应预测放行");
assert!(outer.allow_l2());
assert_eq!(outer.state(), DegradationState::HalfOpen, "外层看似恢复");
let err = backend.get("k").await.expect_err("探测应被内层熔断拒绝");
assert!(
matches!(err, OxCacheError::Degraded(_)),
"内层 Open 拒绝应表现为 Degraded, got {err:?}"
);
assert_eq!(
outer.state(),
DegradationState::Degraded,
"探测失败应打回降级(假恢复)"
);
assert_eq!(inner.admitted(), 0, "内层打开期探测不应触达内层");
let snap = outer.snapshot();
assert_eq!(
snap.consecutive_failures, 2,
"初始失败计数仅 Active 归零,半开探测失败再计 1"
);
assert!(
snap.degraded_elapsed.unwrap() < Duration::from_millis(80),
"重新降级后降级计时重启"
);
std::thread::sleep(Duration::from_millis(340));
assert!(backend.get("k").await.unwrap().is_some());
assert_eq!(
outer.state(),
DegradationState::Active,
"内层放行后应真恢复"
);
assert_eq!(inner.admitted(), 1);
let snap = outer.snapshot();
assert_eq!(snap.state, DegradationState::Active);
assert_eq!(snap.consecutive_failures, 0);
assert_eq!(snap.degraded_elapsed, None);
}
#[tokio::test]
async fn inner_recovery_does_not_shortcut_outer_degraded_window() {
let outer = Arc::new(DegradationController::new(1, Duration::from_millis(200)));
let inner = Arc::new(InnerCircuitBackend::new(Duration::from_millis(50)));
let inner_dyn: Arc<dyn CacheBackend> = inner.clone();
let backend = DegradableBackend::new(inner_dyn, outer.clone());
let _ = backend.get("k").await;
assert_eq!(outer.state(), DegradationState::Degraded);
std::thread::sleep(Duration::from_millis(90));
let err = backend.get("k").await.expect_err("外层窗口内应继续短路");
assert!(matches!(err, OxCacheError::Degraded(_)));
assert_eq!(inner.admitted(), 0, "内层先恢复也不得绕过外层窗口");
std::thread::sleep(Duration::from_millis(150));
assert!(backend.get("k").await.unwrap().is_some());
assert_eq!(outer.state(), DegradationState::Active);
assert_eq!(inner.admitted(), 1);
}
#[cfg(feature = "telemetry")]
mod telemetry_bridge {
use super::*;
use std::fmt;
use tracing::field::{Field, Visit};
use tracing::span::{Attributes, Id, Record};
use tracing::{Event, Metadata, Subscriber};
struct CapturedEvent {
level: tracing::Level,
state: String,
}
struct StateVisitor {
state: String,
}
impl Visit for StateVisitor {
fn record_str(&mut self, field: &Field, value: &str) {
if field.name() == "state" {
self.state = value.to_string();
}
}
fn record_debug(&mut self, field: &Field, value: &dyn fmt::Debug) {
if field.name() == "state" {
self.state = format!("{value:?}");
}
}
}
struct CaptureSubscriber {
events: Arc<Mutex<Vec<CapturedEvent>>>,
}
impl Subscriber for CaptureSubscriber {
fn enabled(&self, _metadata: &Metadata<'_>) -> bool {
true
}
fn new_span(&self, _attrs: &Attributes<'_>) -> Id {
Id::from_u64(1)
}
fn record(&self, _span: &Id, _values: &Record<'_>) {}
fn record_follows_from(&self, _span: &Id, _follows: &Id) {}
fn event(&self, event: &Event<'_>) {
let mut visitor = StateVisitor {
state: String::new(),
};
event.record(&mut visitor);
self.events.lock().unwrap().push(CapturedEvent {
level: *event.metadata().level(),
state: visitor.state,
});
}
fn enter(&self, _span: &Id) {}
fn exit(&self, _span: &Id) {}
}
#[test]
fn tracing_listener_emits_state_transitions() {
let events = Arc::new(Mutex::new(Vec::new()));
let sink = events.clone();
let controller = Arc::new(
DegradationController::new(1, Duration::from_millis(40))
.on_state_change(DegradationTracing::listener()),
);
tracing::subscriber::with_default(Arc::new(CaptureSubscriber { events: sink }), || {
controller.record_failure(); std::thread::sleep(Duration::from_millis(50));
controller.allow_l2(); controller.record_success(); });
let events = events.lock().unwrap();
assert_eq!(events.len(), 3, "三次迁移应各发一条事件");
assert_eq!(events[0].state, "degraded");
assert_eq!(events[0].level, tracing::Level::WARN);
assert_eq!(events[1].state, "half_open");
assert_eq!(events[1].level, tracing::Level::INFO);
assert_eq!(events[2].state, "active");
assert_eq!(events[2].level, tracing::Level::INFO);
}
}
}