
ResilienceTest.pudu
Pudu157 lines10.2 KB
1/** @Test.Handlers.Resilience.Suite — requests sent through resilience pipelines tuned for HTTP */2module PuduLangHttpClient.Handlers.ResilienceTest34import Std.Concurrent.Cancel as Cancel5import Std.Concurrent as Concurrent6import Std.Http as Http7import Std.Io as Io8import Std.Result as Result9import Std.Test as Test10import PuduLangHttpClient.Handler as Handler11import PuduLangHttpClient.Handlers.Resilience as Resilience12import PuduLangHttpClient as HttpClient13import PuduLangHttpClient.Request as Request14import PuduLangHttpClient.Response as Response15import PuduLangHttpClient.Utils.Shared as Shared16import PuduLangResilience.CircuitBreaker as CircuitBreaker17import PuduLangResilience.Context as Guard18import PuduLangResilience as Guarded19import PuduLangResilience.Pipeline as Pipeline20import PuduLangResilience.Predicate as Predicate21import PuduLangResilience.RateLimiter as RateLimiter22import PuduLangResilience.Registry as Registry23import PuduLangResilience.Retry as Retry24import PuduLangResilience.Timeout as Timeout252627fn scripted(calls: &Shared.Shared[Int], statuses: Array[Int]) -> Handler.Send {28 let counter = *calls29 fn(request: Request.Request, _token: Cancel.Token) -> HttpClient.Outcome[Response.Response] {30 let index = Shared.change(&counter, |count: Int| (count + 1, count))31 let code = if index < statuses.length() { statuses[index] } else { 200 }32 if code == 0 { Err(HttpClient.Connection("reset")) } else if code < 0 { Err(HttpClient.InvalidRequest("bad")) } else { Ok(Response.answer(&request, code)) }33 }34}353637fn waiting(request: Request.Request, token: Cancel.Token) -> HttpClient.Outcome[Response.Response] {38 let began = clock()39 match Cancel.pause(&token, 3000) {40 case Ok(_) => Ok(Response.answer(&request, 200))41 case Err(Cancel.Requested(reason)) => Err(HttpClient.Cancelled(reason))42 case Err(Cancel.DeadlineExceeded) => Err(HttpClient.TimedOut(clock() - began))43 }44}454647fn built(outcome: Result[Resilience.Guarding, Pipeline.Invalid]) -> Resilience.Guarding {48 match outcome {49 case Ok(found) => found50 case Err(invalid) => panic(Pipeline.explain(&invalid))51 }52}535455fn sent(handler: Handler.Handler, primary: Handler.Send, request: Request.Request, token: Cancel.Token) -> Result[Int, HttpClient.Failure] {56 Result.map(Handler.around(handler, primary)(request, token), |response: Response.Response| response.status.code)57}585960fn rejected(outcome: &Result[Int, HttpClient.Failure]) -> Bool {61 match outcome {62 case Err(HttpClient.Rejected(_)) => true63 case _ => false64 }65}666768fn main() -> Int {69 let quick = Resilience.Standard{..Resilience.standardOptions(), retry: Retry.Options{..Resilience.standardOptions().retry, delay: 1, useJitter: false}}70 let pipeline = built(Resilience.standardWith(&quick))71 let calls = Shared.shared(0)72 let recovered = sent(Resilience.handler(pipeline), scripted(&calls, [503, 0, 429]), Request.get("http://a/"), Cancel.token())73 let recoveredCalls = Shared.current(&calls)74 let exhaustedCalls = Shared.shared(0)75 let exhausted = sent(Resilience.handler(pipeline), scripted(&exhaustedCalls, [500, 502, 503, 504, 200]), Request.get("http://a/"), Cancel.token())76 let fatalCalls = Shared.shared(0)77 let fatal = sent(Resilience.handler(pipeline), scripted(&fatalCalls, [-1]), Request.get("http://a/"), Cancel.token())78 let notFoundCalls = Shared.shared(0)79 let notFound = sent(Resilience.handler(pipeline), scripted(¬FoundCalls, [404]), Request.get("http://a/"), Cancel.token())8081 let breaker = built(Pipeline.build([CircuitBreaker.strategy(CircuitBreaker.Options{..CircuitBreaker.defaults(), shouldHandle: Resilience.transient(), failureRatio: 0.5d, minimumThroughput: 2, samplingDuration: 1000, breakDuration: 1000})]))82 let breakerCalls = Shared.shared(0)83 let breaking = Resilience.handler(breaker)84 let first = sent(breaking, scripted(&breakerCalls, [503, 503, 200]), Request.get("http://a/"), Cancel.token())85 let second = sent(breaking, scripted(&breakerCalls, [503, 503, 200]), Request.get("http://a/"), Cancel.token())86 let opened = sent(breaking, scripted(&breakerCalls, [503, 503, 200]), Request.get("http://a/"), Cancel.token())8788 let attemptTimeout = sent(Resilience.handler(built(Pipeline.build([Timeout.strategy(Timeout.after(40))]))), waiting, Request.get("http://a/"), Cancel.token())89 let callerDeadline = sent(Resilience.handler(built(Pipeline.build([Timeout.strategy(Timeout.after(5000))]))), waiting, Request.get("http://a/"), Cancel.expiring(30))90 let stopping = Cancel.token()91 let _stopped = Cancel.cancel(&stopping, "caller")92 let callerCancelled = sent(Resilience.handler(built(Pipeline.build([Timeout.strategy(Timeout.after(5000))]))), waiting, Request.get("http://a/"), stopping)9394 let limiter = Resilience.handler(built(Pipeline.build([RateLimiter.strategy(RateLimiter.concurrency(1, 0))])))95 let holder = Cancel.token()96 let held = match Concurrent.start(fn() -> () { let _held = sent(limiter, waiting, Request.get("http://a/"), holder) }) {97 case Ok(found) => found98 case Err(problem) => panic(show(problem))99 }100 let _settle = Concurrent.sleep(30)101 let limited = sent(limiter, waiting, Request.get("http://a/"), Cancel.token())102 let _released = Cancel.cancel(&holder, "done")103 let _joined = Concurrent.join(&held)104105 let readsOnly = Resilience.selecting(|request: Request.Request| if request.method == Http.Get { pipeline } else { Pipeline.empty() })106 let getCalls = Shared.shared(0)107 let postCalls = Shared.shared(0)108 let selectedGet = sent(readsOnly, scripted(&getCalls, [503]), Request.get("http://a/"), Cancel.token())109 let selectedPost = sent(readsOnly, scripted(&postCalls, [503]), Request.create(Http.Post, "http://a/"), Cancel.token())110111 let registry: Registry.Registry[Response.Response, HttpClient.Failure] = Registry.create()112 let _added = Registry.tryAddBuilder(®istry, "api", |_context: Registry.BuilderContext| [Retry.strategy(Retry.Options{..Retry.defaults(), shouldHandle: Resilience.transient(), delay: 1})])113 let registryCalls = Shared.shared(0)114 let fromRegistry = sent(Resilience.fromRegistry(®istry, "api"), scripted(®istryCalls, [500]), Request.get("http://a/"), Cancel.token())115 let missingKey = sent(Resilience.fromRegistry(®istry, "other"), scripted(®istryCalls, []), Request.get("http://a/"), Cancel.token())116117 let onTransientCalls = Shared.shared(0)118 let onTransient = sent(Resilience.handler(built(Resilience.onTransient(|handles: Predicate.Predicate[Response.Response, HttpClient.Failure]| [Retry.strategy(Retry.Options{..Retry.defaults(), shouldHandle: handles, delay: 1})]))), scripted(&onTransientCalls, [503, 429]), Request.get("http://a/"), Cancel.token())119120 let answer = Response.answer(&Request.get("http://a/"), 503).withHeader("retry-after", "2")121 let handled = Resilience.retryAfter(Predicate.arguments(Ok(answer), Guard.create(), 0))122 let unhandled = Resilience.retryAfter(Predicate.arguments(Err(Guarded.TimedOut(1)), Guard.create(), 0))123 let checks = Test.suite("Handlers.Resilience", &[124 Test.equals("transient statuses and failures are retried until success", &(recovered, recoveredCalls), &(Ok(200), 4)),125 Test.equals("when retries run out the last response is answered", &(exhausted, Shared.current(&exhaustedCalls)), &(Ok(504), 4)),126 Test.equals("a failure that is not transient is not retried", &(fatal, Shared.current(&fatalCalls)), &(Err(HttpClient.InvalidRequest("bad")), 1)),127 Test.equals("a status that is not transient is not retried", &(notFound, Shared.current(¬FoundCalls)), &(Ok(404), 1)),128 Test.equals("a breaker passes requests while closed", &(first, second), &(Ok(503), Ok(503))),129 Test.that("an open breaker rejects without calling", rejected(&opened) && Shared.current(&breakerCalls) == 2),130 Test.equals("an attempt past its timeout times out", &attemptTimeout, &Err(HttpClient.TimedOut(40))),131 Test.that("a caller's deadline stays a timeout", match callerDeadline {132 case Err(HttpClient.TimedOut(millis)) => millis >= 25 && millis < 2000133 case _ => false134 }),135 Test.equals("a caller's cancellation stays a cancellation", &callerCancelled, &Err(HttpClient.Cancelled("caller"))),136 Test.that("a full rate limiter rejects", rejected(&limited)),137 Test.equals("a pipeline is chosen per request", &(selectedGet, Shared.current(&getCalls), selectedPost, Shared.current(&postCalls)), &(Ok(200), 2, Ok(503), 1)),138 Test.equals("a registry's pipeline is used", &(fromRegistry, Shared.current(®istryCalls)), &(Ok(200), 2)),139 Test.that("a key the registry cannot build is rejected", rejected(&missingKey)),140 Test.equals("strategies built on transient outcomes retry them", &(onTransient, Shared.current(&onTransientCalls)), &(Ok(200), 3)),141 Test.equals("retry-after sets the delay of a handled response", &(handled, unhandled), &(Some(2000), None)),142 Test.equals("the transient statuses are 408, 429, and every 5xx", &[407, 408, 429, 430, 499, 500, 599].filter(|code: Int| Resilience.isTransientStatus(code)), &[408, 429, 500, 599]),143 Test.that("the standard pipelines build", Result.isOk(&Resilience.standard()) && Result.isOk(&Resilience.standardHedging())),144 Test.equals("the standard hedging pipeline has four layers", &Pipeline.describe(&built(Resilience.standardHedgingWith(&Resilience.standardHedgingOptions()))).strategies.length(), &4),145 Test.equals("resilience failures become request failures", &[Guarded.Raised(HttpClient.ResponseEnded), Guarded.TimedOut(5), Guarded.Cancelled("c"), Guarded.Crashed("x")].map(|failure: Guarded.Failure[HttpClient.Failure]| Resilience.translated(&failure)), &[HttpClient.ResponseEnded, HttpClient.TimedOut(5), HttpClient.Cancelled("c"), HttpClient.Crashed("x")]),146 Test.that("a strategy's rejection becomes a rejection", match Resilience.translated(&Guarded.IsolatedCircuit) {147 case HttpClient.Rejected(_) => true148 case _ => false149 })150 ])151 let ran = Test.run(&checks)152 for failure in Test.failuresOf(&ran) {153 let _reported = Io.writeErrorLine(failure)154 }155 Test.report(&ran)156}157