
Mediator.pudu
Pudu199 lines10.0 KB
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 | Foreign37383940export fn defaults() -> Options { Options{publisher: Publishing.Sequential{}, actionScope: Catalog.ForUnhandled} }414243export fn build(registrations: Array[Registration.Registration]) -> Result[Mediator, Invalid] { buildWith(&defaults(), registrations) }44454647export 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}606162export fn explain(invalid: &Invalid) -> Str { Template.fill(Messages.MEDIATOR_INVALID, &[invalid.problems.join(Messages.SEPARATOR)]) }636465export 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}72737475export 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}828384export 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}878889export 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}939495export fn dispatch(mediator: &Mediator, envelope: Message.Envelope) -> Message.Answer { dispatchWith(mediator, envelope, Context.create()) }96979899export 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}105106107export fn publish[N, E](mediator: &Mediator, subject: &Notification.Kind[N, E], notification: N) -> Messaging.Outcome[(), E] {108 publishWith(mediator, subject, notification, Context.create())109}110111112113export 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}120121122export fn broadcast(mediator: &Mediator, envelope: Message.Envelope) -> Message.Answer { broadcastWith(mediator, envelope, Context.create()) }123124125126export 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}132133134export 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}137138139export 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}146147148export 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}156157158159export 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}164165166167export 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}178179180fn 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}186187188fn 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