Skip to content
File

Blob: src/workerd/api/node/tests/streams-v24-nodejs-test.js

javascript93 lines
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 
7import { Readable, Writable } from 'node:stream';
8import { 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
22export 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};