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

Observable.pudu

Pudu94 lines3.4 KB

GitHub ↗
1/** @Log.Sinks.Observable.Subject — hands events to subscribed observers */2module PuduLangLog.Sinks.Observable34import Std.Sync as Sync5import PuduLangLog as Log6import PuduLangLog.Sink as Sink78/** @Log.Sinks.Observable.Observer — receives events and the end of the stream */9export type Observer = { next: fn(&Log.Event) -> (), completed: fn() -> () }1011/** @Log.Sinks.Observable.Subscription — the ticket that ends one subscription */12export type Subscription = { id: Int }1314/** @Log.Sinks.Observable.Subject — the observers and whether the stream ended */15export type Observable = { observers: Sync.Cell[Array[(Int, Observer)]], issued: Sync.Counter, ended: Sync.Cell[Bool], lock: Sync.Mutex }1617/// A subject with no observers.18export fn create() -> Observable {19  let none: Array[(Int, Observer)] = []20  Observable{observers: Sync.cell(none), issued: Sync.counter(0), ended: Sync.cell(false), lock: Sync.mutex()}21}2223/// An observer calling `next` for each event and ignoring the end of the stream.24export fn observer(next: fn(&Log.Event) -> ()) -> Observer { Observer{next: next, completed: fn() -> () {}} }2526/// Adds an observer and answers the subscription that removes it. An observer added after the sink27/// closed is told at once that the stream ended.28export fn subscribe(subject: &Observable, watcher: Observer) -> Subscription {29  let id = match Sync.increment(&subject.issued, 1) {30    case Ok(next) => next31    case Err(_) => 032  }33  let late = Sync.withLock(&subject.lock, fn() -> Bool {34      if ended(subject) { return true }35      let _stored = Sync.set(&subject.observers, observersOf(subject).push((id, watcher)))36      false37    })38  if late != Ok(false) { (watcher.completed)() }39  Subscription{id: id}40}4142/// Removes the observer of a subscription; removing it twice does nothing.43export fn unsubscribe(subject: &Observable, subscription: Subscription) -> () {44  let _removed = Sync.withLock(&subject.lock, fn() -> () {45      let _stored = Sync.set(&subject.observers, observersOf(subject).filter(|entry: (Int, Observer)| entry[0] != subscription.id))46    })47}4849/// How many observers are subscribed.50export fn count(subject: &Observable) -> Int { observersOf(subject).length() }5152/// The sink handing each event to every observer, in subscription order. Closing it ends the53/// stream: each observer's `completed` runs once and the observers are removed.54export fn sink(subject: &Observable) -> Sink.Sink {55  let shared = *subject56  Sink.Sink {57    emit: fn(event: &Log.Event) -> Result[(), Str] {58      for entry in observersOf(&shared) { (entry[1].next)(event) }59      Ok(())60    },61    flush: fn() -> () {},62    close: fn() -> () {63      let finished = match Sync.withLock(&shared.lock, fn() -> Array[(Int, Observer)] {64          let current = observersOf(&shared)65          let none: Array[(Int, Observer)] = []66          let _stored = Sync.set(&shared.observers, none)67          let _ended = Sync.set(&shared.ended, true)68          current69        }) {70        case Ok(current) => current71        case Err(_) => []72      }73      for entry in finished { (entry[1].completed)() }74    },75    attach: fn(_listener: Sink.Listener) -> () {}76  }77}7879/// The observers subscribed, oldest first.80fn observersOf(subject: &Observable) -> Array[(Int, Observer)] {81  match Sync.get(&subject.observers) {82    case Ok(current) => current83    case Err(_) => []84  }85}8687/// Whether the sink has closed.88fn ended(subject: &Observable) -> Bool {89  match Sync.get(&subject.ended) {90    case Ok(flag) => flag91    case Err(_) => true92  }93}94