1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
17pub struct BreakerAttempt(pub u64);
18
19pub enum ClosedPhase<Route> {
21 Idle { consecutive_failures: u32 },
23 Awaiting {
25 consecutive_failures: u32,
26 attempt: BreakerAttempt,
27 reply_to: Route,
28 },
29}
30
31pub enum ProbePhase<Route> {
33 Available,
35 Awaiting {
37 attempt: BreakerAttempt,
38 reply_to: Route,
39 },
40}
41
42pub enum BreakerPhase<Route> {
44 Closed(ClosedPhase<Route>),
46 Open { generation: TimerGeneration },
48 Probing {
50 generation: TimerGeneration,
51 phase: ProbePhase<Route>,
52 },
53 Exhausted,
55}
56
57#[derive(Debug, Error, Clone, Copy, PartialEq, Eq)]
59pub enum BreakerConfigError {
60 #[error("circuit-breaker reset delay must be non-zero")]
62 ZeroResetDelay,
63}
64
65#[derive(Debug, Clone, Copy, PartialEq, Eq)]
67pub enum BreakerRejection {
68 Busy,
70 Open { generation: TimerGeneration },
72 Exhausted,
74}
75
76#[derive(Debug, Clone, Copy, PartialEq, Eq)]
78pub enum BreakerOutcome {
79 Admitted { attempt: BreakerAttempt },
82 ProbeAdmitted { attempt: BreakerAttempt },
86 Rejected(BreakerRejection),
88 Succeeded { attempt: BreakerAttempt },
90 FailureRecorded {
92 attempt: BreakerAttempt,
93 consecutive_failures: u32,
94 },
95 Opened {
97 attempt: BreakerAttempt,
98 generation: TimerGeneration,
99 },
100}
101
102#[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#[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
124pub enum BreakerMessage<Route> {
126 Admit { reply_to: Route },
128 Succeeded { attempt: BreakerAttempt },
130 Failed { attempt: BreakerAttempt },
132}
133
134#[derive(behavior_macros::SendProduct)]
136pub struct BreakerSends<ReplySends, Schedules> {
137 pub replies: ReplySends,
139 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
147pub 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 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 #[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}