Commit 0e77f91f262 for nodejs
commit 0e77f91f262652665ca9fd7c8c95368c51e6f543
Author: James M Snell <jasnell@gmail.com>
Date: Sat Oct 3 15:00:10 2026 +0000
stream: keep the original error when writer.fail() throws
When pipeTo() or pipeToSync() failed, they called writer.fail(error)
and then rethrew the error. If fail() itself threw, its exception
replaced the error that made the pipe fail, which was then lost.
Call fail() on a best-effort basis and always surface the original
error.
Assisted-by: OpenCode
Signed-off-by: James M Snell <jasnell@gmail.com>
PR-URL: https://github.com/nodejs/node/pull/66483
Reviewed-By: Trivikram Kamat <trivikr.dev@gmail.com>
diff --git a/lib/internal/streams/iter/pull.js b/lib/internal/streams/iter/pull.js
index 792e8d33e9c..878fa420239 100644
--- a/lib/internal/streams/iter/pull.js
+++ b/lib/internal/streams/iter/pull.js
@@ -1083,7 +1083,7 @@ function pipeToSync(source, ...args) {
}
} catch (error) {
if (!options.preventFail) {
- writer.fail?.(error);
+ failWriterQuietly(writer, error);
}
throw error;
}
@@ -1091,6 +1091,17 @@ function pipeToSync(source, ...args) {
return totalBytes;
}
+// Call writer.fail(error) on a best-effort basis. The error that made the
+// pipe fail is what the caller must see; an exception from fail() itself
+// must not replace it.
+function failWriterQuietly(writer, error) {
+ try {
+ writer.fail?.(error);
+ } catch {
+ // Ignored, see above.
+ }
+}
+
/**
* Write an async source through transforms to a writer.
* @param {AsyncIterable<Uint8Array[]>|Iterable<Uint8Array[]>} source
@@ -1108,7 +1119,7 @@ async function pipeTo(source, ...args) {
function failWriter(error) {
if (!options.preventFail) {
- writer.fail?.(error);
+ failWriterQuietly(writer, error);
}
}
diff --git a/test/parallel/test-stream-iter-pipeto-edge.js b/test/parallel/test-stream-iter-pipeto-edge.js
index 13223a226a1..d63171203d6 100644
--- a/test/parallel/test-stream-iter-pipeto-edge.js
+++ b/test/parallel/test-stream-iter-pipeto-edge.js
@@ -5,7 +5,7 @@
const common = require('../common');
const assert = require('assert');
-const { pipeToSync, fromSync } = require('stream/iter');
+const { pipeTo, pipeToSync, fromSync } = require('stream/iter');
// pipeToSync cannot complete when endSync() requires async fallback.
async function testPipeToSyncEndSyncFailure() {
@@ -68,9 +68,46 @@ async function testPipeToSyncPreventClose() {
assert.strictEqual(endCalled, false);
}
+// An exception thrown by writer.fail() must not replace the error that made
+// the pipe fail.
+async function testFailThrowingDoesNotMaskError() {
+ const cause = new Error('write failed');
+ const syncWriter = {
+ writeSync() { throw cause; },
+ endSync: common.mustNotCall(),
+ fail: common.mustCall((error) => {
+ assert.strictEqual(error, cause);
+ throw new Error('fail() threw');
+ }),
+ };
+ assert.throws(() => pipeToSync(fromSync('data'), syncWriter),
+ (error) => error === cause);
+
+ const asyncWriter = {
+ async write() { throw cause; },
+ end: common.mustNotCall(),
+ fail: common.mustCall((error) => {
+ assert.strictEqual(error, cause);
+ throw new Error('fail() threw');
+ }),
+ };
+ await assert.rejects(pipeTo(fromSync('data'), asyncWriter),
+ (error) => error === cause);
+
+ // Same for the already-aborted signal path of pipeTo().
+ const signal = AbortSignal.abort();
+ await assert.rejects(
+ pipeTo(fromSync('data'), {
+ write: common.mustNotCall(),
+ fail() { throw new Error('fail() threw'); },
+ }, { signal }),
+ (error) => error === signal.reason);
+}
+
Promise.all([
testPipeToSyncEndSyncFailure(),
testPipeToSyncNoEndSync(),
testPipeToSyncPreventFail(),
testPipeToSyncPreventClose(),
+ testFailThrowingDoesNotMaskError(),
]).then(common.mustCall());