///|
pub(all) enum Ack {
Continue
Stop
} derive(Debug)
///|
struct Observer[T, E] {
on_next : async (T) -> Ack raise Cancelled
on_complete : async () -> Unit raise Cancelled
on_error : async (E) -> Unit raise Cancelled
}
///|
pub fn[T, E] Observer::Observer(
on_next~ : async (T) -> Ack raise Cancelled,
on_complete~ : async () -> Unit raise Cancelled,
on_error~ : async (E) -> Unit raise Cancelled,
) -> Observer[T, E] {
{ on_next, on_complete, on_error }
}
///|
pub async fn[T, E] Observer::on_next(
self : Self[T, E],
t : T,
) -> Ack raise Cancelled {
(self.on_next)(t)
}
///|
pub async fn[T, E] Observer::on_complete(
self : Self[T, E],
) -> Unit raise Cancelled {
(self.on_complete)()
}
///|
pub async fn[T, E] Observer::on_error(
self : Self[T, E],
e : E,
) -> Unit raise Cancelled {
(self.on_error)(e)
}
///|
pub fn[T, E, U] Observer::replace_on_next(
self : Self[T, E],
on_next : async (U) -> Ack raise Cancelled,
) -> Self[U, E] {
{ on_next, on_complete: self.on_complete, on_error: self.on_error }
}
///|
pub fn[T, E, U] Observer::contramap(
self : Self[T, E],
f : (U) -> T,
) -> Self[U, E] {
let on_next : async (U) -> Ack raise Cancelled = (item : U) => {
let mapped = f(item)
(self.on_next)(mapped)
}
Observer::{ on_next, on_complete: self.on_complete, on_error: self.on_error }
}
///|
pub fn[T, E, U] Observer::contramap_halt(
self : Self[T, E],
f : (U) -> T? noraise,
) -> Self[U, E] {
let on_next : async (U) -> Ack raise Cancelled = (item : U) => {
match f(item) {
Some(mapped) => (self.on_next)(mapped)
None => Ack::Stop
}
}
{ on_next, on_complete: self.on_complete, on_error: self.on_error }
}
///|
pub fn[T, E : Error] Observer::before_on_next(
self : Self[T, E],
effect : async (T) -> Unit raise E,
) -> Self[T, E] {
let on_next = (item : T) => {
effect(item) catch {
e => {
self.on_error(e)
return Ack::Stop
}
}
self.on_next(item)
}
{ ..self, on_next, }
}
///|
pub fn[T, E : Error] Observer::before_on_error(
self : Self[T, E],
effect : async (E) -> Unit raise E,
) -> Self[T, E] {
let on_error = (e : E) => {
effect(e) catch {
e2 => (self.on_error)(e2)
}
self.on_error(e)
}
{ ..self, on_error, }
}
///|
fn[T, E] Observer::ignore() -> Observer[T, E] {
Observer(on_next=_ => Ack::Continue, on_complete=() => (), on_error=_ => ())
}
///|
fn[T, E] Observer::stop() -> Observer[T, E] {
Observer(on_next=_ => Ack::Stop, on_complete=() => (), on_error=_ => ())
}
///|
pub fn[T : Debug, E : Debug] Observer::debug_dump() -> Observer[T, E] {
Observer(
on_next=(item : T) => {
debug("on_next: \{item |> dbg()}")
Ack::Continue
},
on_complete=() => debug("on_complete"),
on_error=(e : E) => debug("on_error: \{e |> dbg()}"),
)
}