///|
pub(all) enum Ack {
  Continue
  Stop
} derive(Debug)

///|
struct Observer[T, E] {
  on_next : async (T) -> Ack raise Cancelled
  on_complete : async () -> Unit raise Cancelled
  on_error : async (E) -> Unit raise Cancelled
}

///|
pub fn[T, E] Observer::Observer(
  on_next~ : async (T) -> Ack raise Cancelled,
  on_complete~ : async () -> Unit raise Cancelled,
  on_error~ : async (E) -> Unit raise Cancelled,
) -> Observer[T, E] {
  { on_next, on_complete, on_error }
}

///|
pub async fn[T, E] Observer::on_next(
  self : Self[T, E],
  t : T,
) -> Ack raise Cancelled {
  (self.on_next)(t)
}

///|
pub async fn[T, E] Observer::on_complete(
  self : Self[T, E],
) -> Unit raise Cancelled {
  (self.on_complete)()
}

///|
pub async fn[T, E] Observer::on_error(
  self : Self[T, E],
  e : E,
) -> Unit raise Cancelled {
  (self.on_error)(e)
}

///|
pub fn[T, E, U] Observer::replace_on_next(
  self : Self[T, E],
  on_next : async (U) -> Ack raise Cancelled,
) -> Self[U, E] {
  { on_next, on_complete: self.on_complete, on_error: self.on_error }
}

///|
pub fn[T, E, U] Observer::contramap(
  self : Self[T, E],
  f : (U) -> T,
) -> Self[U, E] {
  let on_next : async (U) -> Ack raise Cancelled = (item : U) => {
    let mapped = f(item)
    (self.on_next)(mapped)
  }
  Observer::{ on_next, on_complete: self.on_complete, on_error: self.on_error }
}

///|
pub fn[T, E, U] Observer::contramap_halt(
  self : Self[T, E],
  f : (U) -> T? noraise,
) -> Self[U, E] {
  let on_next : async (U) -> Ack raise Cancelled = (item : U) => {
    match f(item) {
      Some(mapped) => (self.on_next)(mapped)
      None => Ack::Stop
    }
  }
  { on_next, on_complete: self.on_complete, on_error: self.on_error }
}

///|
pub fn[T, E : Error] Observer::before_on_next(
  self : Self[T, E],
  effect : async (T) -> Unit raise E,
) -> Self[T, E] {
  let on_next = (item : T) => {
    effect(item) catch {
      e => {
        self.on_error(e)
        return Ack::Stop
      }
    }
    self.on_next(item)
  }
  { ..self, on_next, }
}

///|
pub fn[T, E : Error] Observer::before_on_error(
  self : Self[T, E],
  effect : async (E) -> Unit raise E,
) -> Self[T, E] {
  let on_error = (e : E) => {
    effect(e) catch {
      e2 => (self.on_error)(e2)
    }
    self.on_error(e)
  }
  { ..self, on_error, }
}

///|
fn[T, E] Observer::ignore() -> Observer[T, E] {
  Observer(on_next=_ => Ack::Continue, on_complete=() => (), on_error=_ => ())
}

///|
fn[T, E] Observer::stop() -> Observer[T, E] {
  Observer(on_next=_ => Ack::Stop, on_complete=() => (), on_error=_ => ())
}

///|
pub fn[T : Debug, E : Debug] Observer::debug_dump() -> Observer[T, E] {
  Observer(
    on_next=(item : T) => {
      debug("on_next: \{item |> dbg()}")
      Ack::Continue
    },
    on_complete=() => debug("on_complete"),
    on_error=(e : E) => debug("on_error: \{e |> dbg()}"),
  )
}