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

Mediator.pudu

Pudu199 lines10.0 KB

GitHub ↗
1/** @Mediator.Mediator.Backbone — routes requests, notifications, and streams to their handlers */2module PuduLangMediator.Mediator34import Std.HashMap as HashMap5import PuduLangMediator.Catalog as Catalog6import PuduLangMediator.Compose as Compose7import PuduLangMediator.Constants.Messages as Messages8import PuduLangMediator.Context as Context9import PuduLangMediator.Message as Message10import PuduLangMediator as Messaging11import PuduLangMediator.Notification as Notification12import PuduLangMediator.Publishing as Publishing13import PuduLangMediator.Reader as Reader14import PuduLangMediator.Registration as Registration15import PuduLangMediator.Request as Request16import PuduLangMediator.Stream as Stream17import PuduLangMediator.Utils.Erasure as Erasure18import PuduLangMediator.Utils.Shared as Shared19import PuduLangMediator.Utils.Template as Template2021/** @Mediator.Mediator.Options — how notifications are published and failures observed */22export type Options = { publisher: dynamic Publishing.Strategy, actionScope: Catalog.Scope }2324/** @Mediator.Mediator.Mediator — every kind's composed pipeline, by shape and name */25export type Mediator = {26  requests: HashMap.HashMap[Str, Catalog.Compiled],27  notifications: HashMap.HashMap[Str, Catalog.Compiled],28  streams: HashMap.HashMap[Str, Catalog.Compiled],29  assembly: Catalog.Assembly30}3132/** @Mediator.Mediator.Invalid — every reason a mediator could not be built */33export type Invalid = { problems: Array[Str] }3435/** @Mediator.Mediator.Located — a kind's composed pipeline, its absence, or another kind's */36type Located[H] = Ready(H) | Missing | Foreign3738/// Sequential publishing, and exception actions that observe only failures no exception handler39/// answered.40export fn defaults() -> Options { Options{publisher: Publishing.Sequential{}, actionScope: Catalog.ForUnhandled} }4142/// A mediator of `registrations` with the default options.43export fn build(registrations: Array[Registration.Registration]) -> Result[Mediator, Invalid] { buildWith(&defaults(), registrations) }4445/// A mediator of `registrations`, applied in order, or every problem found in them: a kind without46/// a name, two kinds under one name, and a request or stream kind without exactly one handler.47export fn buildWith(options: &Options, registrations: Array[Registration.Registration]) -> Result[Mediator, Invalid] {48  var catalog = Catalog.empty(options.publisher, options.actionScope)49  for registration in registrations { catalog = (registration.apply)(catalog) }50  let problems = Catalog.problems(&catalog)51  if !problems.isEmpty() { return Err(Invalid{problems: problems}) }52  let assembly = catalog.assembly53  Ok(Mediator {54      requests: Catalog.compile(&catalog.requests, &assembly),55      notifications: Catalog.compile(&catalog.notifications, &assembly),56      streams: Catalog.compile(&catalog.streams, &assembly),57      assembly: assembly58    })59}6061/// Every problem in one sentence.62export fn explain(invalid: &Invalid) -> Str { Template.fill(Messages.MEDIATOR_INVALID, &[invalid.problems.join(Messages.SEPARATOR)]) }6364/// Every registered kind: requests, then notifications, then streams, each in registration order.65export fn kinds(mediator: &Mediator) -> Array[Message.Info] {66  var found: Array[Message.Info] = []67  for table in [mediator.requests, mediator.notifications, mediator.streams] {68    for compiled in HashMap.values(&table) { found = found.push(compiled.info) }69  }70  found71}7273/// The composed pipeline of a request kind, to call many times without looking it up again;74/// `Unhandled` when the kind has no handler and `Mismatched` when another kind took its name.75export fn resolve[Q, R, E](mediator: &Mediator, subject: &Request.Kind[Q, R, E]) -> Result[Request.Handler[Q, R, E], Messaging.Failure[E]] {76  match located(&mediator.requests, subject.info.name, &subject.invokers) {77    case Ready(invoke) => Ok(invoke)78    case Missing => Err(Messaging.Unhandled(subject.info.name))79    case Foreign => Err(Messaging.Mismatched(subject.info.name))80  }81}8283/// Sends a request through its pipeline with a fresh context.84export fn send[Q, R, E](mediator: &Mediator, subject: &Request.Kind[Q, R, E], request: Q) -> Messaging.Outcome[R, E] {85  sendWith(mediator, subject, request, Context.create())86}8788/// Sends a request through its pipeline with the given context.89export fn sendWith[Q, R, E](mediator: &Mediator, subject: &Request.Kind[Q, R, E], request: Q, context: Context.Context) -> Messaging.Outcome[R, E] {90  let invoke = resolve(mediator, subject) ?91  invoke(request, context)92}9394/// Sends a request whose type the caller does not name, with a fresh context.95export fn dispatch(mediator: &Mediator, envelope: Message.Envelope) -> Message.Answer { dispatchWith(mediator, envelope, Context.create()) }9697/// Sends a request whose type the caller does not name, with the given context. Read the answer98/// back with the kind's `answerOf`.99export fn dispatchWith(mediator: &Mediator, envelope: Message.Envelope, context: Context.Context) -> Message.Answer {100  match HashMap.get(&mediator.requests, &envelope.info.name) {101    case Some(compiled) => answering(&compiled, envelope, context)102    case None => Message.undelivered(&envelope.info, Message.NoRoute(envelope.info.name))103  }104}105106/// Publishes a notification to every handler with a fresh context.107export fn publish[N, E](mediator: &Mediator, subject: &Notification.Kind[N, E], notification: N) -> Messaging.Outcome[(), E] {108  publishWith(mediator, subject, notification, Context.create())109}110111/// Publishes a notification to its handlers and the open ones, through the mediator's publishing112/// strategy, with the given context. A kind with no handlers succeeds without running anything.113export fn publishWith[N, E](mediator: &Mediator, subject: &Notification.Kind[N, E], notification: N, context: Context.Context) -> Messaging.Outcome[(), E] {114  match located(&mediator.notifications, subject.info.name, &subject.invokers) {115    case Ready(invoke) => invoke(notification, context)116    case Missing => (Compose.publication(&subject.info, &Notification.route(), &mediator.assembly))(notification, context)117    case Foreign => Err(Messaging.Mismatched(subject.info.name))118  }119}120121/// Publishes a notification whose type the caller does not name, with a fresh context.122export fn broadcast(mediator: &Mediator, envelope: Message.Envelope) -> Message.Answer { broadcastWith(mediator, envelope, Context.create()) }123124/// Publishes a notification whose type the caller does not name, with the given context. A kind125/// the mediator does not know is answered as heard by nobody.126export fn broadcastWith(mediator: &Mediator, envelope: Message.Envelope, context: Context.Context) -> Message.Answer {127  match HashMap.get(&mediator.notifications, &envelope.info.name) {128    case Some(compiled) => answering(&compiled, envelope, context)129    case None => Message.unheard(&envelope.info)130  }131}132133/// Runs a stream request with a fresh context, handing each item to `sink` until it declines one.134export fn stream[Q, T, E](mediator: &Mediator, subject: &Stream.Kind[Q, T, E], request: Q, sink: Stream.Sink[T]) -> Messaging.Outcome[(), E] {135  streamWith(mediator, subject, request, Context.create(), sink)136}137138/// Runs a stream request with the given context, handing each item to `sink` until it declines one.139export fn streamWith[Q, T, E](mediator: &Mediator, subject: &Stream.Kind[Q, T, E], request: Q, context: Context.Context, sink: Stream.Sink[T]) -> Messaging.Outcome[(), E] {140  match located(&mediator.streams, subject.info.name, &subject.invokers) {141    case Ready(invoke) => invoke(request, context, sink)142    case Missing => Err(Messaging.Unhandled(subject.info.name))143    case Foreign => Err(Messaging.Mismatched(subject.info.name))144  }145}146147/// Every item of a stream request, or its failure.148export fn collect[Q, T, E](mediator: &Mediator, subject: &Stream.Kind[Q, T, E], request: Q) -> Messaging.Outcome[Array[T], E] {149  let gathered: Shared.Shared[Array[T]] = Shared.shared([])150  stream(mediator, subject, request, fn(item: T) -> Bool {151      Shared.update(&gathered, fn(held: Array[T]) -> Array[T] { held.push(item) })152      true153    }) ?154  Ok(Shared.current(&gathered))155}156157/// A reader of a stream request produced on its own thread with the given context, holding at158/// most `capacity` items ahead of the consumer. Close it once done with it.159export fn reader[Q, T, E](mediator: &Mediator, subject: &Stream.Kind[Q, T, E], request: Q, context: &Context.Context, capacity: Int) -> Reader.Reader[T, E] {160  let held = *mediator161  let kind = *subject162  Reader.open(capacity, context, fn(passed: Context.Context, sink: Stream.Sink[T]) -> Messaging.Outcome[(), E] { streamWith(&held, &kind, request, passed, sink) })163}164165/// Runs a stream request whose type the caller does not name, handing each packed item to `sink`166/// until it declines one. Read items back with the kind's `itemOf`, and the answer with its `answerOf`.167export fn streamEnvelope(mediator: &Mediator, envelope: Message.Envelope, context: Context.Context, sink: fn(Erasure.Packed) -> Bool) -> Message.Answer {168  match HashMap.get(&mediator.streams, &envelope.info.name) {169    case Some(compiled) => {170      match compiled.dispatch {171        case Catalog.Streaming(run) => run(envelope, context, sink)172        case Catalog.Answering(_) => Message.undelivered(&envelope.info, Message.WrongKind(envelope.info.name))173      }174    }175    case None => Message.undelivered(&envelope.info, Message.NoRoute(envelope.info.name))176  }177}178179/// A compiled request or notification answering an envelope.180fn answering(compiled: &Catalog.Compiled, envelope: Message.Envelope, context: Context.Context) -> Message.Answer {181  match compiled.dispatch {182    case Catalog.Answering(run) => run(envelope, context)183    case Catalog.Streaming(_) => Message.undelivered(&envelope.info, Message.WrongKind(envelope.info.name))184  }185}186187/// The composed pipeline registered under `name`, read through `invokers`.188fn located[H](table: &HashMap.HashMap[Str, Catalog.Compiled], name: Str, invokers: &Erasure.Slot[H]) -> Located[H] {189  match HashMap.get(table, &name) {190    case Some(compiled) => {191      match Erasure.unpack(invokers, &compiled.invoke) {192        case Some(invoke) => Ready(invoke)193        case None => Foreign194      }195    }196    case None => Missing197  }198}199