///|
/// Control-channel failures. `AuthFailed` closes the connection without a
/// response, so a wrong token looks to its presenter like a closed port.
pub suberror IpcError {
  Timeout
  AuthFailed
  Closed(String)
  TooLarge(size~ : Int, limit~ : Int)
} derive(Debug)

///|
pub extend IpcError with @debug.Debug::{to_repr}

///|
/// Hard ceiling for one control message. The receiver allocates exactly the
/// announced length, so this is not a limit on how many instances may be
/// listed — it only stops a hostile peer from forcing an absurd allocation.
pub const MAX_MSG_BYTES : Int = 4 * 1024 * 1024

///|
/// How long either side waits for the auth handshake (and the daemon for a
/// silent client) before dropping the connection.
pub const AUTH_TIMEOUT_MS : Int = 5000

///|
/// Files under the state dir that publish the daemon's contact point.
pub const PORT_FILE : String = "daemon.port"

///|
pub const TOKEN_FILE : String = "daemon.token"

///|
/// 4-byte big-endian length header + payload, over the connection.
pub async fn write_frame(conn : @socket.Tcp, payload : Bytes) -> Unit {
  let len = payload.length()
  if len > MAX_MSG_BYTES {
    raise TooLarge(size=len, limit=MAX_MSG_BYTES)
  }
  let head = Bytes::makei(4, fn(i) { ((len >> ((3 - i) * 8)) & 0xFF).to_byte() })
  conn.write(head)
  conn.write(payload)
}

///|
/// Read one framed message: size the buffer to the announced length, then
/// read the body. Sizing to the payload keeps the accept side free of an
/// arbitrary message-size ceiling.
async fn read_frame(conn : @socket.Tcp) -> Bytes {
  let head = conn.read_exactly(4)
  let len = (head[0].to_int() << 24) |
    (head[1].to_int() << 16) |
    (head[2].to_int() << 8) |
    head[3].to_int()
  if len > MAX_MSG_BYTES {
    raise TooLarge(size=len, limit=MAX_MSG_BYTES)
  }
  conn.read_exactly(len)
}

///|
/// Length-safe equality for token comparison: the loop always consumes the
/// full input, so timing does not reveal a matching prefix.
fn constant_time_eq(a : String, b : String) -> Bool {
  let ab = @utf8.encode(a)
  let bb = @utf8.encode(b)
  let mut diff = ab.length() ^ bb.length()
  let n = if ab.length() < bb.length() { ab.length() } else { bb.length() }
  for i in 0.. String raise {
  hex_encode(@fsx.read_small("/dev/urandom", cap=32))
}

///|
pub fn hex_encode(bytes : Bytes) -> String {
  let digits = "0123456789abcdef"
  let out = StringBuilder()
  for i in 0..> 4:(b >> 4) + 1].to_owned())
    out.write_string(digits[b & 15:(b & 15) + 1].to_owned())
  }
  out.to_string()
}

///|
/// Read the port file the daemon wrote after binding. None when there is
/// nothing readable (no daemon, or a stale file from a dead one).
pub fn read_daemon_port(state_dir : String) -> Int? {
  try {
    let text = @utf8.decode(
      @fsx.read_small(state_dir + "/" + PORT_FILE, cap=32),
    )
    parse_port(text)
  } catch {
    _ => None
  } noraise {
    port => port
  }
}

///|
fn parse_port(text : String) -> Int? {
  let trimmed = trim_ascii(text)
  if trimmed.is_empty() {
    return None
  }
  let mut value = 0
  for i in 0.. 0x39 {
      return None
    }
    value = value * 10 + (ch - 0x30)
  }
  if value <= 0 || value > 65535 {
    return None
  }
  Some(value)
}

///|
fn trim_ascii(text : String) -> String {
  let mut start = 0
  let end = text.length()
  while start < end {
    let ch = text[start].to_int()
    if ch == 0x20 || ch == 0x09 || ch == 0x0A || ch == 0x0D {
      start += 1
    } else {
      break
    }
  }
  let mut stop = end
  while stop > start {
    let ch = text[stop - 1].to_int()
    if ch == 0x20 || ch == 0x09 || ch == 0x0A || ch == 0x0D {
      stop -= 1
    } else {
      break
    }
  }
  if start == 0 && stop == end {
    text
  } else if start >= stop {
    ""
  } else {
    text[start:stop].to_owned()
  }
}

///|
/// Read the shared token. None when missing/unreadable — callers report
/// "no daemon" rather than leaking the reason.
pub fn read_token(state_dir : String) -> String? {
  try {
    let text = @utf8.decode(
      @fsx.read_small(state_dir + "/" + TOKEN_FILE, cap=128),
    )
    let trimmed = trim_ascii(text)
    if trimmed.is_empty() {
      None
    } else {
      Some(trimmed)
    }
  } catch {
    _ => None
  } noraise {
    token => token
  }
}

///|
/// Dial the daemon over loopback TCP using the recorded port. Raises
/// Closed when there is no readable port file (no daemon, or a stale one).
pub async fn dial(state_dir : String) -> @socket.Tcp {
  match read_daemon_port(state_dir) {
    Some(port) => @socket.Tcp::connect_to_host("127.0.0.1", port~)
    None => raise Closed("no daemon port file")
  }
}

///|
/// One authenticated round trip on a fresh connection, over raw bytes:
/// the token frame goes first, the server answers with a single ack byte
/// before the payload frame is accepted (a rejected token therefore
/// surfaces as EOF — AuthFailed — rather than a protocol error), then the
/// payload frame is written and the response frame read. The whole
/// exchange is bounded by `timeout_ms`; silence raises Timeout.
pub async fn round_trip_bytes(
  conn : @socket.Tcp,
  token : String,
  payload : Bytes,
  timeout_ms : Int,
) -> Bytes {
  let exchange = async fn() {
    write_frame(conn, @utf8.encode(token))
    try ignore(conn.read_exactly(1)) catch {
      _ => raise AuthFailed
    } noraise {
      _ => ()
    }
    write_frame(conn, payload)
    read_frame(conn)
  }
  match @async.with_timeout_opt(timeout_ms, exchange) {
    Some(response) => response
    None => raise Timeout
  }
}

///|
/// Accept one connection and authenticate it. Returns the open connection
/// when the peer presented the right token (the ack byte is already sent),
/// None for wrong token / silence / garbage — those connections are closed
/// here without ever reaching the dispatcher. Cancellation propagates so a
/// shutting-down daemon stops accepting.
pub async fn accept_authed(
  server : @socket.TcpServer,
  token : String,
) -> @socket.Tcp? {
  let (conn, _peer) = server.accept()
  let ok = try {
    let exchange = async fn() {
      let presented = @utf8.decode(read_frame(conn))
      write_ack(conn, constant_time_eq(presented, token))
    }
    match @async.with_timeout_opt(AUTH_TIMEOUT_MS, exchange) {
      Some(_) => true
      None => false
    }
  } catch {
    _ => false
  } noraise {
    accepted => accepted
  }
  if !ok {
    conn.close()
    return None
  }
  Some(conn)
}

///|
async fn write_ack(conn : @socket.Tcp, accept : Bool) -> Unit {
  if accept {
    conn.write(b"y")
  } else {
    // Any byte works: the client only needs "the server hung up right
    // after auth" versus "no answer at all" to tell a wrong token from a
    // dead daemon. Closing without the ack is the rejection signal.
    conn.close()
    raise AuthFailed
  }
}

///|
/// Read one framed message from an authenticated connection, tolerating a
/// client that connects and goes silent: None means the connection yielded
/// nothing usable and has been closed.
pub async fn read_frame_timeout(conn : @socket.Tcp, timeout_ms : Int) -> Bytes? {
  let exchange = async fn() { read_frame(conn) }
  let received = try @async.with_timeout_opt(timeout_ms, exchange) catch {
    _ => None
  } noraise {
    received => received
  }
  match received {
    Some(frame) => Some(frame)
    None => {
      conn.close()
      None
    }
  }
}