
Map.pudu
Pudu114 lines4.5 KB
1/** @Log.Sinks.Map.Router — one sink per key, opened on first use */2module PuduLangLog.Sinks.Map34import Std.List as List5import Std.Sync as Sync6import PuduLangLog.Domain.Display as Display7import PuduLangLog.Domain.Properties as Properties8import PuduLangLog.Domain.Recency as Recency9import PuduLangLog as Log10import PuduLangLog.Sink as Sink1112/** @Log.Sinks.Map.Options — the key of an event and the sink for a key */13export type Options = { keyOf: fn(&Log.Event) -> Str, create: fn(Str) -> Sink.Sink, limit: Option[Int] }1415/** @Log.Sinks.Map.Open — a key and the sink opened for it */16type Open = { key: Str, sink: Sink.Sink }1718/** @Log.Sinks.Map.State — the open sinks, least recently used first */19type State = { open: Array[Open], listeners: Array[Sink.Listener], closed: Bool }202122export fn defaults(keyOf: fn(&Log.Event) -> Str, create: fn(Str) -> Sink.Sink) -> Options { Options{keyOf: keyOf, create: create, limit: None} }23242526export fn byProperty(name: Str, fallback: Str) -> fn(&Log.Event) -> Str {27 fn(event: &Log.Event) -> Str {28 match Properties.find(&event.properties, name) {29 case Some(Log.Scalar(Log.Text(text))) => text30 case Some(held) => Display.plain(&Display.value(&held, &None))31 case None => fallback32 }33 }34}3536373839export fn sink(options: Options) -> Sink.Sink {40 let none: Array[Open] = []41 let noListeners: Array[Sink.Listener] = []42 let state = Sync.cell(State{open: none, listeners: noListeners, closed: false})43 let lock = Sync.mutex()44 Sink.Sink {45 emit: fn(event: &Log.Event) -> Result[(), Str] {46 match Sync.withLock(&lock, fn() -> Result[(), Str] { emitLocked(&options, &state, event) }) {47 case Ok(outcome) => outcome48 case Err(_) => Err("The sink map lock failed")49 }50 },51 flush: fn() -> () {52 let _flushed = Sync.withLock(&lock, fn() -> () { for entry in stateOf(&state).open { (entry.sink.flush)() } })53 },54 close: fn() -> () {55 let _closed = Sync.withLock(&lock, fn() -> () {56 let current = stateOf(&state)57 for entry in current.open { retire(&entry) }58 let _stored = Sync.set(&state, State{..current, open: none, closed: true})59 })60 },61 attach: fn(listener: Sink.Listener) -> () {62 let _attached = Sync.withLock(&lock, fn() -> () {63 let current = stateOf(&state)64 for entry in current.open { (entry.sink.attach)(listener) }65 let _stored = Sync.set(&state, State{..current, listeners: current.listeners.push(listener)})66 })67 }68 }69}70717273fn emitLocked(options: &Options, state: &Sync.Cell[State], event: &Log.Event) -> Result[(), Str] {74 let current = stateOf(state)75 if current.closed { return Ok(()) }76 let key = (options.keyOf)(event)77 let target = match List.find(¤t.open, |entry: Open| entry.key == key) {78 case Some(found) => found79 case None => opened(options, ¤t.listeners, key)80 }81 let others = current.open.filter(|entry: Open| entry.key != key)82 let order = Recency.touch(&others.map(|entry: Open| entry.key), key)83 let retired = match options.limit {84 case Some(limit) => Recency.overflow(&order, limit)85 case None => []86 }87 let outcome = (target.sink.emit)(event)88 let kept = others.push(target)89 for entry in kept.filter(|entry: Open| retired.contains(entry.key)) { retire(&entry) }90 let _stored = Sync.set(state, State{..current, open: kept.filter(|entry: Open| !retired.contains(entry.key))})91 outcome92}939495fn opened(options: &Options, listeners: &Array[Sink.Listener], key: Str) -> Open {96 let created = (options.create)(key)97 for listener in listeners { (created.attach)(listener) }98 Open{key: key, sink: created}99}100101102fn retire(entry: &Open) -> () {103 (entry.sink.flush)()104 (entry.sink.close)()105}106107108fn stateOf(state: &Sync.Cell[State]) -> State {109 match Sync.get(state) {110 case Ok(current) => current111 case Err(_) => State{open: [], listeners: [], closed: true}112 }113}114