
Compose.pudu
Pudu248 lines12.4 KB
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 Erasure1718192021export 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}63646566export 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}838485export 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}106107108export 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}120121122export 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}134135136export 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}149150151fn 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}164165166fn 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}172173174fn 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}181182183fn 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}194195196fn 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}205206207fn 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}210211212fn 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}215216217fn 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}220221222fn 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}225226227fn 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(¬ification, &info, context) }}230}231232233fn 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}238239240fn 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}243244245fn 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