
LimiterTest.pudu
Pudu179 lines10.5 KB
1/** @Test.Limiter.Suite — real limiters across threads and manual time */2module PuduLangResilience.Limiter.LimiterTest34import Std.Concurrent.Cancel as Cancel5import Std.Concurrent as Concurrent6import Std.Io as Io7import Std.Sync as Sync8import Std.Test as Test9import PuduLangResilience.Clock as Clock10import PuduLangResilience.Context as Context11import PuduLangResilience.Limiter.Concurrency as Concurrency12import PuduLangResilience.Limiter.FixedWindow as FixedWindow13import PuduLangResilience.Limiter as Limiter14import PuduLangResilience.Limiter.Partitioned as Partitioned15import PuduLangResilience.Limiter.SlidingWindow as SlidingWindow16import PuduLangResilience.Limiter.TokenBucket as TokenBucket171819fn made(result: Result[Limiter.Limiter, Array[Str]]) -> Limiter.Limiter {20 match result {21 case Ok(limiter) => limiter22 case Err(problems) => panic(problems.join("; "))23 }24}252627fn queuedAtLeast(limiter: &Limiter.Limiter, count: Int) -> Bool {28 var waited = 029 while waited < 2000 {30 if Limiter.statistics(limiter).queued >= count { return true }31 let _slept = Concurrent.sleep(2)32 waited = waited + 233 }34 false35}363738fn contend(limiter: &Limiter.Limiter, token: &Cancel.Token, into: &Sync.Cell[Option[Limiter.Grant]]) -> Concurrent.Task {39 let target = *limiter40 let stop = *token41 let slot = *into42 match Concurrent.start(fn() -> () {43 let lease = Limiter.acquire(&target, 1, &stop)44 let _kept = Sync.set(&slot, Some(lease.grant))45 Limiter.release(&lease)46 }) {47 case Ok(started) => started48 case Err(problem) => panic(show(problem))49 }50}515253fn main() -> Int {54 let single = made(Concurrency.create(&Concurrency.Options{..Concurrency.defaults(), permitLimit: 1, queueLimit: 1}))55 let held = Limiter.acquire(&single, 1, &Cancel.token())56 let waiterGrant: Sync.Cell[Option[Limiter.Grant]] = Sync.cell(None)57 let waiter = contend(&single, &Cancel.token(), &waiterGrant)58 let queuedSeen = queuedAtLeast(&single, 1)59 let overflow = Limiter.acquire(&single, 1, &Cancel.token())60 Limiter.release(&held)61 Limiter.release(&held)62 let _joined = Concurrent.join(&waiter)63 let afterRelease = Limiter.statistics(&single)6465 let newest = made(Concurrency.create(&Concurrency.Options{..Concurrency.defaults(), permitLimit: 1, queueLimit: 1, queueOrder: Limiter.NewestFirst}))66 let newestHeld = Limiter.acquire(&newest, 1, &Cancel.token())67 let evictedGrant: Sync.Cell[Option[Limiter.Grant]] = Sync.cell(None)68 let evictedTask = contend(&newest, &Cancel.token(), &evictedGrant)69 let _firstQueued = queuedAtLeast(&newest, 1)70 let newestGrant: Sync.Cell[Option[Limiter.Grant]] = Sync.cell(None)71 let newestTask = contend(&newest, &Cancel.token(), &newestGrant)72 let _evictedDone = Concurrent.join(&evictedTask)73 Limiter.release(&newestHeld)74 let _newestDone = Concurrent.join(&newestTask)7576 let cancelling = made(Concurrency.create(&Concurrency.Options{..Concurrency.defaults(), permitLimit: 1, queueLimit: 5}))77 let blocker = Limiter.acquire(&cancelling, 1, &Cancel.token())78 let stop = Cancel.token()79 let cancelledGrant: Sync.Cell[Option[Limiter.Grant]] = Sync.cell(None)80 let cancelledTask = contend(&cancelling, &stop, &cancelledGrant)81 let _cancelQueued = queuedAtLeast(&cancelling, 1)82 let _fired = Cancel.cancel(&stop, "gave up")83 let _cancelDone = Concurrent.join(&cancelledTask)84 let afterCancel = Limiter.statistics(&cancelling)85 Limiter.release(&blocker)86 let alreadyStopped = Cancel.token()87 let _stoppedFirst = Cancel.cancel(&alreadyStopped, "early")88 let stoppedBefore = Limiter.acquire(&cancelling, 1, &alreadyStopped)8990 let clock = Clock.manual(0)91 let bucket = made(TokenBucket.create(&TokenBucket.Options{..TokenBucket.defaults(), tokenLimit: 2, tokensPerPeriod: 1, replenishmentPeriod: 100, queueLimit: 1, clock: Clock.ofManual(&clock)}))92 let b1 = Limiter.attempt(&bucket, 1)93 let b2 = Limiter.attempt(&bucket, 1)94 let b3 = Limiter.attempt(&bucket, 1)95 let queuedWait = Limiter.acquire(&bucket, 1, &Cancel.token())96 let bucketSleeps = Clock.sleeps(&clock).length()97 let manualBucket = made(TokenBucket.create(&TokenBucket.Options{..TokenBucket.defaults(), tokenLimit: 1, tokensPerPeriod: 1, autoReplenishment: false, clock: Clock.ofManual(&clock)}))98 let _m1 = Limiter.attempt(&manualBucket, 1)99 Clock.advance(&clock, 100000)100 let stillDry = Limiter.attempt(&manualBucket, 1)101 let replenished = Limiter.replenish(&manualBucket)102 let afterReplenish = Limiter.attempt(&manualBucket, 1)103104 let windowClock = Clock.manual(0)105 let window = made(FixedWindow.create(&FixedWindow.Options{..FixedWindow.defaults(), permitLimit: 2, window: 1000, clock: Clock.ofManual(&windowClock)}))106 let _w1 = Limiter.attempt(&window, 2)107 Clock.advance(&windowClock, 250)108 let windowDenied = Limiter.attempt(&window, 1)109 Clock.advance(&windowClock, 750)110 let windowAgain = Limiter.attempt(&window, 1)111112 let slideClock = Clock.manual(0)113 let sliding = made(SlidingWindow.create(&SlidingWindow.Options{..SlidingWindow.defaults(), permitLimit: 2, window: 1000, segmentsPerWindow: 2, clock: Clock.ofManual(&slideClock)}))114 let _s1 = Limiter.attempt(&sliding, 1)115 Clock.advance(&slideClock, 600)116 let _s2 = Limiter.attempt(&sliding, 1)117 let slideDenied = Limiter.attempt(&sliding, 1)118 Clock.advance(&slideClock, 400)119 let slideFreed = Limiter.attempt(&sliding, 1)120121 let chainClock = Clock.manual(0)122 let roomy = made(Concurrency.create(&Concurrency.Options{..Concurrency.defaults(), permitLimit: 5}))123 let tight = made(FixedWindow.create(&FixedWindow.Options{..FixedWindow.defaults(), permitLimit: 1, window: 1000, clock: Clock.ofManual(&chainClock)}))124 let chained = Limiter.chain([roomy, tight])125 let c1 = Limiter.attempt(&chained, 1)126 let c2 = Limiter.attempt(&chained, 1)127 let roomyAfterRefusal = Limiter.statistics(&roomy).available128 Limiter.release(&c1)129 let chainedStats = Limiter.statistics(&chained)130131 let partitions = Partitioned.create(fn(context: Context.Context) -> Partitioned.Partition {132 let key = match context.operationKey { case Some(name) => name case None => "anonymous" }133 Partitioned.Partition{key: key, create: fn() -> Limiter.Limiter { made(Concurrency.create(&Concurrency.Options{..Concurrency.defaults(), permitLimit: 1})) }}134 })135 let alice = Context.keyed("alice")136 let p1 = Partitioned.attempt(&partitions, &alice, 1)137 let p2 = Partitioned.attempt(&partitions, &alice, 1)138 let p3 = Partitioned.acquire(&partitions, &Context.keyed("bob"), 1)139140 let checks = Test.suite("Limiter", &[141 Test.that("a request without permits waits in the queue", queuedSeen),142 Test.equals("a request beyond the queue is denied", &overflow.grant, &Limiter.Denied(None)),143 Test.equals("a released permit goes to the waiter", &Sync.get(&waiterGrant), &Ok(Some(Limiter.Granted))),144 Test.equals("releasing twice returns the permit once", &afterRelease, &Limiter.Statistics{available: 1, queued: 0, failed: 1, succeeded: 2}),145 Test.equals("newest first evicts the older waiter", &Sync.get(&evictedGrant), &Ok(Some(Limiter.Denied(None)))),146 Test.equals("newest first serves the newer waiter", &Sync.get(&newestGrant), &Ok(Some(Limiter.Granted))),147 Test.equals("a cancelled waiter is abandoned", &Sync.get(&cancelledGrant), &Ok(Some(Limiter.Abandoned("cancelled: gave up")))),148 Test.equals("a cancelled waiter leaves the queue", &afterCancel.queued, &0),149 Test.equals("a fired token is abandoned before queueing", &stoppedBefore.grant, &Limiter.Abandoned("cancelled: early")),150 Test.that("a full bucket grants its tokens", Limiter.isGranted(&b1) && Limiter.isGranted(&b2)),151 Test.equals("an empty bucket denies with the wait for one token", &b3.grant, &Limiter.Denied(Some(100))),152 Test.that("a queued request waits on the clock for the next token", Limiter.isGranted(&queuedWait) && bucketSleeps > 0),153 Test.equals("a bucket replenished by hand stays empty over time", &stillDry.grant, &Limiter.Denied(None)),154 Test.that("replenishing by hand refills the bucket", replenished && Limiter.isGranted(&afterReplenish)),155 Test.equals("a spent window denies until it ends", &windowDenied.grant, &Limiter.Denied(Some(750))),156 Test.that("the next window grants again", Limiter.isGranted(&windowAgain)),157 Test.equals("a sliding window denies until its oldest segment leaves", &slideDenied.grant, &Limiter.Denied(Some(400))),158 Test.that("a segment leaving the window frees its permits", Limiter.isGranted(&slideFreed)),159 Test.that("a chain grants when every limiter does", Limiter.isGranted(&c1)),160 Test.equals("a chain answers the first refusal", &c2.grant, &Limiter.Denied(Some(1000))),161 Test.equals("a refused chain releases what it took", &roomyAfterRefusal, &4),162 Test.equals("a chain reports the tightest availability", &chainedStats.available, &0),163 Test.that("each partition has its own permits", Limiter.isGranted(&p1) && !Limiter.isGranted(&p2) && Limiter.isGranted(&p3)),164 Test.equals("partitions are made on first use", &Partitioned.keys(&partitions), &["alice", "bob"]),165 Test.equals("a partition reports its counts", &Partitioned.statistics(&partitions, "alice"), &Some(Limiter.Statistics{available: 0, queued: 0, failed: 1, succeeded: 1})),166 Test.equals("an unused partition has no counts", &Partitioned.statistics(&partitions, "carol"), &None),167 Test.equals("invalid concurrency options are refused", &Concurrency.validate(&Concurrency.Options{..Concurrency.defaults(), permitLimit: 0, queueLimit: -1}), &["permitLimit must be at least 1", "queueLimit must not be negative"]),168 Test.equals("invalid bucket options are refused", &TokenBucket.validate(&TokenBucket.Options{..TokenBucket.defaults(), tokenLimit: 0, tokensPerPeriod: 0, replenishmentPeriod: 0, queueLimit: -1}).length(), &4),169 Test.equals("invalid window options are refused", &FixedWindow.validate(&FixedWindow.Options{..FixedWindow.defaults(), permitLimit: 0, window: 0, queueLimit: -1}).length(), &3),170 Test.equals("a window too short for its segments is refused", &SlidingWindow.validate(&SlidingWindow.Options{..SlidingWindow.defaults(), window: 5, segmentsPerWindow: 10}), &["window must be at least 1 ms for each segment"]),171 Test.that("refused options make no limiter", match SlidingWindow.create(&SlidingWindow.Options{..SlidingWindow.defaults(), segmentsPerWindow: 0}) { case Err(_) => true case Ok(_) => false })172 ])173 let ran = Test.run(&checks)174 for failure in Test.failuresOf(&ran) {175 let _reported = Io.writeErrorLine(failure)176 }177 Test.report(&ran)178}179