
PublishingTest.pudu
Pudu106 lines5.5 KB
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]] }151617fn gauge() -> Gauge { Gauge{running: Shared.shared(0), peak: Shared.shared(0), ran: Shared.shared([])} }181920fn 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}313233fn crashing(name: Str) -> Publishing.Executor[Str] {34 Publishing.Executor{handler: name, run: fn(_context: Context.Context) -> Messaging.Outcome[(), Str] { panic("handler " + name + " crashed") }}35}363738fn 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(¶llel, "a", 60, Ok(())),58 executor(¶llel, "b", 60, Messaging.raise("b failed")),59 executor(¶llel, "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(¶llel.peak) >= 2),90 Test.equals("parallel waits for every handler", &Shared.current(¶llel.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