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

StreamTest.pudu

Pudu43 lines1.9 KB

GitHub ↗
1/** @Test.Stream.Suite — stream kinds, emitting, envelopes, items, and answers */2module PuduLangMediator.StreamTest34import Std.Io as Io5import Std.Sync as Sync6import Std.Test as Test7import PuduLangMediator.Message as Message8import PuduLangMediator as Messaging9import PuduLangMediator.Stream as Stream10import PuduLangMediator.Utils.Erasure as Erasure1112/// Runs the suite.13fn main() -> Int {14  let pages: Stream.Kind[Int, Str, Str] = Stream.kind("pages")15  let seen: Sync.Cell[Array[Str]] = Sync.cell([])16  let emitted: Messaging.Outcome[(), Str] = Stream.each(&["a", "b", "c", "d"], fn(item: Str) -> Bool {17      let held = match Sync.get(&seen) {18        case Ok(found) => found19        case Err(_) => []20      }21      let _stored = Sync.set(&seen, held.push(item))22      item != "b"23    })24  let item = Erasure.pack(&pages.items, "page 1")25  let sealed = Stream.envelope(&pages, 3)26  let answered = Message.answered(&pages.info, &pages.outcomes, Ok(()))27  let checks = Test.suite("Stream", &[28      Test.equals("a stream kind is a sequence", &pages.info.shape, &Message.Sequence),29      Test.equals("tags are carried", &Stream.tagged(&pages, ["paged"]).info.tags, &["paged"]),30      Test.equals("a fresh route is empty", &(Stream.route().handlers.length(), Stream.route().behaviors.length()), &(0, 0)),31      Test.equals("emitting stops after the sink declines", &Sync.get(&seen), &Ok(["a", "b"])),32      Test.equals("a declined stream still succeeds", &emitted, &Ok(())),33      Test.equals("an item reads back through its kind", &Stream.itemOf(&pages, &item), &Some("page 1")),34      Test.equals("an envelope opens for its kind", &Erasure.unpack(&pages.requests, &sealed.packed), &Some(3)),35      Test.equals("the kind reads its own answer", &Stream.answerOf(&pages, &answered), &Some(Ok(())))36    ])37  let ran = Test.run(&checks)38  for failure in Test.failuresOf(&ran) {39    let _reported = Io.writeErrorLine(failure)40  }41  Test.report(&ran)42}43