///|
struct Observable[T, E] {
  subscribe : async (Observer[T, E]) -> Unit raise Cancelled
}

///|
pub fn[T, E] Observable::Observable(
  subscribe : async (Observer[T, E]) -> Unit raise Cancelled,
) -> Observable[T, E] {
  { subscribe, }
}

///|
/// Subscribes the given observer to this observable, blocking until the observable completes,
/// errors, or is cancelled.
pub async fn[T, E] Observable::subscribe(
  self : Self[T, E],
  observer : Observer[T, E],
) -> Unit raise Cancelled {
  (self.subscribe)(observer)
}

///|
pub fn[T] Observable::from_array(
  array : ArrayView[T],
) -> Observable[T, NoError] {
  Observable() <| observer => {
    for item in array {
      if observer.on_next(item) is Ack::Stop {
        break
      }
    } nobreak {
      observer.on_complete()
    }
  }
}

///|
pub fn[T] Observable::from_iter(iter : Iter[T]) -> Observable[T, NoError] {
  Observable() <| observer => {
    for item in iter {
      if observer.on_next(item) is Ack::Stop {
        break
      }
    } nobreak {
      observer.on_complete()
    }
  }
}

///|
pub fn Observable::from_bytes(bytes : BytesView) -> Observable[Byte, NoError] {
  Observable() <| observer => {
    for byte in bytes {
      if observer.on_next(byte) is Ack::Stop {
        break
      }
    } nobreak {
      observer.on_complete()
    }
  }
}

///|
pub fn[T, E] Observable::empty() -> Observable[T, E] {
  Observable() <| observer => { (observer.on_complete)() }
}

///|
pub fn[T, E] Observable::fail(error : E) -> Observable[T, E] {
  Observable() <| observer => { (observer.on_error)(error) }
}

///|
pub fn[T : Compare + Add] Observable::range_between(
  start~ : T,
  end~ : T,
  step~ : T,
) -> Observable[T, NoError] {
  Observable() <| observer => {
    let mut item = start
    while item < end {
      if observer.on_next(item) is Ack::Stop {
        break
      }
      item += step
    } nobreak {
      (observer.on_complete)()
    }
  }
}

///|
pub fn[T : Add] Observable::range_from(
  start~ : T,
  step~ : T,
) -> Observable[T, NoError] {
  Observable() <| observer => {
    let mut item = start
    while true {
      if (observer.on_next)(item) is Ack::Stop {
        break
      }
      item += step
    } nobreak {
      (observer.on_complete)()
    }
  }
}

///|
pub fn[T, E : Error] Observable::from_async_generator(
  generator : async () -> T? raise E,
) -> Observable[T, E] {
  Observable() <| observer => {
    let mut generator_error : E? = None
    async fn run_generator() -> T? noraise {
      generator() catch {
        e => {
          generator_error = Some(e)
          None
        }
      }
    }
    while run_generator() is Some(item) {
      if observer.on_next(item) is Ack::Stop {
        break
      }
    } nobreak {
      match generator_error {
        Some(e) => observer.on_error(e)
        None => observer.on_complete()
      }
    }
  }
}

///|
pub fn[T, E : Error] Observable::from_async(
  task : async () -> T raise E,
) -> Observable[T, E] {
  Observable() <| observer => {
    let maybe_value = Some(task()) catch {
      e => {
        observer.on_error(e)
        None
      }
    }
    if maybe_value is Some(value) {
      if observer.on_next(value) is Ack::Continue {
        observer.on_complete()
      }
    }
  }
}

///|
pub fn[T, E : Error] Observable::from_generator(
  generator : () -> T? raise E,
) -> Observable[T, E] {
  Observable() <| observer => {
    let mut generator_error : E? = None
    fn run_generator() {
      generator() catch {
        e => {
          generator_error = Some(e)
          None
        }
      }
    }
    while run_generator() is Some(item) {
      if observer.on_next(item) is Ack::Stop {
        break
      }
    } nobreak {
      match generator_error {
        Some(e) => observer.on_error(e)
        None => observer.on_complete()
      }
    }
  }
}

///|
pub fn[T, E] Observable::never() -> Observable[T, E] {
  Observable() <| _ => {
    // never completes unless cancelled
    capture_cancellation <| () => {
      while true {
        @async.sleep(@int.MAX_VALUE)
      }
    }
  }
}

///|
pub fn[T, E] Observable::from_value(value : T) -> Observable[T, E] {
  Observable() <| observer => {
    if (observer.on_next)(value) is Ack::Continue {
      (observer.on_complete)()
    }
  }
}

///|
pub fn[T, E, U] Observable::map(
  self : Observable[T, E],
  f : (T) -> U,
) -> Observable[U, E] {
  let subscribe = (observer : Observer[U, E]) => {
    let inner_observer = observer.contramap(f)
    (self.subscribe)(inner_observer)
  }
  { subscribe, }
}

///|
/// Concatenates the inner `Observable`s produced by `f` in order into a single flat `Observable`.
pub fn[T, E, U] Observable::flat_map(
  self : Self[T, E],
  f : (T) -> Self[U, E],
) -> Self[U, E] {
  let subscribe = (observer : Observer[U, E]) => {
    let inner_observer = Observer::{
      on_next: (item : T) => {
        let inner_observable = f(item)
        let mut ack = Ack::Continue
        let flatten_observer = Observer::{
          on_next: item => {
            ack = (observer.on_next)(item)
            ack
          },
          on_complete: () => (),
          on_error: e => {
            ack = Ack::Stop
            (observer.on_error)(e)
          },
        }
        inner_observable.subscribe(flatten_observer)
        ack
      },
      on_complete: observer.on_complete,
      on_error: observer.on_error,
    }
    (self.subscribe)(inner_observer)
  }
  { subscribe, }
}

///|
pub fn[T, E] Observable::filter(
  self : Self[T, E],
  predicate : (T) -> Bool,
) -> Self[T, E] {
  let subscribe = (observer : Observer[T, E]) => {
    let inner_observer = {
      ..observer,
      on_next: item => {
        if predicate(item) {
          (observer.on_next)(item)
        } else {
          Ack::Continue
        }
      },
    }
    self.subscribe(inner_observer)
  }
  { subscribe, }
}

///|
pub fn[T, E : Error] Observable::tap(
  self : Self[T, E],
  effect : async (T) -> Unit raise E,
) -> Self[T, E] {
  let subscribe = (observer : Observer[T, E]) => {
    let inner_observer = observer.before_on_next() <| effect
    (self.subscribe)(inner_observer)
  }
  { subscribe, }
}

///|
pub fn[T, E : Error] Observable::tap_error(
  self : Self[T, E],
  effect : async (E) -> Unit raise E,
) -> Self[T, E] {
  let subscribe = (observer : Observer[T, E]) => {
    let inner_observer = observer.before_on_error() <| effect
    (self.subscribe)(inner_observer)
  }
  { subscribe, }
}

///|
pub fn[T : Debug, E : Debug] Observable::debug_dump(
  self : Self[T, E],
) -> Self[T, E] {
  let subscribe = (observer : Observer[T, E]) => {
    let inner_observer = Observer::{
      on_next: (item : T) => {
        debug("on_next: \{item |> dbg()}")
        (observer.on_next)(item)
      },
      on_complete: () => {
        debug("on_complete")
        (observer.on_complete)()
      },
      on_error: (e : E) => {
        debug("on_error: \{e |> dbg()}")
        (observer.on_error)(e)
      },
    }
    (self.subscribe)(inner_observer)
  }
  { subscribe, }
}

///|
pub fn[T, E] Observable::take(self : Self[T, E], count : UInt) -> Self[T, E] {
  if count == 0 {
    return Observable::empty()
  }
  let subscribe = (observer : Observer[T, E]) => {
    let mut remaining = count
    let inner_observer = {
      ..observer,
      on_next: item => {
        remaining -= 1
        let ack = (observer.on_next)(item)
        if remaining == 0 {
          (observer.on_complete)()
          Ack::Stop
        } else {
          ack
        }
      },
    }
    (self.subscribe)(inner_observer)
  }
  { subscribe, }
}

///|
pub fn[T, E] Observable::drop(self : Self[T, E], count : UInt) -> Self[T, E] {
  if count == 0 {
    return self
  }
  let subscribe = (observer : Observer[T, E]) => {
    let mut remaining = count
    let inner_observer = {
      ..observer,
      on_next: item => {
        if remaining > 0 {
          remaining -= 1
          Continue
        } else {
          (observer.on_next)(item)
        }
      },
    }
    (self.subscribe)(inner_observer)
  }
  { subscribe, }
}

///|
pub fn[T, E1, E2] Observable::map_error(
  self : Self[T, E1],
  f : (E1) -> E2,
) -> Self[T, E2] {
  let subscribe = (observer : Observer[T, E2]) => {
    let inner_observer = {
      on_next: observer.on_next,
      on_complete: observer.on_complete,
      on_error: e => (observer.on_error)(f(e)),
    }
    (self.subscribe)(inner_observer)
  }
  { subscribe, }
}

///|
pub fn[T, E] Observable::catch_error(
  self : Self[T, E],
  f : (E) -> Self[T, E],
) -> Self[T, E] {
  let subscribe = (observer : Observer[T, E]) => {
    let inner_observer = {
      ..observer,
      on_error: e => {
        let new_observable = f(e)
        (new_observable.subscribe)(observer)
      },
    }
    (self.subscribe)(inner_observer)
  }
  { subscribe, }
}

///|
pub fn[T, E] Observable::recover(
  self : Self[T, E],
  f : (E) -> T,
) -> Self[T, NoError] {
  let subscribe = (observer : Observer[T, NoError]) => {
    let inner_observer = {
      on_next: observer.on_next,
      on_complete: observer.on_complete,
      on_error: e => {
        let recovered_value = f(e)
        if (observer.on_next)(recovered_value) is Continue {
          (observer.on_complete)()
        }
      },
    }
    (self.subscribe)(inner_observer)
  }
  { subscribe, }
}

///|
pub fn[T, E1, E2] Observable::refine_error(
  self : Self[T, E1],
  f : (E1) -> E2?,
) -> Self[T, E2] {
  Observable() <| observer => {
    let inner_observer = {
      on_next: observer.on_next,
      on_complete: observer.on_complete,
      on_error: e => {
        match f(e) {
          Some(refined_error) => observer.on_error(refined_error)
          None => panic()
        }
      },
    }
    self.subscribe(inner_observer)
  }
}

///|
pub fn[T, E] Observable::panic_on_error(self : Self[T, E]) -> Self[T, NoError] {
  Observable() <| observer => {
    let inner_observer = {
      on_next: observer.on_next,
      on_complete: observer.on_complete,
      on_error: _ => panic(),
    }
    self.subscribe(inner_observer)
  }
}

///|
pub fn[T, E] Observable::widen_error(self : Self[T, NoError]) -> Self[T, E] {
  Observable() <| observer => {
    let inner_observer = {
      on_next: observer.on_next,
      on_complete: observer.on_complete,
      on_error: _ => panic(),
    }
    self.subscribe(inner_observer)
  }
}

///|
pub fn[T, E] Observable::buffer_tumbling(
  self : Self[T, E],
  size : Int,
) -> Self[Vector[T], E] {
  Observable() <| observer => {
    let mut buffer = Vector::new()
    let inner_observer = {
      on_next: item => {
        buffer = buffer.push(item)
        if buffer.length() == size {
          let to_emit = buffer
          buffer = Vector::new()
          if (observer.on_next)(to_emit) is Stop {
            return Stop
          }
        }
        Continue
      },
      on_complete: () => {
        let ack = if !buffer.is_empty() {
          (observer.on_next)(buffer)
        } else {
          Continue
        }
        if ack is Continue {
          (observer.on_complete)()
        }
      },
      on_error: observer.on_error,
    }
    (self.subscribe)(inner_observer)
  }
}