import { expect, test } from "vitest" import { streamConsumeInMacroTask, streamConsumeInMicroTask, streamFromArrayEager, streamFromArrayLazy, streamTransformInMacroTask, } from "#Source/basic/index.ts" test("streamFromArrayEager creates a readable stream from array values", async () => { const source = [1, 2, 3] const values: number[] = [] const stream = streamFromArrayEager(source) for await (const value of stream) { values.push(value) } expect(values).toEqual(source) }) test("streamFromArrayLazy creates a lazily consumed readable stream", async () => { const source = [1, 2, 3] const stream = streamFromArrayLazy(source) const reader = stream.getReader() const firstChunk = await reader.read() const secondChunk = await reader.read() const thirdChunk = await reader.read() const doneChunk = await reader.read() expect(firstChunk).toEqual({ done: false, value: 1 }) expect(secondChunk).toEqual({ done: false, value: 2 }) expect(thirdChunk).toEqual({ done: false, value: 3 }) expect(doneChunk).toEqual({ done: true, value: undefined }) }) test("streamConsumeInMicroTask consumes stream and handles callback errors", async () => { const consumed: number[] = [] let doneCalled = false await new Promise((resolve, reject) => { streamConsumeInMicroTask({ readableStream: streamFromArrayEager([1, 2, 3]), onValue: (chunk) => { consumed.push(chunk) }, onDone: () => { doneCalled = true resolve() }, onError: reject, }) }) expect(consumed).toEqual([1, 2, 3]) expect(doneCalled).toBe(true) let errorMessage = "" let doneCalledInErrorCase = false await new Promise((resolve) => { streamConsumeInMicroTask({ readableStream: streamFromArrayEager([1, 2]), onValue: (chunk) => { if (chunk === 2) { throw new Error("boom") } }, onDone: () => { doneCalledInErrorCase = true }, onError: (error) => { errorMessage = error.message resolve() }, }) }) expect(errorMessage).toContain("boom") expect(doneCalledInErrorCase).toBe(false) }) test("streamConsumeInMacroTask consumes stream and forwards callback failures", async () => { const consumed: number[] = [] await new Promise((resolve, reject) => { streamConsumeInMacroTask({ readableStream: streamFromArrayEager([1, 2, 3]), onValue: (chunk) => { consumed.push(chunk) }, onDone: () => { resolve() }, onError: (error) => { reject(error) }, }) }) expect(consumed).toEqual([1, 2, 3]) let errorMessage = "" await new Promise((resolve, reject) => { streamConsumeInMacroTask({ readableStream: streamFromArrayEager([1]), onValue: () => { throw new Error("async-boom") }, onDone: () => { reject(new Error("onDone should not be called when onValue throws")) }, onError: (error) => { errorMessage = error.message resolve() }, }) }) expect(errorMessage).toContain("async-boom") }) test("streamTransformInMacroTask transforms values and handles invalid inputs/errors", async () => { expect(() => { streamTransformInMacroTask({}) }).toThrow("Either readableStream or reader must be provided") const transformed = streamTransformInMacroTask({ readableStream: streamFromArrayEager([1, 2, 3]), onChunk: (chunk, controller) => { if (chunk.done === true) { controller.close() return } controller.enqueue(chunk.value * 10) }, }) const values: number[] = [] for await (const value of transformed) { values.push(value) } expect(values).toEqual([10, 20, 30]) const reader = streamFromArrayEager([4]).getReader() const transformedFromReader = streamTransformInMacroTask({ reader, onChunk: (chunk, controller) => { if (chunk.done === true) { controller.close() return } controller.enqueue(chunk.value + 1) }, }) expect(await Array.fromAsync(transformedFromReader)).toEqual([5]) let refinedError: Error | undefined const errored = streamTransformInMacroTask({ readableStream: streamFromArrayEager([1]), onChunk: () => { throw new Error("transform-boom") }, onError: (error) => { refinedError = new Error(`wrapped:${error.message}`) return refinedError }, }) await expect(Array.fromAsync(errored)).rejects.toThrow("wrapped:Error: transform-boom") expect(refinedError?.message).toBe("wrapped:Error: transform-boom") })