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

Async.pudu

Pudu120 lines4.3 KB

GitHub ↗
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}2526/// How long `flush` waits between checks of the queue, in milliseconds.27const POLL: Int = 12829/// Ten thousand queued events, dropping new events when full, and a disabled self-log.30export fn defaults() -> Options { Options{bufferSize: 10000, blockWhenFull: false, selfLog: SelfLog.create()} }3132/// Starts the thread that writes queued events to the inner sink.33export 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}4445/// Writes queued events until the queue is closed and empty.46fn 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}5758/// Reports an inner sink's failure to the self-log and to the attached listener.59fn 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}6364/// The sink that queues events for the thread. Flushing waits until every accepted event was65/// written; closing also stops the thread.66export 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}8586/// Queues an event, waiting for room when blocking and otherwise dropping it when the queue is full.87fn 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}101102/// A counter's value.103fn count(counter: &Sync.Counter) -> Int {104  match Sync.count(counter) {105    case Ok(value) => value106    case Err(_) => 0107  }108}109110/// How many events wait in the queue.111export fn pending(queue: &Async) -> Int {112  match Channel.pending(&queue.queue) {113    case Ok(value) => value114    case Err(_) => 0115  }116}117118/// How many events were dropped because the queue was full.119export fn dropped(queue: &Async) -> Int { count(&queue.dropped) }120