Skip to content
File

Blob: src/workerd/api/tests/streams-error-edge-cases-test.js

javascript280 lines
1// Copyright (c) 2026 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5// Tests for complex error scenarios in streams.
6// These tests focus on error type preservation, error propagation through
7// pipe chains, and race conditions between error and close operations.
8//
9// Test inspirations:
10// - Deno: tests/unit/streams_test.ts (cancel propagation, error type tests)
11// - Bun: test/js/web/streams/streams.test.js (pull rejection, error handling)
12// - Bun: test/js/web/fetch/fetch.stream.test.ts (corrupted data, socket close handling)
13// - Bun: test/js/bun/spawn/spawn-stdin-readable-stream-edge-cases.test.ts (exception in pull)
14 
15import { strictEqual, ok, rejects, deepStrictEqual } from 'node:assert';
16 
17// Custom error class for testing error type preservation
18class CustomStreamError extends Error {
19 constructor(message, code) {
20 super(message);
21 this.name = 'CustomStreamError';
22 this.code = code;
23 }
24}
25 
26// Test error thrown after partial consumption of stream
27// Inspired by: Bun test/js/bun/spawn/spawn-stdin-readable-stream-edge-cases.test.ts (exception in pull)
28export const errorDuringPartialConsumption = {
29 async test() {
30 let chunkCount = 0;
31 
32 const rs = new ReadableStream({
33 pull(controller) {
34 chunkCount++;
35 if (chunkCount <= 3) {
36 controller.enqueue(chunkCount);
37 } else {
38 controller.error(new Error('Error after 3 chunks'));
39 }
40 },
41 });
42 
43 const reader = rs.getReader();
44 const chunks = [];
45 
46 for (let i = 0; i < 3; i++) {
47 const { value, done } = await reader.read();
48 ok(!done);
49 chunks.push(value);
50 }
51 
52 deepStrictEqual(chunks, [1, 2, 3]);
53 
54 await rejects(reader.read(), { message: 'Error after 3 chunks' });
55 
56 await rejects(reader.read(), { message: 'Error after 3 chunks' });
57 },
58};
59 
60// Test that custom error types are preserved through pipeTo
61// Inspired by: Deno tests/unit/streams_test.ts (cancel propagation with "resource closed" reason)
62export const errorTypePreservationPipeTo = {
63 async test() {
64 const customError = new CustomStreamError('Custom error', 'ERR_CUSTOM');
65 
66 const rs = new ReadableStream({
67 start(controller) {
68 controller.error(customError);
69 },
70 });
71 
72 const ws = new WritableStream({
73 write() {},
74 });
75 
76 await rejects(
77 async () => {
78 await rs.pipeTo(ws);
79 },
80 { message: 'Custom error' }
81 );
82 },
83};
84 
85// Test that custom error types are preserved through pipeThrough
86// Inspired by: Deno tests/unit/streams_test.ts (error propagation tests)
87export const errorTypePreservationPipeThrough = {
88 async test() {
89 const customError = new CustomStreamError('Pipe through error', 'ERR_PIPE');
90 
91 const rs = new ReadableStream({
92 pull(controller) {
93 controller.error(customError);
94 },
95 });
96 
97 const transform = new TransformStream();
98 const result = rs.pipeThrough(transform);
99 
100 const reader = result.getReader();
101 
102 await rejects(
103 async () => {
104 await reader.read();
105 },
106 { message: 'Pipe through error' }
107 );
108 },
109};
110 
111// Test race between controller.error() and controller.close() on ReadableStream
112// Inspired by: Bun test/js/web/streams/streams.test.js (error handling edge cases)
113export const errorRaceWithCloseReadable = {
114 async test() {
115 let controller;
116 
117 const rs = new ReadableStream({
118 start(c) {
119 controller = c;
120 },
121 });
122 
123 const reader = rs.getReader();
124 const readPromise = reader.read();
125 
126 controller.error(new Error('Error wins'));
127 try {
128 controller.close();
129 } catch (_e) {
130 // May throw since stream is already errored
131 }
132 
133 await rejects(readPromise, { message: 'Error wins' });
134 },
135};
136 
137// Test race between writer.abort() and writer.close() on WritableStream
138// Inspired by: Bun test/js/web/streams/streams.test.js (abort/close race conditions)
139export const errorRaceWithCloseWritable = {
140 async test() {
141 let writeStarted = false;
142 
143 const ws = new WritableStream({
144 write() {
145 writeStarted = true;
146 // Simulate a slow write that can be aborted
147 return scheduler.wait(100);
148 },
149 });
150 
151 const writer = ws.getWriter();
152 
153 const writePromise = writer.write('data').catch((e) => e);
154 
155 await scheduler.wait(5);
156 ok(writeStarted, 'Write should have started');
157 
158 await writer.abort(new Error('Abort wins'));
159 
160 const writeResult = await writePromise;
161 ok(
162 writeResult === undefined || writeResult instanceof Error,
163 'Write should complete or error'
164 );
165 },
166};
167 
168// Test error thrown in TransformStream transform() callback using controller.error()
169// Inspired by: Bun test/js/web/streams/streams.test.js (TransformStream error handling)
170export const errorInTransformFlush = {
171 async test() {
172 let transformController;
173 
174 const ts = new TransformStream({
175 start(controller) {
176 transformController = controller;
177 },
178 transform(chunk, controller) {
179 controller.enqueue(chunk);
180 },
181 });
182 
183 const reader = ts.readable.getReader();
184 
185 transformController.error(new Error('Transform error'));
186 
187 await rejects(
188 async () => {
189 await reader.read();
190 },
191 { message: 'Transform error' }
192 );
193 },
194};
195 
196// Test error propagation through nested tee branches
197// Inspired by: Bun test/js/web/streams/streams.test.js (tee error handling)
198export const errorPropagationTeeMultiBranch = {
199 async test() {
200 let controller;
201 
202 const rs = new ReadableStream({
203 start(c) {
204 controller = c;
205 },
206 });
207 
208 // Create nested tees: original -> [branch1, temp] -> [branch2, branch3]
209 const [branch1, temp] = rs.tee();
210 const [branch2, branch3] = temp.tee();
211 
212 const reader1 = branch1.getReader();
213 const reader2 = branch2.getReader();
214 const reader3 = branch3.getReader();
215 
216 // Start reads on all branches
217 const read1 = reader1.read();
218 const read2 = reader2.read();
219 const read3 = reader3.read();
220 
221 // Error the source
222 controller.error(new Error('Source error'));
223 
224 // All branches should receive the error
225 const results = await Promise.allSettled([read1, read2, read3]);
226 
227 for (const result of results) {
228 strictEqual(result.status, 'rejected');
229 strictEqual(result.reason.message, 'Source error');
230 }
231 },
232};
233 
234// Test AbortSignal cancellation during active pipeTo
235// Inspired by: Deno tests/unit/streams_test.ts (abort tests), Bun test/js/web/streams/streams.test.js
236export const abortSignalDuringPipe = {
237 async test() {
238 const chunks = [];
239 let pullCount = 0;
240 
241 const rs = new ReadableStream({
242 async pull(controller) {
243 pullCount++;
244 await scheduler.wait(10);
245 if (pullCount <= 10) {
246 controller.enqueue(pullCount);
247 } else {
248 controller.close();
249 }
250 },
251 });
252 
253 const ws = new WritableStream({
254 write(chunk) {
255 chunks.push(chunk);
256 },
257 });
258 
259 const abortController = new AbortController();
260 
261 // Start piping
262 const pipePromise = rs.pipeTo(ws, { signal: abortController.signal });
263 
264 // Wait for some chunks to flow
265 await scheduler.wait(50);
266 
267 // Abort mid-pipe
268 abortController.abort(new Error('User cancelled'));
269 
270 // Pipe should reject with an error (type may vary by implementation)
271 await rejects(async () => {
272 await pipePromise;
273 }, Error);
274 
275 // Some chunks should have been written
276 ok(chunks.length > 0, 'Some chunks written before abort');
277 ok(chunks.length < 10, 'Not all chunks written due to abort');
278 },
279};