Skip to content

Commit 6c454af

Browse files
debadree25danielleadams
authored andcommitted
stream: add pipeline() for webstreams
Refs: #39316 PR-URL: #46307 Reviewed-By: Robert Nagy <ronagy@icloud.com> Reviewed-By: Matteo Collina <matteo.collina@gmail.com> Reviewed-By: Benjamin Gruenbaum <benjamingr@gmail.com>
1 parent 91a550e commit 6c454af

5 files changed

Lines changed: 500 additions & 10 deletions

File tree

‎doc/api/stream.md‎

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2669,6 +2669,9 @@ const cleanup = finished(rs, (err) => {
26692669
<!-- YAML
26702670
added: v10.0.0
26712671
changes:
2672+
- version: REPLACEME
2673+
pr-url: https://github.com/nodejs/node/pull/46307
2674+
description: Added support for webstreams.
26722675
- version: v18.0.0
26732676
pr-url: https://github.com/nodejs/node/pull/41678
26742677
description: Passing an invalid callback to the `callback` argument
@@ -2685,13 +2688,14 @@ changes:
26852688
description: Add support for async generators.
26862689
-->
26872690

2688-
*`streams` {Stream\[]|Iterable\[]|AsyncIterable\[]|Function\[]}
2689-
*`source` {Stream|Iterable|AsyncIterable|Function}
2691+
*`streams` {Stream\[]|Iterable\[]|AsyncIterable\[]|Function\[]|
2692+
ReadableStream\[]|WritableStream\[]|TransformStream\[]}
2693+
*`source` {Stream|Iterable|AsyncIterable|Function|ReadableStream}
26902694
* Returns: {Iterable|AsyncIterable}
2691-
*`...transforms` {Stream|Function}
2695+
*`...transforms` {Stream|Function|TransformStream}
26922696
*`source` {AsyncIterable}
26932697
* Returns: {AsyncIterable}
2694-
*`destination` {Stream|Function}
2698+
*`destination` {Stream|Function|WritableStream}
26952699
*`source` {AsyncIterable}
26962700
* Returns: {AsyncIterable|Promise}
26972701
*`callback` {Function} Called when the pipeline is fully done.

‎lib/internal/streams/pipeline.js‎

Lines changed: 63 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,9 @@ const {
3535
isReadable,
3636
isReadableNodeStream,
3737
isNodeStream,
38+
isTransformStream,
39+
isWebStream,
40+
isReadableStream,
3841
}=require('internal/streams/utils');
3942
const{ AbortController }=require('internal/abort_controller');
4043

@@ -88,7 +91,7 @@ async function* fromReadable(val) {
8891
yield*Readable.prototype[SymbolAsyncIterator].call(val);
8992
}
9093

91-
asyncfunctionpump(iterable,writable,finish,{ end }){
94+
asyncfunctionpumpToNode(iterable,writable,finish,{ end }){
9295
leterror;
9396
letonresolve=null;
9497

@@ -147,6 +150,35 @@ async function pump(iterable, writable, finish, { end }) {
147150
}
148151
}
149152

153+
asyncfunctionpumpToWeb(readable,writable,finish,{ end }){
154+
if(isTransformStream(writable)){
155+
writable=writable.writable;
156+
}
157+
// https://streams.spec.whatwg.org/#example-manual-write-with-backpressure
158+
constwriter=writable.getWriter();
159+
try{
160+
forawait(constchunkofreadable){
161+
awaitwriter.ready;
162+
writer.write(chunk).catch(()=>{});
163+
}
164+
165+
awaitwriter.ready;
166+
167+
if(end){
168+
awaitwriter.close();
169+
}
170+
171+
finish();
172+
}catch(err){
173+
try{
174+
awaitwriter.abort(err);
175+
finish(err);
176+
}catch(err){
177+
finish(err);
178+
}
179+
}
180+
}
181+
150182
functionpipeline(...streams){
151183
returnpipelineImpl(streams,once(popCallback(streams)));
152184
}
@@ -259,7 +291,11 @@ function pipelineImpl(streams, callback, opts) {
259291
ret=Duplex.from(stream);
260292
}
261293
}elseif(typeofstream==='function'){
262-
ret=makeAsyncIterable(ret);
294+
if(isTransformStream(ret)){
295+
ret=makeAsyncIterable(ret?.readable);
296+
}else{
297+
ret=makeAsyncIterable(ret);
298+
}
263299
ret=stream(ret,{ signal });
264300

265301
if(reading){
@@ -303,7 +339,11 @@ function pipelineImpl(streams, callback, opts) {
303339
);
304340
}elseif(isIterable(ret,true)){
305341
finishCount++;
306-
pump(ret,pt,finish,{ end });
342+
pumpToNode(ret,pt,finish,{ end });
343+
}elseif(isReadableStream(ret)||isTransformStream(ret)){
344+
consttoRead=ret.readable||ret;
345+
finishCount++;
346+
pumpToNode(toRead,pt,finish,{ end });
307347
}else{
308348
thrownewERR_INVALID_RETURN_VALUE(
309349
'AsyncIterable or Promise','destination',ret);
@@ -324,12 +364,30 @@ function pipelineImpl(streams, callback, opts) {
324364
if(isReadable(stream)&&isLastStream){
325365
lastStreamCleanup.push(cleanup);
326366
}
367+
}elseif(isTransformStream(ret)||isReadableStream(ret)){
368+
consttoRead=ret.readable||ret;
369+
finishCount++;
370+
pumpToNode(toRead,stream,finish,{ end });
327371
}elseif(isIterable(ret)){
328372
finishCount++;
329-
pump(ret,stream,finish,{ end });
373+
pumpToNode(ret,stream,finish,{ end });
374+
}else{
375+
thrownewERR_INVALID_ARG_TYPE(
376+
'val',['Readable','Iterable','AsyncIterable','ReadableStream','TransformStream'],ret);
377+
}
378+
ret=stream;
379+
}elseif(isWebStream(stream)){
380+
if(isReadableNodeStream(ret)){
381+
finishCount++;
382+
pumpToWeb(makeAsyncIterable(ret),stream,finish,{ end });
383+
}elseif(isReadableStream(ret)||isIterable(ret)){
384+
finishCount++;
385+
pumpToWeb(ret,stream,finish,{ end });
386+
}elseif(isTransformStream(ret)){
387+
pumpToWeb(ret.readable,stream,finish,{ end });
330388
}else{
331389
thrownewERR_INVALID_ARG_TYPE(
332-
'val',['Readable','Iterable','AsyncIterable'],ret);
390+
'val',['Readable','Iterable','AsyncIterable','ReadableStream','TransformStream'],ret);
333391
}
334392
ret=stream;
335393
}else{

‎lib/internal/streams/utils.js‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,19 @@ function isWritableStream(obj) {
7777
);
7878
}
7979

80+
functionisTransformStream(obj){
81+
return!!(
82+
obj&&
83+
!isNodeStream(obj)&&
84+
typeofobj.readable==='object'&&
85+
typeofobj.writable==='object'
86+
);
87+
}
88+
89+
functionisWebStream(obj){
90+
returnisReadableStream(obj)||isWritableStream(obj)||isTransformStream(obj);
91+
}
92+
8093
functionisIterable(obj,isAsync){
8194
if(obj==null)returnfalse;
8295
if(isAsync===true)returntypeofobj[SymbolAsyncIterator]==='function';
@@ -303,6 +316,7 @@ module.exports = {
303316
isReadableFinished,
304317
isReadableErrored,
305318
isNodeStream,
319+
isWebStream,
306320
isWritable,
307321
isWritableNodeStream,
308322
isWritableStream,
@@ -312,4 +326,5 @@ module.exports = {
312326
isServerRequest,
313327
isServerResponse,
314328
willEmitClose,
329+
isTransformStream,
315330
};

‎lib/stream/promises.js‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ const {
88
const{
99
isIterable,
1010
isNodeStream,
11+
isWebStream,
1112
}=require('internal/streams/utils');
1213

1314
const{pipelineImpl: pl}=require('internal/streams/pipeline');
@@ -21,7 +22,7 @@ function pipeline(...streams) {
2122
letend;
2223
constlastArg=streams[streams.length-1];
2324
if(lastArg&&typeoflastArg==='object'&&
24-
!isNodeStream(lastArg)&&!isIterable(lastArg)){
25+
!isNodeStream(lastArg)&&!isIterable(lastArg)&&!isWebStream(lastArg)){
2526
constoptions=ArrayPrototypePop(streams);
2627
signal=options.signal;
2728
end=options.end;

0 commit comments

Comments
 (0)