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

HedgingTest.pudu

Pudu137 lines7.3 KB

GitHub ↗
1/** @Test.Hedging.Suite — racing attempts, choosing a winner, and ending the rest */2module PuduLangResilience.HedgingTest34import Std.Io as Io5import Std.Result as Result6import Std.Sync as Sync7import Std.Test as Test8import Std.Time as Time9import PuduLangResilience.Context as Context10import PuduLangResilience.Hedging as Hedging11import PuduLangResilience.Pipeline as Pipeline12import PuduLangResilience.Predicate as Predicate13import PuduLangResilience as Resilience14import PuduLangResilience.Strategy as Strategy1516/// A pipeline of one hedging strategy.17fn hedged(options: Hedging.Options[Int, Str]) -> Pipeline.Pipeline[Int, Str] {18  match Pipeline.build([Hedging.strategy(options)]) {19    case Ok(pipeline) => pipeline20    case Err(invalid) => panic(Pipeline.explain(&invalid))21  }22}2324/// A callback whose first call waits for cancellation and whose later calls answer at once,25/// counting calls and cancellations.26fn slowThenFast(calls: &Sync.Counter, cancelled: &Sync.Counter) -> Strategy.Callback[Int, Str] {27  let counter = *calls28  let stops = *cancelled29  fn(context: Context.Context) -> Resilience.Outcome[Int, Str] {30    let made = match Sync.increment(&counter, 1) { case Ok(n) => n case Err(_) => 0 }31    if made == 1 {32      if let Err(why) = Context.pause(&context, 5000) {33        let _stopped = Sync.increment(&stops, 1)34        return Err(why)35      }36      return Ok(1)37    }38    Ok(made)39  }40}4142/// Runs the suite.43fn main() -> Int {44  let fastCalls = Sync.counter(0)45  let fast = Pipeline.execute(&hedged(Hedging.defaults()), fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] {46      let _counted = Sync.increment(&fastCalls, 1)47      Ok(7)48    })4950  let calls = Sync.counter(0)51  let cancelled = Sync.counter(0)52  let hedges: Sync.Cell[Array[Int]] = Sync.cell([])53  let started = Time.elapsed()54  let raced = Pipeline.execute(&hedged(Hedging.Options{..Hedging.defaults(), delay: 30, onHedging: Some(fn(given: Hedging.Hedged) -> () {55            let held = match Sync.get(&hedges) { case Ok(list) => list case Err(_) => [] }56            let _kept = Sync.set(&hedges, held.push(given.attempt))57          }) }), slowThenFast(&calls, &cancelled))58  let raceTook = Time.elapsed() - started5960  let fallbackCalls = Sync.counter(0)61  let fallbackMode = Pipeline.execute(&hedged(Hedging.Options{..Hedging.defaults(), delay: Hedging.WAIT_FOR_OUTCOME}), fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] {62      let made = match Sync.increment(&fallbackCalls, 1) { case Ok(n) => n case Err(_) => 0 }63      if made == 1 { Resilience.raise("primary failed") } else { Ok(made) }64    })6566  let failingCalls = Sync.counter(0)67  let allFail = Pipeline.execute(&hedged(Hedging.Options{..Hedging.defaults(), maxHedgedAttempts: 2, delay: 0}), fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] {68      let _counted = Sync.increment(&failingCalls, 1)69      Resilience.raise("down")70    })7172  let declinedCalls = Sync.counter(0)73  let declining = Hedging.Options{..Hedging.defaults(), delay: 0, actionGenerator: Some(fn(_action: Hedging.Action[Int, Str]) -> Option[Strategy.Callback[Int, Str]] { None })}74  let declined = Pipeline.execute(&hedged(declining), fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] {75      let _counted = Sync.increment(&declinedCalls, 1)76      Resilience.raise("only attempt")77    })7879  let generating = Hedging.Options{..Hedging.defaults(), delay: Hedging.WAIT_FOR_OUTCOME, actionGenerator: Some(fn(action: Hedging.Action[Int, Str]) -> Option[Strategy.Callback[Int, Str]] {80        let number = action.attempt81        Some(fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] { Ok(100 + number) })82      }) }83  let generated = Pipeline.execute(&hedged(generating), fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] { Resilience.raise("primary") })8485  let marker = Context.textKey("served-by")86  let caller = Context.create()87  let adoptCalls = Sync.counter(0)88  let adopted = Pipeline.executeWith(&hedged(Hedging.Options{..Hedging.defaults(), delay: Hedging.WAIT_FOR_OUTCOME}), caller, fn(context: Context.Context) -> Resilience.Outcome[Int, Str] {89      let made = match Sync.increment(&adoptCalls, 1) { case Ok(n) => n case Err(_) => 0 }90      if made == 1 {91        Context.set(&context, &marker, "primary")92        Resilience.raise("first")93      } else {94        Context.set(&context, &marker, "hedge")95        Ok(2)96      }97    })9899  let crashed = Pipeline.execute(&hedged(Hedging.Options{..Hedging.defaults(), maxHedgedAttempts: 1, delay: Hedging.WAIT_FOR_OUTCOME, shouldHandle: Predicate.never()}), fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] { panic("boom") })100101  let delays: Sync.Cell[Array[Int]] = Sync.cell([])102  let generatedDelay = Hedging.Options{..Hedging.defaults(), maxHedgedAttempts: 3, delayGenerator: Some(fn(given: Hedging.DelayArguments) -> Int {103        let held = match Sync.get(&delays) { case Ok(list) => list case Err(_) => [] }104        let _kept = Sync.set(&delays, held.push(given.attempt))105        if given.attempt == 2 { Hedging.WAIT_FOR_OUTCOME } else { 0 }106      }) }107  let _generatedRun = Pipeline.execute(&hedged(generatedDelay), fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] { Resilience.raise("down") })108109  let invalid = Hedging.validate(&Hedging.Options{..Hedging.defaults(), maxHedgedAttempts: 11, delay: -2})110  let checks = Test.suite("Hedging", &[111      Test.equals("a fast primary answers alone", &fast, &Ok(7)),112      Test.equals("a fast primary starts no hedge", &Sync.count(&fastCalls), &Ok(1)),113      Test.equals("a slow primary loses to its hedge", &raced, &Ok(2)),114      Test.equals("onHedging sees the hedged attempt", &Sync.get(&hedges), &Ok([1])),115      Test.equals("the losing attempt is cancelled", &Sync.count(&cancelled), &Ok(1)),116      Test.that("the race ends without waiting out the slow primary", raceTook < 4000),117      Test.equals("waiting for the outcome hedges after a failure", &fallbackMode, &Ok(2)),118      Test.equals("when every attempt fails a failure is answered", &allFail, &Resilience.raise("down")),119      Test.equals("every allowed attempt runs", &Sync.count(&failingCalls), &Ok(3)),120      Test.equals("a declined hedge leaves the primary's outcome", &declined, &Resilience.raise("only attempt")),121      Test.equals("a declined hedge runs nothing more", &Sync.count(&declinedCalls), &Ok(1)),122      Test.equals("a generated action runs as the hedge", &generated, &Ok(101)),123      Test.equals("the winner answers", &adopted, &Ok(2)),124      Test.equals("only the winner's properties reach the caller", &Context.get(&caller, &marker), &Some("hedge")),125      Test.that("a crashed attempt is answered as crashed", match crashed { case Err(Resilience.Crashed(_)) => true case _ => false }),126      Test.equals("the delay generator is asked once for each hedge", &Sync.get(&delays), &Ok([1, 2, 3])),127      Test.equals("invalid options are refused", &invalid, &["maxHedgedAttempts must be from 1 to 10", "delay must be from -1 to 86400000 ms"]),128      Test.that("the bounds themselves are accepted", Hedging.validate(&Hedging.Options{..Hedging.defaults(), maxHedgedAttempts: 10, delay: Hedging.WAIT_FOR_OUTCOME}).isEmpty()),129      Test.that("building refuses invalid options", Result.isErr(&Pipeline.build([Hedging.strategy(Hedging.Options{..Hedging.defaults(), maxHedgedAttempts: 0})])))130    ])131  let ran = Test.run(&checks)132  for failure in Test.failuresOf(&ran) {133    let _reported = Io.writeErrorLine(failure)134  }135  Test.report(&ran)136}137