Skip to content

Commit fa50684

Browse files
jasnelladuh95
authored andcommitted
stream: pre-aborted pipeTo now applies dest failure handling
Signed-off-by: James M Snell <jasnell@gmail.com> Assisted-by: Opencode PR-URL: #65658 Reviewed-By: Trivikram Kamat <trivikr.dev@gmail.com>
1 parent 05d1166 commit fa50684

3 files changed

Lines changed: 62 additions & 6 deletions

File tree

‎doc/api/stream_iter.md‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -605,7 +605,8 @@ added: v24.20.0
605605
* `...transforms` {Function|Object} Zero or more transforms to apply.
606606
* `writer` {Object} Destination with `write(chunk)` method.
607607
* `options` {Object}
608-
* `signal` {AbortSignal} Abort the pipeline.
608+
* `signal` {AbortSignal} Abort the pipeline. Aborting fails the destination
609+
writer unless `preventFail` is `true`.
609610
* `preventClose` {boolean} If `true`, do not call `writer.end()` when
610611
the source ends. **Default:** `false`.
611612
* `preventFail` {boolean} If `true`, do not call `writer.fail()` on

‎lib/internal/streams/iter/pull.js‎

Lines changed: 13 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1041,8 +1041,18 @@ async function pipeTo(source, ...args) {
10411041

10421042
const signal = options?.signal;
10431043

1044-
// Check for abort
1045-
signal?.throwIfAborted();
1044+
function failWriter(error) {
1045+
if (!options?.preventFail) {
1046+
writer.fail?.(wrapError(error));
1047+
}
1048+
}
1049+
1050+
try {
1051+
signal?.throwIfAborted();
1052+
} catch (error) {
1053+
failWriter(error);
1054+
throw error;
1055+
}
10461056

10471057
const hasWriteSync = typeof writer.writeSync === 'function';
10481058
const useSyncIterableFastPath =
@@ -1169,9 +1179,7 @@ async function pipeTo(source, ...args) {
11691179
}
11701180
}
11711181
} catch (error) {
1172-
if (!options?.preventFail) {
1173-
writer.fail?.(wrapError(error));
1174-
}
1182+
failWriter(error);
11751183
throw error;
11761184
}
11771185

‎test/parallel/test-stream-iter-pipeto-signal.js‎

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,51 @@ const assert = require('assert');
99
const { setTimeout } = require('timers/promises');
1010
const { pipeTo, from } = require('stream/iter');
1111

12+
async function testPipeToPreAbortedSignalFailsWriter() {
13+
const reason = new Error('already aborted');
14+
let sourceTouched = false;
15+
const source = {
16+
[Symbol.asyncIterator]() {
17+
sourceTouched = true;
18+
return {};
19+
},
20+
};
21+
const writer = {
22+
write: common.mustNotCall(),
23+
fail: common.mustCall((error) => assert.strictEqual(error, reason)),
24+
};
25+
26+
await assert.rejects(
27+
pipeTo(source, writer, { signal: AbortSignal.abort(reason) }),
28+
(error) => error === reason,
29+
);
30+
assert.strictEqual(sourceTouched, false);
31+
}
32+
33+
async function testPipeToPreAbortedSignalPreventFail() {
34+
const reason = new Error('already aborted');
35+
let sourceTouched = false;
36+
const source = {
37+
[Symbol.asyncIterator]() {
38+
sourceTouched = true;
39+
return {};
40+
},
41+
};
42+
const writer = {
43+
write: common.mustNotCall(),
44+
fail: common.mustNotCall(),
45+
};
46+
47+
await assert.rejects(
48+
pipeTo(source, writer, {
49+
signal: AbortSignal.abort(reason),
50+
preventFail: true,
51+
}),
52+
(error) => error === reason,
53+
);
54+
assert.strictEqual(sourceTouched, false);
55+
}
56+
1257
// pipeTo with live signal, no transforms — abort mid-stream
1358
async function testPipeToLiveSignalNoTransforms() {
1459
const ac = new AbortController();
@@ -116,6 +161,8 @@ async function testPipeToLiveSignalWithTransformsCompletes() {
116161
}
117162

118163
Promise.all([
164+
testPipeToPreAbortedSignalFailsWriter(),
165+
testPipeToPreAbortedSignalPreventFail(),
119166
testPipeToLiveSignalNoTransforms(),
120167
testPipeToLiveSignalNoTransformsPendingNext(),
121168
testPipeToLiveSignalWithTransforms(),

0 commit comments

Comments
 (0)