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()
}