import { afterEach, expect, test, vi } from "vitest" import { streamFromArrayEager } from "#Source/basic/index.ts" import { connectTube, readableStreamToTube, Tube, tubeToPromise, tubeToReadableStream, } from "#Source/tube/index.ts" const createTube = ( options: Partial>[0]> = {}, ): Tube => { return new Tube({ historyCount: 3, replayHistory: false, ...options, }) } afterEach(() => { vi.useRealTimers() }) test("connectTube forwards lifecycle, data and error signals until disconnected", async () => { vi.useFakeTimers() const upstream = createTube({ autoEndOnError: false, autoCloseOnError: false }) const downstream = createTube({ autoEndOnError: false, autoCloseOnError: false }) const downstreamValues: number[] = [] downstream.subscribeData({ subscriber: (value) => { downstreamValues.push(value) }, }) const disconnect = connectTube(upstream, downstream) await upstream.pushData(1) await vi.runAllTimersAsync() expect(downstreamValues).toEqual([1]) expect(downstream.isOpen()).toBe(true) expect(downstream.isStart()).toBe(true) await upstream.pushError("boom") await vi.runAllTimersAsync() expect(downstream.getLatestErrorOrThrow()).toBe("boom") disconnect() await upstream.pushData(2) await vi.runAllTimersAsync() expect(downstreamValues).toEqual([1]) }) test("tubeToPromise resolves with the latest data on close and rejects on error or empty close", async () => { vi.useFakeTimers() const successTube = createTube() const successPromise = tubeToPromise(successTube) await successTube.pushData(1) await successTube.pushData(2) await successTube.close() await vi.runAllTimersAsync() await expect(successPromise).resolves.toBe(2) const errorTube = createTube() const errorPromise = tubeToPromise(errorTube) const errorExpectation = expect(errorPromise).rejects.toThrow("boom") await errorTube.pushError(new Error("boom")) await vi.runAllTimersAsync() await errorExpectation const emptyTube = createTube({ autoStartOnOpen: false, autoEndOnClose: false }) const emptyPromise = tubeToPromise(emptyTube) const emptyExpectation = expect(emptyPromise).rejects.toThrow("No data available") await emptyTube.open() await emptyTube.close() await vi.runAllTimersAsync() await emptyExpectation }) test("tubeToReadableStream exposes tube data and propagates stream errors", async () => { vi.useFakeTimers() const successTube = createTube() const successStream = tubeToReadableStream(successTube) const successRead = Array.fromAsync(successStream) await successTube.pushData(1) await successTube.pushData(2) await successTube.close() await vi.runAllTimersAsync() await expect(successRead).resolves.toEqual([1, 2]) const errorTube = createTube() const errorStream = tubeToReadableStream(errorTube) const errorRead = Array.fromAsync(errorStream) const errorExpectation = expect(errorRead).rejects.toThrow("broken") await errorTube.pushError(new Error("broken")) await vi.runAllTimersAsync() await errorExpectation }) test("readableStreamToTube consumes readable streams into a tube and records stream errors", async () => { vi.useFakeTimers() const values: number[] = [] const successTube = readableStreamToTube(streamFromArrayEager([1, 2, 3])) successTube.subscribeData({ subscriber: (value) => { values.push(value) }, }) await vi.runAllTimersAsync() expect(values).toEqual([1, 2, 3]) expect(successTube.hasClosed()).toBe(true) expect(successTube.getLatestDataOrThrow()).toBe(3) const consoleErrorSpy = vi.spyOn(console, "error").mockImplementation(() => undefined) try { const failureTube = readableStreamToTube( new ReadableStream({ start(controller): void { controller.error(new Error("stream-failed")) }, }), ) await vi.runAllTimersAsync() expect(failureTube.getLatestErrorOrThrow()).toBeInstanceOf(Error) expect(failureTube.getLatestErrorOrThrow().message).toBe("Error: stream-failed") } finally { consoleErrorSpy.mockRestore() } })