Pudu programming language
Menu
Package

@chrismichaelps / pudu-lang-resilience

Resilience pipelines for Pudu: retry, circuit breaker, timeout, fallback, hedging, rate limiting, and chaos injection

0.1.0Apache-2.01

InstallClose

CircuitBreaker.pudu

Pudu231 lines11.7 KB

GitHub ↗
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] }5354/// The shortest sampling or break duration an option may name, in milliseconds.55const SHORTEST: Int = 5005657/// The longest sampling or break duration an option may name: one day, in milliseconds.58const LONGEST: Int = 864000005960/// Opens at a tenth of at least 100 executions failing within 30 seconds, for 5 seconds.61export 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}7778/// Every reason the options cannot be followed, not counting an attached state provider.79export 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}8788/// A manual control holding no circuit, not isolated.89export fn manualControl() -> ManualControl { ManualControl{isolated: Shared.shared(false), hooks: Shared.shared([])} }9091/// Isolates every circuit the control holds, and every circuit attached to it later, until closed.92export 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}9697/// Closes every circuit the control holds.98export 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}102103/// Whether the control last isolated its circuits.104export fn isIsolated(control: &ManualControl) -> Bool { Shared.current(&control.isolated) }105106/// A state provider attached to no circuit.107export fn stateProvider() -> StateProvider { StateProvider{reader: Shared.shared(None)} }108109/// The state of the provider's circuit, or `None` before it is attached.110export fn state(provider: &StateProvider) -> Option[CircuitState] {111  let read = Shared.current(&provider.reader) ?112  Some(publicState(&read().phase))113}114115/// The last outcome the provider's circuit handled, while the circuit is not closed.116export fn lastHandled(provider: &StateProvider) -> Option[Str] {117  let read = Shared.current(&provider.reader) ?118  read().lastHandled119}120121/// A circuit breaker strategy following the options. Every pipeline holding the same strategy122/// shares its one circuit.123export 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}152153/// The public name of a circuit phase.154fn 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}162163/// Runs the callback when the circuit admits it and records whether its outcome was handled;164/// answers `BrokenCircuit` or `IsolatedCircuit` without running it otherwise.165fn 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}203204/// Isolates or closes the circuit by hand, reporting the transition through the runtime the205/// strategy was last attached to.206fn 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}217218/// Reports an opened circuit and calls `onOpened`.219fn 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}224225/// Reports a closed circuit and calls `onClosed`.226fn 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