
BatchingTest.pudu
Pudu124 lines5.2 KB
1/** @Test.Sinks.Batching.Suite — batches, timing, retries, and limits */2module PuduLangLog.Sinks.BatchingTest34import 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.Batching as Batching121314fn event(message: Str) -> Log.Event { Event.create(Log.Timestamp{millis: 0, offset: 0}, Log.Information, message, []) }151617fn recording(sizes: &Sync.Cell[Array[Int]], failing: &Sync.Cell[Bool], empties: &Sync.Counter) -> Batching.Target {18 let seen = *sizes19 let down = *failing20 let empty = *empties21 Batching.Target {22 emitBatch: fn(batch: &Array[Log.Event]) -> Result[(), Str] {23 if Sync.get(&down) == Ok(true) { return Err("offline") }24 if let Ok(held) = Sync.get(&seen) { let _stored = Sync.set(&seen, held.push(batch.length())) }25 Ok(())26 },27 onEmptyBatch: fn() -> () { let _counted = Sync.increment(&empty, 1) }28 }29}303132fn sizesOf(sizes: &Sync.Cell[Array[Int]]) -> Array[Int] {33 match Sync.get(sizes) {34 case Ok(held) => held35 case Err(_) => []36 }37}383940fn waitFor(condition: fn() -> Bool) -> Bool {41 var waited = 042 while !condition() && waited < 1000 {43 let _slept = Concurrent.sleep(5)44 waited = waited + 545 }46 condition()47}484950fn main() -> Int {51 let none: Array[Int] = []52 let sizes = Sync.cell(none)53 let failing = Sync.cell(false)54 let empties = Sync.counter(0)55 let slow = Batching.Options{..Batching.defaults(), bufferingTimeLimit: 60000, eagerlyEmitFirstEvent: false, batchSizeLimit: 3}56 let sink = Batching.sink(recording(&sizes, &failing, &empties), slow)57 for number in [1, 2, 3, 4, 5, 6, 7] { let _queued = (sink.emit)(&event(show(number))) }58 let fullBatches = waitFor(fn() -> Bool { sizesOf(&sizes).length() >= 2 })59 (sink.flush)()60 let afterFlush = sizesOf(&sizes)61 (sink.close)()6263 let eagerSizes = Sync.cell(none)64 let eagerSink = Batching.sink(recording(&eagerSizes, &failing, &empties), Batching.Options{..Batching.defaults(), bufferingTimeLimit: 60000})65 let _first = (eagerSink.emit)(&event("first"))66 let eagerly = waitFor(fn() -> Bool { sizesOf(&eagerSizes) == [1] })67 (eagerSink.close)()6869 let timedSizes = Sync.cell(none)70 let timedSink = Batching.sink(recording(&timedSizes, &failing, &empties), Batching.Options{..Batching.defaults(), bufferingTimeLimit: 30, eagerlyEmitFirstEvent: false})71 let _one = (timedSink.emit)(&event("a"))72 let _two = (timedSink.emit)(&event("b"))73 let timely = waitFor(fn() -> Bool { sizesOf(&timedSizes) == [2] })74 let emptyTicks = waitFor(fn() -> Bool { Sync.count(&empties) != Ok(0) })75 (timedSink.close)()7677 let retrySizes = Sync.cell(none)78 let down = Sync.cell(true)79 let lostEvents = Sync.counter(0)80 let retrySink = Batching.sink(recording(&retrySizes, &down, &empties), Batching.Options{..slow, retryTimeLimit: 600000})81 (retrySink.attach)(fn(report: &Sink.Report) -> () { let _counted = Sync.increment(&lostEvents, report.events.length()) })82 let _queued = (retrySink.emit)(&event("kept"))83 (retrySink.flush)()84 let whileDown = sizesOf(&retrySizes)85 let _recovered = Sync.set(&down, false)86 (retrySink.flush)()87 let afterRecovery = sizesOf(&retrySizes)88 (retrySink.close)()8990 let droppingSizes = Sync.cell(none)91 let alwaysDown = Sync.cell(true)92 let dropped = Sync.counter(0)93 let droppingSink = Batching.sink(recording(&droppingSizes, &alwaysDown, &empties), Batching.Options{..slow, retryTimeLimit: 0})94 (droppingSink.attach)(fn(report: &Sink.Report) -> () { let _counted = Sync.increment(&dropped, report.events.length()) })95 let _lost = (droppingSink.emit)(&event("lost"))96 (droppingSink.flush)()97 (droppingSink.close)()9899 let limitedSizes = Sync.cell(none)100 let limitedSink = Batching.sink(recording(&limitedSizes, &failing, &empties), Batching.Options{..slow, queueLimit: Some(2), batchSizeLimit: 10})101 for number in [1, 2, 3, 4] { let _queued = (limitedSink.emit)(&event(show(number))) }102 (limitedSink.close)()103 let refused = (limitedSink.emit)(&event("late"))104105 let checks = Test.suite("Sinks.Batching", &[106 Test.that("full batches go without waiting", fullBatches),107 Test.equals("a flush writes the rest", &afterFlush, &[3, 3, 1]),108 Test.that("the first event goes at once when eager", eagerly),109 Test.that("a batch goes when the buffering time passes", timely),110 Test.that("an empty queue is reported to the target", emptyTicks),111 Test.equals("a failed batch is kept", &whileDown, &[]),112 Test.equals("a kept batch is written after recovery", &afterRecovery, &[1]),113 Test.equals("nothing is lost within the retry time", &Sync.count(&lostEvents), &Ok(0)),114 Test.equals("a batch past its retry time is reported lost", &Sync.count(&dropped), &Ok(1)),115 Test.equals("the queue limit drops the excess", &sizesOf(&limitedSizes), &[2]),116 Test.equals("a closed sink refuses events", &refused, &Err("the batching sink is closed"))117 ])118 let ran = Test.run(&checks)119 for failure in Test.failuresOf(&ran) {120 let _reported = Io.writeErrorLine(failure)121 }122 Test.report(&ran)123}124