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

LimiterTest.pudu

Pudu179 lines10.5 KB

GitHub ↗
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 TokenBucket1718/// The limiter the options make; stops the suite when they are refused.19fn 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}2526/// Waits up to two seconds for `count` permits to be queued; whether they were.27fn 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}3637/// Starts a thread that acquires one permit and stores the grant.38fn 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}5152/// Runs the suite.53fn 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