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

RetryTest.pudu

Pudu139 lines8.0 KB

GitHub ↗
1/** @Test.Retry.Suite — retries, delays, predicates, cancellation, and telemetry */2module PuduLangResilience.RetryTest34import Std.Io as Io5import Std.Sync as Sync6import Std.Test as Test7import PuduLangResilience.Clock as Clock8import PuduLangResilience.Context as Context9import PuduLangResilience.Pipeline as Pipeline10import PuduLangResilience.Predicate as Predicate11import PuduLangResilience.Randomizer as Randomizer12import PuduLangResilience as Resilience13import PuduLangResilience.Retry as Retry14import PuduLangResilience.Telemetry as Telemetry1516/// Options on a manual clock with a fixed randomizer, recording every event.17fn optionsOn(clock: &Clock.Manual, events: &Sync.Cell[Array[Telemetry.Event]]) -> Pipeline.Options {18  let sink = *events19  Pipeline.Options {20    ..Pipeline.defaults(),21    name: Some("orders"),22    clock: Clock.ofManual(clock),23    randomizer: Randomizer.fixed(0.5d),24    listeners: [fn(event: Telemetry.Event) -> () {25        let held = match Sync.get(&sink) { case Ok(list) => list case Err(_) => [] }26        let _kept = Sync.set(&sink, held.push(event))27      }]28  }29}3031/// A pipeline of one retry strategy on a manual clock.32fn retrying(options: Retry.Options[Int, Str], clock: &Clock.Manual, events: &Sync.Cell[Array[Telemetry.Event]]) -> Pipeline.Pipeline[Int, Str] {33  match Pipeline.buildWith(&optionsOn(clock, events), [Retry.strategy(options)]) {34    case Ok(pipeline) => pipeline35    case Err(invalid) => panic(Pipeline.explain(&invalid))36  }37}3839/// A callback failing until the given call, counting its calls.40fn failingUntil(success: Int, calls: &Sync.Counter) -> fn(Context.Context) -> Resilience.Outcome[Int, Str] {41  let counter = *calls42  fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] {43    let made = match Sync.increment(&counter, 1) { case Ok(n) => n case Err(_) => 0 }44    if made < success { Resilience.raise("transient") } else { Ok(made) }45  }46}4748/// The names of the recorded events.49fn names(events: &Sync.Cell[Array[Telemetry.Event]]) -> Array[Str] {50  match Sync.get(events) {51    case Ok(list) => list.map(fn(event: Telemetry.Event) -> Str { event.name })52    case Err(_) => []53  }54}5556/// Runs the suite.57fn main() -> Int {58  let events: Sync.Cell[Array[Telemetry.Event]] = Sync.cell([])59  let clock = Clock.manual(0)60  let calls = Sync.counter(0)61  let recovered = Pipeline.execute(&retrying(Retry.defaults(), &clock, &events), failingUntil(3, &calls))6263  let exhaustedClock = Clock.manual(0)64  let exhaustedCalls = Sync.counter(0)65  let exhausted = Pipeline.execute(&retrying(Retry.Options{..Retry.defaults(), maxRetryAttempts: 2, delay: 10}, &exhaustedClock, &Sync.cell([])), failingUntil(99, &exhaustedCalls))6667  let linearClock = Clock.manual(0)68  let _linear = Pipeline.execute(&retrying(Retry.Options{..Retry.defaults(), backoff: Retry.Linear, delay: 100, maxRetryAttempts: 4}, &linearClock, &Sync.cell([])), failingUntil(99, &Sync.counter(0)))6970  let exponentialClock = Clock.manual(0)71  let _exponential = Pipeline.execute(&retrying(Retry.Options{..Retry.defaults(), backoff: Retry.Exponential, delay: 100, maxRetryAttempts: 5, maxDelay: Some(1000)}, &exponentialClock, &Sync.cell([])), failingUntil(99, &Sync.counter(0)))7273  let jitterClock = Clock.manual(0)74  let _jitter = Pipeline.execute(&retrying(Retry.Options{..Retry.defaults(), useJitter: true, delay: 1000, maxRetryAttempts: 2}, &jitterClock, &Sync.cell([])), failingUntil(99, &Sync.counter(0)))7576  let generatedClock = Clock.manual(0)77  let generated = Retry.Options{..Retry.defaults(), delayGenerator: Some(fn(given: Predicate.Arguments[Int, Str]) -> Option[Int] {78        if given.attempt == 0 { Some(7) } else if given.attempt == 1 { Some(-5) } else { None }79      }) }80  let _generatedRun = Pipeline.execute(&retrying(generated, &generatedClock, &Sync.cell([])), failingUntil(99, &Sync.counter(0)))8182  let unhandledCalls = Sync.counter(0)83  let onlyTimeouts = Retry.Options{..Retry.defaults(), shouldHandle: Predicate.timeouts()}84  let unhandled = Pipeline.execute(&retrying(onlyTimeouts, &Clock.manual(0), &Sync.cell([])), failingUntil(99, &unhandledCalls))8586  let resultCalls = Sync.counter(0)87  let onResult = Retry.Options{..Retry.defaults(), delay: 0, shouldHandle: Predicate.results(fn(value: Int) -> Bool { value < 3 })}88  let byResult = Pipeline.execute(&retrying(onResult, &Clock.manual(0), &Sync.cell([])), fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] {89      match Sync.increment(&resultCalls, 1) { case Ok(n) => Ok(n) case Err(_) => Ok(0) }90    })9192  let cancelledCalls = Sync.counter(0)93  let cancelling = Context.create()94  let cancelled = Pipeline.executeWith(&retrying(Retry.defaults(), &Clock.manual(0), &Sync.cell([])), cancelling, fn(context: Context.Context) -> Resilience.Outcome[Int, Str] {95      let _counted = Sync.increment(&cancelledCalls, 1)96      Context.cancel(&context, "shutdown")97      Resilience.raise("transient")98    })99100  let seen: Sync.Cell[Array[Int]] = Sync.cell([])101  let observed = Retry.Options{..Retry.defaults(), delay: 5, onRetry: Some(fn(retry: Retry.Retrying[Int, Str]) -> () {102        let held = match Sync.get(&seen) { case Ok(list) => list case Err(_) => [] }103        let _kept = Sync.set(&seen, held.push(retry.attempt * 100 + retry.delay))104      }) }105  let _observedRun = Pipeline.execute(&retrying(observed, &Clock.manual(0), &Sync.cell([])), failingUntil(99, &Sync.counter(0)))106107  let unlimitedCalls = Sync.counter(0)108  let unlimited = Pipeline.execute(&retrying(Retry.Options{..Retry.defaults(), maxRetryAttempts: Retry.UNLIMITED, delay: 0}, &Clock.manual(0), &Sync.cell([])), failingUntil(50, &unlimitedCalls))109110  let invalid = Pipeline.build([Retry.strategy(Retry.Options{..Retry.defaults(), maxRetryAttempts: 0, delay: -1, maxDelay: Some(86400001)})])111  let problems = match invalid { case Err(found) => found.problems case Ok(_) => [] }112113  let checks = Test.suite("Retry", &[114      Test.equals("a transient failure is retried until it succeeds", &recovered, &Ok(3)),115      Test.equals("the constant delay is waited before each retry", &Clock.sleeps(&clock), &[2000, 2000]),116      Test.equals("the last failure is answered when retries run out", &exhausted, &Resilience.raise("transient")),117      Test.equals("the callback runs once more than the retry count", &Sync.count(&exhaustedCalls), &Ok(3)),118      Test.equals("linear delays grow by the base delay", &Clock.sleeps(&linearClock), &[100, 200, 300, 400]),119      Test.equals("exponential delays double up to the ceiling", &Clock.sleeps(&exponentialClock), &[100, 200, 400, 800, 1000]),120      Test.equals("jitter at the midpoint keeps the delay", &Clock.sleeps(&jitterClock), &[1000, 1000]),121      Test.equals("a generated delay replaces the computed one; a negative or absent one does not", &Clock.sleeps(&generatedClock), &[7, 2000, 2000]),122      Test.equals("an unhandled failure is answered at once", &unhandled, &Resilience.raise("transient")),123      Test.equals("an unhandled failure is not retried", &Sync.count(&unhandledCalls), &Ok(1)),124      Test.equals("a handled value is retried until one is not handled", &byResult, &Ok(3)),125      Test.equals("a cancelled token stops retrying", &cancelled, &Err(Resilience.Cancelled("cancelled: shutdown"))),126      Test.equals("cancellation ends before the next attempt", &Sync.count(&cancelledCalls), &Ok(1)),127      Test.equals("onRetry sees each attempt and its delay", &Sync.get(&seen), &Ok([5, 105, 205])),128      Test.equals("unlimited retries continue past any small count", &unlimited, &Ok(50)),129      Test.equals("every invalid option is reported", &problems, &["Retry: maxRetryAttempts must be at least 1", "Retry: delay must be from 0 to 86400000 ms", "Retry: maxDelay must be from 0 to 86400000 ms"]),130      Test.that("the defaults are valid", Retry.validate(&Retry.defaults()).isEmpty()),131      Test.equals("attempts, retries, and the pipeline are reported in order", &names(&events), &["PipelineExecuting", "ExecutionAttempt", "OnRetry", "ExecutionAttempt", "OnRetry", "ExecutionAttempt", "PipelineExecuted"])132    ])133  let ran = Test.run(&checks)134  for failure in Test.failuresOf(&ran) {135    let _reported = Io.writeErrorLine(failure)136  }137  Test.report(&ran)138}139