stream: preserve half-open duplexes in async iteration · nodejs/node@efbbb9a · GitHub
Skip to content

Commit efbbb9a

Browse files
efekrsklrichardlau
authored andcommitted
stream: preserve half-open duplexes in async iteration
Signed-off-by: Efe Karasakal <hi@efe.dev> PR-URL: #64275 Reviewed-By: James M Snell <jasnell@gmail.com> Reviewed-By: Jake Yuesong Li <jake.yuesong@gmail.com>
1 parent 9046035 commit efbbb9a

3 files changed

Lines changed: 129 additions & 1 deletion

File tree

lib/internal/streams/readable.js

Lines changed: 8 additions & 1 deletion
Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,80 @@
1+
'use strict';
2+
3+
const common = require('../common');
4+
const assert = require('assert');
5+
const net = require('net');
6+
7+
(async function() {
8+
let resolveServerSocket;
9+
const serverSocketPromise = new Promise((resolve) => {
10+
resolveServerSocket = resolve;
11+
});
12+
13+
const server = net.createServer({
14+
allowHalfOpen: true,
15+
}, common.mustCall((socket) => {
16+
resolveServerSocket(socket);
17+
}));
18+
19+
server.on('error', common.mustNotCall());
20+
server.on('close', common.mustCall());
21+
22+
await new Promise((resolve) => {
23+
server.listen(0, common.localhostIPv4, resolve);
24+
});
25+
26+
const clientSocket = await new Promise((resolve) => {
27+
const socket = net.createConnection({
28+
allowHalfOpen: true,
29+
port: server.address().port,
30+
host: server.address().address,
31+
}, common.mustCall(() => {
32+
resolve(socket);
33+
}));
34+
socket.on('error', common.mustNotCall());
35+
});
36+
37+
const serverSocket = await serverSocketPromise;
38+
serverSocket.on('error', common.mustNotCall());
39+
40+
await new Promise((resolve, reject) => {
41+
clientSocket.write('data written to client socket', (err) => {
42+
if (err) reject(err);
43+
else resolve();
44+
});
45+
});
46+
47+
await new Promise((resolve) => {
48+
clientSocket.end(resolve);
49+
});
50+
51+
let serverRead = '';
52+
for await (const chunk of serverSocket) {
53+
serverRead += chunk;
54+
}
55+
56+
assert.strictEqual(serverRead, 'data written to client socket');
57+
assert.strictEqual(serverSocket.destroyed, false);
58+
59+
await new Promise((resolve, reject) => {
60+
serverSocket.write('data written to server socket', (err) => {
61+
if (err) reject(err);
62+
else resolve();
63+
});
64+
});
65+
66+
await new Promise((resolve) => {
67+
serverSocket.end(resolve);
68+
});
69+
70+
let clientRead = '';
71+
for await (const chunk of clientSocket) {
72+
clientRead += chunk;
73+
}
74+
75+
assert.strictEqual(clientRead, 'data written to server socket');
76+
77+
await new Promise((resolve) => {
78+
server.close(resolve);
79+
});
80+
})().then(common.mustCall());
Lines changed: 41 additions & 0 deletions

0 commit comments

Comments
 (0)