package aceboot::http

import aceboot::web.*
import stdx.net.http.*
import stdx.net.tls.common.*
import std.collection.*
import std.io.*
import std.time.*

// OS 信号处理已移除(Cangjie foreign func 不支持函数指针类型作为 CType 参数)
// 优雅停机通过 ListenOptions.onShutdown 回调由调用方触发

/**
 * catch-all 分发器:把所有路径的请求都交给同一入口,进 ace-web 洋葱处理。
 * 绕开 stdx 默认分发器的精确路径匹配,路由交由 ace-router 在洋葱内完成。
 */
class AceDistributor <: HttpRequestDistributor {
    let app: App
    let streamRequestBodies: Bool
    let streamRequestBodyLimit: Int64
    let streamDrainLimit: Int64

    init(app: App, streamRequestBodies: Bool, streamRequestBodyLimit: Int64, streamDrainLimit: Int64) {
        this.app = app
        this.streamRequestBodies = streamRequestBodies
        this.streamRequestBodyLimit = streamRequestBodyLimit
        this.streamDrainLimit = streamDrainLimit
    }

    public func register(_: String, _: HttpRequestHandler): Unit {}
    public func register(_: String, _: (HttpContext) -> Unit): Unit {}

    public func distribute(_: String): HttpRequestHandler {
        let a = this.app
        let streaming = this.streamRequestBodies
        let limit = this.streamRequestBodyLimit
        let drainLimit = this.streamDrainLimit
        FuncHandler({httpCtx => dispatch(a, httpCtx, streaming, limit, drainLimit)})
    }
}

/** 分块读取请求体为字节(InputStream.read 循环到 EOF;二进制安全,不经 UTF-8)。 */
func readBodyBytes(body: InputStream): Array<UInt8> {
    let out = ArrayList<UInt8>()
    let buf = Array<UInt8>(4096, repeat: 0)
    while (true) {
        let n = body.read(buf)
        if (n <= 0) {
            break
        }
        for (i in 0..n) {
            out.add(buf[i])
        }
    }
    return out.toArray()
}

/** 把 stdx 请求翻译成 ace-web Context,跑洋葱,再写回响应。 */
func dispatch(app: App, httpCtx: HttpContext, streamRequestBodies: Bool, streamRequestBodyLimit: Int64,
    streamDrainLimit: Int64): Unit {
    let req = httpCtx.request
    // WebSocket 升级请求在构造 Context/读 body/跑洋葱之前分流(升级会接管连接)。
    if (isWebSocketUpgrade(req) && handleWsUpgrade(httpCtx, req.url.path)) {
        return
    }
    let headers = HashMap<String, String>()
    for ((k, vs) in req.headers) {
        for (v in vs) {
            headers.add(toLowerAscii(k), v) // 头键统一小写,便于 ctx.header("authorization") 等查找
            break
        }
    }
    let ctx = Context(req.method, req.url.path, headers)
    applyQueryString(ctx, req.url.query ?? "")
    var requestStream: ?RequestBodyStream = None
    if (streamRequestBodies) {
        let stream = RequestBodyStream({buffer: Array<UInt8> => req.body.read(buffer)}, streamRequestBodyLimit)
        requestStream = Some(stream)
        ctx.setRequestBodyStream(stream)
    } else {
        let bodyBytes = readBodyBytes(req.body)
        if (bodyBytes.size > 0) {
            ctx.setRawBodyBytes(bodyBytes)
        }
    }
    // 绑定当前协程的请求 Context,供 Service 等非 Controller 代码静态读取(RequestContextHolder.current())。
    // try/finally 保证请求结束必清理:stdx 每连接一协程,keep-alive 下多请求串行复用同一协程,不清会串到下个请求。
    RequestContextHolder.set(ctx)
    try {
        app.handle(ctx)
        let rb = httpCtx.responseBuilder.status(ctx.status)
        for ((k, v) in ctx.responseHeaders()) {
            rb.header(k, v)
        }
        if (ctx.contentType != "") {
            rb.header("Content-Type", ctx.contentType)
        }
        // Set-Cookie 逐条回写(响应头容器是单值 map,多 cookie 须单独通道避免互相覆盖)。
        for (c in ctx.responseCookies()) {
            rb.header("Set-Cookie", c)
        }
        // 统一响应体一次分流:流式(SSE/大文件)→ chunked + HttpResponseWriter 边产边发;
        // 二进制 → body(Array<UInt8>);文本/空 → 一次性 String body。
        match (ctx.responseBody()) {
            case StreamBody(producer) =>
                // Transfer-Encoding 是 connection 级头,HTTP/2 禁止(RFC 9113 §8.2.2):
                // h2 由 DATA 帧天然分块,仅 HTTP/1.x 需显式 chunked。
                match (req.version) {
                    case HTTP2_0 => ()
                    case _ => rb.header("Transfer-Encoding", "chunked")
                }
                let writer = HttpResponseWriter(httpCtx)
                producer({bytes: Array<UInt8> => writer.write(bytes)})
            case BytesBody(bytes) => rb.body(bytes)
            case TextBody(s) => rb.body(s)
            case NoBody => rb.body("")
        }
    } finally {
        // keep-alive 下必须消费掉业务未读取的请求体,否则下一请求可能读到残留字节。
        match (requestStream) {
            case Some(stream) =>
                if (!stream.drained()) {
                    let drainBuffer = Array<UInt8>(8192, repeat: 0)
                    var drained: Int64 = 0
                    while (streamDrainLimit <= 0 || drained < streamDrainLimit) {
                        let n = req.body.read(drainBuffer)
                        if (n <= 0) { stream.markDrained(); break }
                        drained += n
                    }
                }
            case None => ()
        }
        RequestContextHolder.clear()
    }
}

/**
 * 监听选项:连接级读写超时(stdx 服务端原生,非 racy 中间件)、请求大小上限、TLS、关闭回调。
 * - readTimeoutMs/writeTimeoutMs:0 表无限。
 * - maxRequestBodySize:请求体字节上限,超出 stdx 回 413(仅 HTTP/1.1 非 chunked 生效)。
 *   None 用 stdx 默认(2MB);Some(0) 显式无限。防大文件上传 OOM。
 * - maxRequestHeaderSize:请求头字节上限,超出 stdx 回 431(仅 HTTP/1.1 生效)。
 *   None 用 stdx 默认;Some(0) 显式无限。防恶意大头。
 * - streamRequestBodies:true 时不预读请求体,业务通过 Context.requestBodyStream().readChunk(buffer) 消费。
 *   流式模式下 rawBody/bodyParser 不可用,应用必须显式读取或丢弃请求体。
 * - tls:调用方按 stdx.net.tls 构造 TlsConfig(含证书链/私钥)传入即启用 HTTPS。
 * - onShutdown:服务器关闭(close/closeGracefully)时回调,典型接 shutdownContainer 触发 @PreDestroy。
 */
public struct ListenOptions {
    public var readTimeoutMs: Int64 = 0
    public var writeTimeoutMs: Int64 = 0
    public var maxRequestBodySize: ?Int64 = None
    public var maxRequestHeaderSize: ?Int64 = None
    /** 开启后请求体通过 Context.consumeRequestBody 分块消费,不预先聚合到内存。 */
    public var streamRequestBodies: Bool = false
    /** 应用层流式读取上限;0 表仅使用 stdx maxRequestBodySize。 */
    public var streamRequestBodyLimit: Int64 = 0
    /** 请求处理结束后最多丢弃的未消费字节;默认64KiB,0表示不限制。 */
    public var streamDrainLimit: Int64 = 65536
    public var tls: ?TlsConfig = None
    public var onShutdown: () -> Unit = {=> ()}

    public init() {}
}

/** 用 stdx.net.http 起服务,把 App 的洋葱挂为统一处理入口(默认无超时/无 TLS)。 */
public func listen(app: App, addr: String, port: UInt16): Unit {
    listenWith(app, addr, port, ListenOptions())
}

/** 带选项启动:超时 / TLS / 关闭回调。 */
public func listenWith(app: App, addr: String, port: UInt16, opts: ListenOptions): Unit {
    // enableConnectProtocol:HTTP/2 下发 SETTINGS_ENABLE_CONNECT_PROTOCOL(RFC 8441),
    // 允许客户端经 extended CONNECT 在 h2 流上升级 WebSocket;对 HTTP/1.1 无影响。
    var b = ServerBuilder()
        .addr(addr)
        .port(port)
        .distributor(AceDistributor(app, opts.streamRequestBodies, opts.streamRequestBodyLimit, opts.streamDrainLimit))
        .enableConnectProtocol(true)
    if (opts.readTimeoutMs > 0) {
        b = b.readTimeout(opts.readTimeoutMs * Duration.millisecond)
    }
    if (opts.writeTimeoutMs > 0) {
        b = b.writeTimeout(opts.writeTimeoutMs * Duration.millisecond)
    }
    match (opts.maxRequestBodySize) {
        case Some(n) => b = b.maxRequestBodySize(n)
        case None => ()
    }
    match (opts.maxRequestHeaderSize) {
        case Some(n) => b = b.maxRequestHeaderSize(n)
        case None => ()
    }
    match (opts.tls) {
        case Some(cfg) => b = b.tlsConfig(cfg)
        case None => ()
    }
    let server = b.build()
    server.onShutdown(opts.onShutdown)
    server.serve()
}