///|
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()
}