Asynchronous reactive streams library for Moonbit
Dependencies
[!NOTE] Currently Lunar Reactor only offers "cold" observables, which means that the stream of values is generated anew for every subscription by an observer.
[!NOTE] While the Observer trait is the fundamental way to use an Observable, it is typically only used within the internal implementations of Observables. For application code, Observable is not very convenient to use directly, and it is generally preferable to consume an observable with a Consumer, see below.
async test "consuming an observable" {
let stream = Observable::from_array(["a", "b", "c"])
let result = stream.consume(Consumer::count())
assert_eq(result, 3)
}async test "using a pipeline" {
let pipeline = Pipeline::map(i => i * 2)
.and_then(Pipeline::drop_while(i => i < 7))
let stream = Observable::range_between(start=0, end=10, step=1)
.pipe(pipeline)
let result = stream.consume(Consumer::collect_all())
assert_eq(result, [8, 10, 12, 14, 16, 18])
}pub suberror NoErrortype Consumer[T, E, R]fn[T, E] Consumer::collect_to_fixed(fixed : FixedArray[T]) -> Consumer[T, E, CollectToFixedResult[T]]type Observable[T, E]fn[T, E] Observable::Observable(subscribe : async (Observer[T, E]) -> Unit raise Cancelled) -> Observable[T, E]fn[T, E] Observable::buffer_tumbling(self : Observable[T, E], size : Int) -> Observable[Vector[T], E]fn[T, E] Observable::catch_error(self : Observable[T, E], f : (E) -> Observable[T, E]) -> Observable[T, E]async fn[T, E : Error, R] Observable::consume(self : Observable[T, E], consumer : Consumer[T, E, R]) -> Rasync fn[T, E, R] Observable::consume_result(self : Observable[T, E], consumer : Consumer[T, E, R]) -> Result[R, E] raise Cancelledfn[T, E, U] Observable::flat_map(self : Observable[T, E], f : (T) -> Observable[U, E]) -> Observable[U, E]fn[T, E : Error] Observable::from_async_generator(generator : async () -> T? raise E) -> Observable[T, E]fn[T1, T2, E] Observable::pipe(self : Observable[T1, E], pipeline : Pipeline[T1, T2, E]) -> Observable[T2, E]fn[T : Compare + Add + Eq] Observable::range_between(start~ : T, end~ : T, step~ : T) -> Observable[T, NoError]fn[T, E1, E2] Observable::refine_error(self : Observable[T, E1], f : (E1) -> E2?) -> Observable[T, E2]async fn[T, E] Observable::subscribe(self : Observable[T, E], observer : Observer[T, E]) -> Unit raise Cancelledfn[T, E : Error] Observable::tap(self : Observable[T, E], effect : async (T) -> Unit raise E) -> Observable[T, E]fn[T, E : Error] Observable::tap_error(self : Observable[T, E], effect : async (E) -> Unit raise E) -> Observable[T, E]type Observer[T, E]type Pipeline[T1, T2, E]fn[T1, T2, E] Pipeline::apply(self : Pipeline[T1, T2, E], observable : Observable[T1, E]) -> Observable[T2, E]type Resource[T]Asynchronous reactive streams library for Moonbit
Dependencies