Pudu programming language
Menu
Package

@chrismichaelps / pudu-lang-mediator

In-process messaging for Pudu: requests, notifications, streams, pipeline behaviors, processors, and exception handling

0.1.0Apache-2.01

InstallClose

Publishing.pudu

Pudu135 lines5.3 KB

GitHub ↗
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  /// Runs the executors and answers the publication's outcome.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  /// Runs each executor in order once the context is still live, and answers the first failure.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  /// Runs every executor in order and answers every failure, combined.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  /// Starts every executor on its own thread, waits for all, and answers every failure, combined.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  /// Runs every executor with at most `workers` running at once (one when fewer are asked for),82  /// waits for all, and answers every failure, combined.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}101102/// An empty place for one executor's outcome.103fn outcomeCell[E]() -> Sync.Cell[Option[Messaging.Outcome[(), E]]] { Sync.cell(None) }104105/// The outcome an executor stored, or `Crashed` when its thread stopped before storing one.106fn 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}118119/// The failures among outcomes, in order.120fn 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}127128/// `Ok` when nothing failed, otherwise the failures combined.129fn 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