///|
using @vector {type Vector}
///|
struct Pipeline[T1, T2, E]((Observable[T1, E]) -> Observable[T2, E])
///|
pub fn[T1, T2, E] Pipeline::make(
f : (Observable[T1, E]) -> Observable[T2, E],
) -> Pipeline[T1, T2, E] {
Pipeline(f)
}
///|
pub fn[T1, T2, T3, E] Pipeline::and_then(
self : Pipeline[T1, T2, E],
next : Pipeline[T2, T3, E],
) -> Pipeline[T1, T3, E] {
Pipeline::make((observable : Observable[T1, E]) => {
let intermediate = self.apply(observable)
next.apply(intermediate)
})
}
///|
pub fn[T1, T2, E] Pipeline::apply(
self : Self[T1, T2, E],
observable : Observable[T1, E],
) -> Observable[T2, E] {
let Pipeline(f) = self
f(observable)
}
///|
pub fn[T1, T2, E] Observable::pipe(
self : Observable[T1, E],
pipeline : Pipeline[T1, T2, E],
) -> Observable[T2, E] {
pipeline.apply(self)
}
///|
pub fn[T1, T2, E] Pipeline::map(f : (T1) -> T2) -> Pipeline[T1, T2, E] {
Pipeline::make((observable : Observable[T1, E]) => observable.map(f))
}
///|
pub fn[T1, T2, E] Pipeline::flat_map(
f : (T1) -> Observable[T2, E],
) -> Pipeline[T1, T2, E] {
Pipeline::make((observable : Observable[T1, E]) => observable.flat_map(f))
}
///|
pub fn[T : Eq, E] Pipeline::split_on(
separator : T,
) -> Pipeline[T, Vector[T], E] {
let f = (observable : Observable[T, E]) => {
let subscribe = (observer : Observer[Vector[T], E]) => {
let mut result = Vector::new()
let inner_observer = Observer::{
on_next: item => {
if item == separator {
let ack = (observer.on_next)(result)
result = Vector::new()
ack
} else {
result = result.push(item)
Continue
}
},
on_complete: () => {
let ack = if !result.is_empty() {
(observer.on_next)(result)
} else {
Continue
}
if ack is Ack::Continue {
(observer.on_complete)()
}
},
on_error: observer.on_error,
}
observable.subscribe(inner_observer)
}
{ subscribe, }
}
Pipeline::make(f)
}
///|
pub fn[E] Pipeline::split_lines(
max_line_length? : Int = 1000,
) -> Pipeline[Char, String, E] {
let f = (observable : Observable[Char, E]) => {
let subscribe = (observer : Observer[String, E]) => {
let buffer = FixedArray::make(max_line_length, '\u{CD}')
let mut pos = 0
let inner_observer = Observer::{
on_next: char => {
if char == '\n' {
let line = String::from_array(buffer[0:pos])
pos = 0
(observer.on_next)(line)
} else if pos == max_line_length {
// line too long, output what we have
let ack = (observer.on_next)(String::from_array(buffer))
buffer[0] = char
pos = 1
ack
} else {
buffer[pos] = char
pos += 1
Continue
}
},
on_complete: () => {
let ack = if pos > 0 {
let line = String::from_array(buffer[0:pos])
(observer.on_next)(line)
} else {
Continue
}
if ack is Ack::Continue {
(observer.on_complete)()
}
},
on_error: observer.on_error,
}
observable.subscribe(inner_observer)
}
{ subscribe, }
}
Pipeline::make(f)
}
///|
pub fn[E] Pipeline::unchunk_bytes() -> Pipeline[Bytes, Byte, E] {
Pipeline::make() <| observable => {
Observable() <| observer => {
let byte_chunk_observer = observer.replace_on_next() <| bytes => {
for b in (bytes : Bytes) {
let ack = observer.on_next(b)
if ack is Ack::Stop {
return Ack::Stop
}
}
Ack::Continue
}
observable.subscribe(byte_chunk_observer)
}
}
}
///|
pub fn[E] Pipeline::chunk_bytes(chunk_size : Int) -> Pipeline[Byte, Bytes, E] {
Pipeline::make() <| observable => {
Observable() <| observer => {
let buffer = FixedArray::make(chunk_size, b'\xCD')
let mut pos = 0
let byte_observer = Observer(
on_next=b => {
buffer[pos] = b
pos += 1
if pos == chunk_size {
pos = 0
observer.on_next(Bytes::from_array(buffer))
} else {
Ack::Continue
}
},
on_complete=() => {
let ack = if pos > 0 {
observer.on_next(Bytes::from_array(buffer[0:pos]))
} else {
Ack::Continue
}
if ack is Ack::Continue {
observer.on_complete()
}
},
on_error=observer.on_error,
)
observable.subscribe(byte_observer)
}
}
}
///|
pub fn[E] Pipeline::bytes_to_view() -> Pipeline[Bytes, BytesView, E] {
Pipeline::map(bytes => bytes[:])
}
///|
pub fn[E] Pipeline::string_to_chars() -> Pipeline[String, Char, E] {
Pipeline::flat_map(string => {
Observable::from_iter(string.iter()).widen_error()
})
}
///|
pub fn[T, E] Pipeline::drop_while(predicate : (T) -> Bool) -> Pipeline[T, T, E] {
Pipeline::make() <| observable => {
let subscribe = (observer : Observer[T, E]) => {
let mut dropping = true
let inner_observer = Observer(
on_next=item => {
if dropping && predicate(item) {
Ack::Continue
} else {
dropping = false
observer.on_next(item)
}
},
on_complete=observer.on_complete,
on_error=observer.on_error,
)
observable.subscribe(inner_observer)
}
{ subscribe, }
}
}