fix(trpc): stream backpressure via pull-driven pause and RFC range edge cases
This commit is contained in:
parent
cfddf5bbe1
commit
94e923486e
|
|
@ -8,6 +8,33 @@ export type RemoteFileStream = {
|
||||||
body: ReadableStream<Uint8Array> | null;
|
body: ReadableStream<Uint8Array> | null;
|
||||||
};
|
};
|
||||||
|
|
||||||
|
type ParsedRange = { start: number; end: number } | "unsatisfiable" | null;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 解析单段 bytes Range(RFC 9110):
|
||||||
|
* - null = 语法非法/不支持的形态 → 调用方忽略并回退 200 全量
|
||||||
|
* - "unsatisfiable" = 语义不可满足 → 416
|
||||||
|
*/
|
||||||
|
function parseBytesRange(header: string, size: number): ParsedRange {
|
||||||
|
const m = /^bytes=(\d*)-(\d*)$/i.exec(header.trim());
|
||||||
|
if (!m) return null;
|
||||||
|
const first = m[1] ?? "";
|
||||||
|
const second = m[2] ?? "";
|
||||||
|
if (!first && !second) return null;
|
||||||
|
if (!first) {
|
||||||
|
// 后缀语法 bytes=-N:最后 N 字节;N=0 不可满足
|
||||||
|
const suffix = Number(second);
|
||||||
|
if (suffix === 0) return "unsatisfiable";
|
||||||
|
return { start: Math.max(0, size - suffix), end: size - 1 };
|
||||||
|
}
|
||||||
|
const start = Number(first);
|
||||||
|
if (start >= size) return "unsatisfiable";
|
||||||
|
if (!second) return { start, end: size - 1 };
|
||||||
|
const end = Number(second);
|
||||||
|
if (end < start) return "unsatisfiable";
|
||||||
|
return { start, end: Math.min(end, size - 1) };
|
||||||
|
}
|
||||||
|
|
||||||
/** 打开远端媒体的 Range 流;webdav 细节只在本层出现,web 只拿 Headers + ReadableStream。 */
|
/** 打开远端媒体的 Range 流;webdav 细节只在本层出现,web 只拿 Headers + ReadableStream。 */
|
||||||
export async function openRemoteFileStream(
|
export async function openRemoteFileStream(
|
||||||
userId: string,
|
userId: string,
|
||||||
|
|
@ -36,15 +63,21 @@ export async function openRemoteFileStream(
|
||||||
let end = item.size > 0 ? item.size - 1 : 0;
|
let end = item.size > 0 ? item.size - 1 : 0;
|
||||||
|
|
||||||
if (rangeHeader && item.size > 0) {
|
if (rangeHeader && item.size > 0) {
|
||||||
const m = /bytes=(\d*)-(\d*)/.exec(rangeHeader);
|
const parsed = parseBytesRange(rangeHeader, item.size);
|
||||||
if (m) {
|
if (parsed === "unsatisfiable") {
|
||||||
if (m[1]) start = Number(m[1]);
|
headers.set("Content-Range", `bytes */${item.size}`);
|
||||||
if (m[2]) end = Math.min(Number(m[2]), item.size - 1);
|
return { status: 416, headers, body: null };
|
||||||
|
}
|
||||||
|
if (parsed) {
|
||||||
|
start = parsed.start;
|
||||||
|
end = parsed.end;
|
||||||
status = 206;
|
status = 206;
|
||||||
headers.set("Content-Range", `bytes ${start}-${end}/${item.size}`);
|
headers.set("Content-Range", `bytes ${start}-${end}/${item.size}`);
|
||||||
headers.set("Content-Length", String(end - start + 1));
|
headers.set("Content-Length", String(end - start + 1));
|
||||||
}
|
}
|
||||||
} else if (item.size > 0) {
|
// parsed === null(非法 Range):忽略并回退 200,下方补 Content-Length
|
||||||
|
}
|
||||||
|
if (status === 200 && item.size > 0) {
|
||||||
headers.set("Content-Length", String(item.size));
|
headers.set("Content-Length", String(item.size));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -53,17 +86,24 @@ export async function openRemoteFileStream(
|
||||||
item.size > 0
|
item.size > 0
|
||||||
? client.createReadStream(abs, { range: { start, end } })
|
? client.createReadStream(abs, { range: { start, end } })
|
||||||
: client.createReadStream(abs);
|
: client.createReadStream(abs);
|
||||||
|
// ReadableLike 类型未声明 pause/resume,但实现是 Node Readable(data 监听后进入 flowing 模式)
|
||||||
|
const src = nodeStream as typeof nodeStream & { pause(): void; resume(): void };
|
||||||
|
|
||||||
|
// pull 驱动 + desiredSize 背压:队列满时 pause,pull(消费者取走)时 resume
|
||||||
const body = new ReadableStream<Uint8Array>({
|
const body = new ReadableStream<Uint8Array>({
|
||||||
start(controller) {
|
start(controller) {
|
||||||
nodeStream.on("data", (chunk: Buffer | string) => {
|
src.on("data", (chunk: Buffer | string) => {
|
||||||
controller.enqueue(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
|
controller.enqueue(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
|
||||||
|
if ((controller.desiredSize ?? 1) <= 0) src.pause();
|
||||||
});
|
});
|
||||||
nodeStream.on("end", () => controller.close());
|
src.on("end", () => controller.close());
|
||||||
nodeStream.on("error", (err: unknown) => controller.error(err));
|
src.on("error", (err: unknown) => controller.error(err));
|
||||||
|
},
|
||||||
|
pull() {
|
||||||
|
src.resume();
|
||||||
},
|
},
|
||||||
cancel() {
|
cancel() {
|
||||||
nodeStream.destroy();
|
src.destroy();
|
||||||
},
|
},
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue