rx
Types
pub type Emitter(value, error) {
Emitter(emit: fn(protocol.Notification(value, error)) -> Nil)
}
Constructors
-
Emitter(emit: fn(protocol.Notification(value, error)) -> Nil)
pub opaque type Observable(value, error)
pub type Observer(value, error) {
Observer(
on_next: fn(value) -> Nil,
on_error: fn(error) -> Nil,
on_complete: fn() -> Nil,
)
}
Constructors
-
Observer( on_next: fn(value) -> Nil, on_error: fn(error) -> Nil, on_complete: fn() -> Nil, )
pub opaque type Subscription
Values
pub fn create(
producer: fn(Emitter(value, error)) -> fn() -> Nil,
) -> Observable(value, error)
Create an Observable from a callback producer.
The producer may call the emitter now or at any later time. This is the primitive used to adapt asynchronous push sources such as sockets, queues, actor messages, timers, and callback APIs.
pub fn create_checked(
producer: fn(runtime.Runtime, Emitter(value, error)) -> Result(
fn() -> Nil,
runtime.RuntimeError,
),
) -> Observable(value, error)
Runtime-aware Observable construction with typed subscription failure.
pub fn create_with_runtime(
producer: fn(runtime.Runtime, Emitter(value, error)) -> fn() -> Nil,
) -> Observable(value, error)
Create an Observable whose producer can access the owning Runtime.
This is intended for advanced operators that need to register serialized state with the runtime while still accepting asynchronous source emissions.
pub fn emit(
emitter: Emitter(value, error),
notification: protocol.Notification(value, error),
) -> Nil
pub fn empty() -> Observable(value, error)
pub fn fail(reason: error) -> Observable(value, error)
pub fn filter(
observable: Observable(value, error),
predicate: fn(value) -> Bool,
) -> Observable(value, error)
pub fn from_list(values: List(value)) -> Observable(value, error)
pub fn map(
observable: Observable(a, error),
transform: fn(a) -> b,
) -> Observable(b, error)
pub fn observer(
on_next: fn(value) -> Nil,
on_error: fn(error) -> Nil,
on_complete: fn() -> Nil,
) -> Observer(value, error)
pub fn of(value: value) -> Observable(value, error)
pub fn subscribe(
observable: Observable(value, error),
runtime_: runtime.Runtime,
observer_: Observer(value, error),
) -> Result(Subscription, runtime.RuntimeError)
pub fn tap(
observable: Observable(value, error),
inspect: fn(value) -> Nil,
) -> Observable(value, error)
pub fn unsubscribe(subscription: Subscription) -> Nil