Skip to content

Commit 8a31136

Browse files
mcollinatargos
authored andcommitted
stream: extract Readable.from in its own file
See: nodejs/readable-stream#420 PR-URL: #30140 Reviewed-By: Anna Henningsen <anna@addaleax.net> Reviewed-By: Colin Ihrig <cjihrig@gmail.com> Reviewed-By: Trivikram Kamat <trivikr.dev@gmail.com> Reviewed-By: Gus Caplan <me@gus.host> Reviewed-By: Beth Griggs <Bethany.Griggs@uk.ibm.com>
1 parent 375f349 commit 8a31136

3 files changed

Lines changed: 51 additions & 35 deletions

File tree

‎lib/_stream_readable.js‎

Lines changed: 4 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,7 @@ const {
4747
// Lazy loaded to improve the startup performance.
4848
letStringDecoder;
4949
letcreateReadableStreamAsyncIterator;
50+
letfrom;
5051

5152
Object.setPrototypeOf(Readable.prototype,Stream.prototype);
5253
Object.setPrototypeOf(Readable,Stream);
@@ -1194,40 +1195,8 @@ function endReadableNT(state, stream) {
11941195
}
11951196

11961197
Readable.from=function(iterable,opts){
1197-
letiterator;
1198-
if(iterable&&iterable[Symbol.asyncIterator])
1199-
iterator=iterable[Symbol.asyncIterator]();
1200-
elseif(iterable&&iterable[Symbol.iterator])
1201-
iterator=iterable[Symbol.iterator]();
1202-
else
1203-
thrownewERR_INVALID_ARG_TYPE('iterable',['Iterable'],iterable);
1204-
1205-
constreadable=newReadable({
1206-
objectMode: true,
1207-
...opts
1208-
});
1209-
// Reading boolean to protect against _read
1210-
// being called before last iteration completion.
1211-
letreading=false;
1212-
readable._read=function(){
1213-
if(!reading){
1214-
reading=true;
1215-
next();
1216-
}
1217-
};
1218-
asyncfunctionnext(){
1219-
try{
1220-
const{ value, done }=awaititerator.next();
1221-
if(done){
1222-
readable.push(null);
1223-
}elseif(readable.push(awaitvalue)){
1224-
next();
1225-
}else{
1226-
reading=false;
1227-
}
1228-
}catch(err){
1229-
readable.destroy(err);
1230-
}
1198+
if(from===undefined){
1199+
from=require('internal/streams/from');
12311200
}
1232-
returnreadable;
1201+
returnfrom(Readable,iterable,opts);
12331202
};

‎lib/internal/streams/from.js‎

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,46 @@
1+
'use strict';
2+
3+
const{
4+
ERR_INVALID_ARG_TYPE
5+
}=require('internal/errors').codes;
6+
7+
functionfrom(Readable,iterable,opts){
8+
letiterator;
9+
if(iterable&&iterable[Symbol.asyncIterator])
10+
iterator=iterable[Symbol.asyncIterator]();
11+
elseif(iterable&&iterable[Symbol.iterator])
12+
iterator=iterable[Symbol.iterator]();
13+
else
14+
thrownewERR_INVALID_ARG_TYPE('iterable',['Iterable'],iterable);
15+
16+
constreadable=newReadable({
17+
objectMode: true,
18+
...opts
19+
});
20+
// Reading boolean to protect against _read
21+
// being called before last iteration completion.
22+
letreading=false;
23+
readable._read=function(){
24+
if(!reading){
25+
reading=true;
26+
next();
27+
}
28+
};
29+
asyncfunctionnext(){
30+
try{
31+
const{ value, done }=awaititerator.next();
32+
if(done){
33+
readable.push(null);
34+
}elseif(readable.push(awaitvalue)){
35+
next();
36+
}else{
37+
reading=false;
38+
}
39+
}catch(err){
40+
readable.destroy(err);
41+
}
42+
}
43+
returnreadable;
44+
}
45+
46+
module.exports=from;

‎node.gyp‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -205,6 +205,7 @@
205205
'lib/internal/streams/async_iterator.js',
206206
'lib/internal/streams/buffer_list.js',
207207
'lib/internal/streams/duplexpair.js',
208+
'lib/internal/streams/from.js',
208209
'lib/internal/streams/legacy.js',
209210
'lib/internal/streams/destroy.js',
210211
'lib/internal/streams/state.js',

0 commit comments

Comments
 (0)