
ObservableTest.pudu
Pudu68 lines2.6 KB
1/** @Test.Sinks.Observable.Suite — observers subscribing to the event stream */2module PuduLangLog.Sinks.ObservableTest34import Std.Io as Io5import Std.Sync as Sync6import Std.Test as Test7import PuduLangLog.Configuration as Configuration8import PuduLangLog as Log9import PuduLangLog.Logger as Logger10import PuduLangLog.Sinks.Observable as Observable11import PuduLangLog.Value as Value121314fn recording(journal: &Sync.Cell[Array[Str]], name: Str) -> Observable.Observer {15 let shared = *journal16 Observable.Observer {17 next: fn(event: &Log.Event) -> () { note(&shared, name + ":" + event.template.text) },18 completed: fn() -> () { note(&shared, name + ":done") }19 }20}212223fn note(journal: &Sync.Cell[Array[Str]], line: Str) -> () {24 let _stored = Sync.set(journal, linesOf(journal).push(line))25}262728fn linesOf(journal: &Sync.Cell[Array[Str]]) -> Array[Str] {29 match Sync.get(journal) {30 case Ok(lines) => lines31 case Err(_) => []32 }33}343536fn main() -> Int {37 let none: Array[Str] = []38 let journal = Sync.cell(none)39 let subject = Observable.create()40 let logger = Configuration.create().writeTo(Observable.sink(&subject)).createLogger()41 logger.information("before", [])42 let first = Observable.subscribe(&subject, recording(&journal, "a"))43 let _second = Observable.subscribe(&subject, recording(&journal, "b"))44 let counted = Observable.count(&subject)45 logger.information("both \{N\}", [Value.int(1)])46 Observable.unsubscribe(&subject, first)47 Observable.unsubscribe(&subject, first)48 logger.information("only b", [])49 let seen = Sync.cell(0)50 let plain = Observable.subscribe(&subject, Observable.observer(fn(_event: &Log.Event) -> () { let _counted = Sync.set(&seen, 1) }))51 logger.information("plain", [])52 Observable.unsubscribe(&subject, plain)53 Logger.close(&logger)54 let _late = Observable.subscribe(&subject, recording(&journal, "late"))55 logger.information("after close", [])56 let checks = Test.suite("Sinks.Observable", &[57 Test.equals("subscribers are counted", &counted, &2),58 Test.equals("events reach current observers in subscription order", &linesOf(&journal), &["a:both \{N\}", "b:both \{N\}", "b:only b", "b:plain", "b:done", "late:done"]),59 Test.equals("closing removes every observer", &Observable.count(&subject), &0),60 Test.equals("an observer built from a function ignores the end", &Sync.get(&seen), &Ok(1))61 ])62 let ran = Test.run(&checks)63 for failure in Test.failuresOf(&ran) {64 let _reported = Io.writeErrorLine(failure)65 }66 Test.report(&ran)67}68