
Publishing.pudu
Pudu135 lines5.3 KB
1/** @Mediator.Publishing.Seam — how a notification reaches its handlers */2module PuduLangMediator.Publishing34import Std.Concurrent as Concurrent5import Std.Math as Math6import Std.Result as Result7import Std.Sync as Sync8import PuduLangMediator.Context as Context9import PuduLangMediator as Messaging10import PuduLangMediator.Utils.Threads as Threads1112/** @Mediator.Publishing.Executor — one handler, ready to run with a context */13export type Executor[E] = { handler: Str, run: fn(Context.Context) -> Messaging.Outcome[(), E] }1415/** @Mediator.Publishing.Strategy — runs the executors of one publication */16export trait Strategy {17 18 fn publish[E](self: &Self, executors: &Array[Executor[E]], context: Context.Context) -> Messaging.Outcome[(), E]19}2021/** @Mediator.Publishing.Sequential — one handler at a time, stopping at the first failure */22export type Sequential = {}2324/** @Mediator.Publishing.Continuing — one handler at a time, every handler, every failure */25export type Continuing = {}2627/** @Mediator.Publishing.Parallel — every handler on its own thread at once */28export type Parallel = {}2930/** @Mediator.Publishing.Bounded — every handler, at most `workers` at a time */31export type Bounded = { workers: Int }3233impl Strategy for Sequential {34 35 fn publish[E](self: &Self, executors: &Array[Executor[E]], context: Context.Context) -> Messaging.Outcome[(), E] {36 for executor in executors {37 Context.check(&context) ?38 (executor.run)(context) ?39 }40 Ok(())41 }42}4344impl Strategy for Continuing {45 46 fn publish[E](self: &Self, executors: &Array[Executor[E]], context: Context.Context) -> Messaging.Outcome[(), E] {47 var failures: Array[Messaging.Failure[E]] = []48 for executor in executors {49 match (executor.run)(context) {50 case Ok(_) => {}51 case Err(failure) => { failures = failures.push(failure) }52 }53 }54 settled(failures)55 }56}5758impl Strategy for Parallel {59 60 fn publish[E](self: &Self, executors: &Array[Executor[E]], context: Context.Context) -> Messaging.Outcome[(), E] {61 let cells = executors.map(|_executor: Executor[E]| outcomeCell())62 var started: Array[Result[Concurrent.Task, Concurrent.ConcurrentError]] = []63 var index = 064 for executor in executors {65 let cell = cells[index]66 started = started.push(Concurrent.start(fn() -> () { let _stored = Sync.set(&cell, Some((executor.run)(context))) }))67 index = index + 168 }69 var outcomes: Array[Messaging.Outcome[(), E]] = []70 index = 071 for attempt in started {72 let joined = Result.andThen(attempt, |worker: Concurrent.Task| Concurrent.join(&worker))73 outcomes = outcomes.push(collected(&cells[index], joined))74 index = index + 175 }76 settled(failuresOf(&outcomes))77 }78}7980impl Strategy for Bounded {81 82 83 fn publish[E](self: &Self, executors: &Array[Executor[E]], context: Context.Context) -> Messaging.Outcome[(), E] {84 let cells = executors.map(|_executor: Executor[E]| outcomeCell())85 var actions: Array[fn() -> ()] = []86 var index = 087 for executor in executors {88 let cell = cells[index]89 actions = actions.push(fn() -> () {90 let contained = Concurrent.contain(fn() -> () { let _stored = Sync.set(&cell, Some((executor.run)(context))) })91 if let Err(problem) = contained { let _crashed = Sync.set(&cell, Some(Err(Messaging.Crashed(Threads.reasonOf(&problem))))) }92 })93 index = index + 194 }95 let ran = Concurrent.parallelBounded(&actions, Math.max(self.workers, 1))96 var outcomes: Array[Messaging.Outcome[(), E]] = []97 for cell in cells { outcomes = outcomes.push(collected(&cell, ran)) }98 settled(failuresOf(&outcomes))99 }100}101102103fn outcomeCell[E]() -> Sync.Cell[Option[Messaging.Outcome[(), E]]] { Sync.cell(None) }104105106fn collected[E](cell: &Sync.Cell[Option[Messaging.Outcome[(), E]]], joined: Result[(), Concurrent.ConcurrentError]) -> Messaging.Outcome[(), E] {107 match Sync.get(cell) {108 case Ok(Some(outcome)) => outcome109 case Ok(None) => {110 match joined {111 case Err(problem) => Err(Messaging.Crashed(Threads.reasonOf(&problem)))112 case Ok(_) => Err(Messaging.Crashed("the handler stored no outcome"))113 }114 }115 case Err(problem) => Err(Messaging.Crashed(show(problem)))116 }117}118119120fn failuresOf[E](outcomes: &Array[Messaging.Outcome[(), E]]) -> Array[Messaging.Failure[E]] {121 var failures: Array[Messaging.Failure[E]] = []122 for outcome in outcomes {123 if let Err(failure) = outcome { failures = failures.push(failure) }124 }125 failures126}127128129fn settled[E](failures: Array[Messaging.Failure[E]]) -> Messaging.Outcome[(), E] {130 match Messaging.combine(failures) {131 case Some(failure) => Err(failure)132 case None => Ok(())133 }134}135