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

Sink.pudu

Pudu145 lines5.7 KB

GitHub ↗
1/** @Log.Sink.Seam — where events finally go, and wrappers around them */2module PuduLangLog.Sink34import Std.List as List5import PuduLangLog.Domain.Levels as Levels6import PuduLangLog.LevelSwitch as LevelSwitch7import PuduLangLog as Log8import PuduLangLog.SelfLog as SelfLog910/** @Log.Sink.FailureKind — whether events of a failure may still arrive */11export type FailureKind = Temporary | Permanent | Final1213/** @Log.Sink.Report — a sink failure and the events it lost */14export type Report = { kind: FailureKind, message: Str, events: Array[Log.Event] }1516/** @Log.Sink.Listener — receives reports of sink failures */17export type Listener = fn(&Report) -> ()1819/** @Log.Sink.Destination — emits events, flushes, and closes */20export type Sink = {21  emit: fn(&Log.Event) -> Result[(), Str],22  flush: fn() -> (),23  close: fn() -> (),24  attach: fn(Listener) -> ()25}2627/// A sink that writes each event with a function and needs no flushing or closing.28export fn of(emit: fn(&Log.Event) -> Result[(), Str]) -> Sink {29  Sink{emit: emit, flush: fn() -> () {}, close: fn() -> () {}, attach: fn(_listener: Listener) -> () {}}30}3132/// A sink that formats each event and writes the text.33export fn lines(formatter: Log.Formatter, write: fn(Str) -> Result[(), Str]) -> Sink {34  of(fn(event: &Log.Event) -> Result[(), Str] { write(formatter(event)) })35}3637/// A sink that discards every event.38export fn none() -> Sink { of(fn(_event: &Log.Event) -> Result[(), Str] { Ok(()) }) }3940/// The sink with a different emit, keeping its flush, close, and listener.41fn wrapped(sink: &Sink, emit: fn(&Log.Event) -> Result[(), Str]) -> Sink { Sink{..*sink, emit: emit} }4243/// The sink receiving only events at or above a level.44export fn restricted(sink: Sink, minimum: Log.Level) -> Sink {45  wrapped(&sink, fn(event: &Log.Event) -> Result[(), Str] { if Levels.passes(event.level, minimum) { (sink.emit)(event) } else { Ok(()) } })46}4748/// The sink receiving only events a level switch allows.49export fn controlled(sink: Sink, control: &LevelSwitch.LevelSwitch) -> Sink {50  let shared = *control51  wrapped(&sink, fn(event: &Log.Event) -> Result[(), Str] { if LevelSwitch.allows(&shared, event.level) { (sink.emit)(event) } else { Ok(()) } })52}5354/// The sink receiving only events a condition accepts.55export fn conditional(sink: Sink, condition: fn(&Log.Event) -> Bool) -> Sink {56  wrapped(&sink, fn(event: &Log.Event) -> Result[(), Str] { if condition(event) { (sink.emit)(event) } else { Ok(()) } })57}5859/// A sink writing to every sink. A failing sink is reported to the self-log and the others still60/// receive the event, so the aggregate never fails.61export fn aggregate(sinks: Array[Sink], log: &SelfLog.SelfLog) -> Sink {62  let shared = *log63  Sink {64    emit: fn(event: &Log.Event) -> Result[(), Str] {65      for sink in sinks {66        if let Err(problem) = (sink.emit)(event) { SelfLog.write(&shared, "Caught failure while emitting to sink: " + problem) }67      }68      Ok(())69    },70    flush: fn() -> () { for sink in sinks { (sink.flush)() } },71    close: fn() -> () { for sink in sinks { (sink.close)() } },72    attach: fn(listener: Listener) -> () { for sink in sinks { (sink.attach)(listener) } }73  }74}7576/// A sink writing to every sink and answering the first failure after all have been tried, for77/// audit trails whose writes must be confirmed.78export fn audited(sinks: Array[Sink]) -> Sink {79  Sink {80    emit: fn(event: &Log.Event) -> Result[(), Str] {81      var outcome: Result[(), Str] = Ok(())82      for sink in sinks {83        let result = (sink.emit)(event)84        if outcome == Ok(()) { outcome = result }85      }86      outcome87    },88    flush: fn() -> () { for sink in sinks { (sink.flush)() } },89    close: fn() -> () { for sink in sinks { (sink.close)() } },90    attach: fn(listener: Listener) -> () { for sink in sinks { (sink.attach)(listener) } }91  }92}9394/// The sink with its failures reported to a listener instead of the caller. Failures a sink finds95/// later, such as a batch that could not be sent, reach the same listener.96export fn fallible(sink: Sink, listener: Listener) -> Sink {97  (sink.attach)(listener)98  Sink {99    ..sink,100    emit: fn(event: &Log.Event) -> Result[(), Str] {101      if let Err(problem) = (sink.emit)(event) { listener(&Report{kind: Permanent, message: "Failed to emit event to wrapped sink: " + problem, events: [*event]}) }102      Ok(())103    }104  }105}106107/// A sink trying each sink in turn until one records the event. Events a sink loses later are108/// handed to the sinks after it. The chain fails only when every sink failed.109export fn fallbackChain(sinks: Array[Sink]) -> Sink {110  var index = 0111  for sink in sinks {112    let rest = sinks.slice(index + 1, sinks.length())113    (sink.attach)(fn(report: &Report) -> () {114        for lost in report.events { let _forwarded = emitFirst(&rest, &lost) }115      })116    index = index + 1117  }118  Sink {119    emit: fn(event: &Log.Event) -> Result[(), Str] { emitFirst(&sinks, event) },120    flush: fn() -> () { for sink in sinks { (sink.flush)() } },121    close: fn() -> () { for sink in sinks { (sink.close)() } },122    attach: fn(listener: Listener) -> () {123      if let Some(last) = List.last(&sinks) { (last.attach)(listener) }124    }125  }126}127128/// The event emitted to the first sink that accepts it, or the last failure.129fn emitFirst(sinks: &Array[Sink], event: &Log.Event) -> Result[(), Str] {130  var outcome: Result[(), Str] = Err("No sink in the fallback chain")131  for sink in sinks {132    outcome = (sink.emit)(event)133    if outcome == Ok(()) { return outcome }134  }135  outcome136}137138/// A listener writing each report to the self-log.139export fn reportTo(log: &SelfLog.SelfLog) -> Listener {140  let shared = *log141  fn(report: &Report) -> () {142    SelfLog.write(&shared, "Logging failure (" + show(report.kind) + ", " + show(report.events.length()) + " events): " + report.message)143  }144}145