
Observable.pudu
Pudu94 lines3.4 KB
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 }161718export 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}222324export fn observer(next: fn(&Log.Event) -> ()) -> Observer { Observer{next: next, completed: fn() -> () {}} }25262728export 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}414243export 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}484950export fn count(subject: &Observable) -> Int { observersOf(subject).length() }51525354export 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}787980fn observersOf(subject: &Observable) -> Array[(Int, Observer)] {81 match Sync.get(&subject.observers) {82 case Ok(current) => current83 case Err(_) => []84 }85}868788fn ended(subject: &Observable) -> Bool {89 match Sync.get(&subject.ended) {90 case Ok(flag) => flag91 case Err(_) => true92 }93}94