import { streamConsumeInMacroTask, toError } from "#Source/basic/index.ts" import { Tube } from "./tube.ts" /** * @description 将上游 Tube 的生命周期、错误与数据转发到下游 Tube。 */ export const connectTube = ( upstream: Tube, downstream: Tube, ): (() => void) => { const openEventSubscriber = async (): Promise => { await downstream.open() } upstream.subscribeOpenEvent({ subscriber: openEventSubscriber }) const closeEventSubscriber = async (): Promise => { await downstream.close() } upstream.subscribeCloseEvent({ subscriber: closeEventSubscriber }) const startEventSubscriber = async (): Promise => { await downstream.start() } upstream.subscribeStartEvent({ subscriber: startEventSubscriber }) const endEventSubscriber = async (): Promise => { await downstream.end() } upstream.subscribeEndEvent({ subscriber: endEventSubscriber }) const errorEventSubscriber = async (error: E): Promise => { await downstream.pushError(error) } upstream.subscribeErrorEvent({ subscriber: errorEventSubscriber }) const dataSubscriber = async (data: D): Promise => { await downstream.pushData(data) } upstream.subscribeData({ subscriber: dataSubscriber }) return (): void => { upstream.unsubscribeOpenEvent(openEventSubscriber) upstream.unsubscribeCloseEvent(closeEventSubscriber) upstream.unsubscribeStartEvent(startEventSubscriber) upstream.unsubscribeEndEvent(endEventSubscriber) upstream.unsubscribeErrorEvent(errorEventSubscriber) upstream.unsubscribeData(dataSubscriber) } } /** * @description 把 Tube 收束为一个 Promise,并在通道关闭时以最后一条数据或最新错误结算。 */ export const tubeToPromise = async (tube: Tube): Promise => { const promise = new Promise((resolve, reject) => { const handle = (): void => { if (tube.isError()) { reject(toError(tube.getLatestErrorOrThrow())) unsubscribe() return } if (tube.hasClosed()) { if (tube.isWet()) { resolve(tube.getLatestDataOrThrow()) unsubscribe() return } else { reject(new Error("No data available")) unsubscribe() return } } } const unsubscribe = tube.subscribeCloseEvent({ subscriber: () => { handle() }, }) // invoke once handle() }) return await promise } /** * @description 把 Tube 暴露为一个 `ReadableStream`。 */ export const tubeToReadableStream = (tube: Tube): ReadableStream => { const readableStream = new ReadableStream({ start(controller): void { let hasSettled = false const unsubscribeAll = (): void => { unsubscribeCloseEvent() unsubscribeErrorEvent() unsubscribeDataSubscriber() } const closeEventSubscriber = (): void => { if (hasSettled === true) { return } hasSettled = true controller.close() unsubscribeAll() } const unsubscribeCloseEvent = tube.subscribeCloseEvent({ subscriber: closeEventSubscriber }) const errorEventSubscriber = (error: E): void => { if (hasSettled === true) { return } hasSettled = true controller.error(error) unsubscribeAll() } const unsubscribeErrorEvent = tube.subscribeErrorEvent({ subscriber: errorEventSubscriber }) const dataSubscriber = (data: D): void => { if (hasSettled === true) { return } controller.enqueue(data) } const unsubscribeDataSubscriber = tube.subscribeData({ subscriber: dataSubscriber }) }, }) return readableStream } /** * @description 把 `ReadableStream` 转换为一个立即开始消费上游数据的 Tube。 * 调用后会立刻开始消费 `ReadableStream`,并将数据推入返回的 Tube。 */ export const readableStreamToTube = ( readableStream: ReadableStream, ): Tube => { const tube = new Tube({ historyCount: Infinity, replayHistory: true, }) // NOTE: 这里确保 chunk 的消费和 Tube 的回调穿插在宏任务队列中执行, // 而不是所有 chunk 都获取完毕之后才开始处理 Tube 回调。 streamConsumeInMacroTask({ readableStream, onValue: async (chunk) => await tube.pushData(chunk), onDone: async () => await tube.close(), onError: async (error) => await tube.pushError(error as E), }) return tube }