
CircuitBreaker.pudu
Pudu231 lines11.7 KB
1/** @Resilience.CircuitBreaker.Strategy — stops calling a failing dependency for a while */2module PuduLangResilience.CircuitBreaker34import Std.Decimal as Decimal5import PuduLangResilience.Clock as Clock6import PuduLangResilience.Constants.Events as Events7import PuduLangResilience.Context as Context8import PuduLangResilience.Domain.Circuit as Circuit9import PuduLangResilience.Domain.Health as Health10import PuduLangResilience.Predicate as Predicate11import PuduLangResilience as Resilience12import PuduLangResilience.Strategy as Strategy13import PuduLangResilience.Telemetry as Telemetry14import PuduLangResilience.Utils.Numeric as Numeric15import PuduLangResilience.Utils.Shared as Shared1617/** @Resilience.CircuitBreaker.CircuitState — which executions a circuit admits */18export type CircuitState = Closed | Open | HalfOpen | Isolated1920/** @Resilience.CircuitBreaker.BreakArguments — the health a break duration is chosen from */21export type BreakArguments = { failureRate: Decimal, failureCount: Int, halfOpenAttempts: Int, context: Context.Context }2223/** @Resilience.CircuitBreaker.Opening — why and for how long a circuit opened */24export type Opening[T, E] = { context: Context.Context, outcome: Option[Resilience.Outcome[T, E]], breakDuration: Int, manual: Bool }2526/** @Resilience.CircuitBreaker.Closing — what closed a circuit */27export type Closing[T, E] = { context: Context.Context, outcome: Option[Resilience.Outcome[T, E]], manual: Bool }2829/** @Resilience.CircuitBreaker.ManualControl — isolates and closes circuits by hand */30export type ManualControl = { isolated: Shared.Shared[Bool], hooks: Shared.Shared[Array[fn(Bool) -> ()]] }3132/** @Resilience.CircuitBreaker.StateProvider — reads the state of one circuit */33export type StateProvider = { reader: Shared.Shared[Option[fn() -> Circuit.Machine]] }3435/** @Resilience.CircuitBreaker.Options — when a circuit opens and for how long */36export type Options[T, E] = {37 name: Option[Str],38 shouldHandle: Predicate.Predicate[T, E],39 failureRatio: Decimal,40 minimumThroughput: Int,41 samplingDuration: Int,42 breakDuration: Int,43 breakDurationGenerator: Option[fn(BreakArguments) -> Int],44 manualControl: Option[ManualControl],45 stateProvider: Option[StateProvider],46 onOpened: Option[fn(Opening[T, E]) -> ()],47 onClosed: Option[fn(Closing[T, E]) -> ()],48 onHalfOpened: Option[fn(Context.Context) -> ()]49}5051/** @Resilience.CircuitBreaker.Breaker — one circuit and the runtime it reports through */52type Breaker = { machine: Shared.Shared[Circuit.Machine], runtime: Shared.Shared[Strategy.Runtime] }535455const SHORTEST: Int = 500565758const LONGEST: Int = 86400000596061export fn defaults[T, E]() -> Options[T, E] {62 Options {63 name: None,64 shouldHandle: Predicate.failures(),65 failureRatio: 0.1d,66 minimumThroughput: 100,67 samplingDuration: 30000,68 breakDuration: 5000,69 breakDurationGenerator: None,70 manualControl: None,71 stateProvider: None,72 onOpened: None,73 onClosed: None,74 onHalfOpened: None75 }76}777879export fn validate[T, E](options: &Options[T, E]) -> Array[Str] {80 var problems: Array[Str] = []81 if options.failureRatio < 0d || options.failureRatio > 1d { problems = problems.push("failureRatio must be from 0 to 1") }82 if options.minimumThroughput < 2 { problems = problems.push("minimumThroughput must be at least 2") }83 if options.samplingDuration < SHORTEST || options.samplingDuration > LONGEST { problems = problems.push("samplingDuration must be from 500 to 86400000 ms") }84 if options.breakDuration < SHORTEST || options.breakDuration > LONGEST { problems = problems.push("breakDuration must be from 500 to 86400000 ms") }85 problems86}878889export fn manualControl() -> ManualControl { ManualControl{isolated: Shared.shared(false), hooks: Shared.shared([])} }909192export fn isolate(control: &ManualControl) -> () {93 Shared.update(&control.isolated, fn(_held: Bool) -> Bool { true })94 for hook in Shared.current(&control.hooks) { hook(true) }95}969798export fn close(control: &ManualControl) -> () {99 Shared.update(&control.isolated, fn(_held: Bool) -> Bool { false })100 for hook in Shared.current(&control.hooks) { hook(false) }101}102103104export fn isIsolated(control: &ManualControl) -> Bool { Shared.current(&control.isolated) }105106107export fn stateProvider() -> StateProvider { StateProvider{reader: Shared.shared(None)} }108109110export fn state(provider: &StateProvider) -> Option[CircuitState] {111 let read = Shared.current(&provider.reader) ?112 Some(publicState(&read().phase))113}114115116export fn lastHandled(provider: &StateProvider) -> Option[Str] {117 let read = Shared.current(&provider.reader) ?118 read().lastHandled119}120121122123export fn strategy[T, E](options: Options[T, E]) -> Strategy.Strategy[T, E] {124 let breaker = Breaker{machine: Shared.shared(Circuit.closed(options.samplingDuration)), runtime: Shared.shared(Strategy.standalone())}125 var problems = validate(&options)126 if let Some(provider) = options.stateProvider {127 let machine = breaker.machine128 let attached = Shared.change(&provider.reader, fn(held: Option[fn() -> Circuit.Machine]) -> (Option[fn() -> Circuit.Machine], Bool) {129 match held {130 case Some(_) => (held, false)131 case None => (Some(fn() -> Circuit.Machine { Shared.current(&machine) }), true)132 }133 })134 if !attached { problems = problems.push("the state provider is already attached to another circuit") }135 }136 if let Some(control) = options.manualControl {137 let hook = fn(isolating: Bool) -> () { manual(&options, &breaker, isolating) }138 Shared.update(&control.hooks, fn(hooks: Array[fn(Bool) -> ()]) -> Array[fn(Bool) -> ()] { hooks.push(hook) })139 if isIsolated(&control) { Shared.update(&breaker.machine, fn(held: Circuit.Machine) -> Circuit.Machine { Circuit.isolate(&held) }) }140 }141 Strategy.Strategy {142 name: options.name,143 kind: "CircuitBreaker",144 summary: "opens at " + Decimal.toText(options.failureRatio) + " of " + show(options.minimumThroughput) + " executions in " + show(options.samplingDuration) + " ms for " + show(options.breakDuration) + " ms",145 problems: problems,146 attach: fn(runtime: Strategy.Runtime) -> () { Shared.update(&breaker.runtime, fn(_held: Strategy.Runtime) -> Strategy.Runtime { runtime }) },147 execute: fn(runtime: Strategy.Runtime, context: Context.Context, callback: Strategy.Callback[T, E]) -> Resilience.Outcome[T, E] {148 execute(&options, &breaker, &runtime, context, callback)149 }150 }151}152153154fn publicState(phase: &Circuit.Phase) -> CircuitState {155 match phase {156 case Circuit.Closed => Closed157 case Circuit.Open => Open158 case Circuit.HalfOpen => HalfOpen159 case Circuit.Isolated => Isolated160 }161}162163164165fn execute[T, E](options: &Options[T, E], breaker: &Breaker, runtime: &Strategy.Runtime, context: Context.Context, callback: Strategy.Callback[T, E]) -> Resilience.Outcome[T, E] {166 let admitted = Clock.now(&runtime.clock)167 let breakDuration = options.breakDuration168 let admission = Shared.change(&breaker.machine, fn(held: Circuit.Machine) -> (Circuit.Machine, Circuit.Admission) { Circuit.admit(&held, admitted, breakDuration) })169 match admission {170 case Circuit.Refused => { return Err(Resilience.IsolatedCircuit) }171 case Circuit.Blocked(millis) => { return Err(Resilience.BrokenCircuit(millis)) }172 case Circuit.Probe => {173 Telemetry.report(&runtime.telemetry, &context, Events.ON_CIRCUIT_HALF_OPENED, Telemetry.Warning, None, Telemetry.HalfOpened)174 if let Some(notify) = options.onHalfOpened { notify(context) }175 }176 case Circuit.Admitted => ()177 }178 let outcome = callback(context)179 let finished = Clock.now(&runtime.clock)180 if !(options.shouldHandle)(Predicate.arguments(outcome, context, 0)) {181 let closedNow = Shared.change(&breaker.machine, fn(held: Circuit.Machine) -> (Circuit.Machine, Bool) { Circuit.succeeded(&held, finished) })182 if closedNow { reportClosed(options, runtime, context, Some(outcome), false) }183 return outcome184 }185 let thresholds = Circuit.Thresholds{failureRatio: options.failureRatio, minimumThroughput: options.minimumThroughput}186 let summary = Resilience.summarize(&outcome)187 let generator = options.breakDurationGenerator188 let opened = Shared.change(&breaker.machine, fn(held: Circuit.Machine) -> (Circuit.Machine, Option[Int]) {189 let judged = Circuit.failed(&held, finished, summary, &thresholds)190 if !judged[1] { return (judged[0], None) }191 let duration = match generator {192 case Some(generate) => {193 let totals = Health.info(&judged[0].health, finished)194 generate(BreakArguments{failureRate: Health.failureRate(&totals), failureCount: totals.failures, halfOpenAttempts: judged[0].halfOpenAttempts, context: context})195 }196 case None => breakDuration197 }198 (Circuit.open(&judged[0], finished, duration), Some(duration))199 })200 if let Some(duration) = opened { reportOpened(options, runtime, context, Some(outcome), duration, false) }201 outcome202}203204205206fn manual[T, E](options: &Options[T, E], breaker: &Breaker, isolating: Bool) -> () {207 let runtime = Shared.current(&breaker.runtime)208 let context = Context.create()209 if isolating {210 Shared.update(&breaker.machine, fn(held: Circuit.Machine) -> Circuit.Machine { Circuit.isolate(&held) })211 reportOpened(options, &runtime, context, None, Numeric.LARGEST, true)212 } else {213 let changed = Shared.change(&breaker.machine, fn(held: Circuit.Machine) -> (Circuit.Machine, Bool) { (Circuit.close(&held), held.phase != Circuit.Closed) })214 if changed { reportClosed(options, &runtime, context, None, true) }215 }216}217218219fn reportOpened[T, E](options: &Options[T, E], runtime: &Strategy.Runtime, context: Context.Context, outcome: Option[Resilience.Outcome[T, E]], duration: Int, byHand: Bool) -> () {220 let summary = match outcome { case Some(found) => Some(Resilience.summarize(&found)) case None => None }221 Telemetry.report(&runtime.telemetry, &context, Events.ON_CIRCUIT_OPENED, Telemetry.Error, summary, Telemetry.Opened(duration, byHand))222 if let Some(notify) = options.onOpened { notify(Opening{context: context, outcome: outcome, breakDuration: duration, manual: byHand}) }223}224225226fn reportClosed[T, E](options: &Options[T, E], runtime: &Strategy.Runtime, context: Context.Context, outcome: Option[Resilience.Outcome[T, E]], byHand: Bool) -> () {227 let summary = match outcome { case Some(found) => Some(Resilience.summarize(&found)) case None => None }228 Telemetry.report(&runtime.telemetry, &context, Events.ON_CIRCUIT_CLOSED, Telemetry.Information, summary, Telemetry.Closed(byHand))229 if let Some(notify) = options.onClosed { notify(Closing{context: context, outcome: outcome, manual: byHand}) }230}231