Pudu programming language
Menu
Package

@chrismichaelps / pudu-lang-log

Structured event logging for Pudu: message templates, enrichment, filtering, formatting, and sinks

0.1.0Apache-2.01

InstallClose

AsyncTest.pudu

Pudu71 lines3.1 KB

GitHub ↗
1/** @Test.Sinks.Async.Suite — queued writes on a worker thread */2module PuduLangLog.Sinks.AsyncTest34import Std.Concurrent as Concurrent5import Std.Io as Io6import Std.Sync as Sync7import Std.Test as Test8import PuduLangLog.Event as Event9import PuduLangLog as Log10import PuduLangLog.Sink as Sink11import PuduLangLog.Sinks.Async as Async12import PuduLangLog.Sinks.Memory as Memory1314/// An event carrying a message.15fn event(message: Str) -> Log.Event { Event.create(Log.Timestamp{millis: 0, offset: 0}, Log.Information, message, []) }1617/// Runs the suite.18fn main() -> Int {19  let memory = Memory.create()20  let queue = Async.start(Memory.sink(&memory), Async.defaults())21  let sink = Async.sinkOf(&queue)22  for number in [1, 2, 3, 4, 5] { let _queued = (sink.emit)(&event(show(number))) }23  (sink.flush)()24  let flushed = Memory.messages(&memory)25  (sink.close)()26  let refused = (sink.emit)(&event("late"))2728  let gate = Sync.mutex()29  let _held = Sync.lock(&gate)30  let blocked = Memory.create()31  let slow = Sink.of(fn(written: &Log.Event) -> Result[(), Str] {32      let _waited = Sync.withLock(&gate, fn() -> () {})33      (Memory.sink(&blocked).emit)(written)34    })35  let small = Async.start(slow, Async.Options{..Async.defaults(), bufferSize: 2})36  let smallSink = Async.sinkOf(&small)37  for number in [1, 2, 3, 4, 5, 6] { let _queued = (smallSink.emit)(&event(show(number))) }38  let droppedWhileFull = Async.dropped(&small)39  let _released = Sync.unlock(&gate)40  (smallSink.close)()4142  let failures = Sync.counter(0)43  let failing = Async.start(Sink.of(fn(_written: &Log.Event) -> Result[(), Str] { Err("down") }), Async.defaults())44  let failingSink = Async.sinkOf(&failing)45  (failingSink.attach)(fn(report: &Sink.Report) -> () { let _counted = Sync.increment(&failures, report.events.length()) })46  let _sent = (failingSink.emit)(&event("x"))47  (failingSink.close)()4849  let blocking = Memory.create()50  let blockingQueue = Async.start(Memory.sink(&blocking), Async.Options{..Async.defaults(), bufferSize: 1, blockWhenFull: true})51  let blockingSink = Async.sinkOf(&blockingQueue)52  for number in [1, 2, 3, 4, 5, 6, 7, 8] { let _queued = (blockingSink.emit)(&event(show(number))) }53  (blockingSink.close)()54  let _paused = Concurrent.sleep(1)5556  let checks = Test.suite("Sinks.Async", &[57      Test.equals("flushing waits for every queued event", &flushed, &["1", "2", "3", "4", "5"]),58      Test.equals("a closed queue refuses events", &refused, &Err("the asynchronous sink is closed")),59      Test.that("a full queue drops events rather than waiting", droppedWhileFull >= 3),60      Test.equals("what was queued is written after the queue frees", &(Memory.events(&blocked).length() + droppedWhileFull), &6),61      Test.equals("failures reach the attached listener", &Sync.count(&failures), &Ok(1)),62      Test.equals("a blocking queue loses nothing", &Memory.messages(&blocking), &["1", "2", "3", "4", "5", "6", "7", "8"]),63      Test.equals("nothing is pending after closing", &Async.pending(&blockingQueue), &0)64    ])65  let ran = Test.run(&checks)66  for failure in Test.failuresOf(&ran) {67    let _reported = Io.writeErrorLine(failure)68  }69  Test.report(&ran)70}71