Skip to content

Commit f3a9ea0

Browse files
rluvatontargos
authored andcommitted
stream: call helper function from push and unshift
PR-URL: #50173 Reviewed-By: Matteo Collina <matteo.collina@gmail.com> Reviewed-By: Robert Nagy <ronagy@icloud.com> Reviewed-By: Benjamin Gruenbaum <benjamingr@gmail.com>
1 parent 943047e commit f3a9ea0

1 file changed

Lines changed: 140 additions & 53 deletions

File tree

‎lib/internal/streams/readable.js‎

Lines changed: 140 additions & 53 deletions
Original file line numberDiff line numberDiff line change
@@ -284,77 +284,164 @@ Readable.prototype[SymbolAsyncDispose] = function() {
284284
// similar to how Writable.write() returns true if you should
285285
// write() some more.
286286
Readable.prototype.push=function(chunk,encoding){
287-
returnreadableAddChunk(this,chunk,encoding,false);
287+
debug('push',chunk);
288+
289+
conststate=this._readableState;
290+
return(state[kState]&kObjectMode)===0 ?
291+
readableAddChunkPushByteMode(this,state,chunk,encoding) :
292+
readableAddChunkPushObjectMode(this,state,chunk,encoding);
288293
};
289294

290295
// Unshift should *always* be something directly out of read().
291296
Readable.prototype.unshift=function(chunk,encoding){
292-
returnreadableAddChunk(this,chunk,encoding,true);
297+
debug('unshift',chunk);
298+
conststate=this._readableState;
299+
return(state[kState]&kObjectMode)===0 ?
300+
readableAddChunkUnshiftByteMode(this,state,chunk,encoding) :
301+
readableAddChunkUnshiftObjectMode(this,state,chunk);
293302
};
294303

295-
functionreadableAddChunk(stream,chunk,encoding,addToFront){
296-
debug('readableAddChunk',chunk);
297-
conststate=stream._readableState;
298304

299-
leterr;
300-
if((state[kState]&kObjectMode)===0){
301-
if(typeofchunk==='string'){
302-
encoding=encoding||state.defaultEncoding;
303-
if(state.encoding!==encoding){
304-
if(addToFront&&state.encoding){
305-
// When unshifting, if state.encoding is set, we have to save
306-
// the string in the BufferList with the state encoding.
307-
chunk=Buffer.from(chunk,encoding).toString(state.encoding);
308-
}else{
309-
chunk=Buffer.from(chunk,encoding);
310-
encoding='';
311-
}
305+
functionreadableAddChunkUnshiftByteMode(stream,state,chunk,encoding){
306+
if(chunk===null){
307+
state[kState]&=~kReading;
308+
onEofChunk(stream,state);
309+
310+
returnfalse;
311+
}
312+
313+
if(typeofchunk==='string'){
314+
encoding=encoding||state.defaultEncoding;
315+
if(state.encoding!==encoding){
316+
if(state.encoding){
317+
// When unshifting, if state.encoding is set, we have to save
318+
// the string in the BufferList with the state encoding.
319+
chunk=Buffer.from(chunk,encoding).toString(state.encoding);
320+
}else{
321+
chunk=Buffer.from(chunk,encoding);
312322
}
313-
}elseif(chunkinstanceofBuffer){
314-
encoding='';
315-
}elseif(Stream._isUint8Array(chunk)){
316-
chunk=Stream._uint8ArrayToBuffer(chunk);
317-
encoding='';
318-
}elseif(chunk!=null){
319-
err=newERR_INVALID_ARG_TYPE(
320-
'chunk',['string','Buffer','Uint8Array'],chunk);
321323
}
324+
}elseif(Stream._isUint8Array(chunk)){
325+
chunk=Stream._uint8ArrayToBuffer(chunk);
326+
}elseif(chunk!==undefined&&!(chunkinstanceofBuffer)){
327+
errorOrDestroy(stream,newERR_INVALID_ARG_TYPE(
328+
'chunk',['string','Buffer','Uint8Array'],chunk));
329+
returnfalse;
322330
}
323331

324-
if(err){
325-
errorOrDestroy(stream,err);
326-
}elseif(chunk===null){
332+
333+
if(!(chunk&&chunk.length>0)){
334+
returncanPushMore(state);
335+
}
336+
337+
returnreadableAddChunkUnshiftValue(stream,state,chunk);
338+
}
339+
340+
functionreadableAddChunkUnshiftObjectMode(stream,state,chunk){
341+
if(chunk===null){
327342
state[kState]&=~kReading;
328343
onEofChunk(stream,state);
329-
}elseif(((state[kState]&kObjectMode)!==0)||(chunk&&chunk.length>0)){
330-
if(addToFront){
331-
if((state[kState]&kEndEmitted)!==0)
332-
errorOrDestroy(stream,newERR_STREAM_UNSHIFT_AFTER_END_EVENT());
333-
elseif(state.destroyed||state.errored)
334-
returnfalse;
335-
else
336-
addChunk(stream,state,chunk,true);
337-
}elseif(state.ended){
338-
errorOrDestroy(stream,newERR_STREAM_PUSH_AFTER_EOF());
339-
}elseif(state.destroyed||state.errored){
340-
returnfalse;
341-
}else{
342-
state[kState]&=~kReading;
343-
if(state.decoder&&!encoding){
344-
chunk=state.decoder.write(chunk);
345-
if(state.objectMode||chunk.length!==0)
346-
addChunk(stream,state,chunk,false);
347-
else
348-
maybeReadMore(stream,state);
349-
}else{
350-
addChunk(stream,state,chunk,false);
351-
}
344+
345+
returnfalse;
346+
}
347+
348+
returnreadableAddChunkUnshiftValue(stream,state,chunk);
349+
}
350+
351+
functionreadableAddChunkUnshiftValue(stream,state,chunk){
352+
if((state[kState]&kEndEmitted)!==0)
353+
errorOrDestroy(stream,newERR_STREAM_UNSHIFT_AFTER_END_EVENT());
354+
elseif(state.destroyed||state.errored)
355+
returnfalse;
356+
else
357+
addChunk(stream,state,chunk,true);
358+
359+
returncanPushMore(state);
360+
}
361+
362+
functionreadableAddChunkPushByteMode(stream,state,chunk,encoding){
363+
if(chunk===null){
364+
state[kState]&=~kReading;
365+
onEofChunk(stream,state);
366+
367+
returnfalse;
368+
}
369+
370+
if(typeofchunk==='string'){
371+
encoding=encoding||state.defaultEncoding;
372+
if(state.encoding!==encoding){
373+
chunk=Buffer.from(chunk,encoding);
374+
encoding='';
352375
}
353-
}elseif(!addToFront){
376+
}elseif(chunkinstanceofBuffer){
377+
encoding='';
378+
}elseif(Stream._isUint8Array(chunk)){
379+
chunk=Stream._uint8ArrayToBuffer(chunk);
380+
encoding='';
381+
}elseif(chunk!==undefined){
382+
errorOrDestroy(stream,newERR_INVALID_ARG_TYPE(
383+
'chunk',['string','Buffer','Uint8Array'],chunk));
384+
returnfalse;
385+
}
386+
387+
if(!chunk||chunk.length<=0){
354388
state[kState]&=~kReading;
355389
maybeReadMore(stream,state);
390+
391+
returncanPushMore(state);
392+
}
393+
394+
if(state.ended){
395+
errorOrDestroy(stream,newERR_STREAM_PUSH_AFTER_EOF());
396+
397+
returnfalse;
398+
}
399+
400+
if(state.destroyed||state.errored){
401+
returnfalse;
402+
}
403+
404+
state[kState]&=~kReading;
405+
if(state.decoder&&!encoding){
406+
chunk=state.decoder.write(chunk);
407+
if(chunk.length===0){
408+
maybeReadMore(stream,state);
409+
410+
returncanPushMore(state);
411+
}
356412
}
357413

414+
addChunk(stream,state,chunk,false);
415+
returncanPushMore(state);
416+
}
417+
418+
functionreadableAddChunkPushObjectMode(stream,state,chunk,encoding){
419+
if(chunk===null){
420+
state[kState]&=~kReading;
421+
onEofChunk(stream,state);
422+
423+
returnfalse;
424+
}
425+
426+
if(state.ended){
427+
errorOrDestroy(stream,newERR_STREAM_PUSH_AFTER_EOF());
428+
returnfalse;
429+
}
430+
431+
if(state.destroyed||state.errored){
432+
returnfalse;
433+
}
434+
435+
state[kState]&=~kReading;
436+
if(state.decoder&&!encoding){
437+
chunk=state.decoder.write(chunk);
438+
}
439+
440+
addChunk(stream,state,chunk,false);
441+
returncanPushMore(state);
442+
}
443+
444+
functioncanPushMore(state){
358445
// We can push more data if we are below the highWaterMark.
359446
// Also, if we have no data yet, we can stand some more bytes.
360447
// This is to work around cases where hwm=0, such as the repl.

0 commit comments

Comments
 (0)