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

Attempts.pudu

Pudu96 lines3.8 KB

GitHub ↗
1/** @Hedging.Attempts.Module — concurrent attempts, their outcomes, and their end */2module PuduLangResilience.Hedging.Attempts34import Std.Channel as Channel5import Std.Concurrent.Cancel as Cancel6import Std.Concurrent as Concurrent7import PuduLangResilience.Clock as Clock8import PuduLangResilience.Context as Context9import PuduLangResilience as Resilience10import PuduLangResilience.Utils.Shared as Shared1112/** @Hedging.Attempts.Attempt — one running execution of a hedged callback */13export type Attempt[T, E] = {14  index: Int,15  context: Context.Context,16  started: Int,17  outcome: Shared.Shared[Option[Resilience.Outcome[T, E]]],18  worker: Option[Concurrent.Task]19}2021/// How long a wait for the next finished attempt sleeps before looking again, in milliseconds.22const POLL_MILLIS: Int = 12324/// A wait that ends only when an attempt finishes.25export const FOREVER: Int = -12627/// Starts `callback` on a thread of its own with `context`. When it finishes, its outcome is28/// stored and `index` is sent to `finished`; a callback that stops its thread finishes as29/// `Crashed`, and one that cannot be started finishes at once as `Crashed`.30export fn launch[T, E](index: Int, context: Context.Context, callback: fn(Context.Context) -> Resilience.Outcome[T, E], finished: &Channel.Channel[Int], started: Int) -> Attempt[T, E] {31  let slot: Shared.Shared[Option[Resilience.Outcome[T, E]]] = Shared.shared(None)32  let done = *finished33  let worker = Concurrent.start(fn() -> () {34      let ran = Concurrent.contain(fn() -> () {35          let outcome = callback(context)36          Shared.update(&slot, fn(_held: Option[Resilience.Outcome[T, E]]) -> Option[Resilience.Outcome[T, E]] { Some(outcome) })37        })38      if let Err(problem) = ran {39        Shared.update(&slot, fn(_held: Option[Resilience.Outcome[T, E]]) -> Option[Resilience.Outcome[T, E]] { Some(Err(Resilience.Crashed(show(problem)))) })40      }41      let _sent = Channel.send(&done, index)42    })43  match worker {44    case Ok(running) => Attempt{index: index, context: context, started: started, outcome: slot, worker: Some(running)}45    case Err(problem) => {46      Shared.update(&slot, fn(_held: Option[Resilience.Outcome[T, E]]) -> Option[Resilience.Outcome[T, E]] { Some(Err(Resilience.Crashed(show(problem)))) })47      let _sent = Channel.send(&done, index)48      Attempt{index: index, context: context, started: started, outcome: slot, worker: None}49    }50  }51}5253/// The index of the next attempt to finish, or `None` when `wait` milliseconds pass on the clock54/// first. A `wait` of `FOREVER` waits for as long as it takes.55export fn next(finished: &Channel.Channel[Int], clock: &Clock.Clock, wait: Int) -> Option[Int] {56  if wait == FOREVER {57    return match Channel.receive(finished) {58      case Ok(Some(index)) => Some(index)59      case _ => None60    }61  }62  let until = Clock.now(clock) + wait63  let never = Cancel.token()64  loop {65    match Channel.pending(finished) {66      case Ok(count) => {67        if count > 0 {68          return match Channel.receive(finished) {69            case Ok(Some(index)) => Some(index)70            case _ => None71          }72        }73      }74      case Err(_) => { return None }75    }76    if Clock.now(clock) >= until { return None }77    let _slept = Clock.sleep(clock, POLL_MILLIS, &never)78  }79}8081/// The outcome of a finished attempt.82export fn outcomeOf[T, E](attempt: &Attempt[T, E]) -> Resilience.Outcome[T, E] {83  match Shared.current(&attempt.outcome) {84    case Some(outcome) => outcome85    case None => Err(Resilience.Crashed("the attempt finished without an outcome"))86  }87}8889/// Asks every attempt to stop and waits until each thread has ended.90export fn settle[T, E](attempts: &Array[Attempt[T, E]], why: Str) -> () {91  for attempt in attempts { Context.cancel(&attempt.context, why) }92  for attempt in attempts {93    if let Some(running) = attempt.worker { let _joined = Concurrent.join(&running) }94  }95}96