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

Engine.pudu

Pudu91 lines4.5 KB

GitHub ↗
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 }1213/// How long a queued request waits before asking again, in milliseconds.14const POLL_MILLIS: Int = 11516/// A limiter over `algorithm`, starting from `held`, reading time from `clock`.17export 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(&current, &algorithm, now), queued: Permits.queued(&current), 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(&current, &algorithm) })33    }34  }35}3637/// Leases permits without queueing.38fn 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(&current, algorithm, now, permits, false) })41  answer(book, algorithm, permits, decided)42}4344/// Leases permits, queueing and asking again on the clock until they are granted, the request45/// is evicted, or the token fires.46fn 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(&current, 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}5556/// Waits in the queue holding `ticket`.57fn 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(&current, 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(&current, ticket) })67          return Limiter.refused(Limiter.Abandoned(Cancel.explain(&why)))68        }69      }70    }71  }72}7374/// The lease a decision grants or refuses.75fn 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}8283/// A granted lease that hands `permits` back to the book on release.84fn 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(&current, &counting, permits) })89    })90}91