
RateLimiterTest.pudu
Pudu99 lines5.2 KB
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 Strategy181920fn 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}262728fn 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}343536fn one() -> Strategy.Callback[Int, Str] { fn(_context: Context.Context) -> Resilience.Outcome[Int, Str] { Ok(1) } }373839fn 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