lunar-reactor

Asynchronous reactive streams library for Moonbit

async
reactive
streams
moon add quelgar/lunar-reactor@0.1.0
Download zip
Author
Version
0.1.0
License
Apache-2.0
Last updated
9 hours ago
Downloads
2

Dependencies

README

#Lunar Reactor

Asynchronous reactive streams for Moonbit.

The main thing that Lunar Reactor offers is the Observable type, which provides a push-based asynchronous stream of values with backpressure.

Inspired by the Monix and ZIO Scala libraries.

See examples/fs_examples.mbt for a simple example of using an observable to read and write files with UTF-8 encoding.

#Observables

The type Observable[T, E] represents a stream of values of type T that can fail with an error of type E. Observables can be modified in various ways, such as applying filters or mapping values. Observables are lazy and only start emitting values when they are subscribed to by an Observer.

[!NOTE] Currently Lunar Reactor only offers "cold" observables, which means that the stream of values is generated anew for every subscription by an observer.

#Observers

An observer is used to subscribe to an observable, which activates the observable. The observable will make calls to the three methods of the observer: on_next, on_error and on_complete. All are async.

[!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.

The order of calls is guaranteed to be: 0 or more calls to on_next, followed by either a single call to on_error or a single call to on_complete. No more calls will be made to the observer after either on_error or on_complete has been called. The only exception to this is if the observer returns Stop from on_next, see below.

#on_next

The on_next method accepts a value of type T and returns an Ack. It is called whenever the observable has a value to emit. The Ack return value has two cases: Continue and Stop. If Stop is returned then the subscription is effectively cancelled and the observer will not be called on any of its methods again.

As on_next is async, the observable will wait for the returned Ack before emitting the next value. This allows the observer to apply backpressure to the observable.

#on_error

Provides an error of type E to the observer, indicating the observable failed. After this method is called, the observable will not emit any more values and the observer will not be called again.

#on_complete

Indicates to the observer that the observable has completed successfully. After this method is called, the observable will not emit any more values and the observer will not be called again.

#Consumers

A Consumer[T, E, R] is a type that can consumer from an Observable[T, E] and produces either a result of type R when the observable completes successfully, or an error of type E when the observable or consumer itself fails. A consumer is typically a much more convenient way for application code to use observables than using an observer directly.

Passing a consumer to the consume method of an observable will subscribe to the observable and start the stream. As consume is async, it will block until the observable completes or fails, and will either return the result of type R or raise an error of type E.

There's also the consume_result method that returns a Result[R, E] instead of raising.

async test "consuming an observable" {
let stream = Observable::from_array(["a", "b", "c"])
let result = stream.consume(Consumer::count())
assert_eq(result, 3)
}

#Pipelines

Pipelines are essentially functions from one observable to another. They aren't strictly necessary, as anything they can do could be done directly to an observable, but they provide a convenient way to encapsulate and compose useful transformations.

Pipelines use attached to observers via the pipe method, and can be composed together using and_then.

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])
}

#Cancellation

All the fundamental methods of the system, subscribe, on_next, on_error, on_complete, and consume are all async. This means that they can be cancelled. The only error that these methods are allowed to raise is Cancelled, which is how subscriptions are aborted in response to the task running the subscription being cancelled.

When implementing these methods, whenever you call an async function, you may need to check for cancellation. For functions that can only raise due to cancellation, you can use the capture_cancellation utility function to convert such a raise to a raise of Cancelled.

#To Do

  • Javascript interop

#
Cancelled

pub suberror Cancelled {
Cancelled
}

#
NoError

pub suberror NoError

An error type with no constructors, used to indicate an Observable that never raises an error.

#
Ack

pub(all) enum Ack {
Continue
Stop
} derive(
Debug
)

#
CollectToFixedResult

pub struct CollectToFixedResult[T] {
view : ArrayView[T]
overflowed : Bool
}

#
Consumer

type Consumer[T, E, R]

#
Consumer::Consumer

fn[T, E, R] Consumer::Consumer(observer : Observer[T, E], result : async () -> Result[R, E] raise Cancelled) -> Consumer[T, E, R]

#
Consumer::all

fn[T, E] Consumer::all(predicate : (T) -> Bool) -> Consumer[T, E, Bool]

#
Consumer::any

fn[T, E] Consumer::any(predicate : (T) -> Bool) -> Consumer[T, E, Bool]

#
Consumer::collect_all

fn[T, E] Consumer::collect_all(size_hint? : Int) -> Consumer[T, E, Array[T]]

#
Consumer::collect_to_fixed

fn[T, E] Consumer::collect_to_fixed(fixed : FixedArray[T]) -> Consumer[T, E, CollectToFixedResult[T]]

#
Consumer::contains

fn[T : Eq, E] Consumer::contains(item : T) -> Consumer[T, E, Bool]

#
Consumer::contramap

fn[T, E, R, U] Consumer::contramap(self : Consumer[T, E, R], f : (U) -> T) -> Consumer[U, E, R]

#
Consumer::count

fn[T, E] Consumer::count() -> Consumer[T, E, UInt64]

#
Consumer::drain

fn[T, E] Consumer::drain() -> Consumer[T, E, Unit]

#
Consumer::each

fn[T, E : Error] Consumer::each(effect : async (T) -> Unit raise E) -> Consumer[T, E, Unit]

#
Consumer::find

fn[T, E] Consumer::find(predicate : (T) -> Bool) -> Consumer[T, E, T?]

#
Consumer::from_observer

fn[T, E] Consumer::from_observer(observer : Observer[T, E]) -> Consumer[T, E, Unit]

#
Consumer::hash

fn[T : Hash, E] Consumer::hash(seed? : Int) -> Consumer[T, E, Int]

#
Consumer::make

fn[T, E, R] Consumer::make(on_next~ : async (T) -> Ack raise Cancelled, provide_result~ : async () -> Result[R, E] raise Cancelled) -> Consumer[T, E, R]

#
Consumer::make_nofail

fn[T, E, R] Consumer::make_nofail(on_next~ : async (T) -> Ack raise Cancelled, provide_result~ : async () -> R raise Cancelled) -> Consumer[T, E, R]

#
Consumer::map_result

fn[T, E, R, R2] Consumer::map_result(self : Consumer[T, E, R], f : (R) -> R2) -> Consumer[T, E, R2]

#
Consumer::observer

fn[T, E, R] Consumer::observer(self : Consumer[T, E, R]) -> Observer[T, E]

#
Consumer::result

async fn[T, E, R] Consumer::result(self : Consumer[T, E, R]) -> Result[R, E] raise Cancelled

#
Consumer::stop

fn[T, E] Consumer::stop() -> Consumer[T, E, Unit]

#
Consumer::zip

fn[T, E, R1, R2] Consumer::zip(self : Consumer[T, E, R1], other : Consumer[T, E, R2]) -> Consumer[T, E, (R1, R2)]

#
Observable

type Observable[T, E]

#
Observable::Observable

fn[T, E] Observable::Observable(subscribe : async (Observer[T, E]) -> Unit raise Cancelled) -> Observable[T, E]

#
Observable::buffer_tumbling

fn[T, E] Observable::buffer_tumbling(self : Observable[T, E], size : Int) -> Observable[
Vector
[T], E]

#
Observable::catch_error

fn[T, E] Observable::catch_error(self : Observable[T, E], f : (E) -> Observable[T, E]) -> Observable[T, E]

#
Observable::consume

async fn[T, E : Error, R] Observable::consume(self : Observable[T, E], consumer : Consumer[T, E, R]) -> R

#
Observable::consume_result

async fn[T, E, R] Observable::consume_result(self : Observable[T, E], consumer : Consumer[T, E, R]) -> Result[R, E] raise Cancelled

#
Observable::debug_dump

#
Observable::drop

fn[T, E] Observable::drop(self : Observable[T, E], count : UInt) -> Observable[T, E]

#
Observable::empty

fn[T, E] Observable::empty() -> Observable[T, E]

#
Observable::fail

fn[T, E] Observable::fail(error : E) -> Observable[T, E]

#
Observable::filter

fn[T, E] Observable::filter(self : Observable[T, E], predicate : (T) -> Bool) -> Observable[T, E]

#
Observable::flat_map

fn[T, E, U] Observable::flat_map(self : Observable[T, E], f : (T) -> Observable[U, E]) -> Observable[U, E]

Concatenates the inner Observables produced by f in order into a single flat Observable.

#
Observable::from_array

fn[T] Observable::from_array(array : ArrayView[T]) -> Observable[T, NoError]

#
Observable::from_async

fn[T, E : Error] Observable::from_async(task : async () -> T raise E) -> Observable[T, E]

#
Observable::from_async_generator

fn[T, E : Error] Observable::from_async_generator(generator : async () -> T? raise E) -> Observable[T, E]

#
Observable::from_bytes

fn Observable::from_bytes(bytes : BytesView) -> Observable[Byte, NoError]

#
Observable::from_generator

fn[T, E : Error] Observable::from_generator(generator : () -> T? raise E) -> Observable[T, E]

#
Observable::from_iter

fn[T] Observable::from_iter(iter : Iter[T]) -> Observable[T, NoError]

#
Observable::from_value

fn[T, E] Observable::from_value(value : T) -> Observable[T, E]

#
Observable::map

fn[T, E, U] Observable::map(self : Observable[T, E], f : (T) -> U) -> Observable[U, E]

#
Observable::map_error

fn[T, E1, E2] Observable::map_error(self : Observable[T, E1], f : (E1) -> E2) -> Observable[T, E2]

#
Observable::never

fn[T, E] Observable::never() -> Observable[T, E]

#
Observable::panic_on_error

fn[T, E] Observable::panic_on_error(self : Observable[T, E]) -> Observable[T, NoError]

#
Observable::pipe

fn[T1, T2, E] Observable::pipe(self : Observable[T1, E], pipeline : Pipeline[T1, T2, E]) -> Observable[T2, E]

#
Observable::range_between

fn[T : Compare + Add + Eq] Observable::range_between(start~ : T, end~ : T, step~ : T) -> Observable[T, NoError]

#
Observable::range_from

fn[T : Add] Observable::range_from(start~ : T, step~ : T) -> Observable[T, NoError]

#
Observable::recover

fn[T, E] Observable::recover(self : Observable[T, E], f : (E) -> T) -> Observable[T, NoError]

#
Observable::refine_error

fn[T, E1, E2] Observable::refine_error(self : Observable[T, E1], f : (E1) -> E2?) -> Observable[T, E2]

#
Observable::subscribe

async fn[T, E] Observable::subscribe(self : Observable[T, E], observer : Observer[T, E]) -> Unit raise Cancelled

Subscribes the given observer to this observable, blocking until the observable completes, errors, or is cancelled.

#
Observable::take

fn[T, E] Observable::take(self : Observable[T, E], count : UInt) -> Observable[T, E]

#
Observable::tap

fn[T, E : Error] Observable::tap(self : Observable[T, E], effect : async (T) -> Unit raise E) -> Observable[T, E]

#
Observable::tap_error

fn[T, E : Error] Observable::tap_error(self : Observable[T, E], effect : async (E) -> Unit raise E) -> Observable[T, E]

#
Observable::widen_error

fn[T, E] Observable::widen_error(self : Observable[T, NoError]) -> Observable[T, E]

#
Observer

type Observer[T, E]

#
Observer::Observer

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]

#
Observer::before_on_error

fn[T, E : Error] Observer::before_on_error(self : Observer[T, E], effect : async (E) -> Unit raise E) -> Observer[T, E]

#
Observer::before_on_next

fn[T, E : Error] Observer::before_on_next(self : Observer[T, E], effect : async (T) -> Unit raise E) -> Observer[T, E]

#
Observer::contramap

fn[T, E, U] Observer::contramap(self : Observer[T, E], f : (U) -> T) -> Observer[U, E]

#
Observer::contramap_halt

fn[T, E, U] Observer::contramap_halt(self : Observer[T, E], f : (U) -> T?) -> Observer[U, E]

#
Observer::debug_dump

#
Observer::on_complete

async fn[T, E] Observer::on_complete(self : Observer[T, E]) -> Unit raise Cancelled

#
Observer::on_error

async fn[T, E] Observer::on_error(self : Observer[T, E], e : E) -> Unit raise Cancelled

#
Observer::on_next

async fn[T, E] Observer::on_next(self : Observer[T, E], t : T) -> Ack raise Cancelled

#
Observer::replace_on_next

fn[T, E, U] Observer::replace_on_next(self : Observer[T, E], on_next : async (U) -> Ack raise Cancelled) -> Observer[U, E]

#
Pipeline

type Pipeline[T1, T2, E]

#
Pipeline::and_then

fn[T1, T2, T3, E] Pipeline::and_then(self : Pipeline[T1, T2, E], next : Pipeline[T2, T3, E]) -> Pipeline[T1, T3, E]

#
Pipeline::apply

fn[T1, T2, E] Pipeline::apply(self : Pipeline[T1, T2, E], observable : Observable[T1, E]) -> Observable[T2, E]

#
Pipeline::bytes_to_view

fn[E] Pipeline::bytes_to_view() -> Pipeline[Bytes, BytesView, E]

#
Pipeline::chunk_bytes

fn[E] Pipeline::chunk_bytes(chunk_size : Int) -> Pipeline[Byte, Bytes, E]

#
Pipeline::drop_while

fn[T, E] Pipeline::drop_while(predicate : (T) -> Bool) -> Pipeline[T, T, E]

#
Pipeline::flat_map

fn[T1, T2, E] Pipeline::flat_map(f : (T1) -> Observable[T2, E]) -> Pipeline[T1, T2, E]

#
Pipeline::make

fn[T1, T2, E] Pipeline::make(f : (Observable[T1, E]) -> Observable[T2, E]) -> Pipeline[T1, T2, E]

#
Pipeline::map

fn[T1, T2, E] Pipeline::map(f : (T1) -> T2) -> Pipeline[T1, T2, E]

#
Pipeline::split_lines

fn[E] Pipeline::split_lines(max_line_length? : Int) -> Pipeline[Char, String, E]

#
Pipeline::split_on

fn[T : Eq, E] Pipeline::split_on(separator : T) -> Pipeline[T,
Vector
[T], E]

#
Pipeline::string_to_chars

fn[E] Pipeline::string_to_chars() -> Pipeline[String, Char, E]

#
Pipeline::unchunk_bytes

fn[E] Pipeline::unchunk_bytes() -> Pipeline[Bytes, Byte, E]

#
Resource

type Resource[T]

#
Resource::Resource

fn[T] Resource::Resource(acquire : async () -> T, release : (T) -> Unit) -> Resource[T]

#
Resource::acquire

async fn[T] Resource::acquire(self : Resource[T]) -> T

#
Resource::release

fn[T] Resource::release(self : Resource[T], resource : T) -> Unit

#
Resource::use_resource

async fn[T, U] Resource::use_resource(self : Resource[T], f : async (T) -> U) -> U

#
capture_cancellation

async fn[T] capture_cancellation(effect : async () -> T) -> T raise Cancelled

Powered by MoonBit

Site sourceReport issuePackagesBuild queueSkillsStatistics

© 2026 mooncakes.io