
Stream.pudu
Pudu70 lines2.9 KB
1/** @Mediator.Stream.Kind — requests answered by a sequence of items */2module PuduLangMediator.Stream34import PuduLangMediator.Context as Context5import PuduLangMediator.Message as Message6import PuduLangMediator as Messaging7import PuduLangMediator.Utils.Erasure as Erasure89/** @Mediator.Stream.Sink — takes one item and answers whether it wants another */10export type Sink[T] = fn(T) -> Bool1112/** @Mediator.Stream.Next — the rest of the stream pipeline */13export type Next[T, E] = fn(Context.Context, Sink[T]) -> Messaging.Outcome[(), E]1415/** @Mediator.Stream.Handler — emits the items of one stream request */16export type Handler[Q, T, E] = fn(Q, Context.Context, Sink[T]) -> Messaging.Outcome[(), E]1718/** @Mediator.Stream.Behavior — wraps the rest of the stream pipeline for one kind */19export type Behavior[Q, T, E] = fn(Q, Context.Context, Sink[T], Next[T, E]) -> Messaging.Outcome[(), E]2021/** @Mediator.Stream.Route — every component registered for one kind */22export type Route[Q, T, E] = { handlers: Array[Handler[Q, T, E]], behaviors: Array[Behavior[Q, T, E]] }2324/** @Mediator.Stream.Kind — the typed identity of one stream request */25export type Kind[Q, T, E] = {26 info: Message.Info,27 routes: Erasure.Slot[Route[Q, T, E]],28 invokers: Erasure.Slot[Handler[Q, T, E]],29 requests: Erasure.Slot[Q],30 items: Erasure.Slot[T],31 outcomes: Erasure.Slot[Messaging.Outcome[(), E]]32}333435export fn kind[Q, T, E](name: Str) -> Kind[Q, T, E] {36 Kind {37 info: Message.Info{name: name, shape: Message.Sequence, tags: []},38 routes: Erasure.slot(),39 invokers: Erasure.slot(),40 requests: Erasure.slot(),41 items: Erasure.slot(),42 outcomes: Erasure.slot()43 }44}454647export fn tagged[Q, T, E](subject: &Kind[Q, T, E], tags: Array[Str]) -> Kind[Q, T, E] {48 Kind{..*subject, info: Message.Info{..subject.info, tags: subject.info.tags.concat(tags)}}49}505152export fn route[Q, T, E]() -> Route[Q, T, E] { Route{handlers: [], behaviors: []} }535455export fn each[T, E](items: &Array[T], sink: Sink[T]) -> Messaging.Outcome[(), E] {56 for item in items {57 if !sink(item) { return Ok(()) }58 }59 Ok(())60}616263export fn envelope[Q, T, E](subject: &Kind[Q, T, E], request: Q) -> Message.Envelope { Message.sealed(&subject.info, &subject.requests, request) }646566export fn itemOf[Q, T, E](subject: &Kind[Q, T, E], packed: &Erasure.Packed) -> Option[T] { Erasure.unpack(&subject.items, packed) }676869export fn answerOf[Q, T, E](subject: &Kind[Q, T, E], answer: &Message.Answer) -> Option[Messaging.Outcome[(), E]] { Message.read(&subject.outcomes, answer) }70