Skip to content
File

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

javascript185 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
4import { Duplex, Readable, Writable } from 'node:stream';
5import { strictEqual, deepStrictEqual, ok } from 'node:assert';
6import { mock } from 'node:test';
7 
8export const testStreamDuplexDestroy = {
9 async test() {
10 const { promise, resolve } = Promise.withResolvers();
11 const duplex = new Duplex({
12 read() {},
13 write(chunk, enc, cb) {
14 cb();
15 },
16 });
17 
18 duplex.cork();
19 duplex.write('foo', (err) => {
20 strictEqual(err.code, 'ERR_STREAM_DESTROYED');
21 resolve();
22 });
23 duplex.destroy();
24 await promise;
25 },
26};
27 
28// Prevents stream unexpected pause when highWaterMark set to 0
29// Ref: https://github.com/nodejs/node/commit/50695e5de14ccd8255537972181bbc9a1f44368e
30export const testStreamsHighwatermark = {
31 async test() {
32 const { promise, resolve } = Promise.withResolvers();
33 const res = [];
34 const r = new Readable({
35 read() {},
36 });
37 const w = new Writable({
38 highWaterMark: 0,
39 write(chunk, encoding, callback) {
40 res.push(chunk.toString());
41 callback();
42 },
43 });
44 
45 r.pipe(w);
46 r.push('a');
47 r.push('b');
48 r.push('c');
49 r.push(null);
50 
51 r.on('end', () => {
52 deepStrictEqual(res, ['a', 'b', 'c']);
53 resolve();
54 });
55 await promise;
56 },
57};
58 
59// Tests are taken from
60// https://github.com/nodejs/node/blob/9cc019575961ad9fcc18883993c3f9056699908d/test/parallel/test-stream-readable-data.js
61 
62export const testStreamReadableData = {
63 async test() {
64 const { promise, resolve } = Promise.withResolvers();
65 const readable = new Readable({
66 read() {},
67 });
68 
69 const onReadableFn = mock.fn();
70 readable.setEncoding('utf8');
71 readable.on('readable', onReadableFn);
72 readable.removeListener('readable', onReadableFn);
73 readable.on('end', resolve);
74 
75 const onDataFn = mock.fn();
76 queueMicrotask(function () {
77 readable.on('data', onDataFn);
78 readable.push('hello');
79 
80 queueMicrotask(() => {
81 readable.push(null);
82 });
83 });
84 
85 await promise;
86 strictEqual(onDataFn.mock.callCount(), 1);
87 strictEqual(onReadableFn.mock.callCount(), 0);
88 },
89};
90 
91// Tests are taken from
92// https://github.com/nodejs/node/blob/9cc019575961ad9fcc18883993c3f9056699908d/test/parallel/test-stream-readable-ended.js
93export const testStreamReadableEnded = {
94 async test() {
95 // basic
96 {
97 // Find it on Readable.prototype
98 ok(Object.hasOwn(Readable.prototype, 'readableEnded'));
99 }
100 
101 // event
102 {
103 const { promise, resolve } = Promise.withResolvers();
104 const readable = new Readable();
105 
106 readable._read = () => {
107 // The state ended should start in false.
108 strictEqual(readable.readableEnded, false);
109 readable.push('asd');
110 strictEqual(readable.readableEnded, false);
111 readable.push(null);
112 strictEqual(readable.readableEnded, false);
113 };
114 
115 const onDataFn = mock.fn(() => {
116 strictEqual(readable.readableEnded, false);
117 });
118 
119 readable.on('end', () => {
120 strictEqual(readable.readableEnded, true);
121 strictEqual(onDataFn.mock.callCount(), 1);
122 resolve();
123 });
124 
125 readable.on('data', onDataFn);
126 
127 await promise;
128 }
129 
130 // Verifies no `error` triggered on multiple .push(null) invocations
131 {
132 const { promise, resolve, reject } = Promise.withResolvers();
133 const readable = new Readable();
134 
135 readable.on('readable', () => {
136 readable.read();
137 });
138 readable.on('error', reject);
139 readable.on('end', resolve);
140 
141 readable.push('a');
142 readable.push(null);
143 readable.push(null);
144 
145 await promise;
146 }
147 },
148};
149 
150// Tests are taken from
151// https://github.com/nodejs/node/blob/9cc019575961ad9fcc18883993c3f9056699908d/test/parallel/test-stream-readable-hwm-0.js
152export const testStreamReadableHwm0 = {
153 async test() {
154 // This test ensures that Readable stream will call _read() for streams
155 // with highWaterMark === 0 upon .read(0) instead of just trying to
156 // emit 'readable' event.
157 const { promise, resolve } = Promise.withResolvers();
158 const onReadFn = mock.fn();
159 const r = new Readable({
160 // Must be called only once upon setting 'readable' listener
161 read: onReadFn,
162 highWaterMark: 0,
163 });
164 
165 let pushedNull = false;
166 // This will trigger read(0) but must only be called after push(null)
167 // because the we haven't pushed any data
168 const onReadableFn = mock.fn(() => {
169 strictEqual(r.read(), null);
170 strictEqual(pushedNull, true);
171 });
172 r.on('readable', onReadableFn);
173 r.on('end', resolve);
174 queueMicrotask(() => {
175 strictEqual(r.read(), null);
176 pushedNull = true;
177 r.push(null);
178 });
179 
180 await promise;
181 strictEqual(onReadFn.mock.callCount(), 1);
182 strictEqual(onReadableFn.mock.callCount(), 1);
183 },
184};