// stream.mbt 拆出来的另一半:「读全量」三件套(drain / read_all_partial / read_all)。
//
// 为什么单独一个文件:AGENTS.md RL-04 对子包文件的上限是 300 行,加上下载进度
// 报告后 stream.mbt 会超限;而「读全量」正好是能整体搬走的一块——它只用
// `read_or_fail` / `close` 这些同包可见的原语,与「按块读」也确实是两组语义。
//
// 下载进度的契约(谁触发、粒度、total 从哪来)见 `docs/10-progress.md`。
///|
/// 内存体上模拟的「一次到货多少字节」,只影响进度回调的粒度。
///
/// 内存体不带 `max_len` 时其实一次就能把剩余的字节全给出去,真实连接则是
/// 对端发多少到多少。按固定大小报告,内存体的回调节奏才与真实连接一致——
/// Mock 的用例也才能确定地观察到「多块递增」。
const PROGRESS_CHUNK_SIZE : Int = 65536
///|
/// 把连接上的剩余字节读进 `buffer`,直到 EOF,逐块报告下载进度。
///
/// 单独抽出来是为了让它作为**一个整体**交给 `read_or_fail`:读全量的超时
/// 口径与逐次读取不同(整段读完受一次 `timeout` 约束,见 `read_all_partial`),
/// 分成多次读取会把这个口径改掉。
///
/// `loaded` 从 0 起算,只统计**这次读取**取到的字节;`total` 是响应体的总长度
/// (可能未知)。
async fn drain_into(
client : @http.Client,
buffer : @buffer.Buffer,
total : Int?,
on_progress : @config.ProgressCallback?,
) -> Unit {
let mut loaded = 0
while client.read_some() is Some(chunk) {
buffer.write_bytes(chunk)
loaded = loaded + chunk.length()
match on_progress {
Some(report) => report({ loaded, total, })
None => ()
}
}
}
///|
/// 读到 EOF,返回剩下的全部字节;**失败时不抛错**,把已读到的部分与错误一起返回。
///
/// 为什么要它:网络在读到一半时断掉(超时、连接被重置)是常态,此时已经到手的
/// 字节往往正是现场——服务端错误响应的正文、JSON 的开头、下载进度。上层要把它们
/// 挂到错误上(`HttpError::response`),所以不能像 `read_all` 那样在中途失败时
/// 把字节丢掉。
///
/// 返回 `(全部字节, None)` 表示正常读完,`(已读到的部分, Some(错误))` 表示中途失败。
/// 两种情况流都已关闭,与 `read_all` 一致。超时口径也与 `read_all` 一致:
/// 整段读取受一次 `timeout` 约束,而不是每次分块各算一次。
///
/// `on_progress` 有值时逐块报告下载进度。它是**这次读取**的进度而不是整个
/// 响应体的:`loaded` 从 0 起算,先按块读过一段再调用本函数时不会接着累加。
/// 回调在读取过程中同步执行,占用的也是这次读取的时间预算(`timeout`)。
pub async fn ResponseBody::read_all_partial(
self : ResponseBody,
on_progress? : @config.ProgressCallback,
) -> (Bytes, TransportError?) noraise {
// 取消发生在「响应头到手」与「开始读全量」之间时,这里报取消,而不是落到
// 下面的 `Closed => (b"", None)` —— 后者会被上层当成「读完了但一个字节都没有」。
match cancelled_failure(self.cancel_token) {
Some(error) => return (b"", Some(error))
None => ()
}
match self.inner {
Closed => (b"", None)
Memory(memory) => {
let start = memory.offset
let end = memory.data.length()
match on_progress {
Some(report) => {
// 内存体其实一次就能给完,这里按 PROGRESS_CHUNK_SIZE 模拟「分块到货」,
// 让 Mock 观察到的回调节奏与真实连接一致(真实连接报告的是
// 每次 read_some 的实际块大小)。
let mut loaded = 0
while loaded < end - start {
let mut next = loaded + PROGRESS_CHUNK_SIZE
if next > end - start {
next = end - start
}
loaded = next
report({ loaded, total: self.total, })
}
}
None => ()
}
let tail = memory.data.exact_view(start~, end~).to_owned()
memory.offset = end
self.close()
(tail, None)
}
Http(client) => {
let buffer = @buffer.Buffer::Buffer()
// 失败不当异常走,而是作为值返回:调用方要的是「已读到的字节 + 错误」。
let failure : TransportError? = try {
read_or_fail(self.timeout, self.cancel_token, () => {
drain_into(client, buffer, self.total, on_progress)
})
None
} catch {
error => Some(error)
}
// 读完与失败都一样:这条连接没有别的用处,释放掉(`read_some` 的
// errdefer 已经关过一次,`close` 幂等)。
self.close()
(buffer.to_bytes(), failure)
}
}
}
///|
/// 读到 EOF,返回剩下的全部字节,并关闭流。
///
/// 这是「非流式」用法的入口:`Client::request` 就是先读完再走原有的 JSON 解析
/// 与状态码校验。中途失败时抛 `TransportError`——已经读到的部分在这一层丢掉,
/// 需要它请用 `read_all_partial`(本函数就是它的「失败即抛」包装)。
///
/// `on_progress` 原样转交给 `read_all_partial`。
pub async fn ResponseBody::read_all(
self : ResponseBody,
on_progress? : @config.ProgressCallback,
) -> Bytes raise TransportError {
let (bytes, failure) = self.read_all_partial(on_progress?)
match failure {
None => bytes
Some(error) => raise error
}
}