
Resilience.pudu
Pudu191 lines8.9 KB
1/** @Handlers.Resilience.Module — requests sent through resilience pipelines tuned for HTTP */2module PuduLangHttpClient.Handlers.Resilience34import Std.Concurrent.Cancel as Cancel5import Std.Time as Time6import PuduLangHttpClient.Handler as Handler7import PuduLangHttpClient as HttpClient8import PuduLangHttpClient.Request as Request9import PuduLangHttpClient.Response as Response10import PuduLangResilience.CircuitBreaker as CircuitBreaker11import PuduLangResilience.Context as Guard12import PuduLangResilience as Guarded13import PuduLangResilience.Hedging as Hedging14import PuduLangResilience.Pipeline as Pipeline15import PuduLangResilience.Predicate as Predicate16import PuduLangResilience.RateLimiter as RateLimiter17import PuduLangResilience.Registry as Registry18import PuduLangResilience.Retry as Retry19import PuduLangResilience.Strategy as Strategy20import PuduLangResilience.Timeout as Timeout2122/** @Handlers.Resilience.Guarding — a resilience pipeline over HTTP responses and failures */23export type Guarding = Pipeline.Pipeline[Response.Response, HttpClient.Failure]2425/** @Handlers.Resilience.Standard — the five layers of the standard pipeline, outermost first */26export type Standard = {27 rateLimiter: RateLimiter.Options,28 totalRequestTimeout: Timeout.Options,29 retry: Retry.Options[Response.Response, HttpClient.Failure],30 circuitBreaker: CircuitBreaker.Options[Response.Response, HttpClient.Failure],31 attemptTimeout: Timeout.Options32}3334/** @Handlers.Resilience.StandardHedging — the four layers of the standard hedging pipeline */35export type StandardHedging = {36 totalRequestTimeout: Timeout.Options,37 hedging: Hedging.Options[Response.Response, HttpClient.Failure],38 circuitBreaker: CircuitBreaker.Options[Response.Response, HttpClient.Failure],39 attemptTimeout: Timeout.Options40}414243const TRANSIENT_STATUSES: Set[Int] = #{408, 429}44454647export fn handler(pipeline: Guarding) -> Handler.Handler { selecting(|_request: Request.Request| pipeline) }48495051export fn selecting(choose: fn(Request.Request) -> Guarding) -> Handler.Handler {52 fn(request: Request.Request, token: Cancel.Token, next: Handler.Send) -> HttpClient.Outcome[Response.Response] {53 guarded(&choose(request), request, &token, next)54 }55}56575859export fn fromRegistry(registry: &Registry.Registry[Response.Response, HttpClient.Failure], key: Str) -> Handler.Handler {60 let held = *registry61 fn(request: Request.Request, token: Cancel.Token, next: Handler.Send) -> HttpClient.Outcome[Response.Response] {62 match Registry.get(&held, key) {63 case Ok(pipeline) => guarded(&pipeline, request, &token, next)64 case Err(problem) => Err(HttpClient.Rejected(Registry.explain(&problem)))65 }66 }67}686970717273export fn standardOptions() -> Standard {74 Standard {75 rateLimiter: RateLimiter.concurrency(1000, 0),76 totalRequestTimeout: Timeout.after(30000),77 retry: Retry.Options{..Retry.defaults(), shouldHandle: transient(), backoff: Retry.Exponential, useJitter: true, delayGenerator: Some(retryAfter)},78 circuitBreaker: CircuitBreaker.Options{..CircuitBreaker.defaults(), shouldHandle: transient()},79 attemptTimeout: Timeout.after(10000)80 }81}828384export fn standard() -> Result[Guarding, Pipeline.Invalid] { standardWith(&standardOptions()) }858687export fn standardWith(options: &Standard) -> Result[Guarding, Pipeline.Invalid] {88 Pipeline.build([89 RateLimiter.strategy(options.rateLimiter),90 Timeout.strategy(options.totalRequestTimeout),91 Retry.strategy(options.retry),92 CircuitBreaker.strategy(options.circuitBreaker),93 Timeout.strategy(options.attemptTimeout)94 ])95}96979899export fn standardHedgingOptions() -> StandardHedging {100 StandardHedging {101 totalRequestTimeout: Timeout.after(30000),102 hedging: Hedging.Options{..Hedging.defaults(), shouldHandle: transient()},103 circuitBreaker: CircuitBreaker.Options{..CircuitBreaker.defaults(), shouldHandle: transient()},104 attemptTimeout: Timeout.after(10000)105 }106}107108109export fn standardHedging() -> Result[Guarding, Pipeline.Invalid] { standardHedgingWith(&standardHedgingOptions()) }110111112export fn standardHedgingWith(options: &StandardHedging) -> Result[Guarding, Pipeline.Invalid] {113 Pipeline.build([114 Timeout.strategy(options.totalRequestTimeout),115 Hedging.strategy(options.hedging),116 CircuitBreaker.strategy(options.circuitBreaker),117 Timeout.strategy(options.attemptTimeout)118 ])119}120121122123export fn onTransient(build: fn(Predicate.Predicate[Response.Response, HttpClient.Failure]) -> Array[Strategy.Strategy[Response.Response, HttpClient.Failure]]) -> Result[Guarding, Pipeline.Invalid] {124 Pipeline.build(build(transient()))125}126127128export fn transient() -> Predicate.Predicate[Response.Response, HttpClient.Failure] {129 Predicate.anyOf(&[130 Predicate.results(|response: Response.Response| isTransientStatus(response.status.code)),131 Predicate.raised(|failure: HttpClient.Failure| HttpClient.isTransient(&failure)),132 Predicate.timeouts()133 ])134}135136137export fn isTransientStatus(code: Int) -> Bool { code >= 500 || code in TRANSIENT_STATUSES }138139140export fn retryAfter(given: Predicate.Arguments[Response.Response, HttpClient.Failure]) -> Option[Int] {141 match given.outcome {142 case Ok(response) => Response.retryAfter(&response, Time.toMillis(&Time.currentInstant()))143 case Err(_) => None144 }145}146147148149export fn translated(failure: &Guarded.Failure[HttpClient.Failure]) -> HttpClient.Failure {150 match failure {151 case Guarded.Raised(inner) => inner152 case Guarded.TimedOut(millis) => HttpClient.TimedOut(millis)153 case Guarded.Cancelled(reason) => HttpClient.Cancelled(reason)154 case Guarded.Crashed(reason) => HttpClient.Crashed(reason)155 case rejection => HttpClient.Rejected(Guarded.describe(&rejection))156 }157}158159160fn reasonOf(why: &Cancel.Reason) -> Str {161 match why {162 case Cancel.Requested(reason) => reason163 case Cancel.DeadlineExceeded => Cancel.explain(why)164 }165}166167168169170fn guarded(pipeline: &Guarding, request: Request.Request, token: &Cancel.Token, next: Handler.Send) -> HttpClient.Outcome[Response.Response] {171 let began = clock()172 let outcome = Pipeline.executeWith(pipeline, Guard.cancellable(*token), fn(attempt: Guard.Context) -> Guarded.Outcome[Response.Response, HttpClient.Failure] {173 match next(request, attempt.token) {174 case Ok(response) => Ok(response)175 case Err(failure) => {176 match Cancel.reason(&attempt.token) {177 case Some(why) => if HttpClient.isTimeout(&failure) || HttpClient.isCancellation(&failure) { Err(Guarded.Cancelled(reasonOf(&why))) } else { Err(Guarded.Raised(failure)) }178 case None => Err(Guarded.Raised(failure))179 }180 }181 }182 })183 match outcome {184 case Ok(response) => Ok(response)185 case Err(Guarded.Cancelled(reason)) => {186 if Cancel.reason(token) == Some(Cancel.DeadlineExceeded) { Err(HttpClient.TimedOut(clock() - began)) } else { Err(HttpClient.Cancelled(reason)) }187 }188 case Err(failure) => Err(translated(&failure))189 }190}191