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

Hedging.pudu

Pudu169 lines7.9 KB

GitHub ↗
1/** @Resilience.Hedging.Strategy — races extra attempts against a slow one */2module PuduLangResilience.Hedging34import Std.Channel as Channel5import Std.Concurrent.Cancel as Cancel6import Std.Option as Option7import PuduLangResilience.Clock as Clock8import PuduLangResilience.Constants.Events as Events9import PuduLangResilience.Context as Context10import PuduLangResilience.Hedging.Attempts as Attempts11import PuduLangResilience.Predicate as Predicate12import PuduLangResilience as Resilience13import PuduLangResilience.Strategy as Strategy14import PuduLangResilience.Telemetry as Telemetry1516/** @Resilience.Hedging.Action — what a hedged attempt is generated from */17export type Action[T, E] = {18  primaryContext: Context.Context,19  actionContext: Context.Context,20  attempt: Int,21  callback: Strategy.Callback[T, E]22}2324/** @Resilience.Hedging.DelayArguments — the attempt a hedging delay is chosen for */25export type DelayArguments = { context: Context.Context, attempt: Int }2627/** @Resilience.Hedging.Hedged — the hedged attempt about to start */28export type Hedged = { primaryContext: Context.Context, actionContext: Context.Context, attempt: Int }2930/** @Resilience.Hedging.Options — how many attempts to race and when to start them */31export type Options[T, E] = {32  name: Option[Str],33  shouldHandle: Predicate.Predicate[T, E],34  maxHedgedAttempts: Int,35  delay: Int,36  actionGenerator: Option[fn(Action[T, E]) -> Option[Strategy.Callback[T, E]]],37  delayGenerator: Option[fn(DelayArguments) -> Int],38  onHedging: Option[fn(Hedged) -> ()]39}4041/// A delay that starts the next hedged attempt only once an earlier attempt's outcome is handled.42export const WAIT_FOR_OUTCOME: Int = -14344/// The most hedged attempts an option may name.45const MOST_HEDGED: Int = 104647/// The longest delay an option may name: one day, in milliseconds.48const LONGEST_DELAY: Int = 864000004950/// One hedged attempt, started after two seconds, for every failure except a cancellation.51export fn defaults[T, E]() -> Options[T, E] {52  Options{name: None, shouldHandle: Predicate.failures(), maxHedgedAttempts: 1, delay: 2000, actionGenerator: None, delayGenerator: None, onHedging: None}53}5455/// Every reason the options cannot be followed.56export fn validate[T, E](options: &Options[T, E]) -> Array[Str] {57  var problems: Array[Str] = []58  if options.maxHedgedAttempts < 1 || options.maxHedgedAttempts > MOST_HEDGED { problems = problems.push("maxHedgedAttempts must be from 1 to 10") }59  if options.delay < WAIT_FOR_OUTCOME || options.delay > LONGEST_DELAY { problems = problems.push("delay must be from -1 to 86400000 ms") }60  problems61}6263/// A hedging strategy following the options.64export fn strategy[T, E](options: Options[T, E]) -> Strategy.Strategy[T, E] {65  let pace = if options.delay == WAIT_FOR_OUTCOME { "after each handled outcome" } else { "every " + show(options.delay) + " ms" }66  Strategy.Strategy {67    name: options.name,68    kind: "Hedging",69    summary: "up to " + show(options.maxHedgedAttempts) + " hedged attempts " + pace,70    problems: validate(&options),71    attach: Strategy.detached,72    execute: fn(runtime: Strategy.Runtime, context: Context.Context, callback: Strategy.Callback[T, E]) -> Resilience.Outcome[T, E] {73      execute(&options, &runtime, context, callback)74    }75  }76}7778/// Runs the callback and, each time the delay passes or an attempt's outcome is handled, one more79/// attempt, until an attempt answers an outcome that is not handled or every attempt has80/// answered. The first unhandled outcome wins; without one, the last outcome to arrive is81/// answered. The other attempts are cancelled and waited for, and the winner's properties are82/// copied into the caller's context.83fn execute[T, E](options: &Options[T, E], runtime: &Strategy.Runtime, context: Context.Context, callback: Strategy.Callback[T, E]) -> Resilience.Outcome[T, E] {84  let total = options.maxHedgedAttempts + 185  let finished: Channel.Channel[Int] = Channel.channel(total)86  var attempts: Array[Attempts.Attempt[T, E]] = Option.toArray(start(options, runtime, &context, callback, 0, &finished))87  var exhausted = false88  var answered = 089  var delayFor = attempts.length()90  var delay = waitFor(options, context, delayFor)91  loop {92    if !exhausted && delayFor != attempts.length() {93      delayFor = attempts.length()94      delay = waitFor(options, context, delayFor)95    }96    let wait = if exhausted { Attempts.FOREVER } else { delay }97    match Attempts.next(&finished, &runtime.clock, wait) {98      case None => {99        let grown = launchNext(options, runtime, &context, callback, &attempts, total, &finished)100        attempts = grown[0]101        exhausted = grown[1]102      }103      case Some(index) => {104        answered = answered + 1105        let attempt = attempts[index]106        let outcome = Attempts.outcomeOf(&attempt)107        let handled = (options.shouldHandle)(Predicate.arguments(outcome, attempt.context, index))108        let severity = if handled { Telemetry.Warning } else { Telemetry.Information }109        let duration = Clock.now(&runtime.clock) - attempt.started110        Telemetry.report(&runtime.telemetry, &context, Events.EXECUTION_ATTEMPT, severity, Some(Resilience.summarize(&outcome)), Telemetry.Attempted(index, duration, handled))111        if handled && !exhausted {112          let grown = launchNext(options, runtime, &context, callback, &attempts, total, &finished)113          attempts = grown[0]114          exhausted = grown[1]115        }116        if !handled || (exhausted && answered == attempts.length()) {117          Attempts.settle(&attempts, "another hedged attempt was accepted")118          Context.adopt(&context, &attempt.context)119          return outcome120        }121      }122    }123  }124}125126/// How long to wait for an attempt to finish before starting the attempt numbered `launched`.127fn waitFor[T, E](options: &Options[T, E], context: Context.Context, launched: Int) -> Int {128  let delay = nextDelay(options, context, launched)129  if delay == WAIT_FOR_OUTCOME { Attempts.FOREVER } else if delay < 0 { 0 } else { delay }130}131132/// The attempts with the next one started, and whether no further attempt can start.133fn launchNext[T, E](options: &Options[T, E], runtime: &Strategy.Runtime, context: &Context.Context, callback: Strategy.Callback[T, E], attempts: &Array[Attempts.Attempt[T, E]], total: Int, finished: &Channel.Channel[Int]) -> (Array[Attempts.Attempt[T, E]], Bool) {134  if attempts.length() >= total { return (*attempts, true) }135  match start(options, runtime, context, callback, attempts.length(), finished) {136    case Some(attempt) => {137      let grown = attempts.push(attempt)138      (grown, grown.length() >= total)139    }140    case None => (*attempts, true)141  }142}143144/// The delay before the attempt numbered `attempt`, asked once per attempt.145fn nextDelay[T, E](options: &Options[T, E], context: Context.Context, attempt: Int) -> Int {146  match options.delayGenerator {147    case Some(generate) => generate(DelayArguments{context: context, attempt: attempt})148    case None => options.delay149  }150}151152/// Starts the attempt numbered `index` with its own copy of the context, or `None` when the153/// action generator declines to make one.154fn start[T, E](options: &Options[T, E], runtime: &Strategy.Runtime, primary: &Context.Context, callback: Strategy.Callback[T, E], index: Int, finished: &Channel.Channel[Int]) -> Option[Attempts.Attempt[T, E]] {155  let forked = Context.fork(primary, Cancel.child(&primary.token))156  let action = if index == 0 { Some(callback) } else {157    match options.actionGenerator {158      case Some(generate) => generate(Action{primaryContext: *primary, actionContext: forked, attempt: index, callback: callback})159      case None => Some(callback)160    }161  }162  let chosen = action ?163  if index > 0 {164    Telemetry.report(&runtime.telemetry, primary, Events.ON_HEDGING, Telemetry.Warning, None, Telemetry.Hedged(index))165    if let Some(notify) = options.onHedging { notify(Hedged{primaryContext: *primary, actionContext: forked, attempt: index}) }166  }167  Some(Attempts.launch(index, forked, chosen, finished, Clock.now(&runtime.clock)))168}169