///|
struct Observable[T, E] {
subscribe : async (Observer[T, E]) -> Unit raise Cancelled
}
///|
pub fn[T, E] Observable::Observable(
subscribe : async (Observer[T, E]) -> Unit raise Cancelled,
) -> Observable[T, E] {
{ subscribe, }
}
///|
/// Subscribes the given observer to this observable, blocking until the observable completes,
/// errors, or is cancelled.
pub async fn[T, E] Observable::subscribe(
self : Self[T, E],
observer : Observer[T, E],
) -> Unit raise Cancelled {
(self.subscribe)(observer)
}
///|
pub fn[T] Observable::from_array(
array : ArrayView[T],
) -> Observable[T, NoError] {
Observable() <| observer => {
for item in array {
if observer.on_next(item) is Ack::Stop {
break
}
} nobreak {
observer.on_complete()
}
}
}
///|
pub fn[T] Observable::from_iter(iter : Iter[T]) -> Observable[T, NoError] {
Observable() <| observer => {
for item in iter {
if observer.on_next(item) is Ack::Stop {
break
}
} nobreak {
observer.on_complete()
}
}
}
///|
pub fn Observable::from_bytes(bytes : BytesView) -> Observable[Byte, NoError] {
Observable() <| observer => {
for byte in bytes {
if observer.on_next(byte) is Ack::Stop {
break
}
} nobreak {
observer.on_complete()
}
}
}
///|
pub fn[T, E] Observable::empty() -> Observable[T, E] {
Observable() <| observer => { (observer.on_complete)() }
}
///|
pub fn[T, E] Observable::fail(error : E) -> Observable[T, E] {
Observable() <| observer => { (observer.on_error)(error) }
}
///|
pub fn[T : Compare + Add] Observable::range_between(
start~ : T,
end~ : T,
step~ : T,
) -> Observable[T, NoError] {
Observable() <| observer => {
let mut item = start
while item < end {
if observer.on_next(item) is Ack::Stop {
break
}
item += step
} nobreak {
(observer.on_complete)()
}
}
}
///|
pub fn[T : Add] Observable::range_from(
start~ : T,
step~ : T,
) -> Observable[T, NoError] {
Observable() <| observer => {
let mut item = start
while true {
if (observer.on_next)(item) is Ack::Stop {
break
}
item += step
} nobreak {
(observer.on_complete)()
}
}
}
///|
pub fn[T, E : Error] Observable::from_async_generator(
generator : async () -> T? raise E,
) -> Observable[T, E] {
Observable() <| observer => {
let mut generator_error : E? = None
async fn run_generator() -> T? noraise {
generator() catch {
e => {
generator_error = Some(e)
None
}
}
}
while run_generator() is Some(item) {
if observer.on_next(item) is Ack::Stop {
break
}
} nobreak {
match generator_error {
Some(e) => observer.on_error(e)
None => observer.on_complete()
}
}
}
}
///|
pub fn[T, E : Error] Observable::from_async(
task : async () -> T raise E,
) -> Observable[T, E] {
Observable() <| observer => {
let maybe_value = Some(task()) catch {
e => {
observer.on_error(e)
None
}
}
if maybe_value is Some(value) {
if observer.on_next(value) is Ack::Continue {
observer.on_complete()
}
}
}
}
///|
pub fn[T, E : Error] Observable::from_generator(
generator : () -> T? raise E,
) -> Observable[T, E] {
Observable() <| observer => {
let mut generator_error : E? = None
fn run_generator() {
generator() catch {
e => {
generator_error = Some(e)
None
}
}
}
while run_generator() is Some(item) {
if observer.on_next(item) is Ack::Stop {
break
}
} nobreak {
match generator_error {
Some(e) => observer.on_error(e)
None => observer.on_complete()
}
}
}
}
///|
pub fn[T, E] Observable::never() -> Observable[T, E] {
Observable() <| _ => {
// never completes unless cancelled
capture_cancellation <| () => {
while true {
@async.sleep(@int.MAX_VALUE)
}
}
}
}
///|
pub fn[T, E] Observable::from_value(value : T) -> Observable[T, E] {
Observable() <| observer => {
if (observer.on_next)(value) is Ack::Continue {
(observer.on_complete)()
}
}
}
///|
pub fn[T, E, U] Observable::map(
self : Observable[T, E],
f : (T) -> U,
) -> Observable[U, E] {
let subscribe = (observer : Observer[U, E]) => {
let inner_observer = observer.contramap(f)
(self.subscribe)(inner_observer)
}
{ subscribe, }
}
///|
/// Concatenates the inner `Observable`s produced by `f` in order into a single flat `Observable`.
pub fn[T, E, U] Observable::flat_map(
self : Self[T, E],
f : (T) -> Self[U, E],
) -> Self[U, E] {
let subscribe = (observer : Observer[U, E]) => {
let inner_observer = Observer::{
on_next: (item : T) => {
let inner_observable = f(item)
let mut ack = Ack::Continue
let flatten_observer = Observer::{
on_next: item => {
ack = (observer.on_next)(item)
ack
},
on_complete: () => (),
on_error: e => {
ack = Ack::Stop
(observer.on_error)(e)
},
}
inner_observable.subscribe(flatten_observer)
ack
},
on_complete: observer.on_complete,
on_error: observer.on_error,
}
(self.subscribe)(inner_observer)
}
{ subscribe, }
}
///|
pub fn[T, E] Observable::filter(
self : Self[T, E],
predicate : (T) -> Bool,
) -> Self[T, E] {
let subscribe = (observer : Observer[T, E]) => {
let inner_observer = {
..observer,
on_next: item => {
if predicate(item) {
(observer.on_next)(item)
} else {
Ack::Continue
}
},
}
self.subscribe(inner_observer)
}
{ subscribe, }
}
///|
pub fn[T, E : Error] Observable::tap(
self : Self[T, E],
effect : async (T) -> Unit raise E,
) -> Self[T, E] {
let subscribe = (observer : Observer[T, E]) => {
let inner_observer = observer.before_on_next() <| effect
(self.subscribe)(inner_observer)
}
{ subscribe, }
}
///|
pub fn[T, E : Error] Observable::tap_error(
self : Self[T, E],
effect : async (E) -> Unit raise E,
) -> Self[T, E] {
let subscribe = (observer : Observer[T, E]) => {
let inner_observer = observer.before_on_error() <| effect
(self.subscribe)(inner_observer)
}
{ subscribe, }
}
///|
pub fn[T : Debug, E : Debug] Observable::debug_dump(
self : Self[T, E],
) -> Self[T, E] {
let subscribe = (observer : Observer[T, E]) => {
let inner_observer = Observer::{
on_next: (item : T) => {
debug("on_next: \{item |> dbg()}")
(observer.on_next)(item)
},
on_complete: () => {
debug("on_complete")
(observer.on_complete)()
},
on_error: (e : E) => {
debug("on_error: \{e |> dbg()}")
(observer.on_error)(e)
},
}
(self.subscribe)(inner_observer)
}
{ subscribe, }
}
///|
pub fn[T, E] Observable::take(self : Self[T, E], count : UInt) -> Self[T, E] {
if count == 0 {
return Observable::empty()
}
let subscribe = (observer : Observer[T, E]) => {
let mut remaining = count
let inner_observer = {
..observer,
on_next: item => {
remaining -= 1
let ack = (observer.on_next)(item)
if remaining == 0 {
(observer.on_complete)()
Ack::Stop
} else {
ack
}
},
}
(self.subscribe)(inner_observer)
}
{ subscribe, }
}
///|
pub fn[T, E] Observable::drop(self : Self[T, E], count : UInt) -> Self[T, E] {
if count == 0 {
return self
}
let subscribe = (observer : Observer[T, E]) => {
let mut remaining = count
let inner_observer = {
..observer,
on_next: item => {
if remaining > 0 {
remaining -= 1
Continue
} else {
(observer.on_next)(item)
}
},
}
(self.subscribe)(inner_observer)
}
{ subscribe, }
}
///|
pub fn[T, E1, E2] Observable::map_error(
self : Self[T, E1],
f : (E1) -> E2,
) -> Self[T, E2] {
let subscribe = (observer : Observer[T, E2]) => {
let inner_observer = {
on_next: observer.on_next,
on_complete: observer.on_complete,
on_error: e => (observer.on_error)(f(e)),
}
(self.subscribe)(inner_observer)
}
{ subscribe, }
}
///|
pub fn[T, E] Observable::catch_error(
self : Self[T, E],
f : (E) -> Self[T, E],
) -> Self[T, E] {
let subscribe = (observer : Observer[T, E]) => {
let inner_observer = {
..observer,
on_error: e => {
let new_observable = f(e)
(new_observable.subscribe)(observer)
},
}
(self.subscribe)(inner_observer)
}
{ subscribe, }
}
///|
pub fn[T, E] Observable::recover(
self : Self[T, E],
f : (E) -> T,
) -> Self[T, NoError] {
let subscribe = (observer : Observer[T, NoError]) => {
let inner_observer = {
on_next: observer.on_next,
on_complete: observer.on_complete,
on_error: e => {
let recovered_value = f(e)
if (observer.on_next)(recovered_value) is Continue {
(observer.on_complete)()
}
},
}
(self.subscribe)(inner_observer)
}
{ subscribe, }
}
///|
pub fn[T, E1, E2] Observable::refine_error(
self : Self[T, E1],
f : (E1) -> E2?,
) -> Self[T, E2] {
Observable() <| observer => {
let inner_observer = {
on_next: observer.on_next,
on_complete: observer.on_complete,
on_error: e => {
match f(e) {
Some(refined_error) => observer.on_error(refined_error)
None => panic()
}
},
}
self.subscribe(inner_observer)
}
}
///|
pub fn[T, E] Observable::panic_on_error(self : Self[T, E]) -> Self[T, NoError] {
Observable() <| observer => {
let inner_observer = {
on_next: observer.on_next,
on_complete: observer.on_complete,
on_error: _ => panic(),
}
self.subscribe(inner_observer)
}
}
///|
pub fn[T, E] Observable::widen_error(self : Self[T, NoError]) -> Self[T, E] {
Observable() <| observer => {
let inner_observer = {
on_next: observer.on_next,
on_complete: observer.on_complete,
on_error: _ => panic(),
}
self.subscribe(inner_observer)
}
}
///|
pub fn[T, E] Observable::buffer_tumbling(
self : Self[T, E],
size : Int,
) -> Self[Vector[T], E] {
Observable() <| observer => {
let mut buffer = Vector::new()
let inner_observer = {
on_next: item => {
buffer = buffer.push(item)
if buffer.length() == size {
let to_emit = buffer
buffer = Vector::new()
if (observer.on_next)(to_emit) is Stop {
return Stop
}
}
Continue
},
on_complete: () => {
let ack = if !buffer.is_empty() {
(observer.on_next)(buffer)
} else {
Continue
}
if ack is Continue {
(observer.on_complete)()
}
},
on_error: observer.on_error,
}
(self.subscribe)(inner_observer)
}
}