Skip to main content

sz_orm_pool/
circuit_breaker.rs

1//! 断路器抽象 trait
2//!
3//! P1-4 修复:核心层定义断路器接口,消除对 sz-orm-health 的反向依赖。
4//!
5//! ## 设计动机
6//!
7//! 之前 `sz-orm-core` 通过 optional dependency 反向依赖 `sz-orm-health`,
8//! 违反了"核心 ← 扩展 ← 集成"的单向依赖原则。本模块将断路器的抽象 trait
9//! 提升至核心层,扩展包(如 sz-orm-health)实现该 trait,从而反转依赖方向。
10//!
11//! ## 使用方式
12//!
13//! - 池内部默认使用 [`DefaultCircuitBreaker`](核心层自带实现)
14//! - 调用方可通过实现 [`CircuitBreaker`] trait 自定义断路器逻辑
15//! - sz-orm-health 包的 `CircuitBreaker` 结构体也实现了本 trait
16
17use std::time::{Duration, Instant};
18
19/// 断路器状态
20///
21/// 三态有限状态机:
22/// ```text
23///   Closed ──(失败次数达阈值)──> Open
24///   Open ──(reset_timeout 到达)──> HalfOpen
25///   HalfOpen ──(一次成功)──> Closed
26///   HalfOpen ──(一次失败)──> Open
27/// ```
28#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29pub enum CircuitState {
30    /// 正常运行:请求放行
31    Closed,
32    /// 跳闸:请求被阻塞,直到 `reset_timeout` 到达
33    Open,
34    /// 探测:放行一次试探请求(reset_timeout 到达后进入)
35    HalfOpen,
36}
37
38/// 断路器抽象 trait
39///
40/// P1-4:核心层定义的断路器接口,消除对 sz-orm-health 的反向依赖。
41///
42/// 实现方需保证:
43/// - `can_execute()` 在 `Open` 状态下根据 `reset_timeout` 自动迁移到 `HalfOpen`
44/// - `record_success()` 将状态重置为 `Closed` 并清空失败计数
45/// - `record_failure()` 累计失败次数,达阈值时迁移到 `Open`
46pub trait CircuitBreaker: Send + Sync {
47    /// 获取当前断路器状态
48    fn state(&self) -> CircuitState;
49
50    /// 记录一次成功请求
51    ///
52    /// 通常将状态重置为 `Closed` 并清空连续失败计数。
53    fn record_success(&mut self);
54
55    /// 记录一次失败请求
56    ///
57    /// 累计连续失败次数,达到阈值时将状态迁移到 `Open`。
58    fn record_failure(&mut self);
59
60    /// 判断当前是否允许请求通过
61    ///
62    /// 在 `Open` 状态下,若 `reset_timeout` 已到达,应迁移到 `HalfOpen` 并返回 `true`。
63    fn can_execute(&mut self) -> bool;
64
65    /// 手动重置断路器到 `Closed` 状态
66    ///
67    /// 与 `record_success` 的区别:`reset` 是运维主动强制重置,
68    /// 无视当前状态(含 `Open`),常用于故障排除后手动恢复。
69    ///
70    /// 返回是否实际发生了状态变更。
71    fn reset(&mut self) -> bool;
72}
73
74/// 默认断路器实现
75///
76/// 提供 `failure_threshold`(连续失败阈值)+ `reset_timeout`(重置超时)的经典三态机。
77///
78/// # 示例
79///
80/// ```ignore
81/// use sz_orm_pool::circuit_breaker::{CircuitBreaker, CircuitState, DefaultCircuitBreaker};
82/// use std::time::Duration;
83///
84/// let mut cb = DefaultCircuitBreaker::new(3, Duration::from_secs(60));
85/// assert_eq!(cb.state(), CircuitState::Closed);
86/// assert!(cb.can_execute());
87///
88/// // 3 次失败后跳闸
89/// cb.record_failure();
90/// cb.record_failure();
91/// cb.record_failure();
92/// assert_eq!(cb.state(), CircuitState::Open);
93/// assert!(!cb.can_execute());
94///
95/// // 一次成功后恢复
96/// cb.record_success();
97/// assert_eq!(cb.state(), CircuitState::Closed);
98/// ```
99#[derive(Debug)]
100pub struct DefaultCircuitBreaker {
101    /// 连续失败达到此阈值时跳闸(>= 比较)
102    failure_threshold: usize,
103    /// Open 状态持续时间,到达后进入 HalfOpen
104    reset_timeout: Duration,
105    /// 当前状态
106    state: CircuitState,
107    /// 连续失败次数
108    consecutive_failures: usize,
109    /// 最后一次失败的时间(用于判断 reset_timeout 是否到达)
110    last_failure_at: Option<Instant>,
111}
112
113impl DefaultCircuitBreaker {
114    /// 创建默认断路器
115    ///
116    /// - `failure_threshold`:连续失败次数阈值(>= 时跳闸)
117    /// - `reset_timeout`:Open 状态持续时间,到达后进入 HalfOpen
118    pub fn new(failure_threshold: usize, reset_timeout: Duration) -> Self {
119        Self {
120            failure_threshold,
121            reset_timeout,
122            state: CircuitState::Closed,
123            consecutive_failures: 0,
124            last_failure_at: None,
125        }
126    }
127
128    /// 获取连续失败次数(主要用于测试与监控)
129    pub fn consecutive_failures(&self) -> usize {
130        self.consecutive_failures
131    }
132
133    /// 获取失败阈值
134    pub fn failure_threshold(&self) -> usize {
135        self.failure_threshold
136    }
137
138    /// 获取重置超时
139    pub fn reset_timeout(&self) -> Duration {
140        self.reset_timeout
141    }
142}
143
144impl Default for DefaultCircuitBreaker {
145    /// 默认配置:5 次连续失败跳闸,30 秒后进入 HalfOpen
146    fn default() -> Self {
147        Self::new(5, Duration::from_secs(30))
148    }
149}
150
151impl CircuitBreaker for DefaultCircuitBreaker {
152    fn state(&self) -> CircuitState {
153        self.state
154    }
155
156    fn record_success(&mut self) {
157        self.consecutive_failures = 0;
158        self.state = CircuitState::Closed;
159        self.last_failure_at = None;
160    }
161
162    fn record_failure(&mut self) {
163        self.consecutive_failures += 1;
164        self.last_failure_at = Some(Instant::now());
165        if self.consecutive_failures >= self.failure_threshold {
166            self.state = CircuitState::Open;
167        }
168    }
169
170    fn can_execute(&mut self) -> bool {
171        match self.state {
172            CircuitState::Closed => true,
173            CircuitState::HalfOpen => true,
174            CircuitState::Open => {
175                let elapsed = self
176                    .last_failure_at
177                    .map(|t| t.elapsed())
178                    .unwrap_or_else(|| Duration::ZERO);
179                if elapsed >= self.reset_timeout {
180                    self.state = CircuitState::HalfOpen;
181                    true
182                } else {
183                    false
184                }
185            }
186        }
187    }
188
189    fn reset(&mut self) -> bool {
190        let changed = self.state != CircuitState::Closed || self.consecutive_failures != 0;
191        self.state = CircuitState::Closed;
192        self.consecutive_failures = 0;
193        self.last_failure_at = None;
194        changed
195    }
196}
197
198#[cfg(test)]
199mod tests {
200    use super::*;
201
202    #[test]
203    fn test_circuit_breaker_starts_closed() {
204        let mut cb = DefaultCircuitBreaker::new(3, Duration::from_millis(100));
205        assert_eq!(cb.state(), CircuitState::Closed);
206        assert!(cb.can_execute());
207    }
208
209    #[test]
210    fn test_circuit_breaker_trips_after_threshold() {
211        let mut cb = DefaultCircuitBreaker::new(3, Duration::from_secs(60));
212        assert!(cb.can_execute());
213        cb.record_failure();
214        cb.record_failure();
215        assert_eq!(cb.state(), CircuitState::Closed);
216        cb.record_failure();
217        assert_eq!(cb.state(), CircuitState::Open);
218        assert!(!cb.can_execute());
219    }
220
221    #[test]
222    fn test_circuit_breaker_success_resets() {
223        let mut cb = DefaultCircuitBreaker::new(3, Duration::from_secs(60));
224        cb.record_failure();
225        cb.record_failure();
226        cb.record_success();
227        assert_eq!(cb.state(), CircuitState::Closed);
228        // After success, should need 3 failures again to trip.
229        cb.record_failure();
230        cb.record_failure();
231        assert_eq!(cb.state(), CircuitState::Closed);
232    }
233
234    #[test]
235    fn test_circuit_breaker_half_open_after_timeout() {
236        let mut cb = DefaultCircuitBreaker::new(1, Duration::from_millis(10));
237        cb.record_failure();
238        assert_eq!(cb.state(), CircuitState::Open);
239        // Immediately: still open.
240        assert!(!cb.can_execute());
241        // Wait for reset timeout.
242        std::thread::sleep(Duration::from_millis(30));
243        assert!(cb.can_execute());
244        assert_eq!(cb.state(), CircuitState::HalfOpen);
245    }
246
247    #[test]
248    fn test_circuit_breaker_half_open_success_closes() {
249        let mut cb = DefaultCircuitBreaker::new(1, Duration::from_millis(10));
250        cb.record_failure();
251        std::thread::sleep(Duration::from_millis(20));
252        assert!(cb.can_execute());
253        assert_eq!(cb.state(), CircuitState::HalfOpen);
254        cb.record_success();
255        assert_eq!(cb.state(), CircuitState::Closed);
256    }
257
258    #[test]
259    fn test_circuit_breaker_half_open_failure_reopens() {
260        let mut cb = DefaultCircuitBreaker::new(1, Duration::from_millis(10));
261        cb.record_failure();
262        std::thread::sleep(Duration::from_millis(20));
263        assert!(cb.can_execute());
264        assert_eq!(cb.state(), CircuitState::HalfOpen);
265        cb.record_failure();
266        assert_eq!(cb.state(), CircuitState::Open);
267    }
268
269    #[test]
270    fn test_circuit_breaker_boundary_exactly_threshold() {
271        // threshold = 3 means 3 failures should trip (>= comparison).
272        let mut cb = DefaultCircuitBreaker::new(3, Duration::from_secs(60));
273        cb.record_failure();
274        cb.record_failure();
275        assert_eq!(cb.state(), CircuitState::Closed);
276        cb.record_failure();
277        assert_eq!(cb.state(), CircuitState::Open);
278    }
279
280    #[test]
281    fn test_circuit_breaker_reset_from_open() {
282        let mut cb = DefaultCircuitBreaker::new(2, Duration::from_secs(60));
283        cb.record_failure();
284        cb.record_failure();
285        assert_eq!(cb.state(), CircuitState::Open);
286        assert!(cb.reset());
287        assert_eq!(cb.state(), CircuitState::Closed);
288        assert_eq!(cb.consecutive_failures(), 0);
289        // 重置后立即可执行
290        assert!(cb.can_execute());
291        // 仍需累计到阈值才会再次跳闸
292        cb.record_failure();
293        assert_eq!(cb.state(), CircuitState::Closed);
294    }
295
296    #[test]
297    fn test_circuit_breaker_reset_from_half_open() {
298        let mut cb = DefaultCircuitBreaker::new(1, Duration::from_millis(10));
299        cb.record_failure();
300        assert_eq!(cb.state(), CircuitState::Open);
301        std::thread::sleep(Duration::from_millis(20));
302        assert!(cb.can_execute());
303        assert_eq!(cb.state(), CircuitState::HalfOpen);
304        // HalfOpen 状态下手动重置
305        assert!(cb.reset());
306        assert_eq!(cb.state(), CircuitState::Closed);
307    }
308
309    #[test]
310    fn test_circuit_breaker_reset_idempotent_when_closed() {
311        let mut cb = DefaultCircuitBreaker::new(3, Duration::from_secs(60));
312        assert!(!cb.reset());
313        assert_eq!(cb.state(), CircuitState::Closed);
314        // 存在失败计数但未跳闸时,reset 视为变更
315        cb.record_failure();
316        cb.record_failure();
317        assert!(cb.reset());
318        assert_eq!(cb.consecutive_failures(), 0);
319        // 再次 reset 无变更
320        assert!(!cb.reset());
321    }
322
323    #[test]
324    fn test_circuit_state_variants_distinct() {
325        assert_ne!(CircuitState::Closed, CircuitState::Open);
326        assert_ne!(CircuitState::Open, CircuitState::HalfOpen);
327        assert_ne!(CircuitState::Closed, CircuitState::HalfOpen);
328    }
329
330    #[test]
331    fn test_default_circuit_breaker_default_config() {
332        let cb = DefaultCircuitBreaker::default();
333        assert_eq!(cb.failure_threshold(), 5);
334        assert_eq!(cb.reset_timeout(), Duration::from_secs(30));
335        assert_eq!(cb.state(), CircuitState::Closed);
336    }
337
338    #[test]
339    fn test_circuit_breaker_send_sync() {
340        fn assert_send_sync<T: Send + Sync>() {}
341        assert_send_sync::<DefaultCircuitBreaker>();
342        assert_send_sync::<CircuitState>();
343    }
344
345    /// 验证 trait object 可用(动态分发)
346    #[test]
347    fn test_circuit_breaker_via_trait_object() {
348        let cb: Box<dyn CircuitBreaker> =
349            Box::new(DefaultCircuitBreaker::new(2, Duration::from_secs(60)));
350        // 只能调用 trait 方法
351        assert_eq!(cb.state(), CircuitState::Closed);
352    }
353}