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

RateLimiter.pudu

Pudu89 lines3.8 KB

GitHub ↗
1/** @Resilience.RateLimiter.Strategy — admits executions only while permits last */2module PuduLangResilience.RateLimiter34import PuduLangResilience.Constants.Events as Events5import PuduLangResilience.Context as Context6import PuduLangResilience.Limiter.Concurrency as Concurrency7import PuduLangResilience.Limiter as Limiter8import PuduLangResilience.Limiter.Partitioned as Partitioned9import PuduLangResilience as Resilience10import PuduLangResilience.Strategy as Strategy11import PuduLangResilience.Telemetry as Telemetry1213/** @Resilience.RateLimiter.Rejected — the execution a limiter turned away */14export type Rejected = { context: Context.Context, retryAfter: Option[Int] }1516/** @Resilience.RateLimiter.Options — where permits come from and who hears of rejections */17export type Options = {18  name: Option[Str],19  acquire: Option[fn(Context.Context) -> Limiter.Lease],20  defaultLimiter: Concurrency.Options,21  onRejected: Option[fn(Rejected) -> ()]22}2324/// One permit per execution from a concurrency limiter of a thousand permits and no queue.25export fn defaults() -> Options { Options{name: None, acquire: None, defaultLimiter: Concurrency.defaults(), onRejected: None} }2627/// One permit per execution from a concurrency limiter of `permitLimit` permits and a queue of28/// `queueLimit`.29export fn concurrency(permitLimit: Int, queueLimit: Int) -> Options {30  Options{..defaults(), defaultLimiter: Concurrency.Options{..Concurrency.defaults(), permitLimit: permitLimit, queueLimit: queueLimit}}31}3233/// One permit per execution from `limiter`, waiting as it allows.34export fn using(limiter: Limiter.Limiter) -> Options {35  Options{..defaults(), acquire: Some(fn(context: Context.Context) -> Limiter.Lease { Limiter.acquire(&limiter, 1, &context.token) })}36}3738/// One permit per execution from the limiter of its partition.39export fn partitioned(limiters: Partitioned.Partitioned) -> Options {40  Options{..defaults(), acquire: Some(fn(context: Context.Context) -> Limiter.Lease { Partitioned.acquire(&limiters, &context, 1) })}41}4243/// A rate limiter strategy following the options. Every pipeline holding the same strategy44/// shares its default limiter.45export fn strategy[T, E](options: Options) -> Strategy.Strategy[T, E] {46  let made = match options.acquire {47    case Some(given) => Ok(given)48    case None => match Concurrency.create(&options.defaultLimiter) {49      case Ok(limiter) => Ok(fn(context: Context.Context) -> Limiter.Lease { Limiter.acquire(&limiter, 1, &context.token) })50      case Err(problems) => Err(problems)51    }52  }53  let problems = match made {54    case Ok(_) => []55    case Err(found) => found56  }57  let take = match made {58    case Ok(given) => given59    case Err(_) => fn(_context: Context.Context) -> Limiter.Lease { Limiter.refused(Limiter.Denied(None)) }60  }61  let summary = match options.acquire {62    case Some(_) => "custom limiter"63    case None => show(options.defaultLimiter.permitLimit) + " permits, queue of " + show(options.defaultLimiter.queueLimit)64  }65  Strategy.Strategy {66    name: options.name,67    kind: "RateLimiter",68    summary: summary,69    problems: problems,70    attach: Strategy.detached,71    execute: fn(runtime: Strategy.Runtime, context: Context.Context, callback: Strategy.Callback[T, E]) -> Resilience.Outcome[T, E] {72      let lease = take(context)73      match lease.grant {74        case Limiter.Granted => {75          let outcome = callback(context)76          Limiter.release(&lease)77          outcome78        }79        case Limiter.Denied(retryAfter) => {80          Telemetry.report(&runtime.telemetry, &context, Events.ON_RATE_LIMITER_REJECTED, Telemetry.Error, None, Telemetry.Rejected(retryAfter))81          if let Some(notify) = options.onRejected { notify(Rejected{context: context, retryAfter: retryAfter}) }82          Err(Resilience.RateLimited(retryAfter))83        }84        case Limiter.Abandoned(why) => Err(Resilience.Cancelled(why))85      }86    }87  }88}89