Skip to content

Commit fcbff00

Browse files
efekrskladuh95
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 9a20866 commit fcbff00

3 files changed

Lines changed: 129 additions & 1 deletion

File tree

‎lib/internal/streams/readable.js‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1401,9 +1401,16 @@ async function* createAsyncIterator(stream, options) {
14011401
error=aggregateTwoErrors(error,err);
14021402
throwerror;
14031403
}finally{
1404+
constpreserveHalfOpenDuplex=
1405+
error===null&&
1406+
stream.allowHalfOpen===true&&
1407+
stream.writable===true&&
1408+
stream.writableEnded!==true;
1409+
14041410
if(
14051411
(error||options?.destroyOnReturn!==false)&&
1406-
(error===undefined||stream._readableState.autoDestroy)
1412+
(error===undefined||stream._readableState.autoDestroy)&&
1413+
!preserveHalfOpenDuplex
14071414
){
14081415
destroyImpl.destroyer(stream,null);
14091416
}else{
Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,80 @@
1+
'use strict';
2+
3+
constcommon=require('../common');
4+
constassert=require('assert');
5+
constnet=require('net');
6+
7+
(asyncfunction(){
8+
letresolveServerSocket;
9+
constserverSocketPromise=newPromise((resolve)=>{
10+
resolveServerSocket=resolve;
11+
});
12+
13+
constserver=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+
awaitnewPromise((resolve)=>{
23+
server.listen(0,common.localhostIPv4,resolve);
24+
});
25+
26+
constclientSocket=awaitnewPromise((resolve)=>{
27+
constsocket=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+
constserverSocket=awaitserverSocketPromise;
38+
serverSocket.on('error',common.mustNotCall());
39+
40+
awaitnewPromise((resolve,reject)=>{
41+
clientSocket.write('data written to client socket',(err)=>{
42+
if(err)reject(err);
43+
elseresolve();
44+
});
45+
});
46+
47+
awaitnewPromise((resolve)=>{
48+
clientSocket.end(resolve);
49+
});
50+
51+
letserverRead='';
52+
forawait(constchunkofserverSocket){
53+
serverRead+=chunk;
54+
}
55+
56+
assert.strictEqual(serverRead,'data written to client socket');
57+
assert.strictEqual(serverSocket.destroyed,false);
58+
59+
awaitnewPromise((resolve,reject)=>{
60+
serverSocket.write('data written to server socket',(err)=>{
61+
if(err)reject(err);
62+
elseresolve();
63+
});
64+
});
65+
66+
awaitnewPromise((resolve)=>{
67+
serverSocket.end(resolve);
68+
});
69+
70+
letclientRead='';
71+
forawait(constchunkofclientSocket){
72+
clientRead+=chunk;
73+
}
74+
75+
assert.strictEqual(clientRead,'data written to server socket');
76+
77+
awaitnewPromise((resolve)=>{
78+
server.close(resolve);
79+
});
80+
})().then(common.mustCall());
Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,41 @@
1+
'use strict';
2+
3+
constcommon=require('../common');
4+
constassert=require('assert');
5+
const{
6+
Duplex,
7+
}=require('stream');
8+
9+
{
10+
letwritten='';
11+
12+
constduplex=newDuplex({
13+
allowHalfOpen: true,
14+
read(){
15+
this.push('hello');
16+
this.push(null);
17+
},
18+
write(chunk,encoding,callback){
19+
written+=chunk;
20+
callback();
21+
},
22+
});
23+
24+
duplex.on('error',common.mustNotCall());
25+
duplex.on('close',common.mustCall());
26+
27+
(async()=>{
28+
letread='';
29+
forawait(constchunkofduplex){
30+
read+=chunk;
31+
}
32+
33+
assert.strictEqual(read,'hello');
34+
assert.strictEqual(duplex.destroyed,false);
35+
36+
duplex.write('world',common.mustSucceed(()=>{
37+
assert.strictEqual(written,'world');
38+
duplex.end();
39+
}));
40+
})().then(common.mustCall());
41+
}

0 commit comments

Comments
 (0)