File
Blob: src/workerd/api/node/tests/streams-v24-nodejs-test.js
| 1 | // Copyright (c) 2025 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 | // Test for Node.js v24 streams compatibility |
| 5 | // These tests specifically verify the v24 behavior changes work correctly |
| 6 | |
| 7 | import { Readable, Writable } from 'node:stream'; |
| 8 | import { strictEqual } from 'node:assert'; |
| 9 | |
| 10 | // Regression test for v24 compat: large stream piping with backpressure |
| 11 | // This tests that needDrain is calculated AFTER write() completes, not before. |
| 12 | // |
| 13 | // Bug context: When needDrain was calculated before _write(), synchronous writes that modified |
| 14 | // state.length would cause the backpressure flag to be set based on stale buffer state. This |
| 15 | // created a deadlock where the stream would never emit 'drain', causing pipes to hang. |
| 16 | // |
| 17 | // Symptoms if broken: |
| 18 | // - Test will timeout (stream never completes) |
| 19 | // - totalWritten will be stuck at approximately 2.3MB (around the buffer threshold) |
| 20 | // - The stream hangs waiting for a 'drain' event that never fires |
| 21 | // - Worker will show "code had hung and would never generate a response" error |
| 22 | export const testLargeStreamPipeBackpressureV24 = { |
| 23 | async test() { |
| 24 | const { promise, resolve, reject } = Promise.withResolvers(); |
| 25 | |
| 26 | // Create a 6MB buffer to ensure we exceed highWaterMark multiple times |
| 27 | const chunkSize = 64 * 1024; // 64KB chunks |
| 28 | const numChunks = 100; // ~6.4MB total |
| 29 | const expectedTotal = chunkSize * numChunks; |
| 30 | |
| 31 | let totalWritten = 0; |
| 32 | let writesCompleted = 0; |
| 33 | |
| 34 | const readable = new Readable({ |
| 35 | read() { |
| 36 | // Push data in chunks |
| 37 | if (writesCompleted < numChunks) { |
| 38 | const chunk = Buffer.alloc(chunkSize, writesCompleted % 256); |
| 39 | this.push(chunk); |
| 40 | writesCompleted++; |
| 41 | } else { |
| 42 | this.push(null); // End the stream |
| 43 | } |
| 44 | }, |
| 45 | }); |
| 46 | |
| 47 | const writable = new Writable({ |
| 48 | // Set highWaterMark low to trigger backpressure |
| 49 | highWaterMark: 128 * 1024, // 128KB |
| 50 | write(chunk, encoding, callback) { |
| 51 | totalWritten += chunk.length; |
| 52 | // Simulate async write to ensure backpressure logic is exercised |
| 53 | setImmediate(callback); |
| 54 | }, |
| 55 | }); |
| 56 | |
| 57 | writable.on('finish', () => { |
| 58 | try { |
| 59 | strictEqual( |
| 60 | totalWritten, |
| 61 | expectedTotal, |
| 62 | `Expected ${expectedTotal} bytes but got ${totalWritten}` |
| 63 | ); |
| 64 | resolve(); |
| 65 | } catch (err) { |
| 66 | reject(err); |
| 67 | } |
| 68 | }); |
| 69 | |
| 70 | writable.on('error', reject); |
| 71 | readable.on('error', reject); |
| 72 | |
| 73 | // Pipe with backpressure handling |
| 74 | readable.pipe(writable); |
| 75 | |
| 76 | // Add timeout to catch the hang condition if the bug regresses |
| 77 | const timeout = setTimeout(() => { |
| 78 | reject( |
| 79 | new Error( |
| 80 | `Stream hung after ${totalWritten} bytes (expected ${expectedTotal}). ` + |
| 81 | `This indicates needDrain backpressure calculation is broken.` |
| 82 | ) |
| 83 | ); |
| 84 | }, 10000); |
| 85 | |
| 86 | try { |
| 87 | await promise; |
| 88 | } finally { |
| 89 | clearTimeout(timeout); |
| 90 | } |
| 91 | }, |
| 92 | }; |