stream: validate writer options signal · nodejs/node@2a3babb · GitHub
Skip to content

Commit 2a3babb

Browse files
trivikraduh95
authored andcommitted
stream: validate writer options signal
Validate options.signal for stream/iter writer write(), writev(), and end() methods across push, broadcast, and fromWritable. Signed-off-by: Kamat, Trivikram <16024985+trivikr@users.noreply.github.com> Assisted-by: openai:gpt-5.5 PR-URL: #64385 Fixes: #64384 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent 0ab8f26 commit 2a3babb

5 files changed

Lines changed: 93 additions & 22 deletions

File tree

lib/internal/streams/iter/broadcast.js

Lines changed: 11 additions & 9 deletions

lib/internal/streams/iter/classic.js

Lines changed: 12 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,7 @@ const {
6161
} = require('internal/streams/iter/types');
6262

6363
const {
64+
getWriterSignal,
6465
validateBackpressure,
6566
toUint8Array,
6667
} = require('internal/streams/iter/utils');
@@ -572,10 +573,11 @@ function fromWritable(writable, options = kNullPrototype) {
572573
// as 'error' events caught by our generic error handler, rejecting
573574
// the next pending operation rather than the already-resolved one.
574575
//
575-
// The options.signal parameter from the Writer interface is ignored.
576-
// Classic stream.Writable has no per-write abort signal support;
577-
// cancellation should be handled at the pipeline level instead.
578-
write(chunk) {
576+
// The options.signal parameter from the Writer interface is validated but
577+
// otherwise ignored. Classic stream.Writable has no per-write abort signal
578+
// support; cancellation should be handled at the pipeline level instead.
579+
write(chunk, options) {
580+
getWriterSignal(options);
579581
if (!isWritable()) {
580582
return PromiseReject(new ERR_STREAM_WRITE_AFTER_END());
581583
}
@@ -617,10 +619,11 @@ function fromWritable(writable, options = kNullPrototype) {
617619
return PromiseResolve();
618620
},
619621

620-
writev(chunks) {
622+
writev(chunks, options) {
621623
if (!ArrayIsArray(chunks)) {
622624
throw new ERR_INVALID_ARG_TYPE('chunks', 'Array', chunks);
623625
}
626+
getWriterSignal(options);
624627
if (!isWritable()) {
625628
return PromiseReject(new ERR_STREAM_WRITE_AFTER_END());
626629
}
@@ -666,8 +669,10 @@ function fromWritable(writable, options = kNullPrototype) {
666669
return -1;
667670
},
668671

669-
// options.signal is ignored for the same reason as write().
670-
end() {
672+
// options.signal is validated but otherwise ignored for the same reason as
673+
// write().
674+
end(options) {
675+
getWriterSignal(options);
671676
if ((writable.writableFinished ?? false) ||
672677
(writable.destroyed ?? false)) {
673678
cleanup();

lib/internal/streams/iter/push.js

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ const {
4242
onSignalAbort,
4343
toUint8Array,
4444
convertChunks,
45+
getWriterSignal,
4546
parsePullArgs,
4647
validateBackpressure,
4748
} = require('internal/streams/iter/utils');
@@ -565,26 +566,28 @@ class PushWriter {
565566
}
566567

567568
write(chunk, options) {
568-
if (!options?.signal && this.#queue.canWriteSync()) {
569+
const signal = getWriterSignal(options);
570+
if (!signal && this.#queue.canWriteSync()) {
569571
const bytes = toUint8Array(chunk);
570572
this.#queue.writeSync([bytes]);
571573
return kResolvedPromise;
572574
}
573575
const bytes = toUint8Array(chunk);
574-
return this.#queue.writeAsync([bytes], options?.signal);
576+
return this.#queue.writeAsync([bytes], signal);
575577
}
576578

577579
writev(chunks, options) {
578580
if (!ArrayIsArray(chunks)) {
579581
throw new ERR_INVALID_ARG_TYPE('chunks', 'Array', chunks);
580582
}
581-
if (!options?.signal && this.#queue.canWriteSync()) {
583+
const signal = getWriterSignal(options);
584+
if (!signal && this.#queue.canWriteSync()) {
582585
const bytes = convertChunks(chunks);
583586
this.#queue.writeSync(bytes);
584587
return kResolvedPromise;
585588
}
586589
const bytes = convertChunks(chunks);
587-
return this.#queue.writeAsync(bytes, options?.signal);
590+
return this.#queue.writeAsync(bytes, signal);
588591
}
589592

590593
writeSync(chunk) {
@@ -601,6 +604,7 @@ class PushWriter {
601604
}
602605

603606
end(options) {
607+
getWriterSignal(options);
604608
const result = this.#queue.end();
605609
if (result === -2) {
606610
// Errored: reject with stored error

lib/internal/streams/iter/utils.js

Lines changed: 16 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,7 +35,10 @@ const { isError } = require('internal/util');
3535

3636
const { isSharedArrayBuffer, isUint8Array } = require('internal/util/types');
3737

38-
const { validateOneOf } = require('internal/validators');
38+
const {
39+
validateAbortSignal,
40+
validateOneOf,
41+
} = require('internal/validators');
3942

4043
// Cached resolved promise to avoid allocating a new one on every sync fast-path.
4144
const kResolvedPromise = PromiseResolve();
@@ -267,6 +270,17 @@ function convertChunks(chunks) {
267270
return result;
268271
}
269272

273+
/**
274+
* Validate Writer options and return options.signal.
275+
* @param {object|undefined} options
276+
* @returns {AbortSignal|undefined}
277+
*/
278+
function getWriterSignal(options) {
279+
const signal = options?.signal;
280+
validateAbortSignal(signal, 'options.signal');
281+
return signal;
282+
}
283+
270284
/**
271285
* Wrap a caught value as an Error, converting non-Error values.
272286
* @param {unknown} error
@@ -378,6 +392,7 @@ module.exports = {
378392
clampHWM,
379393
concatBytes,
380394
convertChunks,
395+
getWriterSignal,
381396
getMinCursor,
382397
hasProtocol,
383398
isPullOptions,

test/parallel/test-stream-iter-validation.js

Lines changed: 46 additions & 1 deletion

0 commit comments

Comments
 (0)