
Retry.pudu
Pudu144 lines5.6 KB
1/** @Resilience.Retry.Strategy — runs a callback again after handled failures */2module PuduLangResilience.Retry34import Std.Concurrent.Cancel as Cancel5import PuduLangResilience.Clock as Clock6import PuduLangResilience.Constants.Events as Events7import PuduLangResilience.Context as Context8import PuduLangResilience.Domain.Backoff as Delays9import PuduLangResilience.Predicate as Predicate10import PuduLangResilience.Randomizer as Randomizer11import PuduLangResilience as Resilience12import PuduLangResilience.Strategy as Strategy13import PuduLangResilience.Telemetry as Telemetry1415/** @Resilience.Retry.Backoff — how the delay grows with each retry */16export type Backoff = Constant | Linear | Exponential1718/** @Resilience.Retry.Retrying — the attempt that failed and the wait before the next */19export type Retrying[T, E] = {20 outcome: Resilience.Outcome[T, E],21 context: Context.Context,22 attempt: Int,23 delay: Int,24 duration: Int25}2627/** @Resilience.Retry.Options — when to retry and how long to wait */28export type Options[T, E] = {29 name: Option[Str],30 shouldHandle: Predicate.Predicate[T, E],31 maxRetryAttempts: Int,32 backoff: Backoff,33 delay: Int,34 maxDelay: Option[Int],35 useJitter: Bool,36 delayGenerator: Option[fn(Predicate.Arguments[T, E]) -> Option[Int]],37 onRetry: Option[fn(Retrying[T, E]) -> ()]38}394041export const UNLIMITED: Int = 9223372036854775807424344const LONGEST_DELAY: Int = 86400000454647export fn defaults[T, E]() -> Options[T, E] {48 Options {49 name: None,50 shouldHandle: Predicate.failures(),51 maxRetryAttempts: 3,52 backoff: Constant,53 delay: 2000,54 maxDelay: None,55 useJitter: false,56 delayGenerator: None,57 onRetry: None58 }59}606162export fn validate[T, E](options: &Options[T, E]) -> Array[Str] {63 var problems: Array[Str] = []64 if options.maxRetryAttempts < 1 { problems = problems.push("maxRetryAttempts must be at least 1") }65 if options.delay < 0 || options.delay > LONGEST_DELAY { problems = problems.push("delay must be from 0 to 86400000 ms") }66 if let Some(ceiling) = options.maxDelay {67 if ceiling < 0 || ceiling > LONGEST_DELAY { problems = problems.push("maxDelay must be from 0 to 86400000 ms") }68 }69 problems70}717273export fn strategy[T, E](options: Options[T, E]) -> Strategy.Strategy[T, E] {74 Strategy.Strategy {75 name: options.name,76 kind: "Retry",77 summary: describe(&options),78 problems: validate(&options),79 attach: Strategy.detached,80 execute: fn(runtime: Strategy.Runtime, context: Context.Context, callback: Strategy.Callback[T, E]) -> Resilience.Outcome[T, E] {81 execute(&options, &runtime, context, callback)82 }83 }84}858687fn describe[T, E](options: &Options[T, E]) -> Str {88 let count = if options.maxRetryAttempts == UNLIMITED { "unlimited" } else { show(options.maxRetryAttempts) }89 let shape = match options.backoff {90 case Constant => "constant"91 case Linear => "linear"92 case Exponential => "exponential"93 }94 let jitter = if options.useJitter { " with jitter" } else { "" }95 count + " retries, " + shape + " backoff from " + show(options.delay) + " ms" + jitter96}979899fn plan[T, E](options: &Options[T, E]) -> Delays.Plan {100 let shape = match options.backoff {101 case Constant => Delays.Constant102 case Linear => Delays.Linear103 case Exponential => Delays.Exponential104 }105 Delays.Plan{shape: shape, delay: options.delay, maxDelay: options.maxDelay, useJitter: options.useJitter}106}107108109110fn execute[T, E](options: &Options[T, E], runtime: &Strategy.Runtime, context: Context.Context, callback: Strategy.Callback[T, E]) -> Resilience.Outcome[T, E] {111 var attempt = 0112 var state = 0.0113 loop {114 let started = Clock.now(&runtime.clock)115 let outcome = callback(context)116 let judged = Predicate.arguments(outcome, context, attempt)117 let handled = (options.shouldHandle)(judged)118 let duration = Clock.now(&runtime.clock) - started119 let last = attempt != UNLIMITED && attempt >= options.maxRetryAttempts120 let severity = if !handled { Telemetry.Information } else if last { Telemetry.Error } else { Telemetry.Warning }121 Telemetry.report(&runtime.telemetry, &context, Events.EXECUTION_ATTEMPT, severity, Some(Resilience.summarize(&outcome)), Telemetry.Attempted(attempt, duration, handled))122 if last || !handled { return outcome }123 let computed = Delays.next(&plan(options), attempt, state, Randomizer.unit(&runtime.randomizer))124 state = computed[1]125 var delay = computed[0]126 if let Some(generate) = options.delayGenerator {127 if let Some(generated) = generate(judged) {128 if generated >= 0 { delay = generated }129 }130 }131 Telemetry.report(&runtime.telemetry, &context, Events.ON_RETRY, Telemetry.Warning, Some(Resilience.summarize(&outcome)), Telemetry.Retrying(attempt, delay))132 if let Some(notify) = options.onRetry {133 notify(Retrying{outcome: outcome, context: context, attempt: attempt, delay: delay, duration: duration})134 }135 if let Err(stopped) = Context.check(&context) { return Err(stopped) }136 if delay > 0 {137 if let Err(why) = Clock.sleep(&runtime.clock, delay, &context.token) {138 return Err(Resilience.Cancelled(Cancel.explain(&why)))139 }140 }141 if attempt != UNLIMITED { attempt = attempt + 1 }142 }143}144