
Hedging.pudu
Pudu169 lines7.9 KB
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}404142export const WAIT_FOR_OUTCOME: Int = -1434445const MOST_HEDGED: Int = 10464748const LONGEST_DELAY: Int = 86400000495051export 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}545556export 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}626364export 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}77787980818283fn 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}125126127fn 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}131132133fn 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}143144145fn 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}151152153154fn 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