
RateLimiter.pudu
Pudu89 lines3.8 KB
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}232425export fn defaults() -> Options { Options{name: None, acquire: None, defaultLimiter: Concurrency.defaults(), onRejected: None} }26272829export fn concurrency(permitLimit: Int, queueLimit: Int) -> Options {30 Options{..defaults(), defaultLimiter: Concurrency.Options{..Concurrency.defaults(), permitLimit: permitLimit, queueLimit: queueLimit}}31}323334export fn using(limiter: Limiter.Limiter) -> Options {35 Options{..defaults(), acquire: Some(fn(context: Context.Context) -> Limiter.Lease { Limiter.acquire(&limiter, 1, &context.token) })}36}373839export fn partitioned(limiters: Partitioned.Partitioned) -> Options {40 Options{..defaults(), acquire: Some(fn(context: Context.Context) -> Limiter.Lease { Partitioned.acquire(&limiters, &context, 1) })}41}42434445export 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