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

PublishingTest.pudu

Pudu106 lines5.5 KB

GitHub ↗
1/** @Test.Publishing.Suite — sequential, continuing, parallel, and bounded publication */2module PuduLangMediator.PublishingTest34import Std.Concurrent as Concurrent5import Std.Io as Io6import Std.Math as Math7import Std.Test as Test8import PuduLangMediator.Context as Context9import PuduLangMediator as Messaging10import PuduLangMediator.Publishing as Publishing11import PuduLangMediator.Utils.Shared as Shared1213/** @Test.Publishing.Gauge — how many executors run now, and the most that ever did */14type Gauge = { running: Shared.Shared[Int], peak: Shared.Shared[Int], ran: Shared.Shared[Array[Str]] }1516/// A fresh gauge.17fn gauge() -> Gauge { Gauge{running: Shared.shared(0), peak: Shared.shared(0), ran: Shared.shared([])} }1819/// An executor that records itself, holds for `millis`, and answers `outcome`.20fn executor(meter: &Gauge, name: Str, millis: Int, outcome: Messaging.Outcome[(), Str]) -> Publishing.Executor[Str] {21  let held = *meter22  Publishing.Executor{handler: name, run: fn(_context: Context.Context) -> Messaging.Outcome[(), Str] {23      let now = Shared.change(&held.running, fn(count: Int) -> (Int, Int) { (count + 1, count + 1) })24      Shared.update(&held.peak, |peak: Int| Math.max(peak, now))25      let _held = Concurrent.sleep(millis)26      Shared.update(&held.ran, |names: Array[Str]| names.push(name))27      Shared.update(&held.running, |count: Int| count - 1)28      outcome29    } }30}3132/// An executor that stops its thread.33fn crashing(name: Str) -> Publishing.Executor[Str] {34  Publishing.Executor{handler: name, run: fn(_context: Context.Context) -> Messaging.Outcome[(), Str] { panic("handler " + name + " crashed") }}35}3637/// Runs the suite.38fn main() -> Int {39  let sequential = gauge()40  let stopped = Publishing.Sequential{}.publish(&[41      executor(&sequential, "a", 0, Ok(())),42      executor(&sequential, "b", 0, Messaging.raise("b failed")),43      executor(&sequential, "c", 0, Ok(()))44    ], Context.create())45  let cancelledContext = Context.create()46  Context.cancel(&cancelledContext, "shutdown")47  let cancelledGauge = gauge()48  let cancelled = Publishing.Sequential{}.publish(&[executor(&cancelledGauge, "a", 0, Ok(()))], cancelledContext)49  let continuing = gauge()50  let continued = Publishing.Continuing{}.publish(&[51      executor(&continuing, "a", 0, Messaging.raise("a failed")),52      executor(&continuing, "b", 0, Ok(())),53      executor(&continuing, "c", 0, Messaging.raise("c failed"))54    ], Context.create())55  let parallel = gauge()56  let together = Publishing.Parallel{}.publish(&[57      executor(&parallel, "a", 60, Ok(())),58      executor(&parallel, "b", 60, Messaging.raise("b failed")),59      executor(&parallel, "c", 60, Ok(())),60      crashing("d")61    ], Context.create())62  let bounded = gauge()63  let limited = Publishing.Bounded{workers: 2}.publish(&[64      executor(&bounded, "a", 30, Ok(())),65      executor(&bounded, "b", 30, Ok(())),66      executor(&bounded, "c", 30, Messaging.raise("c failed")),67      executor(&bounded, "d", 30, Ok(())),68      crashing("e")69    ], Context.create())70  let atLeastOne = gauge()71  let single = Publishing.Bounded{workers: 0}.publish(&[executor(&atLeastOne, "a", 0, Ok(())), executor(&atLeastOne, "b", 0, Ok(()))], Context.create())72  let empty: Messaging.Outcome[(), Str] = Publishing.Parallel{}.publish(&[], Context.create())73  let strategies: Array[dynamic Publishing.Strategy] = [Publishing.Sequential{}, Publishing.Continuing{}, Publishing.Parallel{}, Publishing.Bounded{workers: 4}]74  var allSucceed = true75  for strategy in strategies {76    let meter = gauge()77    if strategy.publish(&[executor(&meter, "a", 0, Ok(())), executor(&meter, "b", 0, Ok(()))], Context.create()) != Ok(()) { allSucceed = false }78  }79  let checks = Test.suite("Publishing", &[80      Test.equals("sequential stops at the first failure", &stopped, &Err(Messaging.Raised("b failed"))),81      Test.equals("handlers after a failure never run", &Shared.current(&sequential.ran), &["a", "b"]),82      Test.equals("sequential does not start on a cancelled context", &(cancelled, Shared.current(&cancelledGauge.ran)), &(Err(Messaging.Cancelled("cancelled: shutdown")), [])),83      Test.equals("continuing runs every handler and answers every failure", &(continued, Shared.current(&continuing.ran)), &(84          Err(Messaging.Aggregate([Messaging.Raised("a failed"), Messaging.Raised("c failed")])), ["a", "b", "c"]85        )),86      Test.equals("parallel answers every failure in handler order, a crash included", &together, &Err(87          Messaging.Aggregate([Messaging.Raised("b failed"), Messaging.Crashed("handler d crashed")])88        )),89      Test.that("parallel runs the handlers at once", Shared.current(&parallel.peak) >= 2),90      Test.equals("parallel waits for every handler", &Shared.current(&parallel.ran).length(), &3),91      Test.equals("bounded answers every failure, a crash included", &limited, &Err(92          Messaging.Aggregate([Messaging.Raised("c failed"), Messaging.Crashed("handler e crashed")])93        )),94      Test.equals("bounded never runs more than its workers", &Shared.current(&bounded.peak), &2),95      Test.equals("bounded runs every handler", &Shared.current(&bounded.ran).length(), &4),96      Test.equals("bounded with no workers still runs one at a time", &(single, Shared.current(&atLeastOne.peak)), &(Ok(()), 1)),97      Test.equals("nothing to run succeeds", &empty, &Ok(())),98      Test.that("every strategy succeeds when every handler does", allSucceed)99    ])100  let ran = Test.run(&checks)101  for failure in Test.failuresOf(&ran) {102    let _reported = Io.writeErrorLine(failure)103  }104  Test.report(&ran)105}106