
Pipeline.pudu
Pudu124 lines4.2 KB
1/** @Log.Pipeline.Core — the stages every event of one logger passes */2module PuduLangLog.Pipeline34import Std.List as List5import Std.Map as Map6import Std.Sync as Sync7import PuduLangLog.Clock as Clock8import PuduLangLog.Domain.Capture as Capture9import PuduLangLog.Domain.Levels as Levels10import PuduLangLog.Domain.Parser as Parser11import PuduLangLog.Domain.Sources as Sources12import PuduLangLog.Enricher as Enricher13import PuduLangLog.Filter as Filter14import PuduLangLog.LevelSwitch as LevelSwitch15import PuduLangLog as Log16import PuduLangLog.SelfLog as SelfLog17import PuduLangLog.Sink as Sink1819/** @Log.Pipeline.Stages — levels, enrichment, filters, and sinks */20export type Pipeline = {21 minimum: Log.Level,22 control: Option[LevelSwitch.LevelSwitch],23 overrides: Array[Override],24 enrichers: Array[Enricher.Enricher],25 filters: Array[Filter.Filter],26 sink: Sink.Sink,27 audit: Option[Sink.Sink],28 policy: Capture.Policy,29 selfLog: SelfLog.SelfLog,30 clock: Clock.Clock,31 templates: Templates,32 silent: Bool33}3435/** @Log.Pipeline.Override — a level switch for one source prefix */36export type Override = { source: Str, control: LevelSwitch.LevelSwitch }3738/** @Log.Pipeline.Templates — parsed templates kept for reuse */39export type Templates = { entries: Sync.Cell[Map[Str, Log.Template]], lock: Sync.Mutex }404142const MOST_TEMPLATES: Int = 1000434445const LONGEST_TEMPLATE: Int = 1024464748export fn templates() -> Templates {49 let none: Map[Str, Log.Template] = Map.empty()50 Templates{entries: Sync.cell(none), lock: Sync.mutex()}51}525354export fn silent() -> Pipeline {55 Pipeline {56 minimum: Log.Fatal,57 control: None,58 overrides: [],59 enrichers: [],60 filters: [],61 sink: Sink.none(),62 audit: None,63 policy: Capture.defaults(),64 selfLog: SelfLog.create(),65 clock: Clock.system(),66 templates: templates(),67 silent: true68 }69}70717273export fn templateOf(pipeline: &Pipeline, text: Str) -> Log.Template {74 if let Ok(entries) = Sync.get(&pipeline.templates.entries) {75 if let Some(found) = Map.get(&entries, text) { return found }76 }77 let parsed = Parser.parse(text)78 if text.length() <= LONGEST_TEMPLATE {79 let _kept = Sync.withLock(&pipeline.templates.lock, fn() -> () {80 if let Ok(entries) = Sync.get(&pipeline.templates.entries) {81 if Map.size(&entries) < MOST_TEMPLATES { let _stored = Sync.set(&pipeline.templates.entries, Map.insert(&entries, text, parsed)) }82 }83 })84 }85 parsed86}87888990export fn enabled(pipeline: &Pipeline, source: &Option[Str], level: Log.Level) -> Bool {91 if pipeline.silent { return false }92 if let Some(context) = source {93 if let Some(found) = List.find(&pipeline.overrides, |candidate: Override| Sources.covers(candidate.source, context)) {94 return LevelSwitch.allows(&found.control, level)95 }96 }97 if !Levels.passes(level, pipeline.minimum) { return false }98 match pipeline.control {99 case Some(control) => LevelSwitch.allows(&control, level)100 case None => true101 }102}103104105106export fn process(pipeline: &Pipeline, event: Log.Event, context: &Array[Enricher.Enricher]) -> Result[(), Str] {107 var current = event108 for enricher in context.reverse() { current = enricher(current, &pipeline.policy) }109 for enricher in pipeline.enrichers { current = enricher(current, &pipeline.policy) }110 for filter in pipeline.filters {111 if !filter(¤t) { return Ok(()) }112 }113 let _emitted = (pipeline.sink.emit)(¤t)114 match pipeline.audit {115 case Some(audit) => (audit.emit)(¤t)116 case None => Ok(())117 }118}119120121export fn report(pipeline: &Pipeline, problems: &Array[Str]) -> () {122 for problem in problems { SelfLog.write(&pipeline.selfLog, problem) }123}124