// The main Connection type (port of h11/_connection.py). Everything in h11
// revolves around this.

///|
/// If we ever have this much buffered without it making a complete parseable
/// event, we error out. The only time we really buffer is when reading the
/// request/response line + headers together, so this is effectively the
/// limit on the size of that.
pub const DEFAULT_MAX_INCOMPLETE_EVENT_SIZE : Int = 16 * 1024

///|
/// The result of `Connection::next_event`.
pub(all) enum NextEvent {
  /// A parsed event.
  Event(Event)
  /// You need to read more data from your socket and pass it to
  /// `receive_data` before more events can be returned.
  NeedData
  /// We are not in a state where we can process incoming data (usually
  /// because the peer has finished their part of the current
  /// request/response cycle, and you have not yet called
  /// `start_next_cycle`).
  Paused
} derive(Eq, Debug)

///|
/// How a message body is framed.
priv enum Framing {
  /// `declared` keeps the digits from the header: lengths beyond Int64 are
  /// saturated, but error messages should still quote what the peer said.
  ContentLength(Int64, declared~ : Bytes?)
  Chunked
  Http10
}

///|
/// RFC 7230's rules for connection lifecycles, simplified: if someone says
/// `Connection: close` we will close, and if someone uses HTTP/1.0 we will
/// close (we don't support keep-alive with HTTP/1.0 peers).
fn keep_alive(headers : Headers, http_version : Bytes) -> Bool {
  let connection = get_comma_header(headers, b"connection")
  if connection.contains(b"close") {
    return false
  }
  if bytes_lt(http_version, b"1.1") {
    return false
  }
  true
}

///|
/// Steps 2 and 3 of RFC 7230 section 3.3.3: Transfer-Encoding beats
/// Content-Length. `None` if neither is present.
fn headers_framing(headers : Headers) -> Framing? {
  if headers.chunked {
    Some(Chunked)
  } else if headers.content_length is Some((n, declared)) {
    Some(ContentLength(n, declared=Some(declared)))
  } else {
    None
  }
}

///|
/// Framing of a request body: without framing headers there is no body.
fn request_framing(req : Request) -> Framing {
  headers_framing(req.headers).unwrap_or(ContentLength(0L, declared=None))
}

///|
/// Framing of a response body, given the method of the request it answers.
/// Reference: RFC 7230 section 3.3.3.
fn response_framing(request_method : Bytes?, resp : Response) -> Framing {
  // Step 1: some responses always have an empty body, regardless of what the
  // headers say. (Section 3.3.3 also lists responses with status_code < 200;
  // for us these are InformationalResponses, which never have a body.)
  if resp.status_code is (204 | 304) ||
    request_method is Some(b"HEAD") ||
    (
      request_method is Some(b"CONNECT") &&
      resp.status_code >= 200 &&
      resp.status_code < 300
    ) {
    return ContentLength(0L, declared=None)
  }
  // Step 4: no applicable headers: read until the connection closes.
  headers_framing(resp.headers).unwrap_or(Http10)
}

///|
/// An object encapsulating the state of an HTTP connection.
///
/// It performs no I/O: feed received bytes with `receive_data`, pull parsed
/// events with `next_event`, and turn your own events into bytes with
/// `send`.
pub struct Connection {
  priv our_role : Role
  priv their_role : Role
  priv max_incomplete_event_size : Int
  priv cstate : ConnectionState
  // Reader for their messages, and writer for our message bodies, given the
  // current states
  priv mut body_writer : BodyWriter?
  priv mut reader : Reader?
  // Holds any unprocessed received data
  priv receive_buffer : ReceiveBuffer
  // If this is true, then it indicates that the incoming connection was
  // closed *after* the end of whatever's in receive_buffer
  priv mut receive_buffer_closed : Bool
  // Only used to interpret framing headers for figuring out how to
  // read/write response bodies. their_http_version is also public.
  priv mut their_http_version : Bytes?
  priv mut request_method : Bytes?
  // This is pure flow-control and doesn't at all affect the set of legal
  // transitions, so no need to bother ConnectionState with it
  priv mut client_is_waiting_for_100_continue : Bool
}

///|
/// Create a connection playing `our_role`.
///
/// `max_incomplete_event_size` is the maximum number of bytes we're willing
/// to buffer of an incomplete event. In practice this mostly sets a limit on
/// the maximum size of the request/response line + headers. If this is
/// exceeded, then `next_event` will raise `RemoteProtocolError`.
pub fn Connection::new(
  our_role : Role,
  max_incomplete_event_size? : Int = DEFAULT_MAX_INCOMPLETE_EVENT_SIZE,
) -> Connection {
  let their_role = our_role.other()
  {
    our_role,
    their_role,
    max_incomplete_event_size,
    cstate: ConnectionState::new(),
    body_writer: None,
    reader: reader_for_state(their_role, Idle),
    receive_buffer: ReceiveBuffer::new(),
    receive_buffer_closed: false,
    their_http_version: None,
    request_method: None,
    client_is_waiting_for_100_continue: false,
  }
}

///|
/// The role we are playing.
pub fn Connection::our_role(self : Connection) -> Role {
  self.our_role
}

///|
/// The role our peer is playing.
pub fn Connection::their_role(self : Connection) -> Role {
  self.their_role
}

///|
/// The current state of both the client and the server.
pub fn Connection::states(self : Connection) -> States {
  self.cstate.states
}

///|
/// The current state of whichever role we are playing.
pub fn Connection::our_state(self : Connection) -> State {
  self.cstate.states.get(self.our_role)
}

///|
/// The current state of whichever role we are NOT playing.
pub fn Connection::their_state(self : Connection) -> State {
  self.cstate.states.get(self.their_role)
}

///|
/// The HTTP version our peer used in their last request or response, if
/// any has been received.
pub fn Connection::their_http_version(self : Connection) -> Bytes? {
  self.their_http_version
}

///|
/// Whether the client has sent `Expect: 100-continue` and is still waiting
/// for a response.
pub fn Connection::client_is_waiting_for_100_continue(
  self : Connection,
) -> Bool {
  self.client_is_waiting_for_100_continue
}

///|
/// Whether our peer is a client waiting for a `100 Continue`.
pub fn Connection::they_are_waiting_for_100_continue(self : Connection) -> Bool {
  self.their_role is Client && self.client_is_waiting_for_100_continue
}

///|
/// Attempt to reset our connection state for a new request/response cycle.
///
/// If both client and server are in `Done` state, then resets them both to
/// `Idle` in preparation for a new request/response cycle on this same
/// connection. Otherwise, raises a `LocalProtocolError`.
pub fn Connection::start_next_cycle(
  self : Connection,
) -> Unit raise ProtocolError {
  let old_states = self.cstate.states
  self.cstate.start_next_cycle()
  self.request_method = None
  // their_http_version gets left alone, since it presumably lasts beyond a
  // single request/response cycle. (client_is_waiting_for_100_continue is
  // already false: the client's EndOfMessage cleared it.)
  self.respond_to_state_changes(old_states, None)
}

///|
fn Connection::process_error(self : Connection, role : Role) -> Unit {
  let old_states = self.cstate.states
  self.cstate.process_error(role)
  self.respond_to_state_changes(old_states, None)
}

///|
fn Connection::server_switch_event(
  self : Connection,
  event : Event,
) -> SwitchType? {
  match event {
    InformationalResponse(r) if r.status_code == 101 => Some(SwitchUpgrade)
    Response(r) if self.cstate.pending_switch_proposals.contains(SwitchConnect) &&
      r.status_code >= 200 &&
      r.status_code < 300 => Some(SwitchConnect)
    _ => None
  }
}

///|
/// All events go through here.
fn Connection::process_event(
  self : Connection,
  role : Role,
  event : Event,
) -> Unit raise ProtocolError {
  // First, pass the event through the state machine to make sure it
  // succeeds.
  let old_states = self.cstate.states
  if role is Client && event is Request(req) {
    if req.method_ == b"CONNECT" {
      self.cstate.process_client_switch_proposal(SwitchConnect)
    }
    if !get_comma_header(req.headers, b"upgrade").is_empty() {
      self.cstate.process_client_switch_proposal(SwitchUpgrade)
    }
  }
  let server_switch_event = if role is Server {
    self.server_switch_event(event)
  } else {
    None
  }
  self.cstate.process_event(role, event.event_type(), server_switch_event?)
  // Then perform the updates triggered by it.
  if event is Request(req) {
    self.request_method = Some(req.method_)
  }
  if role == self.their_role {
    match event {
      Request(r) => self.their_http_version = Some(r.http_version)
      Response(r) => self.their_http_version = Some(r.http_version)
      InformationalResponse(r) => self.their_http_version = Some(r.http_version)
      _ => ()
    }
  }
  // Keep alive handling
  //
  // RFC 7230 doesn't really say what one should do if Connection: close
  // shows up on a 1xx InformationalResponse. I think the idea is that this
  // is not supposed to happen. In any case, if it does happen, we ignore it.
  match event {
    Request(r) if !keep_alive(r.headers, r.http_version) =>
      self.cstate.process_keep_alive_disabled()
    Response(r) if !keep_alive(r.headers, r.http_version) =>
      self.cstate.process_keep_alive_disabled()
    _ => ()
  }
  // 100-continue
  if event is Request(req) && has_expect_100_continue(req) {
    self.client_is_waiting_for_100_continue = true
  }
  if event is (InformationalResponse(_) | Response(_)) {
    self.client_is_waiting_for_100_continue = false
  }
  if role is Client && event is (Data(_) | EndOfMessage(_)) {
    self.client_is_waiting_for_100_continue = false
  }
  self.respond_to_state_changes(old_states, Some(event))
}

///|
/// This must be called after any action that might have caused the states
/// to change. `event` is only used when entering SEND_BODY.
fn Connection::respond_to_state_changes(
  self : Connection,
  old_states : States,
  event : Event?,
) -> Unit {
  // Update reader/writer
  if self.our_state() != old_states.get(self.our_role) {
    self.body_writer = match self.our_state() {
      SendBody => Some(writer_for_framing(self.framing_for(event)))
      _ => None
    }
  }
  if self.their_state() != old_states.get(self.their_role) {
    self.reader = match self.their_state() {
      SendBody => Some(reader_for_framing(self.framing_for(event)))
      state => reader_for_state(self.their_role, state)
    }
  }
}

///|
fn Connection::framing_for(self : Connection, event : Event?) -> Framing {
  match event {
    Some(Request(req)) => request_framing(req)
    Some(Response(resp)) => response_framing(self.request_method, resp)
    // Only a Request or a Response can move a party into SEND_BODY.
    _ => abort("entered SEND_BODY without a Request or Response")
  }
}

///|
/// Data that has been received, but not yet processed, together with a flag
/// that is true if the receive connection was closed.
///
/// See the h11 docs on switching protocols for why you'd want this.
pub fn Connection::trailing_data(self : Connection) -> (Bytes, Bool) {
  (self.receive_buffer.to_bytes(), self.receive_buffer_closed)
}

///|
/// Add data to our internal receive buffer.
///
/// This does not actually do any processing on the data, just stores it. To
/// trigger processing, you have to call `next_event`.
///
/// Special case: if `data` is empty, then this indicates that the remote
/// side has closed the connection (end of file). Calling
/// `receive_data(b"")` multiple times is fine, and equivalent to calling it
/// once.
///
/// Raises `RuntimeError` if you pass an empty `data`, indicating EOF, and
/// then pass a non-empty `data`, indicating more data that somehow arrived
/// after the EOF.
pub fn Connection::receive_data(
  self : Connection,
  data : BytesView,
) -> Unit raise RuntimeError {
  if !data.is_empty() {
    if self.receive_buffer_closed {
      raise RuntimeError("received close, then received more data?")
    }
    self.receive_buffer.append(data)
  } else {
    self.receive_buffer_closed = true
  }
}

///|
fn Connection::extract_next_receive_event(
  self : Connection,
) -> NextEvent raise ProtocolError {
  let state = self.their_state()
  // We don't pause immediately when they enter DONE, because even in DONE
  // state we can still process a ConnectionClosed() event. But if we have
  // data in our buffer, then we definitely aren't getting a
  // ConnectionClosed() immediately and we need to pause.
  if state is Done && !self.receive_buffer.is_empty() {
    return Paused
  }
  if state is (MightSwitchProtocol | SwitchedProtocol) {
    return Paused
  }
  guard! self.reader is Some(reader)
  let mut event = reader.read(self.receive_buffer)
  if event is None {
    if self.receive_buffer.is_empty() && self.receive_buffer_closed {
      // In some unusual cases (basically just HTTP/1.0 bodies), EOF
      // triggers an actual protocol event; in that case, we want to return
      // that event, and then the state will change and we'll get called
      // again to generate the actual ConnectionClosed().
      event = match reader.read_eof() {
        Some(e) => Some(e)
        None => Some(ConnectionClosed)
      }
    }
  }
  match event {
    Some(e) => Event(e)
    None => NeedData
  }
}

///|
/// Parse the next event out of our receive buffer, update our internal
/// state, and return it.
///
/// This is a mutating operation -- think of it like calling `next` on an
/// iterator. Returns one of:
///
/// 1. `Event(event)`: an event object.
/// 2. `NeedData`: you need to read more data from your socket and pass it to
///    `receive_data` before this method will be able to return any more
///    events.
/// 3. `Paused`: we are not in a state where we can process incoming data
///    (usually because the peer has finished their part of the current
///    request/response cycle, and you have not yet called
///    `start_next_cycle`).
///
/// Raises `RemoteProtocolError` if the peer has misbehaved. You should close
/// the connection (possibly after sending some kind of 4xx response).
///
/// Once this method returns `ConnectionClosed` once, then all subsequent
/// calls will also return `ConnectionClosed`.
///
/// If this method raises then it also sets `their_state` to `Error`.
pub fn Connection::next_event(
  self : Connection,
) -> NextEvent raise ProtocolError {
  if self.their_state() is Error {
    raise RemoteProtocolError(
      "Can't receive data when peer state is ERROR",
      error_status_hint=400,
    )
  }
  try {
    let event = self.extract_next_receive_event()
    if event is Event(e) {
      self.process_event(self.their_role, e)
    }
    if event is NeedData {
      if self.receive_buffer.length() > self.max_incomplete_event_size {
        // 431 is "Request header fields too large" which is pretty much the
        // only situation where we can get here
        raise RemoteProtocolError(
          "Receive buffer too long",
          error_status_hint=431,
        )
      }
      if self.receive_buffer_closed {
        // We're still trying to complete some event, but that's never going
        // to happen because no more data is coming
        raise RemoteProtocolError(
          "peer unexpectedly closed connection",
          error_status_hint=400,
        )
      }
    }
    event
  } catch {
    exc => {
      self.process_error(self.their_role)
      match exc {
        LocalProtocolError(msg, error_status_hint~) =>
          raise RemoteProtocolError(msg, error_status_hint~)
        exc => raise exc
      }
    }
  }
}

///|
/// Convert a high-level event into bytes that can be sent to the peer, while
/// updating our internal state machine.
///
/// Returns `None` if `event` is `ConnectionClosed`, and the bytes to send
/// otherwise.
///
/// Raises `LocalProtocolError` if sending this event at this time would
/// violate our understanding of the HTTP/1.1 protocol. If this method
/// raises then it also sets `our_state` to `Error`.
pub fn Connection::send(
  self : Connection,
  event : Event,
) -> Bytes? raise ProtocolError {
  match self.send_with_data_passthrough(event) {
    None => None
    Some(data_list) => Some(bytes_concat(data_list))
  }
}

///|
/// Identical to `send`, except that in situations where `send` returns a
/// single byte string, this instead returns a list of them -- and when
/// sending a `Data` event, this list is guaranteed to contain the exact
/// `Bytes` object you passed in as `Data.data`.
pub fn Connection::send_with_data_passthrough(
  self : Connection,
  event : Event,
) -> Array[Bytes]? raise ProtocolError {
  if self.our_state() is Error {
    raise local_error("Can't send data when our state is ERROR")
  }
  errdefer self.process_error(self.our_role)
  let event = match event {
    Response(resp) =>
      Event::Response(self.clean_up_response_headers_for_sending(resp))
    e => e
  }
  // We want to call process_event before writing anything, because if
  // someone tries to do something invalid then this will give a sensible
  // error message, while our writers all just assume they will only receive
  // valid events. But, process_event might replace self.body_writer (e.g.
  // EndOfMessage leaves SEND_BODY), so we grab it first:
  let body_writer = self.body_writer
  self.process_event(self.our_role, event)
  let data_list = []
  match (event, body_writer) {
    (ConnectionClosed, _) => return None
    (Request(r), _) => write_request(r, data_list)
    (InformationalResponse(r), _) =>
      write_any_response(
        r.status_code,
        r.headers,
        r.http_version,
        r.reason,
        data_list,
      )
    (Response(r), _) =>
      write_any_response(
        r.status_code,
        r.headers,
        r.http_version,
        r.reason,
        data_list,
      )
    (Data(d), Some(w)) => w.send_data(d.data, data_list)
    (EndOfMessage(eom), Some(w)) => w.send_eom(eom.headers, data_list)
    // process_event only accepts Data/EndOfMessage in SEND_BODY, where the
    // body writer is always set.
    (Data(_) | EndOfMessage(_), None) => abort("body event outside SEND_BODY")
  }
  Some(data_list)
}

///|
/// Notify the state machine that we failed to send the data it gave us.
///
/// This causes `our_state` to immediately become `Error`.
pub fn Connection::send_failed(self : Connection) -> Unit {
  self.process_error(self.our_role)
}

///|
/// When sending a Response, we take responsibility for a few things:
///
/// - Sometimes you MUST set Connection: close. We take care of those times.
///   (You can also set it yourself if you want, and if you do then we'll
///   respect that and close the connection at the right time.)
///
/// - The user has to set Content-Length if they want it. Otherwise, for
///   responses that have bodies (e.g. not HEAD), we will automatically
///   select the right mechanism for streaming a body of unknown length,
///   which depends on the peer's HTTP version.
///
/// This function's *only* responsibility is making sure headers are set up
/// right -- everything downstream just looks at the headers.
fn Connection::clean_up_response_headers_for_sending(
  self : Connection,
  response : Response,
) -> Response raise ProtocolError {
  let mut headers = response.headers
  let mut need_close = false
  // HEAD requests need some special handling: they always act like they
  // have Content-Length: 0, and that's how response_framing treats them. But
  // their headers are supposed to match what we would send if the request
  // was a GET.
  let method_for_choosing_headers = match self.request_method {
    Some(b"HEAD") => Some(b"GET")
    m => m
  }
  let framing = response_framing(method_for_choosing_headers, response)
  if framing is (Chunked | Http10) {
    // This response has a body of unknown length.
    // If our peer is HTTP/1.1, we use Transfer-Encoding: chunked
    // If our peer is HTTP/1.0, we use no framing headers, and close the
    // connection afterwards.
    //
    // Make sure to clear Content-Length (in principle user could have set
    // both and then we ignored Content-Length b/c Transfer-Encoding
    // overwrote it -- this would be naughty of them, but the HTTP spec says
    // that if our peer does this then we have to fix it instead of erroring
    // out, so we'll accord the user the same respect).
    headers = set_comma_header(headers, b"content-length", [])
    let peer_is_http10 = match self.their_http_version {
      None => true
      Some(v) => bytes_lt(v, b"1.1")
    }
    if peer_is_http10 {
      // Either we never got a valid request and are sending back an error
      // (their_http_version is None), so we assume the worst; or else we did
      // get a valid HTTP/1.0 request, so we know that they don't understand
      // chunked encoding.
      headers = set_comma_header(headers, b"transfer-encoding", [])
      // This is actually redundant ATM, since currently we unconditionally
      // disable keep-alive when talking to HTTP/1.0 peers. But let's be
      // defensive just in case we add Connection: keep-alive support later:
      if !(self.request_method is Some(b"HEAD")) {
        need_close = true
      }
    } else {
      headers = set_comma_header(headers, b"transfer-encoding", [b"chunked"])
    }
  }
  if !self.cstate.keep_alive || need_close {
    // Make sure Connection: close is set
    let connection_set : Map[Bytes, Unit] = Map([])
    for v in get_comma_header(headers, b"connection") {
      connection_set[v] = ()
    }
    connection_set.remove(b"keep-alive")
    connection_set[b"close"] = ()
    let connection = connection_set.keys().to_array()
    connection.sort_by((a, b) => a.lexical_compare(b))
    headers = set_comma_header(headers, b"connection", connection)
  }
  Response::make(
    response.status_code,
    headers,
    response.http_version,
    response.reason,
  )
}