///|
struct Consumer[T, E, R] {
observer : Observer[T, E]
result : async () -> Result[R, E] raise Cancelled
}
///|
pub fn[T, E, R] Consumer::Consumer(
observer : Observer[T, E],
result : async () -> Result[R, E] raise Cancelled,
) -> Consumer[T, E, R] {
{ observer, result }
}
///|
pub fn[T, E, R] Consumer::observer(self : Self[T, E, R]) -> Observer[T, E] {
self.observer
}
///|
pub async fn[T, E, R] Consumer::result(
self : Self[T, E, R],
) -> Result[R, E] raise Cancelled {
(self.result)()
}
///|
pub fn[T, E, R, U] Consumer::contramap(
self : Self[T, E, R],
f : (U) -> T,
) -> Self[U, E, R] {
let observer = self.observer.contramap(f)
let result = self.result
Consumer::{ observer, result }
}
///|
pub fn[T, E, R, R2] Consumer::map_result(
self : Self[T, E, R],
f : (R) -> R2,
) -> Self[T, E, R2] {
let observer = self.observer
let result = () => self.result().map(f)
Consumer::{ observer, result }
}
///|
pub 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] {
let result_queue = @aqueue.Queue(kind=Blocking(1))
let observer = Observer(
on_next=item => {
let ack = on_next(item)
if ack is Ack::Stop {
capture_cancellation <| () => { result_queue.put(provide_result()) }
}
ack
},
on_complete=() => {
capture_cancellation <| () => { result_queue.put(provide_result()) }
},
on_error=e => capture_cancellation <| () => { result_queue.put(Err(e)) },
)
let result = () => capture_cancellation <| () => { result_queue.get() }
{ observer, result }
}
///|
pub 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::make(on_next~, provide_result=() => Ok(provide_result()))
}
///|
pub fn[T, E] Consumer::collect_all(
size_hint? : Int,
) -> Consumer[T, E, Array[T]] {
let array = Array::new(capacity?=size_hint)
Consumer::make_nofail(
on_next=item => {
array.push(item)
Ack::Continue
},
provide_result=() => array,
)
}
///|
pub struct CollectToFixedResult[T] {
view : ArrayView[T]
overflowed : Bool
}
///|
pub fn[T, E] Consumer::collect_to_fixed(
fixed : FixedArray[T],
) -> Consumer[T, E, CollectToFixedResult[T]] {
let mut index : Int = 0
let mut overflowed = false
Consumer::make_nofail(
on_next=item => {
if index < fixed.length() {
fixed[index] = item
index += 1
Ack::Continue
} else {
overflowed = true
Ack::Stop
}
},
provide_result=() => { view: fixed.view()[0:index], overflowed },
)
}
///|
pub fn[T, E] Consumer::count() -> Consumer[T, E, UInt64] {
let mut count : UInt64 = 0
Consumer::make_nofail(
on_next=_ => {
count += 1
Ack::Continue
},
provide_result=() => count,
)
}
///|
pub fn[T, E] Consumer::all(predicate : (T) -> Bool) -> Consumer[T, E, Bool] {
let mut all = true
Consumer::make_nofail(
on_next=item => {
if !predicate(item) {
all = false
Ack::Stop
} else {
Ack::Continue
}
},
provide_result=() => all,
)
}
///|
pub fn[T, E] Consumer::any(predicate : (T) -> Bool) -> Consumer[T, E, Bool] {
let mut any = false
Consumer::make_nofail(
on_next=item => {
if predicate(item) {
any = true
Ack::Stop
} else {
Ack::Continue
}
},
provide_result=() => any,
)
}
///|
pub fn[T, E] Consumer::find(predicate : (T) -> Bool) -> Consumer[T, E, T?] {
let mut found : T? = None
Consumer::make_nofail(
on_next=item => {
if predicate(item) {
found = Some(item)
Ack::Stop
} else {
Ack::Continue
}
},
provide_result=() => found,
)
}
///|
pub fn[T : Eq, E] Consumer::contains(item : T) -> Consumer[T, E, Bool] {
Consumer::any(x => x == item)
}
///|
pub fn[T : Hash, E] Consumer::hash(seed? : Int = 0) -> Consumer[T, E, Int] {
let hasher = Hasher(seed~)
Consumer::make_nofail(
on_next=item => {
hasher.combine(item)
Ack::Continue
},
provide_result=() => hasher.finalize(),
)
}
///|
pub fn[T, E, R1, R2] Consumer::zip(
self : Self[T, E, R1],
other : Self[T, E, R2],
) -> Self[T, E, (R1, R2)] {
async fn zip_results() {
let other_result = other.result()
self.result().bind() <| r1 => { other_result.map() <| r2 => { (r1, r2) } }
}
let mut self_finished = false
let mut other_finished = false
let result_queue = @aqueue.Queue(kind=Blocking(1))
let observer = {
on_next: (item : T) => {
if !self_finished {
let ack = (self.observer.on_next)(item)
if ack is Stop {
self_finished = true
}
}
if !other_finished {
let ack = (other.observer.on_next)(item)
if ack is Stop {
other_finished = true
}
}
if self_finished && other_finished {
capture_cancellation <| () => { result_queue.put(zip_results()) }
Ack::Stop
} else {
Ack::Continue
}
},
on_complete: () => {
if !self_finished {
(self.observer.on_complete)()
}
if !other_finished {
(other.observer.on_complete)()
}
capture_cancellation <| () => { result_queue.put(zip_results()) }
},
on_error: (e : E) => {
if !self_finished {
(self.observer.on_error)(e)
}
if !other_finished {
(other.observer.on_error)(e)
}
capture_cancellation <| () => { result_queue.put(Err(e)) }
},
}
let result = () => capture_cancellation() <| () => { result_queue.get() }
{ observer, result }
}
///|
pub fn[T, E] Consumer::from_observer(
observer : Observer[T, E],
) -> Consumer[T, E, Unit] {
let result_queue = @aqueue.Queue(kind=Blocking(1))
let inner_observer = Observer::{
on_next: (item : T) => {
let ack = (observer.on_next)(item)
if ack is Ack::Stop {
capture_cancellation <| () => { result_queue.put(Ok(())) }
}
ack
},
on_complete: () => {
(observer.on_complete)()
capture_cancellation <| () => { result_queue.put(Ok(())) }
},
on_error: (e : E) => {
(observer.on_error)(e)
capture_cancellation <| () => { result_queue.put(Err(e)) }
},
}
let result = () => capture_cancellation() <| () => { result_queue.get() }
{ observer: inner_observer, result }
}
///|
pub fn[T, E] Consumer::drain() -> Consumer[T, E, Unit] {
Consumer::from_observer(Observer::ignore())
}
///|
pub fn[T, E] Consumer::stop() -> Consumer[T, E, Unit] {
Consumer::from_observer(Observer::stop())
}
///|
pub fn[T, E : Error] Consumer::each(
effect : async (T) -> Unit raise E,
) -> Consumer[T, E, Unit] {
let observer = Observer::ignore().before_on_next(effect)
Consumer::from_observer(observer)
}
///|
pub async fn[T, E, R] Observable::consume_result(
self : Self[T, E],
consumer : Consumer[T, E, R],
) -> Result[R, E] raise Cancelled {
self.subscribe(consumer.observer)
consumer.result()
}
///|
pub async fn[T, E : Error, R] Observable::consume(
self : Self[T, E],
consumer : Consumer[T, E, R],
) -> R {
self.consume_result(consumer).unwrap_or_error()
}