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

RateLimiterTest.pudu

Pudu99 lines5.2 KB

GitHub ↗
1/** @Test.RateLimiter.Suite — admission, rejection, release, and cancellation */2module PuduLangResilience.RateLimiterTest34import Std.Io as Io5import Std.Result as Result6import Std.Sync as Sync7import Std.Test as Test8import PuduLangResilience.Clock as Clock9import PuduLangResilience.Context as Context10import PuduLangResilience.Limiter.Concurrency as Concurrency11import PuduLangResilience.Limiter.FixedWindow as FixedWindow12import PuduLangResilience.Limiter as Limiter13import PuduLangResilience.Limiter.Partitioned as Partitioned14import PuduLangResilience.Pipeline as Pipeline15import PuduLangResilience.RateLimiter as RateLimiter16import PuduLangResilience as Resilience17import PuduLangResilience.Strategy as Strategy1819/// A pipeline of one rate limiter strategy.20fn limited(options: RateLimiter.Options) -> Pipeline.Pipeline[Int, Str] {21  match Pipeline.build([RateLimiter.strategy(options)]) {22    case Ok(pipeline) => pipeline23    case Err(invalid) => panic(Pipeline.explain(&invalid))24  }25}2627/// The limiter the options make; stops the suite when they are refused.28fn made(result: Result[Limiter.Limiter, Array[Str]]) -> Limiter.Limiter {29  match result {30    case Ok(limiter) => limiter31    case Err(problems) => panic(problems.join("; "))32  }33}3435/// A callback answering 1.36fn one() -> Strategy.Callback[Int, Str] { fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] { Ok(1) } }3738/// Runs the suite.39fn main() -> Int {40  let clock = Clock.manual(0)41  let window = made(FixedWindow.create(&FixedWindow.Options{..FixedWindow.defaults(), permitLimit: 1, window: 1000, clock: Clock.ofManual(&clock)}))42  let rejections: Sync.Cell[Array[Option[Int]]] = Sync.cell([])43  let windowed = limited(RateLimiter.Options{..RateLimiter.using(window), onRejected: Some(fn(rejected: RateLimiter.Rejected) -> () {44          let _kept = Sync.set(&rejections, [rejected.retryAfter])45        }) })46  let admitted = Pipeline.execute(&windowed, one())47  Clock.advance(&clock, 400)48  let rejected = Pipeline.execute(&windowed, one())4950  let single = limited(RateLimiter.concurrency(1, 0))51  let nested = Pipeline.execute(&single, fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] {52      match Pipeline.execute(&single, one()) {53        case Err(Resilience.RateLimited(_)) => Ok(-1)54        case other => other55      }56    })57  let afterRelease = Pipeline.execute(&single, one())58  let afterFailure = Pipeline.execute(&single, fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] { Resilience.raise("bad") })59  let permitReturned = Pipeline.execute(&single, one())6061  let queued = made(Concurrency.create(&Concurrency.Options{..Concurrency.defaults(), permitLimit: 1, queueLimit: 1}))62  let blocker = Limiter.acquire(&queued, 1, &Context.create().token)63  let stopped = Context.create()64  Context.cancel(&stopped, "caller left")65  let abandoned = Pipeline.executeWith(&limited(RateLimiter.using(queued)), stopped, one())66  Limiter.release(&blocker)6768  let partitions = Partitioned.create(fn(context: Context.Context) -> Partitioned.Partition {69      Partitioned.Partition{key: match context.operationKey { case Some(key) => key case None => "" }, create: fn() -> Limiter.Limiter {70          made(FixedWindow.create(&FixedWindow.Options{..FixedWindow.defaults(), permitLimit: 1, window: 60000}))71        } }72    })73  let byTenant = limited(RateLimiter.partitioned(partitions))74  let tenantA = Pipeline.executeWith(&byTenant, Context.keyed("a"), one())75  let tenantAAgain = Pipeline.executeWith(&byTenant, Context.keyed("a"), one())76  let tenantB = Pipeline.executeWith(&byTenant, Context.keyed("b"), one())7778  let invalid = Pipeline.build([RateLimiter.strategy(RateLimiter.concurrency(0, -1))])79  let checks = Test.suite("RateLimiter", &[80      Test.equals("an execution with a permit runs", &admitted, &Ok(1)),81      Test.equals("an execution without a permit is rejected with the wait", &rejected, &Err(Resilience.RateLimited(Some(600)))),82      Test.equals("onRejected sees the wait", &Sync.get(&rejections), &Ok([Some(600)])),83      Test.equals("a concurrency limit rejects a second execution at once", &nested, &Ok(-1)),84      Test.equals("the permit returns when the execution ends", &afterRelease, &Ok(1)),85      Test.equals("a failed execution passes its failure through", &afterFailure, &Resilience.raise("bad")),86      Test.equals("a failed execution returns its permit", &permitReturned, &Ok(1)),87      Test.equals("a fired token abandons the wait as a cancellation", &abandoned, &Err(Resilience.Cancelled("cancelled: caller left"))),88      Test.equals("each partition admits its own first execution", &(tenantA, tenantB), &(Ok(1), Ok(1))),89      Test.that("a partition over its limit rejects", match tenantAAgain { case Err(Resilience.RateLimited(Some(_))) => true case _ => false }),90      Test.equals("the default limiter holds a thousand permits and no queue", &(RateLimiter.defaults().defaultLimiter.permitLimit, RateLimiter.defaults().defaultLimiter.queueLimit), &(1000, 0)),91      Test.equals("invalid default limiter options are refused", &Result.err(invalid), &Some(Pipeline.Invalid{problems: ["RateLimiter: permitLimit must be at least 1", "RateLimiter: queueLimit must not be negative"]}))92    ])93  let ran = Test.run(&checks)94  for failure in Test.failuresOf(&ran) {95    let _reported = Io.writeErrorLine(failure)96  }97  Test.report(&ran)98}99