Skip to main content

behavior_actors/routing/
circuit_breaker.rs

1//! Single-flight circuit admission with explicit reset-timer evidence.
2
3use std::{num::NonZeroU32, time::Duration};
4
5#[cfg(test)]
6use behavior::Recipient;
7use behavior::{
8    Actions, Address, Behavior, BehaviorActed, BehaviorBase, EventLayer, InterpreterRequests,
9    Never, NoBirths, SendEffects, User,
10};
11use thiserror::Error;
12
13use crate::{DeliveryRoute, ScheduleAfter, TimedEvent, TimerGeneration, TimerId};
14
15/// Identity of one admitted operation.
16#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
17pub struct BreakerAttempt(pub u64);
18
19/// Closed-state subphase; only one operation may own the breaker at a time.
20pub enum ClosedPhase<Route> {
21    /// The breaker may admit one operation.
22    Idle { consecutive_failures: u32 },
23    /// One admitted operation owns its terminal report capability.
24    Awaiting {
25        consecutive_failures: u32,
26        attempt: BreakerAttempt,
27        reply_to: Route,
28    },
29}
30
31/// Half-open probe subphase.
32pub enum ProbePhase<Route> {
33    /// Exactly one probe may be admitted.
34    Available,
35    /// The admitted probe owns its terminal report capability.
36    Awaiting {
37        attempt: BreakerAttempt,
38        reply_to: Route,
39    },
40}
41
42/// Complete circuit-breaker phase sum.
43pub enum BreakerPhase<Route> {
44    /// Ordinary admission with a consecutive-failure count.
45    Closed(ClosedPhase<Route>),
46    /// Admission is denied until matching timer evidence arrives.
47    Open { generation: TimerGeneration },
48    /// One recovery probe is available or in progress.
49    Probing {
50        generation: TimerGeneration,
51        phase: ProbePhase<Route>,
52    },
53    /// No fresh attempt or timer generation remains representable.
54    Exhausted,
55}
56
57/// Typed construction failure.
58#[derive(Debug, Error, Clone, Copy, PartialEq, Eq)]
59pub enum BreakerConfigError {
60    /// Reset delay must be non-zero.
61    #[error("circuit-breaker reset delay must be non-zero")]
62    ZeroResetDelay,
63}
64
65/// Why an admission request was declined.
66#[derive(Debug, Clone, Copy, PartialEq, Eq)]
67pub enum BreakerRejection {
68    /// A previously admitted operation has not reported a result.
69    Busy,
70    /// The breaker is open.
71    Open { generation: TimerGeneration },
72    /// Numeric freshness is exhausted.
73    Exhausted,
74}
75
76/// Facts emitted to the recipient supplied with an admission request.
77#[derive(Debug, Clone, Copy, PartialEq, Eq)]
78pub enum BreakerOutcome {
79    /// The caller may begin this uniquely identified operation while the
80    /// circuit is closed.
81    Admitted { attempt: BreakerAttempt },
82    /// The caller's request is the single half-open recovery probe: a later
83    /// [`BreakerOutcome::Succeeded`] closes the circuit, while a failed
84    /// completion reopens it with a fresh timer generation.
85    ProbeAdmitted { attempt: BreakerAttempt },
86    /// A request was not admitted.
87    Rejected(BreakerRejection),
88    /// The admitted operation succeeded and closed the circuit.
89    Succeeded { attempt: BreakerAttempt },
90    /// Failure was recorded while the circuit remained closed.
91    FailureRecorded {
92        attempt: BreakerAttempt,
93        consecutive_failures: u32,
94    },
95    /// Failure reached the threshold and opened the circuit.
96    Opened {
97        attempt: BreakerAttempt,
98        generation: TimerGeneration,
99    },
100}
101
102/// Complete terminal fact submitted for one admitted operation.
103#[derive(Debug, Clone, Copy, PartialEq, Eq)]
104pub enum BreakerCompletion {
105    Succeeded { attempt: BreakerAttempt },
106    Failed { attempt: BreakerAttempt },
107}
108
109impl BreakerCompletion {
110    const fn attempt(self) -> BreakerAttempt {
111        match self {
112            Self::Succeeded { attempt } | Self::Failed { attempt } => attempt,
113        }
114    }
115}
116
117/// Typed transition failure that retains an unowned completion command.
118#[derive(Debug, Error, Clone, Copy, PartialEq, Eq)]
119pub enum BreakerError {
120    #[error("completion does not belong to the operation currently admitted")]
121    UnexpectedCompletion(BreakerCompletion),
122}
123
124/// Closed user-command protocol.
125pub enum BreakerMessage<Route> {
126    /// Request one operation admission.
127    Admit { reply_to: Route },
128    /// Report successful completion of an admitted operation.
129    Succeeded { attempt: BreakerAttempt },
130    /// Report failed completion of an admitted operation.
131    Failed { attempt: BreakerAttempt },
132}
133
134/// Named output lanes for circuit facts and reset scheduling.
135#[derive(behavior_macros::SendProduct)]
136pub struct BreakerSends<ReplySends, Schedules> {
137    /// Admission and completion facts.
138    pub replies: ReplySends,
139    /// Relative reset requests interpreted by Bombay Timers.
140    pub schedules: Schedules,
141}
142
143type BreakerEvent<A, Route> = TimedEvent<User<A, BreakerMessage<Route>>>;
144type BreakerActions<A, ReplySends> =
145    Actions<A, Never, BreakerSends<ReplySends, InterpreterRequests<ScheduleAfter>>, NoBirths>;
146
147/// Pure single-flight closed/open/probing circuit-breaker fold.
148///
149/// Admission and completion are distinct typed events. Closed and probing
150/// phases admit at most one operation, so a completion cannot be attributed by
151/// timing or adjacency. Consecutive closed failures open the circuit at the
152/// configured threshold; matching reset evidence makes exactly one probe
153/// available. Probe success closes the breaker and probe failure reopens it
154/// with a fresh generation. A stale or duplicate application completion is
155/// returned as [`BreakerError::UnexpectedCompletion`]; stale timer evidence is
156/// consumed as correlation evidence without changing state. Empty
157/// initialization, single-flight policy, failure counting, and reset ordering
158/// are Bombay policy. Timer scheduling is interpreted by Bombay Timers; the
159/// protected operation remains ordinary domain behavior. No transition panics.
160pub struct CircuitBreaker<A, Route>
161where
162    A: Address,
163    Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = BreakerOutcome>>,
164{
165    threshold: NonZeroU32,
166    reset_after: Duration,
167    timer_id: TimerId,
168    next_attempt: u64,
169    phase: BreakerPhase<Route>,
170    marker: core::marker::PhantomData<fn() -> A>,
171}
172
173impl<A, Route> CircuitBreaker<A, Route>
174where
175    A: Address,
176    Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = BreakerOutcome>> + Clone,
177{
178    /// Construct a closed breaker.
179    ///
180    /// # Errors
181    ///
182    /// Returns [`BreakerConfigError::ZeroResetDelay`] when no positive reset
183    /// interval was supplied.
184    pub fn new(
185        threshold: NonZeroU32,
186        reset_after: Duration,
187        timer_id: TimerId,
188    ) -> Result<Self, BreakerConfigError> {
189        if reset_after.is_zero() {
190            return Err(BreakerConfigError::ZeroResetDelay);
191        }
192        Ok(Self {
193            threshold,
194            reset_after,
195            timer_id,
196            next_attempt: 0,
197            phase: BreakerPhase::Closed(ClosedPhase::Idle {
198                consecutive_failures: 0,
199            }),
200            marker: core::marker::PhantomData,
201        })
202    }
203
204    /// Borrow the complete current phase.
205    #[must_use]
206    pub const fn phase(&self) -> &BreakerPhase<Route> {
207        &self.phase
208    }
209
210    fn reply(reply_to: Route, outcome: BreakerOutcome) -> BreakerActions<A, Route::Sends> {
211        Actions::send(BreakerSends {
212            replies: reply_to.deliver(outcome),
213            schedules: InterpreterRequests::empty(),
214        })
215    }
216
217    fn admit(&mut self, reply_to: Route) -> BreakerActions<A, Route::Sends> {
218        let attempt = BreakerAttempt(self.next_attempt);
219        match &self.phase {
220            BreakerPhase::Closed(ClosedPhase::Idle {
221                consecutive_failures,
222            }) => {
223                let Some(next) = self.next_attempt.checked_add(1) else {
224                    self.phase = BreakerPhase::Exhausted;
225                    return Self::reply(
226                        reply_to,
227                        BreakerOutcome::Rejected(BreakerRejection::Exhausted),
228                    );
229                };
230                let failures = *consecutive_failures;
231                self.next_attempt = next;
232                self.phase = BreakerPhase::Closed(ClosedPhase::Awaiting {
233                    consecutive_failures: failures,
234                    attempt,
235                    reply_to: reply_to.clone(),
236                });
237                Self::reply(reply_to, BreakerOutcome::Admitted { attempt })
238            }
239            BreakerPhase::Probing {
240                generation,
241                phase: ProbePhase::Available,
242            } => {
243                let Some(next) = self.next_attempt.checked_add(1) else {
244                    self.phase = BreakerPhase::Exhausted;
245                    return Self::reply(
246                        reply_to,
247                        BreakerOutcome::Rejected(BreakerRejection::Exhausted),
248                    );
249                };
250                let generation = *generation;
251                self.next_attempt = next;
252                self.phase = BreakerPhase::Probing {
253                    generation,
254                    phase: ProbePhase::Awaiting {
255                        attempt,
256                        reply_to: reply_to.clone(),
257                    },
258                };
259                Self::reply(reply_to, BreakerOutcome::ProbeAdmitted { attempt })
260            }
261            BreakerPhase::Open { generation } => Self::reply(
262                reply_to,
263                BreakerOutcome::Rejected(BreakerRejection::Open {
264                    generation: *generation,
265                }),
266            ),
267            BreakerPhase::Exhausted => Self::reply(
268                reply_to,
269                BreakerOutcome::Rejected(BreakerRejection::Exhausted),
270            ),
271            _ => Self::reply(reply_to, BreakerOutcome::Rejected(BreakerRejection::Busy)),
272        }
273    }
274
275    fn complete(
276        &mut self,
277        completion: BreakerCompletion,
278    ) -> Result<BreakerActions<A, Route::Sends>, BreakerError> {
279        let attempt = completion.attempt();
280        let ownership = match &self.phase {
281            BreakerPhase::Closed(ClosedPhase::Awaiting {
282                consecutive_failures,
283                attempt: current,
284                reply_to,
285            }) if *current == attempt => Some((reply_to.clone(), *consecutive_failures, None)),
286            BreakerPhase::Probing {
287                generation,
288                phase:
289                    ProbePhase::Awaiting {
290                        attempt: current,
291                        reply_to,
292                    },
293            } if *current == attempt => Some((reply_to.clone(), 0, Some(*generation))),
294            _ => None,
295        };
296        let Some((reply_to, failures, probe_generation)) = ownership else {
297            return Err(BreakerError::UnexpectedCompletion(completion));
298        };
299        if matches!(completion, BreakerCompletion::Succeeded { .. }) {
300            self.phase = BreakerPhase::Closed(ClosedPhase::Idle {
301                consecutive_failures: 0,
302            });
303            return Ok(Self::reply(reply_to, BreakerOutcome::Succeeded { attempt }));
304        }
305        let failures = failures.saturating_add(1);
306        if probe_generation.is_none() && failures < self.threshold.get() {
307            self.phase = BreakerPhase::Closed(ClosedPhase::Idle {
308                consecutive_failures: failures,
309            });
310            return Ok(Self::reply(
311                reply_to,
312                BreakerOutcome::FailureRecorded {
313                    attempt,
314                    consecutive_failures: failures,
315                },
316            ));
317        }
318        let next_generation = match probe_generation {
319            Some(generation) => generation.0.checked_add(1).map(TimerGeneration),
320            None => Some(TimerGeneration(0)),
321        };
322        let Some(generation) = next_generation else {
323            self.phase = BreakerPhase::Exhausted;
324            return Ok(Self::reply(
325                reply_to,
326                BreakerOutcome::Rejected(BreakerRejection::Exhausted),
327            ));
328        };
329        self.phase = BreakerPhase::Open { generation };
330        Ok(Actions::send(BreakerSends {
331            replies: reply_to.deliver(BreakerOutcome::Opened {
332                attempt,
333                generation,
334            }),
335            schedules: InterpreterRequests::one(ScheduleAfter::new(
336                self.timer_id,
337                generation,
338                self.reset_after,
339            )),
340        }))
341    }
342}
343
344impl<A, Route> BehaviorBase for CircuitBreaker<A, Route>
345where
346    A: Address,
347    Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = BreakerOutcome>>,
348{
349    type Base = Self;
350    fn base(&self) -> &Self {
351        self
352    }
353}
354
355impl<A, Route> behavior::Protocol for CircuitBreaker<A, Route>
356where
357    A: Address,
358    Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = BreakerOutcome>>,
359{
360    type Addr = A;
361    type Msg = BreakerMessage<Route>;
362}
363
364impl<A, Route> Behavior for CircuitBreaker<A, Route>
365where
366    A: Address,
367    Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = BreakerOutcome>> + Clone,
368    Route::Sends: behavior::SendsFor<BreakerEvent<A, Route>>,
369{
370    type Protocol = Self;
371    type Event = BreakerEvent<A, Route>;
372    type Sends = BreakerSends<Route::Sends, InterpreterRequests<ScheduleAfter>>;
373    type Ph = Never;
374    type Error = BreakerError;
375    type Birth = NoBirths;
376    fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
377        match event {
378            EventLayer::Inner(event) => match event.message {
379                BreakerMessage::Admit { reply_to } => Ok(self.admit(reply_to)),
380                BreakerMessage::Succeeded { attempt } => {
381                    self.complete(BreakerCompletion::Succeeded { attempt })
382                }
383                BreakerMessage::Failed { attempt } => {
384                    self.complete(BreakerCompletion::Failed { attempt })
385                }
386            },
387            EventLayer::Owned(elapsed) => {
388                if let BreakerPhase::Open { generation } = self.phase
389                    && elapsed.id == self.timer_id
390                    && elapsed.generation == generation
391                {
392                    self.phase = BreakerPhase::Probing {
393                        generation,
394                        phase: ProbePhase::Available,
395                    };
396                }
397                Ok(Actions::cont())
398            }
399        }
400    }
401}
402
403#[cfg(test)]
404mod tests {
405    use super::*;
406    use crate::{Activate as _, TimerElapsed};
407    use behavior::MailAddr;
408
409    struct Reply;
410    impl behavior::Protocol for Reply {
411        type Addr = MailAddr;
412        type Msg = BreakerOutcome;
413    }
414
415    impl Behavior for Reply {
416        type Protocol = Self;
417        type Event = User<MailAddr, behavior::BehaviorMessage<Self>>;
418        type Sends = Vec<Never>;
419        type Ph = Never;
420        type Error = Never;
421        type Birth = NoBirths;
422        fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
423            Ok(Actions::cont())
424        }
425    }
426
427    fn subject() -> crate::Active<CircuitBreaker<MailAddr, Recipient<Reply>>> {
428        (CircuitBreaker::new(
429            NonZeroU32::new(2).unwrap(),
430            Duration::from_secs(1),
431            TimerId(8),
432        )
433        .unwrap())
434        .initialize()
435        .unwrap()
436        .behavior
437    }
438    fn reply() -> Recipient<Reply> {
439        Recipient::global(MailAddr(9))
440    }
441    fn admit(
442        subject: &mut crate::Active<CircuitBreaker<MailAddr, Recipient<Reply>>>,
443    ) -> BreakerAttempt {
444        let actions = subject
445            .receive(MailAddr(0), BreakerMessage::Admit { reply_to: reply() })
446            .unwrap();
447        match actions.sends.replies[0].message {
448            BreakerOutcome::Admitted { attempt } | BreakerOutcome::ProbeAdmitted { attempt } => {
449                attempt
450            }
451            _ => unreachable!("test expected admission"),
452        }
453    }
454
455    #[test]
456    fn threshold_opens_matching_elapsed_allows_one_probe() {
457        let mut subject = subject();
458        let first = admit(&mut subject);
459        let first_failure = subject
460            .receive(MailAddr(0), BreakerMessage::Failed { attempt: first })
461            .unwrap();
462        assert!(matches!(
463            first_failure.sends.replies.as_slice(),
464            [behavior::Delivery {
465                message: BreakerOutcome::FailureRecorded {
466                    attempt,
467                    consecutive_failures: 1,
468                },
469                ..
470            }] if *attempt == first
471        ));
472        assert!(first_failure.sends.schedules.is_empty());
473        assert!(first_failure.creates.is_empty());
474        assert_eq!(first_failure.become_, behavior::Step::Continue);
475        let second = admit(&mut subject);
476        let opened = subject
477            .receive(MailAddr(0), BreakerMessage::Failed { attempt: second })
478            .unwrap();
479        assert_eq!(
480            opened.sends.schedules.as_slice()[0].generation,
481            TimerGeneration(0)
482        );
483        let denied = subject
484            .receive(MailAddr(0), BreakerMessage::Admit { reply_to: reply() })
485            .unwrap();
486        assert!(matches!(
487            denied.sends.replies[0].message,
488            BreakerOutcome::Rejected(BreakerRejection::Open { .. })
489        ));
490        let elapsed = subject
491            .on_path(TimerElapsed::new(TimerId(8), TimerGeneration(0)))
492            .unwrap();
493        assert!(elapsed.sends.replies.is_empty());
494        assert!(elapsed.sends.schedules.is_empty());
495        assert!(elapsed.creates.is_empty());
496        assert_eq!(elapsed.become_, behavior::Step::Continue);
497        let probe = admit(&mut subject);
498        assert_eq!(probe, BreakerAttempt(2));
499        let busy = subject
500            .receive(MailAddr(0), BreakerMessage::Admit { reply_to: reply() })
501            .unwrap();
502        assert!(matches!(
503            busy.sends.replies[0].message,
504            BreakerOutcome::Rejected(BreakerRejection::Busy)
505        ));
506    }
507
508    #[test]
509    fn stale_completion_is_returned_while_stale_timer_evidence_is_consumed() {
510        let mut subject = subject();
511        let attempt = admit(&mut subject);
512        let stale = BreakerCompletion::Succeeded {
513            attempt: BreakerAttempt(99),
514        };
515        let rejection = subject.receive(
516            MailAddr(0),
517            BreakerMessage::Succeeded {
518                attempt: BreakerAttempt(99),
519            },
520        );
521        assert!(matches!(
522            rejection,
523            Err(BreakerError::UnexpectedCompletion(returned)) if returned == stale
524        ));
525        assert!(
526            subject
527                .on_path(TimerElapsed::new(TimerId(8), TimerGeneration(7)))
528                .unwrap()
529                .sends
530                .replies
531                .is_empty()
532        );
533        let success = subject
534            .receive(MailAddr(0), BreakerMessage::Succeeded { attempt })
535            .unwrap();
536        assert!(matches!(
537            success.sends.replies[0].message,
538            BreakerOutcome::Succeeded { .. }
539        ));
540    }
541
542    fn breaker() -> CircuitBreaker<MailAddr, Recipient<Reply>> {
543        CircuitBreaker::new(
544            NonZeroU32::new(2).unwrap(),
545            Duration::from_secs(1),
546            TimerId(8),
547        )
548        .unwrap()
549    }
550
551    #[test]
552    fn attempt_counter_exhaustion_is_typed_not_wrapped() {
553        let mut breaker = breaker();
554        breaker.next_attempt = u64::MAX;
555        let mut subject = (breaker).initialize().unwrap().behavior;
556        let actions = subject
557            .receive(MailAddr(0), BreakerMessage::Admit { reply_to: reply() })
558            .unwrap();
559        assert!(matches!(*subject.phase(), BreakerPhase::Exhausted));
560        assert!(matches!(
561            actions.sends.replies[0].message,
562            BreakerOutcome::Rejected(BreakerRejection::Exhausted)
563        ));
564    }
565
566    #[test]
567    fn timer_generation_exhaustion_is_typed_not_wrapped() {
568        let mut breaker = breaker();
569        breaker.phase = BreakerPhase::Probing {
570            generation: TimerGeneration(u64::MAX),
571            phase: ProbePhase::Awaiting {
572                attempt: BreakerAttempt(0),
573                reply_to: reply(),
574            },
575        };
576        let mut subject = (breaker).initialize().unwrap().behavior;
577        let actions = subject
578            .receive(
579                MailAddr(0),
580                BreakerMessage::Failed {
581                    attempt: BreakerAttempt(0),
582                },
583            )
584            .unwrap();
585        assert!(matches!(*subject.phase(), BreakerPhase::Exhausted));
586        assert!(matches!(
587            actions.sends.replies[0].message,
588            BreakerOutcome::Rejected(BreakerRejection::Exhausted)
589        ));
590    }
591}