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

Stream.pudu

Pudu70 lines2.9 KB

GitHub ↗
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}3334/// A kind of stream request emitting `T` or failing with `E`. Make each kind once and share it.35export 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}4546/// The same kind carrying `tags` as well, which open components may read.47export 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}5051/// A route with no components.52export fn route[Q, T, E]() -> Route[Q, T, E] { Route{handlers: [], behaviors: []} }5354/// Hands each item to the sink in order until it declines one.55export 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}6162/// An envelope carrying `request`, for a sender that does not name its type.63export fn envelope[Q, T, E](subject: &Kind[Q, T, E], request: Q) -> Message.Envelope { Message.sealed(&subject.info, &subject.requests, request) }6465/// An item an untyped stream emitted, when this kind emitted it.66export fn itemOf[Q, T, E](subject: &Kind[Q, T, E], packed: &Erasure.Packed) -> Option[T] { Erasure.unpack(&subject.items, packed) }6768/// The outcome an answer holds, when this kind answered it.69export fn answerOf[Q, T, E](subject: &Kind[Q, T, E], answer: &Message.Answer) -> Option[Messaging.Outcome[(), E]] { Message.read(&subject.outcomes, answer) }70