///|
/// 真实传输层:把请求交给 `moonbitlang/async/http`,真的走网络。
///
/// 这是整个模块唯一碰网络的地方。上层拿到的 `RawResponse` 已经与底层类型解耦,
/// 所以底层库升级或替换不会波及配置合并、URL 拼接等逻辑。
///
/// 当前实现每次都新建连接(用底层的 `@http.Client` 手动走
/// 「connect → request → write → end_request」四步),不做连接复用;
/// TLS 校验开关暂未暴露,见 README 的「暂不支持」清单。
/// 手动四步而非 `@http.request` 便捷函数,是为了不在传输层就把响应体读光。
///
/// 代理(`PreparedRequest::proxy`)由底层的 CONNECT 隧道实现:给它一个干净的
/// 代理客户端,它就负责发 `CONNECT host:port`、验收并进入隧道,https 目标
/// 再在隧道上叠一层 TLS。逐请求各建一条隧道(与「不复用连接」一致),
/// 契约见 `docs/09-proxy.md`。
pub struct AsyncHttpTransport {}
///|
/// 创建真实传输层。
pub fn AsyncHttpTransport::new() -> AsyncHttpTransport {
AsyncHttpTransport::{ }
}
///|
/// 本项目的 `Method` 到底层 `RequestMethod` 的一一映射。
///
/// 两个枚举的构造子完全对应,映射写在这里而不是做成类型别名,
/// 是为了让上层完全不认识底层类型——上层的纯逻辑包不需要依赖 async。
fn to_request_method(meth : @config.Method) -> @http.RequestMethod {
match meth {
Get => @http.RequestMethod::Get
Head => @http.RequestMethod::Head
Post => @http.RequestMethod::Post
Put => @http.RequestMethod::Put
Delete => @http.RequestMethod::Delete
Connect => @http.RequestMethod::Connect
Options => @http.RequestMethod::Options
Trace => @http.RequestMethod::Trace
Patch => @http.RequestMethod::Patch
}
}
///|
/// 把我们的头集合转成底层要求的 `Map[CaseInsensitiveString, String]`。
///
/// 两者都是大小写不敏感的容器,这里只做类型搬运,不改变任何语义。
fn to_http_headers(
headers : @headers.Headers,
) -> Map[@http.CaseInsensitiveString, String] {
let result : Map[@http.CaseInsensitiveString, String] = Map([])
for entry in headers.entries() {
result[@http.CaseInsensitiveString(entry.0)] = entry.1
}
result
}
///|
/// 把底层响应头转回我们自己的 `Headers`。
///
/// 注意 `Set-Cookie` 不在这里:底层把响应 cookie 单独放在 `Response::cookies`
/// 里,不混进 headers,所以本项目也不提供 `headers.get("Set-Cookie")`
/// (见 README 的「暂不支持」清单)。
fn from_http_headers(
headers : Map[@http.CaseInsensitiveString, String],
) -> @headers.Headers {
let mut result = @headers.Headers::new()
for name, value in headers {
result = result.set(name.to_string(), value)
}
result
}
///|
/// 把完整地址切成「连接用的根地址」和「请求路径」两段。
///
/// 底层 `@http.Client::Client(uri)` 要求 uri 的 path 恰好是 `/`,
/// 路径要留给 `Client::request(meth, path)`,所以这里得自己切一刀;
/// 也正因为这个限制,不能用只覆盖 GET 的 `@http.get_stream`。
/// 切法镜像底层的 `resolve_url`:先按 `://` 分出协议,再取之后第一个 `/`。
///
/// 相对地址与 http/https 之外的协议都抛 `Unsupported`:
/// 它们的共同点是「请求还没发出去就失败了」,正是这个分类的定义。
fn split_url(url : String) -> (String, String) raise TransportError {
let scheme_end = match url.find("://") {
Some(index) => index
None =>
raise TransportError::Unsupported("地址不是绝对地址:" + url)
}
let scheme = url[:scheme_end].to_owned()
guard scheme is ("http" | "https") else {
raise TransportError::Unsupported("只支持 http/https 协议:" + scheme)
}
let rest = url[scheme_end + 3:]
// 切不开路径时(例如 `https://host`)按 HTTP 约定补成 `/`;
// 切到的 path 一定以 `/` 开头,正好满足 `Client::request` 对路径的要求。
let (host, path) = match rest.find("/") {
Some(index) => (rest[:index].to_owned(), rest[index:].to_owned())
None => (rest.to_owned(), "/")
}
(scheme + "://" + host, path)
}
///|
/// 建一个「处于干净状态」的代理客户端。
///
/// 底层的 `Client::Client(uri, proxy~)` 要求 proxy 是一个干净的 `@http.Client`
/// 并且**会接管它的所有权**(进入隧道后调用方不能再碰它),所以每次请求都得
/// 新建一个——这正好与本实现「每次请求新建连接」的现状一致。
///
/// 凭据以持久头的形式交给它:`CONNECT` 请求带的就是代理客户端自己的这批头
/// (见底层 `Client::connect`),而不是目标请求的头——后者是隧道打通之后
/// 发给目标服务器的。`Proxy-Authorization` 因此不会泄漏给目标服务器。
///
/// 已知微瑕:目标地址非法(例如端口越界)时,底层在接管这个客户端**之前**
/// 就抛错,那条到代理的连接不会被关掉。它只出现在「地址本来就发不出去」的
/// 失败路径上,代价是一条连接,不影响请求结果。
async fn open_proxy(endpoint : ProxyEndpoint) -> @http.Client {
let headers : Map[@http.CaseInsensitiveString, String] = Map([])
match endpoint.authorization {
Some(authorization) =>
headers[@http.CaseInsensitiveString("Proxy-Authorization")] = authorization
None => ()
}
@http.Client::Client(endpoint.url, headers~)
}
///|
/// 上传进度的分块大小:每写完这么多字节报告一次进度。
///
/// 它不是协议要求(请求体本来就是 chunked 编码),只是「进度回调的粒度」:
/// 太小会让回调次数暴涨、太大则进度条一顿一顿。64 KiB 是常见的折中。
const UPLOAD_CHUNK_SIZE : Int = 65536
///|
/// 把请求体写出去;有上传进度回调时按 `UPLOAD_CHUNK_SIZE` 分块写、逐块报告。
///
/// 为什么分块不改变线上行为:底层不传 `Content-Length` 时本来就把请求体编码成
/// `Transfer-Encoding: chunked`,它的发送缓冲只有 1 KiB——也就是说现在的整块
/// `write` 在底层早就被切成很多个 chunk 了。这里只是换个切法,并对齐
/// 「报告了 = 已经交给内核」这个口径:每块写完 `flush()` 一次,
/// 免得进度报完了数据还堆在库的缓冲里。
///
/// 报告不保证对端已经收到:`flush` 只把字节交给内核,TCP 缓冲区里还剩多少
/// 只有对端知道(见 `docs/10-progress.md`)。
///
/// 空 body(或 `None`)不触发任何回调——没有字节可报。
async fn write_body(
client : @http.Client,
data : Bytes,
on_progress : @config.ProgressCallback?,
) -> Unit {
match on_progress {
None => {
let body : &@io.Data = data
client.write(body)
}
Some(report) => {
let total = data.length()
let mut sent = 0
while sent < total {
let mut end = sent + UPLOAD_CHUNK_SIZE
if end > total {
end = total
}
let chunk : &@io.Data = data.exact_view(start=sent, end~)
client.write(chunk)
client.flush()
sent = end
report({ loaded: sent, total: Some(total), })
}
}
}
}
///|
/// 从响应头里取 `Content-Length`,作为下载进度的 `total`。
///
/// 取不到(chunked 响应没有这个头、值不是合法整数)就是 `None`:进度事件据此
/// 表达「长度未知」,而不是编一个 0 出来。注意压缩响应里这个值是**压缩后**的
/// 长度,而读出来的是解压后的字节,所以 `loaded` 可能超过 `total`
/// (见 `docs/10-progress.md`)。
fn content_length(headers : Map[@http.CaseInsensitiveString, String]) -> Int? {
let raw = match headers.get(@http.CaseInsensitiveString("Content-Length")) {
Some(value) => value
None => return None
}
let parsed : Int? = Some(@string.parse_int(raw.trim())) catch { _ => None }
match parsed {
Some(length) if length >= 0 => Some(length)
_ => None
}
}
///|
/// 建连、发出请求头与请求体、等回响应头,然后把「还开着的连接」与响应头
/// 一起交出去。
///
/// 这一段(含等响应头)整体受 `timeout` 约束;响应体还没开始读,
/// 读取阶段的超时由 `ResponseBody` 按**单次读取**来管,见 `stream.mbt`。
/// 用 `@http.Client` 而不是 `@http.request` 的便捷函数,正是为了不在
/// 这里就把响应体读光——`end_request()` 返回时响应体还在连接上。
///
/// 走代理时,与代理的 CONNECT 握手也在这一段里,因此同样受 `timeout` 约束。
async fn send_head(
request : PreparedRequest,
) -> (@http.Client, @http.Response) raise TransportError {
let (root, path) = split_url(request.url)
let attempt = () => {
// 代理客户端在这里(而不是外层)新建,握手代价才落进 timeout;
// 无代理时 `proxy?` 传 `None`,与不传该参数等价。
let proxy = match request.proxy {
Some(endpoint) => Some(open_proxy(endpoint))
None => None
}
// 头一次性交给 connect(与便捷函数 @http.request 的用法一致:
// 它同样是把 headers 交给 connect,body 单独 write)。
let client = @http.Client::Client(
root,
headers=to_http_headers(request.headers),
proxy?,
)
// 失败路径必须关连接,否则会漏下一条没人管的半开连接。
errdefer client.close()
client.request(to_request_method(request.http_method), path)
match request.body {
Some(bytes) => write_body(client, bytes, request.on_upload_progress)
// 无 body 时什么都不写:底层据此发 Content-Length: 0 或 chunked。
None => ()
}
(client, client.end_request())
}
try {
// timeout 为 None 或非正数时不加限制,避免把「不限时」误当成「立即超时」。
match request.timeout {
Some(milliseconds) if milliseconds > 0 =>
@async.with_timeout(milliseconds, attempt)
_ => attempt()
}
} catch {
@async.TimeoutError => raise TransportError::Timeout
// 代理拒绝建立隧道:CONNECT 的响应不是 2xx(`407` 需要认证最常见)。
// 单独翻成人话,比把底层的 `ProxyError({code: 407, ...})` 原样塞进网络错误有用。
@http.ProxyError(response) =>
raise TransportError::Network(
"代理拒绝建立隧道:HTTP " +
response.code.to_string() +
" " +
response.reason,
)
// 其余错误(连接失败、DNS、TLS、协议错误等)统一归到网络层失败,
// 并保留原始错误文本,便于排查。
error => raise TransportError::Network(error.to_string())
}
}
///|
/// 真正发请求:拿回响应头与一条可读的响应体流。
///
/// 错误分类的意义在于:上层只需要认识 `TransportError` 就能决定
/// 对外报 `ErrorCode::Timeout` 还是 `ErrorCode::Network`,
/// 不必知道底层用的是哪套 HTTP 实现、抛了哪些具体错误类型。
///
/// 取消作用域包住**整跳**:建连、写头、写 body(含上传进度回调)、等响应头。
/// 从上传进度回调里取消也落在同一个子任务上,所以下一块写出去之前就会断。
// 实现里的方法写 `fn` 而不是 `async fn`:异步性是 trait 声明的一部分,
// 由 trait 决定,实现侧不再重复标注(与 moonbitlang/async 自身的写法一致)。
pub impl Transport for AsyncHttpTransport with fn send(self, request) {
ignore(self)
// 取消句柄随请求带过来(`PreparedRequest::cancel_token`);没有就直跑,
// 不建任务组(见 with_cancel_scope)。
let sent = with_cancel_scope(request.cancel_token, () => send_head(request))
let (client, response) = sent
{
status: response.code,
status_text: response.reason,
headers: from_http_headers(response.headers),
// 响应体还没读:连接的所有权转移给这个流,读到 EOF 或 close() 时关闭。
// 不在这里读全,是为了让 chunked / SSE 这类协议能边到边读。
// `Content-Length` 一并交出去,它是下载进度的 total(可能没有,见 content_length);
// token 也交出去:每次读取都靠它可被取消,取消时它还会顺手关掉这条连接。
body: ResponseBody::open(
client,
request.timeout,
content_length(response.headers),
request.cancel_token,
),
}
}
///|
pub extend AsyncHttpTransport with Transport::{send}