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

Registration.pudu

Pudu187 lines9.4 KB

GitHub ↗
1/** @Mediator.Registration.Builder — the handlers and components a mediator is built from */2module PuduLangMediator.Registration34import PuduLangMediator.Catalog as Catalog5import PuduLangMediator.Compose as Compose6import PuduLangMediator.Domain.Order as Order7import PuduLangMediator.Notification as Notification8import PuduLangMediator.Open as Open9import PuduLangMediator.Request as Request10import PuduLangMediator.Stream as Stream1112/** @Mediator.Registration.Registration — one change to the catalog a mediator is built from */13export type Registration = { apply: fn(Catalog.Catalog) -> Catalog.Catalog }1415/// The handler of a request kind. A request kind takes exactly one.16export fn handler[Q, R, E](subject: &Request.Kind[Q, R, E], handle: Request.Handler[Q, R, E]) -> Registration {17  requestChange(subject, 1, 0, None, fn(route: Request.Route[Q, R, E]) -> Request.Route[Q, R, E] {18      Request.Route{..route, handlers: route.handlers.push(handle)}19    })20}2122/// A behavior around the pipeline of one request kind, ordered among every behavior by registration.23export fn behavior[Q, R, E](subject: &Request.Kind[Q, R, E], wrap: Request.Behavior[Q, R, E]) -> Registration {24  requestChange(subject, 0, 1, Some(Order.Behaviors), fn(route: Request.Route[Q, R, E]) -> Request.Route[Q, R, E] {25      Request.Route{..route, behaviors: route.behaviors.push(wrap)}26    })27}2829/// A pre-processor of one request kind, ordered among every pre-processor by registration.30export fn preProcessor[Q, R, E](subject: &Request.Kind[Q, R, E], process: Request.PreProcessor[Q, E]) -> Registration {31  requestChange(subject, 0, 1, Some(Order.PreProcessors), fn(route: Request.Route[Q, R, E]) -> Request.Route[Q, R, E] {32      Request.Route{..route, preProcessors: route.preProcessors.push(process)}33    })34}3536/// A post-processor of one request kind, ordered among every post-processor by registration.37export fn postProcessor[Q, R, E](subject: &Request.Kind[Q, R, E], process: Request.PostProcessor[Q, R, E]) -> Registration {38  requestChange(subject, 0, 1, Some(Order.PostProcessors), fn(route: Request.Route[Q, R, E]) -> Request.Route[Q, R, E] {39      Request.Route{..route, postProcessors: route.postProcessors.push(process)}40    })41}4243/// An exception handler of one request kind; the kind's exception handlers are tried in registration order.44export fn exceptionHandler[Q, R, E](subject: &Request.Kind[Q, R, E], recover: Request.ExceptionHandler[Q, R, E]) -> Registration {45  requestChange(subject, 0, 1, None, fn(route: Request.Route[Q, R, E]) -> Request.Route[Q, R, E] {46      Request.Route{..route, exceptionHandlers: route.exceptionHandlers.push(recover)}47    })48}4950/// An exception action of one request kind, ordered among every exception action by registration.51export fn exceptionAction[Q, R, E](subject: &Request.Kind[Q, R, E], act: Request.ExceptionAction[Q, E]) -> Registration {52  requestChange(subject, 0, 1, Some(Order.ExceptionActions), fn(route: Request.Route[Q, R, E]) -> Request.Route[Q, R, E] {53      Request.Route{..route, exceptionActions: route.exceptionActions.push(act)}54    })55}5657/// A handler of a notification kind, reported to publishing strategies as `name`. A notification58/// kind takes any number, ordered among every notification handler by registration.59export fn notificationHandler[N, E](subject: &Notification.Kind[N, E], name: Str, handle: Notification.Handler[N, E]) -> Registration {60  let held = *subject61  let named = Notification.Named{name: name, handle: handle}62  Registration{apply: fn(catalog: Catalog.Catalog) -> Catalog.Catalog {63      let table = Catalog.touch(&catalog.notifications, &held.info, &held.routes, Notification.route(), fn(route: Notification.Route[N, E]) -> Notification.Route[N, E] {64          Notification.Route{handlers: route.handlers.push(named)}65        }, 1, 0, Compose.notificationFinish(held))66      Catalog.marked(Catalog.Catalog{..catalog, notifications: table}, Order.NotificationHandlers, held.info.name)67    } }68}6970/// The handler of a stream kind. A stream kind takes exactly one.71export fn streamHandler[Q, T, E](subject: &Stream.Kind[Q, T, E], handle: Stream.Handler[Q, T, E]) -> Registration {72  streamChange(subject, 1, 0, None, fn(route: Stream.Route[Q, T, E]) -> Stream.Route[Q, T, E] {73      Stream.Route{..route, handlers: route.handlers.push(handle)}74    })75}7677/// A behavior around the pipeline of one stream kind, ordered among every stream behavior by registration.78export fn streamBehavior[Q, T, E](subject: &Stream.Kind[Q, T, E], wrap: Stream.Behavior[Q, T, E]) -> Registration {79  streamChange(subject, 0, 1, Some(Order.StreamBehaviors), fn(route: Stream.Route[Q, T, E]) -> Stream.Route[Q, T, E] {80      Stream.Route{..route, behaviors: route.behaviors.push(wrap)}81    })82}8384/// A behavior around the pipeline of every request kind.85export fn openBehavior(open: dynamic Open.Behavior) -> Registration {86  Registration{apply: fn(catalog: Catalog.Catalog) -> Catalog.Catalog {87      let held = catalog.assembly88      let mark = Order.Open(Order.Behaviors, held.behaviors.length())89      Catalog.Catalog{..catalog, assembly: Catalog.Assembly{..held, behaviors: held.behaviors.push(open), marks: held.marks.push(mark)}}90    } }91}9293/// A pre-processor of every request kind.94export fn openPreProcessor(open: dynamic Open.PreProcessor) -> Registration {95  Registration{apply: fn(catalog: Catalog.Catalog) -> Catalog.Catalog {96      let held = catalog.assembly97      let mark = Order.Open(Order.PreProcessors, held.preProcessors.length())98      Catalog.Catalog{..catalog, assembly: Catalog.Assembly{..held, preProcessors: held.preProcessors.push(open), marks: held.marks.push(mark)}}99    } }100}101102/// A post-processor of every request kind.103export fn openPostProcessor(open: dynamic Open.PostProcessor) -> Registration {104  Registration{apply: fn(catalog: Catalog.Catalog) -> Catalog.Catalog {105      let held = catalog.assembly106      let mark = Order.Open(Order.PostProcessors, held.postProcessors.length())107      Catalog.Catalog{..catalog, assembly: Catalog.Assembly{..held, postProcessors: held.postProcessors.push(open), marks: held.marks.push(mark)}}108    } }109}110111/// An exception action of every request kind.112export fn openExceptionAction(open: dynamic Open.ExceptionAction) -> Registration {113  Registration{apply: fn(catalog: Catalog.Catalog) -> Catalog.Catalog {114      let held = catalog.assembly115      let mark = Order.Open(Order.ExceptionActions, held.exceptionActions.length())116      Catalog.Catalog{..catalog, assembly: Catalog.Assembly{..held, exceptionActions: held.exceptionActions.push(open), marks: held.marks.push(mark)}}117    } }118}119120/// A handler of every notification kind, reported to publishing strategies as `name`. An envelope121/// reaches it only when its kind is registered with the mediator.122export fn openNotificationHandler(name: Str, open: dynamic Open.NotificationHandler) -> Registration {123  let listener = Open.Listener{name: name, handler: open}124  Registration{apply: fn(catalog: Catalog.Catalog) -> Catalog.Catalog {125      let held = catalog.assembly126      let mark = Order.Open(Order.NotificationHandlers, held.listeners.length())127      Catalog.Catalog{..catalog, assembly: Catalog.Assembly{..held, listeners: held.listeners.push(listener), marks: held.marks.push(mark)}}128    } }129}130131/// A behavior around the pipeline of every stream kind.132export fn openStreamBehavior(open: dynamic Open.StreamBehavior) -> Registration {133  Registration{apply: fn(catalog: Catalog.Catalog) -> Catalog.Catalog {134      let held = catalog.assembly135      let mark = Order.Open(Order.StreamBehaviors, held.streamBehaviors.length())136      Catalog.Catalog{..catalog, assembly: Catalog.Assembly{..held, streamBehaviors: held.streamBehaviors.push(open), marks: held.marks.push(mark)}}137    } }138}139140/// Several registrations as one, applied in order.141export fn group(registrations: Array[Registration]) -> Registration {142  Registration{apply: fn(catalog: Catalog.Catalog) -> Catalog.Catalog {143      var changed = catalog144      for registration in registrations { changed = (registration.apply)(changed) }145      changed146    } }147}148149/// A change to one request kind's route that adds `handlers` and `components`, marked in `layer`150/// when its order among other kinds' components matters.151fn requestChange[Q, R, E](152  subject: &Request.Kind[Q, R, E],153  handlers: Int,154  components: Int,155  layer: Option[Order.Layer],156  change: fn(Request.Route[Q, R, E]) -> Request.Route[Q, R, E]157) -> Registration {158  let held = *subject159  Registration{apply: fn(catalog: Catalog.Catalog) -> Catalog.Catalog {160      let table = Catalog.touch(&catalog.requests, &held.info, &held.routes, Request.route(), change, handlers, components, Compose.requestFinish(held))161      let updated = Catalog.Catalog{..catalog, requests: table}162      match layer {163        case Some(at) => Catalog.marked(updated, at, held.info.name)164        case None => updated165      }166    } }167}168169/// A change to one stream kind's route, as `requestChange` is for requests.170fn streamChange[Q, T, E](171  subject: &Stream.Kind[Q, T, E],172  handlers: Int,173  components: Int,174  layer: Option[Order.Layer],175  change: fn(Stream.Route[Q, T, E]) -> Stream.Route[Q, T, E]176) -> Registration {177  let held = *subject178  Registration{apply: fn(catalog: Catalog.Catalog) -> Catalog.Catalog {179      let table = Catalog.touch(&catalog.streams, &held.info, &held.routes, Stream.route(), change, handlers, components, Compose.streamFinish(held))180      let updated = Catalog.Catalog{..catalog, streams: table}181      match layer {182        case Some(at) => Catalog.marked(updated, at, held.info.name)183        case None => updated184      }185    } }186}187