
Attempts.pudu
Pudu96 lines3.8 KB
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}202122const POLL_MILLIS: Int = 1232425export const FOREVER: Int = -12627282930export 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}52535455export 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}808182export 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}888990export 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