
Engine.pudu
Pudu91 lines4.5 KB
1/** @Limiter.Engine.Module — thread-safe leasing over one permit algorithm */2module PuduLangResilience.Limiter.Engine34import Std.Concurrent.Cancel as Cancel5import PuduLangResilience.Clock as Clock6import PuduLangResilience.Domain.Permits as Permits7import PuduLangResilience.Limiter as Limiter8import PuduLangResilience.Utils.Shared as Shared910/** @Limiter.Engine.Queueing — the limits every engine enforces */11export type Queueing = { permitLimit: Int, queueLimit: Int, order: Limiter.QueueOrder }121314const POLL_MILLIS: Int = 1151617export fn limiter[S](held: S, algorithm: Permits.Algorithm[S], queueing: &Queueing, clock: Clock.Clock) -> Limiter.Limiter {18 let order = match queueing.order {19 case Limiter.OldestFirst => Permits.Oldest20 case Limiter.NewestFirst => Permits.Newest21 }22 let book = Shared.shared(Permits.open(held, queueing.permitLimit, queueing.queueLimit, order))23 Limiter.Limiter {24 acquire: fn(permits: Int, token: Cancel.Token) -> Limiter.Lease { acquire(&book, &algorithm, &clock, permits, &token) },25 attempt: fn(permits: Int) -> Limiter.Lease { attempt(&book, &algorithm, &clock, permits) },26 statistics: fn() -> Limiter.Statistics {27 let now = Clock.now(&clock)28 let current = Shared.current(&book)29 Limiter.Statistics{available: Permits.available(¤t, &algorithm, now), queued: Permits.queued(¤t), failed: current.failed, succeeded: current.succeeded}30 },31 replenish: fn() -> Bool {32 Shared.change(&book, fn(current: Permits.Book[S]) -> (Permits.Book[S], Bool) { Permits.replenish(¤t, &algorithm) })33 }34 }35}363738fn attempt[S](book: &Shared.Shared[Permits.Book[S]], algorithm: &Permits.Algorithm[S], clock: &Clock.Clock, permits: Int) -> Limiter.Lease {39 let now = Clock.now(clock)40 let decided = Shared.change(book, fn(current: Permits.Book[S]) -> (Permits.Book[S], Permits.Decision) { Permits.request(¤t, algorithm, now, permits, false) })41 answer(book, algorithm, permits, decided)42}43444546fn acquire[S](book: &Shared.Shared[Permits.Book[S]], algorithm: &Permits.Algorithm[S], clock: &Clock.Clock, permits: Int, token: &Cancel.Token) -> Limiter.Lease {47 if let Err(why) = Cancel.check(token) { return Limiter.refused(Limiter.Abandoned(Cancel.explain(&why))) }48 let now = Clock.now(clock)49 let decided = Shared.change(book, fn(current: Permits.Book[S]) -> (Permits.Book[S], Permits.Decision) { Permits.request(¤t, algorithm, now, permits, true) })50 match decided {51 case Permits.Queue(ticket) => wait(book, algorithm, clock, permits, ticket, token)52 case _ => answer(book, algorithm, permits, decided)53 }54}555657fn wait[S](book: &Shared.Shared[Permits.Book[S]], algorithm: &Permits.Algorithm[S], clock: &Clock.Clock, permits: Int, ticket: Int, token: &Cancel.Token) -> Limiter.Lease {58 loop {59 let now = Clock.now(clock)60 let turn = Shared.change(book, fn(current: Permits.Book[S]) -> (Permits.Book[S], Permits.Turn) { Permits.poll(¤t, algorithm, now, ticket) })61 match turn {62 case Permits.Taken => { return lease(book, algorithm, permits) }63 case Permits.Evicted => { return Limiter.refused(Limiter.Denied(None)) }64 case Permits.Waiting => {65 if let Err(why) = Clock.sleep(clock, POLL_MILLIS, token) {66 Shared.update(book, fn(current: Permits.Book[S]) -> Permits.Book[S] { Permits.withdraw(¤t, ticket) })67 return Limiter.refused(Limiter.Abandoned(Cancel.explain(&why)))68 }69 }70 }71 }72}737475fn answer[S](book: &Shared.Shared[Permits.Book[S]], algorithm: &Permits.Algorithm[S], permits: Int, decided: Permits.Decision) -> Limiter.Lease {76 match decided {77 case Permits.Grant => lease(book, algorithm, permits)78 case Permits.Deny(retryAfter) => Limiter.refused(Limiter.Denied(retryAfter))79 case Permits.Queue(_) => Limiter.refused(Limiter.Denied(None))80 }81}828384fn lease[S](book: &Shared.Shared[Permits.Book[S]], algorithm: &Permits.Algorithm[S], permits: Int) -> Limiter.Lease {85 let target = *book86 let counting = *algorithm87 Limiter.granted(fn() -> () {88 Shared.update(&target, fn(current: Permits.Book[S]) -> Permits.Book[S] { Permits.giveBack(¤t, &counting, permits) })89 })90}91