
RetryTest.pudu
Pudu139 lines8.0 KB
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 Telemetry151617fn 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}303132fn 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}383940fn 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}474849fn 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}555657fn 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