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

CircuitBreakerTest.pudu

Pudu156 lines9.5 KB

GitHub ↗
1/** @Test.CircuitBreaker.Suite — opening, probing, closing, and manual control */2module PuduLangResilience.CircuitBreakerTest34import Std.Io as Io5import Std.Result as Result6import Std.Sync as Sync7import Std.Test as Test8import PuduLangResilience.CircuitBreaker as CircuitBreaker9import PuduLangResilience.Clock as Clock10import PuduLangResilience.Context as Context11import PuduLangResilience.Pipeline as Pipeline12import PuduLangResilience as Resilience13import PuduLangResilience.Telemetry as Telemetry1415/// Opens at half of at least two executions in one second, for half a second.16fn sensitive() -> CircuitBreaker.Options[Int, Str] {17  CircuitBreaker.Options{..CircuitBreaker.defaults(), failureRatio: 0.5d, minimumThroughput: 2, samplingDuration: 1000, breakDuration: 500}18}1920/// A pipeline of one circuit breaker on a manual clock, recording event names.21fn breaking(options: CircuitBreaker.Options[Int, Str], clock: &Clock.Manual, events: &Sync.Cell[Array[Str]]) -> Pipeline.Pipeline[Int, Str] {22  let sink = *events23  let settings = Pipeline.Options{..Pipeline.defaults(), clock: Clock.ofManual(clock), listeners: [fn(event: Telemetry.Event) -> () {24        if event.name != "PipelineExecuting" && event.name != "PipelineExecuted" {25          let held = match Sync.get(&sink) { case Ok(list) => list case Err(_) => [] }26          let _kept = Sync.set(&sink, held.push(event.name))27        }28      }] }29  match Pipeline.buildWith(&settings, [CircuitBreaker.strategy(options)]) {30    case Ok(pipeline) => pipeline31    case Err(invalid) => panic(Pipeline.explain(&invalid))32  }33}3435/// A callback answering `outcome` and counting its calls.36fn answering(outcome: Resilience.Outcome[Int, Str], calls: &Sync.Counter) -> fn(Context.Context) -> Resilience.Outcome[Int, Str] {37  let counter = *calls38  fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] {39    let _counted = Sync.increment(&counter, 1)40    outcome41  }42}4344/// Runs the suite.45fn main() -> Int {46  let events: Sync.Cell[Array[Str]] = Sync.cell([])47  let clock = Clock.manual(0)48  let calls = Sync.counter(0)49  let provider = CircuitBreaker.stateProvider()50  let beforeAttached = CircuitBreaker.state(&provider)51  let pipeline = breaking(CircuitBreaker.Options{..sensitive(), stateProvider: Some(provider)}, &clock, &events)52  let fail = answering(Resilience.raise("down"), &calls)53  let succeed = answering(Ok(1), &calls)54  let first = Pipeline.execute(&pipeline, fail)55  let stillClosed = CircuitBreaker.state(&provider)56  let second = Pipeline.execute(&pipeline, fail)57  let openState = CircuitBreaker.state(&provider)58  let lastFailure = CircuitBreaker.lastHandled(&provider)59  Clock.advance(&clock, 200)60  let rejected = Pipeline.execute(&pipeline, succeed)61  let callsWhileOpen = Sync.count(&calls)62  Clock.advance(&clock, 300)63  let probe = Pipeline.execute(&pipeline, succeed)64  let afterProbe = CircuitBreaker.state(&provider)65  let clearedFailure = CircuitBreaker.lastHandled(&provider)6667  let reopenClock = Clock.manual(0)68  let attempts: Sync.Cell[Array[Int]] = Sync.cell([])69  let generated = CircuitBreaker.Options{..sensitive(), breakDurationGenerator: Some(fn(given: CircuitBreaker.BreakArguments) -> Int {70        let held = match Sync.get(&attempts) { case Ok(list) => list case Err(_) => [] }71        let _kept = Sync.set(&attempts, held.push(given.halfOpenAttempts))72        1000 * (given.halfOpenAttempts + 1)73      }) }74  let reopening = breaking(generated, &reopenClock, &Sync.cell([]))75  let reopenCalls = Sync.counter(0)76  let _a = Pipeline.execute(&reopening, answering(Resilience.raise("down"), &reopenCalls))77  let _b = Pipeline.execute(&reopening, answering(Resilience.raise("down"), &reopenCalls))78  Clock.advance(&reopenClock, 1000)79  let failedProbe = Pipeline.execute(&reopening, answering(Resilience.raise("still down"), &reopenCalls))80  Clock.advance(&reopenClock, 1999)81  let beforeLongerBreak = Pipeline.execute(&reopening, answering(Ok(1), &reopenCalls))8283  let windowClock = Clock.manual(0)84  let windowed = breaking(sensitive(), &windowClock, &Sync.cell([]))85  let windowCalls = Sync.counter(0)86  let _old = Pipeline.execute(&windowed, answering(Resilience.raise("down"), &windowCalls))87  Clock.advance(&windowClock, 1000)88  let _fresh = Pipeline.execute(&windowed, answering(Resilience.raise("down"), &windowCalls))89  let afterExpiry = Pipeline.execute(&windowed, answering(Ok(3), &windowCalls))9091  let cancelClock = Clock.manual(0)92  let cancelling = breaking(sensitive(), &cancelClock, &Sync.cell([]))93  for _round in [1, 2, 3] {94    let _cancelled = Pipeline.execute(&cancelling, answering(Err(Resilience.Cancelled("stop")), &Sync.counter(0)))95  }96  let afterCancellations = Pipeline.execute(&cancelling, answering(Ok(4), &Sync.counter(0)))9798  let control = CircuitBreaker.manualControl()99  let manualEvents: Sync.Cell[Array[Str]] = Sync.cell([])100  let manualOpened: Sync.Cell[Array[Bool]] = Sync.cell([])101  let controlled = breaking(CircuitBreaker.Options{..sensitive(), manualControl: Some(control), onOpened: Some(fn(opening: CircuitBreaker.Opening[Int, Str]) -> () {102          let _kept = Sync.set(&manualOpened, [opening.manual, opening.outcome == None])103        }) }, &Clock.manual(0), &manualEvents)104  CircuitBreaker.isolate(&control)105  let isolated = Pipeline.execute(&controlled, answering(Ok(1), &Sync.counter(0)))106  let isolatedFlag = CircuitBreaker.isIsolated(&control)107  CircuitBreaker.close(&control)108  let reopened = Pipeline.execute(&controlled, answering(Ok(5), &Sync.counter(0)))109  CircuitBreaker.close(&control)110111  let preIsolated = CircuitBreaker.manualControl()112  CircuitBreaker.isolate(&preIsolated)113  let lateProvider = CircuitBreaker.stateProvider()114  let _late = breaking(CircuitBreaker.Options{..sensitive(), manualControl: Some(preIsolated), stateProvider: Some(lateProvider)}, &Clock.manual(0), &Sync.cell([]))115116  let shared = CircuitBreaker.stateProvider()117  let _owner = CircuitBreaker.strategy(CircuitBreaker.Options{..sensitive(), stateProvider: Some(shared)})118  let reused = Pipeline.build([CircuitBreaker.strategy(CircuitBreaker.Options{..sensitive(), stateProvider: Some(shared)})])119  let invalid = Pipeline.build([CircuitBreaker.strategy(CircuitBreaker.Options{..sensitive(), failureRatio: 1.5d, minimumThroughput: 1, samplingDuration: 499, breakDuration: 86400001})])120  let problems = match invalid { case Err(found) => found.problems case Ok(_) => [] }121122  let checks = Test.suite("CircuitBreaker", &[123      Test.equals("a provider reads nothing before it is attached", &beforeAttached, &None),124      Test.equals("a failure below the threshold passes through", &first, &Resilience.raise("down")),125      Test.equals("the circuit stays closed below the minimum throughput", &stillClosed, &Some(CircuitBreaker.Closed)),126      Test.equals("the failure that meets the threshold passes through", &second, &Resilience.raise("down")),127      Test.equals("the circuit opens at the threshold", &openState, &Some(CircuitBreaker.Open)),128      Test.equals("the provider reads the last handled outcome", &lastFailure, &Some("The operation failed: \"down\"")),129      Test.equals("an open circuit rejects with the time left", &rejected, &Err(Resilience.BrokenCircuit(300))),130      Test.equals("an open circuit does not run the callback", &callsWhileOpen, &Ok(2)),131      Test.equals("after the break a probe runs", &probe, &Ok(1)),132      Test.equals("a successful probe closes the circuit", &afterProbe, &Some(CircuitBreaker.Closed)),133      Test.equals("a closed circuit forgets its last failure", &clearedFailure, &None),134      Test.equals("transitions are reported", &Sync.get(&events), &Ok(["OnCircuitOpened", "OnCircuitHalfOpened", "OnCircuitClosed"])),135      Test.equals("a failed probe passes its failure through", &failedProbe, &Resilience.raise("still down")),136      Test.equals("the break generator sees the half-open attempts", &Sync.get(&attempts), &Ok([0, 1])),137      Test.equals("a failed probe reopens for the generated duration", &beforeLongerBreak, &Err(Resilience.BrokenCircuit(1))),138      Test.equals("failures older than the sampling period are not counted", &afterExpiry, &Ok(3)),139      Test.equals("cancellations are not handled and never open the circuit", &afterCancellations, &Ok(4)),140      Test.equals("an isolated circuit refuses every execution", &isolated, &Err(Resilience.IsolatedCircuit)),141      Test.that("the control reports isolation", isolatedFlag),142      Test.equals("onOpened sees a manual opening without an outcome", &Sync.get(&manualOpened), &Ok([true, true])),143      Test.equals("closing by hand admits executions again", &reopened, &Ok(5)),144      Test.equals("closing an already closed circuit reports nothing", &Sync.get(&manualEvents), &Ok(["OnCircuitOpened", "OnCircuitClosed"])),145      Test.equals("a circuit attached to an isolated control starts isolated", &CircuitBreaker.state(&lateProvider), &Some(CircuitBreaker.Isolated)),146      Test.that("a state provider serves one circuit only", Result.isErr(&reused)),147      Test.equals("every invalid option is reported", &problems, &["CircuitBreaker: failureRatio must be from 0 to 1", "CircuitBreaker: minimumThroughput must be at least 2", "CircuitBreaker: samplingDuration must be from 500 to 86400000 ms", "CircuitBreaker: breakDuration must be from 500 to 86400000 ms"]),148      Test.that("the bounds themselves are accepted", CircuitBreaker.validate(&CircuitBreaker.Options{..sensitive(), failureRatio: 0d, samplingDuration: 500, breakDuration: 86400000}).isEmpty() && CircuitBreaker.validate(&CircuitBreaker.Options{..sensitive(), failureRatio: 1d}).isEmpty())149    ])150  let ran = Test.run(&checks)151  for failure in Test.failuresOf(&ran) {152    let _reported = Io.writeErrorLine(failure)153  }154  Test.report(&ran)155}156