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

Compose.pudu

Pudu248 lines12.4 KB

GitHub ↗
1/** @Mediator.Compose.Builder — one kind's pipeline from its route and the open components */2module PuduLangMediator.Compose34import Std.List as List5import Std.Option as Option6import PuduLangMediator.Catalog as Catalog7import PuduLangMediator.Context as Context8import PuduLangMediator.Domain.Order as Order9import PuduLangMediator.Message as Message10import PuduLangMediator as Messaging11import PuduLangMediator.Notification as Notification12import PuduLangMediator.Open as Open13import PuduLangMediator.Publishing as Publishing14import PuduLangMediator.Request as Request15import PuduLangMediator.Stream as Stream16import PuduLangMediator.Utils.Erasure as Erasure1718/// A request kind's handler wrapped, outermost first, in its exception layers (actions outside19/// handlers for `ForUnhandled`, inside them for `ForAll`), its pre-processors, its20/// post-processors, and its behaviors in registration order. A layer with no components is left out.21export fn requestHandler[Q, R, E](info: &Message.Info, route: &Request.Route[Q, R, E], assembly: &Catalog.Assembly) -> Request.Handler[Q, R, E] {22  let seen = *info23  let marks = assembly.marks24  var actions: Array[Request.ExceptionAction[Q, E]] = []25  for pick in Order.picks(&marks, Order.ExceptionActions, seen.name) {26    actions = actions.push(match pick {27        case Order.Shared(index) => openAction(assembly.exceptionActions[index], seen)28        case Order.Own(index) => route.exceptionActions[index]29      })30  }31  var pres: Array[Request.PreProcessor[Q, E]] = []32  for pick in Order.picks(&marks, Order.PreProcessors, seen.name) {33    pres = pres.push(match pick {34        case Order.Shared(index) => openPre(assembly.preProcessors[index], seen)35        case Order.Own(index) => route.preProcessors[index]36      })37  }38  var posts: Array[Request.PostProcessor[Q, R, E]] = []39  for pick in Order.picks(&marks, Order.PostProcessors, seen.name) {40    posts = posts.push(match pick {41        case Order.Shared(index) => openPost(assembly.postProcessors[index], seen)42        case Order.Own(index) => route.postProcessors[index]43      })44  }45  var observing: Array[Request.Behavior[Q, R, E]] = []46  if !actions.isEmpty() { observing = observing.push(observed(actions)) }47  var recovering: Array[Request.Behavior[Q, R, E]] = []48  if !route.exceptionHandlers.isEmpty() { recovering = recovering.push(recovered(route.exceptionHandlers)) }49  var layers = match assembly.actionScope {50    case Catalog.ForUnhandled => observing.concat(recovering)51    case Catalog.ForAll => recovering.concat(observing)52  }53  if !pres.isEmpty() { layers = layers.push(preprocessed(pres)) }54  if !posts.isEmpty() { layers = layers.push(postprocessed(posts)) }55  for pick in Order.picks(&marks, Order.Behaviors, seen.name) {56    layers = layers.push(match pick {57        case Order.Shared(index) => openBehavior(assembly.behaviors[index], seen)58        case Order.Own(index) => route.behaviors[index]59      })60  }61  wrap(layers, Option.unwrapOr(List.first(&route.handlers), unhandled(seen.name)))62}6364/// A notification kind's publication: its own handlers and the open ones in registration order,65/// run by the assembly's publishing strategy.66export fn publication[N, E](info: &Message.Info, route: &Notification.Route[N, E], assembly: &Catalog.Assembly) -> Notification.Handler[N, E] {67  let seen = *info68  var handlers: Array[Notification.Named[N, E]] = []69  for pick in Order.picks(&assembly.marks, Order.NotificationHandlers, seen.name) {70    handlers = handlers.push(match pick {71        case Order.Shared(index) => listened(assembly.listeners[index], seen)72        case Order.Own(index) => route.handlers[index]73      })74  }75  let publisher = assembly.publisher76  fn(notification: N, context: Context.Context) -> Messaging.Outcome[(), E] {77    let executors = handlers.map(fn(named: Notification.Named[N, E]) -> Publishing.Executor[E] {78        Publishing.Executor{handler: named.name, run: fn(passed: Context.Context) -> Messaging.Outcome[(), E] { (named.handle)(notification, passed) }}79      })80    publisher.publish(&executors, context)81  }82}8384/// A stream kind's handler wrapped in its behaviors, open and its own, in registration order.85export fn streamHandler[Q, T, E](info: &Message.Info, route: &Stream.Route[Q, T, E], assembly: &Catalog.Assembly) -> Stream.Handler[Q, T, E] {86  let seen = *info87  var layers: Array[Stream.Behavior[Q, T, E]] = []88  for pick in Order.picks(&assembly.marks, Order.StreamBehaviors, seen.name) {89    layers = layers.push(match pick {90        case Order.Shared(index) => openStream(assembly.streamBehaviors[index], seen)91        case Order.Own(index) => route.behaviors[index]92      })93  }94  var composed = Option.unwrapOr(List.first(&route.handlers), unhandledStream(seen.name))95  var index = layers.length() - 196  while index >= 0 {97    let layer = layers[index]98    let inner = composed99    composed = fn(request: Q, context: Context.Context, sink: Stream.Sink[T]) -> Messaging.Outcome[(), E] {100      layer(request, context, sink, fn(passed: Context.Context, passing: Stream.Sink[T]) -> Messaging.Outcome[(), E] { inner(request, passed, passing) })101    }102    index = index - 1103  }104  composed105}106107/// How a request kind's entry is composed once the catalog is complete.108export fn requestFinish[Q, R, E](subject: Request.Kind[Q, R, E]) -> fn(Erasure.Packed, Catalog.Assembly) -> Catalog.Compiled {109  fn(packed: Erasure.Packed, assembly: Catalog.Assembly) -> Catalog.Compiled {110    let info = subject.info111    let invoke = requestHandler(&info, &Option.unwrapOr(Erasure.unpack(&subject.routes, &packed), Request.route()), &assembly)112    Catalog.Compiled{info: info, invoke: Erasure.pack(&subject.invokers, invoke), dispatch: Catalog.Answering(fn(envelope: Message.Envelope, context: Context.Context) -> Message.Answer {113          match Erasure.unpack(&subject.requests, &envelope.packed) {114            case Some(found) => Message.answered(&info, &subject.outcomes, invoke(found, context))115            case None => Message.undelivered(&info, Message.WrongKind(info.name))116          }117        }) }118  }119}120121/// How a notification kind's entry is composed once the catalog is complete.122export fn notificationFinish[N, E](subject: Notification.Kind[N, E]) -> fn(Erasure.Packed, Catalog.Assembly) -> Catalog.Compiled {123  fn(packed: Erasure.Packed, assembly: Catalog.Assembly) -> Catalog.Compiled {124    let info = subject.info125    let invoke = publication(&info, &Option.unwrapOr(Erasure.unpack(&subject.routes, &packed), Notification.route()), &assembly)126    Catalog.Compiled{info: info, invoke: Erasure.pack(&subject.invokers, invoke), dispatch: Catalog.Answering(fn(envelope: Message.Envelope, context: Context.Context) -> Message.Answer {127          match Erasure.unpack(&subject.notifications, &envelope.packed) {128            case Some(found) => Message.answered(&info, &subject.outcomes, invoke(found, context))129            case None => Message.undelivered(&info, Message.WrongKind(info.name))130          }131        }) }132  }133}134135/// How a stream kind's entry is composed once the catalog is complete.136export fn streamFinish[Q, T, E](subject: Stream.Kind[Q, T, E]) -> fn(Erasure.Packed, Catalog.Assembly) -> Catalog.Compiled {137  fn(packed: Erasure.Packed, assembly: Catalog.Assembly) -> Catalog.Compiled {138    let info = subject.info139    let items = subject.items140    let invoke = streamHandler(&info, &Option.unwrapOr(Erasure.unpack(&subject.routes, &packed), Stream.route()), &assembly)141    Catalog.Compiled{info: info, invoke: Erasure.pack(&subject.invokers, invoke), dispatch: Catalog.Streaming(fn(envelope: Message.Envelope, context: Context.Context, sink: fn(Erasure.Packed) -> Bool) -> Message.Answer {142          match Erasure.unpack(&subject.requests, &envelope.packed) {143            case Some(found) => Message.answered(&info, &subject.outcomes, invoke(found, context, fn(item: T) -> Bool { sink(Erasure.pack(&items, item)) }))144            case None => Message.undelivered(&info, Message.WrongKind(info.name))145          }146        }) }147  }148}149150/// The handler wrapped in `layers`, the first outermost.151fn wrap[Q, R, E](layers: Array[Request.Behavior[Q, R, E]], handler: Request.Handler[Q, R, E]) -> Request.Handler[Q, R, E] {152  var composed = handler153  var index = layers.length() - 1154  while index >= 0 {155    let layer = layers[index]156    let inner = composed157    composed = fn(request: Q, context: Context.Context) -> Messaging.Outcome[R, E] {158      layer(request, context, fn(passed: Context.Context) -> Messaging.Outcome[R, E] { inner(request, passed) })159    }160    index = index - 1161  }162  composed163}164165/// Runs every pre-processor in order, then the rest of the pipeline; a failure stops the request.166fn preprocessed[Q, R, E](processors: Array[Request.PreProcessor[Q, E]]) -> Request.Behavior[Q, R, E] {167  fn(request: Q, context: Context.Context, next: Request.Next[R, E]) -> Messaging.Outcome[R, E] {168    for processor in processors { processor(request, context) ? }169    next(context)170  }171}172173/// Runs the rest of the pipeline, then every post-processor in order on a value it answered.174fn postprocessed[Q, R, E](processors: Array[Request.PostProcessor[Q, R, E]]) -> Request.Behavior[Q, R, E] {175  fn(request: Q, context: Context.Context, next: Request.Next[R, E]) -> Messaging.Outcome[R, E] {176    let response = next(context) ?177    for processor in processors { processor(request, response, context) ? }178    Ok(response)179  }180}181182/// Answers the first value an exception handler gives in place of a failure, in order.183fn recovered[Q, R, E](handlers: Array[Request.ExceptionHandler[Q, R, E]]) -> Request.Behavior[Q, R, E] {184  fn(request: Q, context: Context.Context, next: Request.Next[R, E]) -> Messaging.Outcome[R, E] {185    let outcome = next(context)186    if let Err(failure) = outcome {187      for handler in handlers {188        if let Some(value) = handler(request, failure, context) { return Ok(value) }189      }190    }191    outcome192  }193}194195/// Hands a failure to every exception action in order and answers it unchanged.196fn observed[Q, R, E](actions: Array[Request.ExceptionAction[Q, E]]) -> Request.Behavior[Q, R, E] {197  fn(request: Q, context: Context.Context, next: Request.Next[R, E]) -> Messaging.Outcome[R, E] {198    let outcome = next(context)199    if let Err(failure) = outcome {200      for action in actions { action(request, failure, context) }201    }202    outcome203  }204}205206/// An open behavior applied to one kind.207fn openBehavior[Q, R, E](open: dynamic Open.Behavior, info: Message.Info) -> Request.Behavior[Q, R, E] {208  fn(request: Q, context: Context.Context, next: Request.Next[R, E]) -> Messaging.Outcome[R, E] { open.handleRequest(&request, &info, context, next) }209}210211/// An open pre-processor applied to one kind.212fn openPre[Q, E](open: dynamic Open.PreProcessor, info: Message.Info) -> Request.PreProcessor[Q, E] {213  fn(request: Q, context: Context.Context) -> Messaging.Outcome[(), E] { open.beforeRequest(&request, &info, context) }214}215216/// An open post-processor applied to one kind.217fn openPost[Q, R, E](open: dynamic Open.PostProcessor, info: Message.Info) -> Request.PostProcessor[Q, R, E] {218  fn(request: Q, response: R, context: Context.Context) -> Messaging.Outcome[(), E] { open.afterRequest(&request, &response, &info, context) }219}220221/// An open exception action applied to one kind.222fn openAction[Q, E](open: dynamic Open.ExceptionAction, info: Message.Info) -> Request.ExceptionAction[Q, E] {223  fn(request: Q, failure: Messaging.Failure[E], context: Context.Context) -> () { open.onFailure(&request, &failure, &info, context) }224}225226/// An open notification handler applied to one kind.227fn listened[N, E](listener: Open.Listener, info: Message.Info) -> Notification.Named[N, E] {228  let open = listener.handler229  Notification.Named{name: listener.name, handle: fn(notification: N, context: Context.Context) -> Messaging.Outcome[(), E] { open.handleNotification(&notification, &info, context) }}230}231232/// An open stream behavior applied to one kind.233fn openStream[Q, T, E](open: dynamic Open.StreamBehavior, info: Message.Info) -> Stream.Behavior[Q, T, E] {234  fn(request: Q, context: Context.Context, sink: Stream.Sink[T], next: Stream.Next[T, E]) -> Messaging.Outcome[(), E] {235    open.handleStream(&request, &info, context, sink, next)236  }237}238239/// A handler answering `Unhandled`, for a kind whose route has none.240fn unhandled[Q, R, E](name: Str) -> Request.Handler[Q, R, E] {241  fn(_request: Q, _context: Context.Context) -> Messaging.Outcome[R, E] { Err(Messaging.Unhandled(name)) }242}243244/// A stream handler answering `Unhandled`, for a kind whose route has none.245fn unhandledStream[Q, T, E](name: Str) -> Stream.Handler[Q, T, E] {246  fn(_request: Q, _context: Context.Context, _sink: Stream.Sink[T]) -> Messaging.Outcome[(), E] { Err(Messaging.Unhandled(name)) }247}248