Sign in

    #transport —— 传输层契约

    传输层:把「真正把字节发出去」抽象成 Transport trait,上层只认 PreparedRequest / RawResponse 两个契约。真实的网络实现是 httpconn 包的自研 HTTP/1.1 栈(HttpConnTransport,Client::new 的缺省传输);本包自带 MockTransport 供测试替换,也可以换成自己的实现。

    import { "q2316367743/moonhttp/transport", }

    #契约

    类型内容
    Transport(trait)只有一个方法:async fn send(Self, PreparedRequest) -> RawResponse raise TransportError
    PreparedRequest方法、完整 URL、拍平后的头、body(RequestBody?:Buffered(Bytes) 或 Stream(StreamBody))、超时、代理端点、上传进度回调、取消信号
    RawResponse状态码、状态短语、响应头、响应体流 ResponseBody
    ProxyEndpoint{ url, authorization? }:隧道地址与 CONNECT 的凭据
    TransportErrorTimeout / Network(String) / Unsupported(String) / Malformed(String) / Cancelled(String?);上层看到的 HttpError 分类就来自它(Malformed 是「响应体与它声称的 Content-Encoding 不符」,报 ERR_BAD_RESPONSE),取消的载荷是取消理由

    #ResponseBody 的读法

    方法说明
    from_bytes(bytes)手里已有完整响应体时包一层(Mock 与测试用)
    read_all(on_progress?)读到 EOF;传了回调就逐块报下载进度
    read_all_partial(on_progress?)同上但不抛错:返回 (已读到的字节, 失败原因?),「读到一半失败」时保住现场
    read_some(max_len?)读一块,None 表示 EOF
    read_until(delim)读到分隔符(含分隔符)为止
    close()释放连接

    超时在这里分两种口径:read_all / read_all_partial 是「整段读完」一个时限,read_some / read_until 是「每次等待」一个时限。

    #解压与响应头

    缓冲路径的解压工具在这个包:decode_gzip(bytes) 把一整份 gzip 字节解成实体字节,declares_gzip(headers) 判断响应是否声称自己是 gzip。解压归属按「谁解压谁声明」:库自己声明 Accept-Encoding 的响应,交到你手里前已解掉(Content-Encoding / Content-Length 一并摘除);用户自己声明 Accept-Encoding 的,流式路径字节与头原样交出(头与体说同一件事)。自定义传输实现同样守这条口径。谁声明、谁解压、失败怎么报见 docs/15-response-compression.md。

    #两个实现

    • HttpConnTransport(定义在 httpconn 包,Client::new 的缺省传输):自研 HTTP/1.1 栈——TCP / TLS 直连与 CONNECT 隧道代理、gzip 自持、自定义方法原样落线;每次请求一条连接(发 Connection: close,连接池在路线图上),见 docs/18-httpconn-transport.md。
    • MockTransport:测试用,不碰网络。记录收到的每一份请求(received() / request_count() / last_request()),返回预置响应(new(response) 或 from_responses([...]),取完之后一直复用最后一个),也可以固定失败(failing(error) / with_failure(error))。

    #用法

    // 在测试里替换掉真实网络

    ///|
    let mock = @transport.MockTransport::new(response)

    ///|
    let transport : &@transport.Transport = mock

    ///|
    let client = @moonhttp.Client::new(transport~)

    自定义传输只要实现一个方法,底层类型完全被挡在上层之外:

    ///|
    pub impl Transport for MyTransport with fn send(self, request) {
    // request : PreparedRequest(方法、完整地址、已拍平的头、body(缓冲或流式)、超时)
    // 返回 : RawResponse(状态码、状态短语、响应头、响应体流)
    ...
    }

    字段含义、超时语义与自定义传输的注意事项见 docs/05-transport.md。

    AbortSignal

    取消信号:一个可以传给任意多次请求的「取消观察口」。

    它自己不能被取消——叫停只能由持有 AbortController 的一方发起 (对应 JS 的 controller.abort())。这条分工正是这套机制比 CancelToken 好的地方:拿到信号的代码只能观察,「误调用取消」在类型上就写不出来。

    几条刻意的语义(与 JS 一致):
    • 一次性:取消之后永远是取消状态,不能复用;需要「一次请求一个信号」 就每次 AbortController::new()。
    • 可共享,但只能按请求显式共享:同一个信号可以挂在任意多次调用上, abort 一次全部生效。它不会被实例默认值继承——信号是一次性的, 放进实例默认值会让那个实例之后的所有请求一起失效,所以本项目的 三个入口只从 signal? 参数收信号(理由见 docs/12-cancellation.md)。
    • 取消晚一步也算数:取消发生之后才开始的请求会立刻失败,不会 「取消晚了一步就照常发出去」——靠的是每个可取消的 I/O 入口先查 aborted()(见 attach 的说明与 docs/12-cancellation.md 的「两条兜底」)。
    • 幂等:重复 abort 只生效第一次,第二次连理由都不会改。

    Headers

    大小写不敏感的 HTTP 头集合,对应 axios 的 AxiosHeaders。

    内部以小写化后的头名作为 Map 的 key,值里再记住写入时的原始拼写, 于是 get("Content-Type") 与 get("content-type") 命中同一条记录, 而遍历/打印时仍然保留调用方书写时的形状(与 axios「首次出现的拼写胜出」一致)。

    与 axios 相比的两处有意简化(README 里有完整清单):
    • 不支持 axios 用 false 表示「这条头禁止被默认值覆盖」的哨兵值;
    • 同一个头名只能有一个字符串值,不展开多值头(axios 允许数组)。

    所有变更方法都返回新实例而不是就地修改:Config 是值语义的, 合并配置时同一份 Headers 可能被多个 Config 引用, 返回新值可以避免「改一个实例的头,另一个实例跟着变」这类共享可变状态 bug。

    Method

    这几个类型出现在 PreparedRequest / RawResponse 的公开签名里, 因此随本包一起再导出,使用者不必为了构造一份请求再 import 两个包。 AbortSignal 同理:它由 PreparedRequest 携带到这一层(定义在 abort 包)。 StreamBody 是请求体流式形态(docs/20)的载荷,自定义传输实现要按 RequestBody::Stream 解构它。

    ProgressCallback

    进度回调。

    noraise:回调是在请求发送 / 响应读取的中途被调用的,它抛出的错误没有 合理的归属方(既不属于这次请求的传输错误,也不该让整条请求失败), 所以类型上就要求它不抛。想在回调里记日志、推 UI 状态都没问题。

    回调是同步执行的:它占用请求本身的时间预算(上传回调落在 timeout 覆盖的发送阶段里),所以别在里面做耗时的事。

    StreamBody

    这几个类型出现在 PreparedRequest / RawResponse 的公开签名里, 因此随本包一起再导出,使用者不必为了构造一份请求再 import 两个包。 AbortSignal 同理:它由 PreparedRequest 携带到这一层(定义在 abort 包)。 StreamBody 是请求体流式形态(docs/20)的载荷,自定义传输实现要按 RequestBody::Stream 解构它。

    Transport

    pub(open) trait Transport {
    async fn send(Self, PreparedRequest) -> RawResponse raise TransportError
    }

    传输层抽象:把一份准备好的请求发出去,拿回原始响应。

    用 trait 而不是直接调用异步 HTTP 库有两个目的:
    1. 把整个 async 运行时依赖关在实现里,上层(配置合并、URL 拼接、 错误映射)都能用普通同步测试覆盖;
    2. 使用方和测试可以注入自己的实现,不必真的发网络请求。

    TransportError

    pub(all) suberror TransportError {
    Timeout
    Network(String)
    Unsupported(String)
    Malformed(String)
    Cancelled(String?)
    }

    传输层错误:把「网络世界里可能出什么事」收敛成五种情况, 上层再映射成对外的 ErrorCode。

    有了这层收敛,上层不需要 import 任何 async 相关的包, 自定义传输实现也不需要知道底层用的是哪套 HTTP 库。

    MockTransport

    pub struct MockTransport {
    // private fields
    }

    测试用传输层:记录收到的每一份请求,并返回预置的响应。

    它本身就是「传输层可替换」这个设计的示范——不碰网络就能把 合并 → 拼接 → 拍平 → 校验的完整流程跑一遍。 使用方也可以用它在自己的测试里替换掉真实网络。

    流式请求体(RequestBody::Stream)只记录引用、不消费——Mock 不真的 发送数据(与「不报告上传进度」同一契约);要验证「流真的被泵出去、 字节正确」,得用真实传输(httpconn)打靶子。

    MockTransport::failing

    构造一个总是失败的 Mock,用来验证错误分类与传播。

    MockTransport::from_responses

    fn MockTransport::from_responses(responses : Array[RawResponse]) -> MockTransport

    构造一个按队列依次返回响应的 Mock;队列用完后一直复用最后一个响应。

    MockTransport::last_request

    fn MockTransport::last_request(self : MockTransport) -> PreparedRequest?

    最近一次收到的请求;一次都没发过时为 None。

    MockTransport::new

    fn MockTransport::new(response : RawResponse) -> MockTransport

    构造一个始终返回同一响应的 Mock。

    MockTransport::received

    收到过的请求,按顺序。

    MockTransport::request_count

    fn MockTransport::request_count(self : MockTransport) -> Int

    断言用的便捷入口:接收一个 &Transport 也能读出请求记录。

    MockTransport::send

    async fn MockTransport::send(self : MockTransport, request : PreparedRequest) -> RawResponse raise TransportError

    MockTransport::with_failure

    fn MockTransport::with_failure(self : MockTransport, error : TransportError) -> MockTransport

    设置失败错误,返回新的 Mock(failure 是可变字段,这里刻意返回新实例, 让「构造—配置—使用」的写法保持一致)。

    PreparedRequest

    pub(all) struct PreparedRequest {
    http_method :
    Method

    url : String
    headers :
    Headers

    body : RequestBody?
    timeout : Int?
    proxy : ProxyEndpoint?
    on_upload_progress : (
    ProgressEvent
    ) -> Unit?
    signal :
    AbortSignal
    ?
    }

    交给传输层去发送的一份请求。

    到达这一步时,配置合并、base_url 拼接、query 序列化、头拍平都已经完成, 所以它只包含「把字节发出去」真正需要的信息,不含任何配置语义。 这样任何实现——真实的 HTTP、测试用的 Mock、将来可能的连接池—— 都只需要关心这一层,替换传输实现不会影响上层语义。

    PreparedRequest::to_repr

    ProxyEndpoint

    pub(all) struct ProxyEndpoint {
    url : String
    authorization : String?
    } derive(
    Debug
    )

    走代理需要的两样东西:连到哪儿、以及 CONNECT 请求带什么凭据。

    已经是解析完的最小形态:协议名与端口默认值在拼请求时就写进了 url, 凭据也已经编码成可直接落头的字符串。这样传输层不必认识 ProxyProtocol 这类配置枚举,也不必知道 Basic 认证怎么编码——与 PreparedRequest 「不含任何配置语义」的口径一致。

    RawResponse

    pub(all) struct RawResponse {
    status : Int
    status_text : String
    headers :
    Headers

    body : ResponseBody
    set_cookies : Array[String]
    } derive(
    Debug
    )

    传输层拿到的原始响应。

    刻意不叫 Response:上层对外的 Response 还要承担 JSON 解析、 状态码校验等语义,那些不属于传输层的职责。

    body 是流而不是字节:send 返回时只保证状态行与响应头已到手, 响应体按需读取(ResponseBody)。底层原语只保留最弱的能力, 一次性读全是上层的组合结果——这样 chunked / SSE 这类「边到边读」 的协议才能表达出来。

    RequestBody

    pub(all) enum RequestBody {
    Buffered(Bytes)
    Stream(
    StreamBody
    )
    }

    请求体在传输层的形态:一块完整字节,或一个待泵出的读取流。

    四种缓冲形态(文本 / JSON / 表单 / urlencoded)序列化完都落 Buffered; Stream 来自 with_data_from_stream(docs/20)——content_length 有值 时按 Content-Length 定长发送(进度 total 已知),没值时按 Transfer-Encoding: chunked 发送。流是一次性资源:内置实现(httpconn) 的泵循环消费它;Mock 与自定义实现可以不消费,但绝不能假设它能读第二次。

    ResponseBody

    pub struct ResponseBody {
    // private fields
    }

    响应体的可读流。

    Transport::send 只保证「状态行与响应头已经到手,响应体按需读」, 一次性读全是上层的组合结果(Client::request 就是 read_all() 之后 走原来的解析流程)。这样 chunked / SSE 这类「边到边读」的协议才有落点: 响应头先用于判断状态码与内容类型,数据到了再一段段取。

    内部有两种来源:自研传输层(httpconn)的连接流(open_wire 包进来的任意 @io.Reader),以及内存字节(Mock 与测试用)。 两者的读语义刻意保持一致(对齐 @io.Reader),Mock 才能真实代表网络侧:
    • read_some 到 EOF 返回 None;
    • read_until 消费掉分隔符,且不把分隔符放进返回值;
    • 读到 EOF 会自动关闭底层连接,之后继续读仍然返回 None。

    拿到流之后必须读到 EOF 或调用 close()。本项目不复用连接,也没有析构器, 忘记关闭就会漏掉一条 TCP 连接。

    本文件只放读这一半(读语义、超时、取消作用域);构造与释放 (from_bytes / open / open_wire / rewind / close)在 stream_lifecycle.mbt,读全量那三件套在 stream_all.mbt—— 都是 RL-04 的 300 行上限逼出来的拆分。

    ResponseBody::close

    fn ResponseBody::close(self : ResponseBody) -> Unit

    关闭流并释放底层连接。幂等:重复调用只生效一次。

    只读了半截就停止(例如 SSE 收到想结束就断开)时必须显式调用它, 否则连接会一直挂着。

    ResponseBody::from_bytes

    fn ResponseBody::from_bytes(data : Bytes) -> ResponseBody

    用内存字节构造响应体。Mock 传输层与自定义实现用它造出 「不是网络来的」响应体,读语义与真实连接一致。

    ResponseBody::open_wire

    fn ResponseBody::open_wire(reader : &
    Reader
    , close : () -> Unit, timeout~ : Int?, total~ : Int?, signal~ :
    AbortSignal
    ?) -> ResponseBody

    用自研传输层(httpconn)的连接流构造响应体:reader 是已按 HTTP 分帧 切好的读端(chunked / content-length / close-delimited 都在它底下解好), close 负责释放底层连接(TLS + TCP)。连接的生命周期从这里交给响应体, 读到 EOF 或 close() 时才会关闭。

    timeout 约束单次读取、signal 登记「取消即关闭」、进来就已取消的 信号立刻关流——随后任何读取都会报 Cancelled 而不是安静地 EOF。这份检查 必须自己写——信号上的登记在已取消时不会补触发,见 src/abort/abort.mbt 的 attach。close 应当幂等:取消触发的关闭与显式 close() 可能只差一拍。

    ResponseBody::read_all

    async fn ResponseBody::read_all(self : ResponseBody, on_progress? : (
    ProgressEvent
    ) -> Unit) -> Bytes raise TransportError

    读到 EOF,返回剩下的全部字节,并关闭流。

    这是「非流式」用法的入口:Client::request 就是先读完再走原有的 JSON 解析 与状态码校验。中途失败时抛 TransportError——已经读到的部分在这一层丢掉, 需要它请用 read_all_partial(本函数就是它的「失败即抛」包装)。

    on_progress 原样转交给 read_all_partial。

    ResponseBody::read_all_partial

    async fn ResponseBody::read_all_partial(self : ResponseBody, on_progress? : (
    ProgressEvent
    ) -> Unit) -> (Bytes, TransportError?) noraise

    读到 EOF,返回剩下的全部字节;失败时不抛错,把已读到的部分与错误一起返回。

    为什么要它:网络在读到一半时断掉(超时、连接被重置)是常态,此时已经到手的 字节往往正是现场——服务端错误响应的正文、JSON 的开头、下载进度。上层要把它们 挂到错误上(HttpError::response),所以不能像 read_all 那样在中途失败时 把字节丢掉。

    返回 (全部字节, None) 表示正常读完,(已读到的部分, Some(错误)) 表示中途失败。 两种情况流都已关闭,与 read_all 一致。超时口径也与 read_all 一致: 整段读取受一次 timeout 约束,而不是每次分块各算一次。

    on_progress 有值时逐块报告下载进度。它是这次读取的进度而不是整个 响应体的:loaded 从 0 起算,先按块读过一段再调用本函数时不会接着累加。 回调在读取过程中同步执行,占用的也是这次读取的时间预算(timeout)。

    ResponseBody::read_some

    async fn ResponseBody::read_some(self : ResponseBody, max_len? : Int) -> Bytes? raise TransportError

    读取一段响应体;到 EOF 返回 None,此时连接已经关闭。

    max_len 限制单次返回的字节数,缺省时能取多少取多少—— 与 @io.Reader::read_some 一样,返回的块可能小于 max_len, 需要按长度区分的协议请自己缓冲。

    ResponseBody::read_until

    async fn ResponseBody::read_until(self : ResponseBody, sep : String) -> String? raise TransportError

    读到分隔符 sep 为止,返回 sep 之前的内容;sep 被消费掉但不返回。

    到 EOF 时把剩余内容当作最后一段返回,再读一次才返回 None—— 与 @io.Reader::read_until 一致。SSE 这类「按空行切事件」的协议 可以直接 read_until("\n\n") 取一个事件,不必自己处理跨块的边界。

    ResponseBody::to_repr

    declares_gzip

    fn declares_gzip(headers :
    Headers
    ) -> Bool

    响应头是否声明了 gzip 内容编码(Content-Encoding: gzip)。

    只认「整份值恰好是 gzip」这一种:按 RFC 9110,Content-Encoding 是逗号分隔的 编码列表,多层编码(gzip, gzip)与别的编码(br / deflate)这里都返回 false——没有对应的解码器,认得更宽只会把没解压的字节当成实体字节交出去。 值比较按大小写不敏感 + 裁空白(与 httpconn 流式解压的判定口径一致),头名比较由 Headers::get 负责。

    判定只有这一处:根包据此决定缓冲路径要不要解压(见 docs/15-response-compression.md)。

    decode_gzip

    async fn decode_gzip(data : Bytes) -> Bytes raise TransportError

    解压一整份 gzip 字节。

    只做内存里的数据变换:不碰网络、不做 IO、不看 HTTP 头——输入是已经读全的 响应体。缓冲路径(Client::request)本来就先把响应体读全,解压因此能和 「读」完全解耦,不必在流式读取的每一跳上挂解码器。

    实现借一个内存管道把两边接起来:@gzip.Decoder 只吃 @io.Reader, 而我们手上是完整字节。写端放进后台任务一次性写完就关闭,读端交给解码器, 这是 async 库官方自测里同一件事的写法。不假手 ResponseBody:它刻意不实现 @io.Reader(要命名 async 的内部类型,理由见 docs/05-transport.md)。

    失败一律报 Malformed(响应体与它声称的编码不符):gzip 头非法、数据损坏、 被截断(CRC32 / 长度校验不过、NeedMoreInput)都落在这一条上。宁可显式失败, 也不把半截字节当正文交出去——上层据此报 ERR_BAD_RESPONSE,并把原始字节 留在错误里(现场保真)。

    with_abort_scope

    async fn[T] with_abort_scope(signal :
    AbortSignal
    ?, f : async () -> T raise TransportError) -> T raise TransportError

    跑 f,并在 signal 取消时中断它。

    • signal 是 None:直接跑,不建任务组、零额外开销;
    • f 正常结束:原样返回它的结果;f 自己抛的错原样向外抛,不做包装;
    • 进来时信号已经取消:一步都不跑,直接抛 TransportError::Cancelled(带着理由);
    • 跑的过程中被取消:同样抛 TransportError::Cancelled(带着理由)。

    进来时那次检查是「取消晚一步」的兜底:信号上的登记在已取消时不会补触发 (与 JS 的 addEventListener 一致,见 src/abort/abort.mbt 的 attach), 所以每个登记点都要自己先查一眼。它管的正是「取消落在两段 I/O 之间」这个窗口 ——重定向的两跳之间就是真实的一例:这一跳必须立刻断,不能照常发出去。 检查与登记之间没有挂起点(协作式模型下另一条协程只在挂起点才可能运行), 所以「查完再登记」不会把取消漏掉。

    为什么不是「执行到检查点时看一眼信号」:检查只在执行到检查点时有效,而请求 卡在等首字节 / 建连时根本没有下一个检查点——那正是最需要取消的场景。 协程级取消的落点是挂起点,所以挂起中的 socket 动作会被真的中断 (@async.with_timeout 用的就是同一套机制,超时就是「定时取消」)。

    一条关键性质:被取消的是这里 spawn 的子任务,不是调用方协程。所以取消 之后调用方不会被「粘住」(取消是被取消任务的粘性状态),它拿到的只是一个 普通错误,可以接着发下一个请求。

    开放为 pub 是自研传输层落地(docs/17)的一部分:取消作用域是 Transport 契约的一半(另一半是 ResponseBody 上的「取消即关闭」登记),任何传输 实现都应该走同一套机制,取消语义才能在实现之间保持一致。