Pudu programming language
Menu
Package

@chrismichaelps / pudu-lang-httpclient

Named HTTP clients for Pudu: a client factory, delegating handlers, pooled keep-alive connections, handler lifetimes, logging, and resilience

0.1.1Apache-2.01

InstallClose

Resilience.pudu

Pudu191 lines8.9 KB

GitHub ↗
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}4142/// The statuses worth another attempt besides every 5xx: request timeout and too many requests.43const TRANSIENT_STATUSES: Set[Int] = #{408, 429}4445/// A handler running the rest of the pipeline under a resilience pipeline: each attempt observes46/// the strategies' token, and a strategy's rejection is answered as `Rejected`.47export fn handler(pipeline: Guarding) -> Handler.Handler { selecting(|_request: Request.Request| pipeline) }4849/// A handler choosing the resilience pipeline for each request, such as one for reads and another50/// for writes.51export 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}5657/// A handler running each request under the registry's pipeline for `key`, built on first use and58/// rebuilt when the registry reloads it; a key the registry cannot build is `Rejected`.59export 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}6869/// The standard layers: a concurrency limiter of 1000 permits, a 30-second total timeout, three70/// exponential retries from two seconds with jitter honouring `retry-after`, a circuit breaker71/// opening at a tenth of at least 100 requests failing in 30 seconds for 5 seconds, and a 10-second72/// timeout per attempt; the retry and the breaker handle transient outcomes.73export 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}8283/// The standard pipeline.84export fn standard() -> Result[Guarding, Pipeline.Invalid] { standardWith(&standardOptions()) }8586/// The standard pipeline with its layers as given.87export 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}9697/// The standard hedging layers: a 30-second total timeout, one hedged attempt after two seconds for98/// transient outcomes, a circuit breaker for transient outcomes, and a 10-second timeout per attempt.99export 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}107108/// The standard hedging pipeline.109export fn standardHedging() -> Result[Guarding, Pipeline.Invalid] { standardHedgingWith(&standardHedgingOptions()) }110111/// The standard hedging pipeline with its layers as given.112export 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}120121/// A pipeline of the strategies `build` makes when given the transient predicate, so each strategy122/// handles exactly the outcomes worth another attempt.123export 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}126127/// Outcomes worth another attempt: a 5xx, 408, or 429 response, a transient failure, and a timeout.128export 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}135136/// Whether a status is worth another attempt.137export fn isTransientStatus(code: Int) -> Bool { code >= 500 || code in TRANSIENT_STATUSES }138139/// The wait a handled response's `retry-after` header asks for, or `None` to keep the backoff.140export 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}146147/// A resilience failure as a request failure: the request's own failure unchanged, a timeout or148/// cancellation as one, and a strategy's rejection as `Rejected`.149export 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}158159/// The words a cancellation carries: the reason given, or that the deadline passed.160fn reasonOf(why: &Cancel.Reason) -> Str {161  match why {162    case Cancel.Requested(reason) => reason163    case Cancel.DeadlineExceeded => Cancel.explain(why)164  }165}166167/// The response the rest of the pipeline answers under the resilience pipeline. An attempt stopped168/// by its own token answers a cancellation so the strategy that fired it can tell; a deadline of the169/// caller's that passed stays a timeout.170fn 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