
ReaderTest.pudu
Pudu87 lines3.5 KB
1/** @Test.Reader.Suite — pulled streams, back-pressure, failures, crashes, and closing */2module PuduLangMediator.ReaderTest34import Std.Concurrent as Concurrent5import Std.Io as Io6import Std.Test as Test7import PuduLangMediator.Context as Context8import PuduLangMediator as Messaging9import PuduLangMediator.Reader as Reader10import PuduLangMediator.Stream as Stream11import PuduLangMediator.Utils.Shared as Shared121314fn numbers(count: Int, produced: Shared.Shared[Int]) -> fn(Context.Context, Stream.Sink[Int]) -> Messaging.Outcome[(), Str] {15 fn(_context: Context.Context, sink: Stream.Sink[Int]) -> Messaging.Outcome[(), Str] {16 var next = 117 while next <= count {18 if !sink(next) { return Ok(()) }19 Shared.update(&produced, |held: Int| held + 1)20 next = next + 121 }22 Ok(())23 }24}252627fn main() -> Int {28 let wholeCount = Shared.shared(0)29 let whole = Reader.open(2, &Context.create(), numbers(5, wholeCount))30 let first = Reader.next(&whole)31 let others = Reader.rest(&whole)32 let afterEnd = Reader.next(&whole)33 Reader.close(&whole)3435 let pressureCount = Shared.shared(0)36 let pressured = Reader.open(1, &Context.create(), numbers(1000, pressureCount))37 let _taken = Reader.next(&pressured)38 let _settled = Concurrent.sleep(100)39 let producedAhead = Shared.current(&pressureCount)40 Reader.close(&pressured)41 let afterClose = Reader.next(&pressured)4243 let failing = Reader.open(4, &Context.create(), fn(_context: Context.Context, sink: Stream.Sink[Int]) -> Messaging.Outcome[(), Str] {44 let _first = sink(1)45 Messaging.raise("disk full")46 })47 let failingItem = Reader.next(&failing)48 let failure = Reader.next(&failing)49 let afterFailure = Reader.next(&failing)50 Reader.close(&failing)5152 let crashing = Reader.open(4, &Context.create(), fn(_context: Context.Context, _sink: Stream.Sink[Int]) -> Messaging.Outcome[(), Str] { panic("producer crashed") })53 let crashed = Reader.rest(&crashing)54 Reader.close(&crashing)5556 let parent = Context.create()57 let observed = Shared.shared(false)58 let watching = Reader.open(1, &parent, fn(context: Context.Context, sink: Stream.Sink[Int]) -> Messaging.Outcome[(), Str] {59 let _first = sink(1)60 Context.pause(&context, 60000) ?61 Shared.update(&observed, |_held: Bool| true)62 Ok(())63 })64 let _one = Reader.next(&watching)65 Context.cancel(&parent, "caller left")66 let cancelled = Reader.next(&watching)67 Reader.close(&watching)6869 let checks = Test.suite("Reader", &[70 Test.equals("the first item is pulled", &first, &Some(Ok(1))),71 Test.equals("the rest follow in order", &others, &Ok([2, 3, 4, 5])),72 Test.equals("a finished reader answers nothing", &afterEnd, &None),73 Test.that("a full reader holds the producer back", producedAhead <= 3),74 Test.equals("a closed reader answers nothing", &afterClose, &None),75 Test.equals("items before a failure are delivered", &failingItem, &Some(Ok(1))),76 Test.equals("the failure is delivered once", &(failure, afterFailure), &(Some(Err(Messaging.Raised("disk full"))), None)),77 Test.equals("a crashing producer is reported, not hung on", &crashed, &Err(Messaging.Crashed("producer crashed"))),78 Test.equals("cancelling the caller's context stops the producer", &cancelled, &Some(Err(Messaging.Cancelled("cancelled: caller left")))),79 Test.not("the producer never ran past its cancellation", Shared.current(&observed))80 ])81 let ran = Test.run(&checks)82 for failure in Test.failuresOf(&ran) {83 let _reported = Io.writeErrorLine(failure)84 }85 Test.report(&ran)86}87