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

PipelineScenarioTest.pudu

Pudu73 lines4.1 KB

GitHub ↗
1/** @Test.Integration.Scenario — every strategy stacked in one pipeline */2module Integration.PipelineScenarioTest34import Std.Io as Io5import Std.Sync as Sync6import Std.Test as Test7import PuduLangResilience.CircuitBreaker as CircuitBreaker8import PuduLangResilience.Clock as Clock9import PuduLangResilience.Context as Context10import PuduLangResilience.Fallback as Fallback11import PuduLangResilience.Pipeline as Pipeline12import PuduLangResilience.Predicate as Predicate13import PuduLangResilience.RateLimiter as RateLimiter14import PuduLangResilience as Resilience15import PuduLangResilience.Retry as Retry16import PuduLangResilience.Telemetry.Meter as Meter17import PuduLangResilience.Timeout as Timeout1819/// Runs the suite.20fn main() -> Int {21  let clock = Clock.manual(0)22  let meter = Meter.create()23  let provider = CircuitBreaker.stateProvider()24  let settings = Pipeline.Options{..Pipeline.defaults(), name: Some("inventory"), clock: Clock.ofManual(&clock), listeners: [Meter.listener(&meter)]}25  let built = Pipeline.buildWith(&settings, [26      Fallback.strategy(Fallback.Options{..Fallback.withValue(-1), shouldHandle: Predicate.brokenCircuits()}),27      RateLimiter.strategy(RateLimiter.concurrency(4, 0)),28      Retry.strategy(Retry.Options{..Retry.defaults(), maxRetryAttempts: 2, delay: 100, shouldHandle: Predicate.anyOf(&[Predicate.raised(|problem: Str| problem == "transient"), Predicate.timeouts()])}),29      CircuitBreaker.strategy(CircuitBreaker.Options{..CircuitBreaker.defaults(), failureRatio: 0.5d, minimumThroughput: 4, samplingDuration: 10000, breakDuration: 1000, stateProvider: Some(provider)}),30      Timeout.strategy(Timeout.after(200))31    ])32  let pipeline = match built {33    case Ok(found) => found34    case Err(invalid) => panic(Pipeline.explain(&invalid))35  }36  let calls = Sync.counter(0)37  let flaky = fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] {38    let made = match Sync.increment(&calls, 1) { case Ok(n) => n case Err(_) => 0 }39    if made % 2 == 1 { Resilience.raise("transient") } else { Ok(made) }40  }41  let recovered = Pipeline.execute(&pipeline, flaky)42  let down = fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] { Resilience.raise("transient") }43  let exhausted = Pipeline.execute(&pipeline, down)44  let stateAfterFailures = CircuitBreaker.state(&provider)45  let fellBack = Pipeline.execute(&pipeline, flaky)46  let fatal = Pipeline.execute(&pipeline, fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] { Resilience.raise("fatal") })47  let timedOut = Pipeline.execute(&match Pipeline.build([Retry.strategy(Retry.Options{..Retry.defaults(), maxRetryAttempts: 1, delay: 0, shouldHandle: Predicate.timeouts()}), Timeout.strategy(Timeout.after(20))]) {48      case Ok(found) => found49      case Err(invalid) => panic(Pipeline.explain(&invalid))50    }, fn(context: Context.Context) -> Resilience.Outcome[Int, Str] {51      Context.pause(&context, 1000) ?52      Ok(0)53    })54  let checks = Test.suite("Integration.PipelineScenario", &[55      Test.equals("a transient failure is retried to success", &recovered, &Ok(2)),56      Test.equals("a circuit opening during retries hands its rejection to the fallback", &exhausted, &Ok(-1)),57      Test.equals("the failures open the circuit", &stateAfterFailures, &Some(CircuitBreaker.Open)),58      Test.equals("an open circuit falls back to the substitute", &fellBack, &Ok(-1)),59      Test.equals("an open circuit keeps rejecting whatever the callback would do", &fatal, &Ok(-1)),60      Test.equals("each timeout is retried and the last one answered", &timedOut, &Err(Resilience.TimedOut(20))),61      Test.equals("the retries waited on the pipeline clock", &Clock.sleeps(&clock), &[100, 100, 100]),62      Test.equals("every execution is metered", &Meter.count(&meter, "PipelineExecuted"), &4),63      Test.equals("the circuit opened once", &Meter.count(&meter, "OnCircuitOpened"), &1),64      Test.equals("every fallback is metered", &Meter.count(&meter, "OnFallback"), &3),65      Test.equals("a rejection by the circuit ends the retries after one attempt", &Meter.count(&meter, "ExecutionAttempt"), &7)66    ])67  let ran = Test.run(&checks)68  for failure in Test.failuresOf(&ran) {69    let _reported = Io.writeErrorLine(failure)70  }71  Test.report(&ran)72}73