
HedgingTest.pudu
Pudu137 lines7.3 KB
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 Strategy151617fn 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}23242526fn 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}414243fn 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