
Async.pudu
Pudu120 lines4.3 KB
1/** @Log.Sinks.Async.Wrapper — hands events to a sink on another thread */2module PuduLangLog.Sinks.Async34import Std.Channel as Channel5import Std.Concurrent as Concurrent6import Std.Sync as Sync7import PuduLangLog as Log8import PuduLangLog.SelfLog as SelfLog9import PuduLangLog.Sink as Sink1011/** @Log.Sinks.Async.Options — queue size and behaviour when full */12export type Options = { bufferSize: Int, blockWhenFull: Bool, selfLog: SelfLog.SelfLog }1314/** @Log.Sinks.Async.Queue — a running wrapper and its counters */15export type Async = {16 inner: Sink.Sink,17 options: Options,18 queue: Channel.Channel[Log.Event],19 accepted: Sync.Counter,20 written: Sync.Counter,21 dropped: Sync.Counter,22 listener: Sync.Cell[Option[fn(&Sink.Report) -> ()]],23 worker: Option[Concurrent.Task]24}252627const POLL: Int = 1282930export fn defaults() -> Options { Options{bufferSize: 10000, blockWhenFull: false, selfLog: SelfLog.create()} }313233export fn start(inner: Sink.Sink, options: Options) -> Async {34 let none: Option[fn(&Sink.Report) -> ()] = None35 let base = Async{inner: inner, options: options, queue: Channel.channel(options.bufferSize), accepted: Sync.counter(0), written: Sync.counter(0), dropped: Sync.counter(0), listener: Sync.cell(none), worker: None}36 match Concurrent.start(fn() -> () { drain(&base) }) {37 case Ok(started) => Async{..base, worker: Some(started)}38 case Err(_) => {39 SelfLog.write(&options.selfLog, "The asynchronous sink could not start its thread; events are dropped")40 base41 }42 }43}444546fn drain(queue: &Async) -> () {47 loop {48 match Channel.receive(&queue.queue) {49 case Ok(Some(event)) => {50 if let Err(problem) = (queue.inner.emit)(&event) { report(queue, problem, event) }51 let _counted = Sync.increment(&queue.written, 1)52 }53 case _ => { return () }54 }55 }56}575859fn report(queue: &Async, problem: Str, event: Log.Event) -> () {60 SelfLog.write(&queue.options.selfLog, "The asynchronous sink could not write an event: " + problem)61 if let Ok(Some(listener)) = Sync.get(&queue.listener) { listener(&Sink.Report{kind: Sink.Permanent, message: problem, events: [event]}) }62}63646566export fn sinkOf(queue: &Async) -> Sink.Sink {67 let shared = *queue68 Sink.Sink {69 emit: fn(event: &Log.Event) -> Result[(), Str] { enqueue(&shared, event) },70 flush: fn() -> () {71 while count(&shared.written) < count(&shared.accepted) { let _slept = Concurrent.sleep(POLL) }72 (shared.inner.flush)()73 },74 close: fn() -> () {75 let _closed = Channel.close(&shared.queue)76 if let Some(worker) = shared.worker { let _joined = Concurrent.join(&worker) }77 (shared.inner.close)()78 },79 attach: fn(listener: Sink.Listener) -> () {80 let _stored = Sync.set(&shared.listener, Some(listener))81 (shared.inner.attach)(listener)82 }83 }84}858687fn enqueue(queue: &Async, event: &Log.Event) -> Result[(), Str] {88 if !queue.options.blockWhenFull && pending(queue) >= queue.options.bufferSize {89 let _counted = Sync.increment(&queue.dropped, 1)90 SelfLog.write(&queue.options.selfLog, "The asynchronous sink's queue is full; an event was dropped")91 return Ok(())92 }93 match Channel.send(&queue.queue, *event) {94 case Ok(_) => {95 let _counted = Sync.increment(&queue.accepted, 1)96 Ok(())97 }98 case Err(_) => Err("the asynchronous sink is closed")99 }100}101102103fn count(counter: &Sync.Counter) -> Int {104 match Sync.count(counter) {105 case Ok(value) => value106 case Err(_) => 0107 }108}109110111export fn pending(queue: &Async) -> Int {112 match Channel.pending(&queue.queue) {113 case Ok(value) => value114 case Err(_) => 0115 }116}117118119export fn dropped(queue: &Async) -> Int { count(&queue.dropped) }120