///|
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, }
  }
}