Skip to content

Commit 89a35d9

Browse files
committed
stream: add pipeToSync() failOnIncompleteClose option
pipeToSync() does not fail the writer when endSync() returns -1, so that a caller can still close it asynchronously. A caller that cannot, e.g. because the writer is sync-only and has no end(), would be left with a writer that is neither closed nor failed. Add a failOnIncompleteClose option (a Node.js extension) that fails the writer with the thrown ERR_INVALID_STATE error in that case. preventFail takes precedence over it. Assisted-by: OpenCode Signed-off-by: James M Snell <jasnell@gmail.com>
1 parent 7989e49 commit 89a35d9

4 files changed

Lines changed: 71 additions & 5 deletions

File tree

‎doc/api/stream_iter.md‎

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -716,6 +716,10 @@ added:
716716
* `...transforms` {Function|Object} Zero or more sync transforms.
717717
* `writer` {Object} Destination with `write(chunk)` method.
718718
* `options` {Object}
719+
* `failOnIncompleteClose` {boolean} If `true`, call `writer.fail()` when
720+
`writer.endSync()` cannot close the writer synchronously. Ignored when
721+
`preventFail` is `true`. This option is a Node.js extension.
722+
**Default:** `false`.
719723
* `preventClose` {boolean} **Default:** `false`.
720724
* `preventFail` {boolean} **Default:** `false`.
721725
* Returns: {number} Total bytes written.
@@ -730,8 +734,10 @@ The `writer` must have the `*Sync` methods (`writeSync`, `writevSync`,
730734
`writer.endSync()` returns `-1` because the writer cannot close synchronously
731735
(for example, a `push()` writer whose consumer has not read all of the data
732736
yet), `pipeToSync()` throws `ERR_INVALID_STATE`. All of the data was accepted
733-
by then, so the writer is not failed: it can still be closed, for example with
734-
`await writer.end()`.
737+
by then, so by default the writer is not failed: it can still be closed, for
738+
example with `await writer.end()`. If the writer cannot be closed any other
739+
way (for example, it has no `end()` method), or the caller will not close it,
740+
set `failOnIncompleteClose` to fail it with the thrown error instead.
735741

736742
### `pull(source[, ...transforms][, options])`
737743

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

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1088,11 +1088,16 @@ function pipeToSync(source, ...args) {
10881088

10891089
// endSync() returning -1 only means that the writer cannot close
10901090
// synchronously; every chunk was accepted. pipeToSync() never falls back to
1091-
// the async end(), so report it, but leave the writer as it is: the caller
1092-
// can still close it, e.g. with `await writer.end()`.
1091+
// the async end(), so report it. By default the writer is left as it is so
1092+
// that the caller can still close it (e.g. with `await writer.end()`);
1093+
// `failOnIncompleteClose` fails it instead, for callers that cannot.
10931094
if (!closedSync) {
1094-
throw new ERR_INVALID_STATE(
1095+
const error = new ERR_INVALID_STATE(
10951096
'Writer could not be closed synchronously');
1097+
if (options.failOnIncompleteClose && !options.preventFail) {
1098+
failWriterQuietly(writer, error);
1099+
}
1100+
throw error;
10961101
}
10971102

10981103
return totalBytes;

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

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -104,6 +104,13 @@ converters.PipeToOptions = createDictionaryConverter('PipeToOptions', [
104104
]);
105105
converters.PipeToSyncOptions = createDictionaryConverter(
106106
'PipeToSyncOptions', [
107+
// Node.js extension.
108+
{
109+
__proto__: null,
110+
key: 'failOnIncompleteClose',
111+
converter: baseConverters.boolean,
112+
defaultValue: () => false,
113+
},
107114
{
108115
__proto__: null,
109116
key: 'preventClose',

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

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -129,9 +129,57 @@ async function testFailThrowingDoesNotMaskError() {
129129
(error) => error === signal.reason);
130130
}
131131

132+
// failOnIncompleteClose fails a writer that cannot be closed synchronously,
133+
// e.g. a sync-only writer that has no end().
134+
async function testPipeToSyncFailOnIncompleteClose() {
135+
let failReason;
136+
const writer = {
137+
writeSync() { return true; },
138+
endSync: common.mustCall(() => -1),
139+
fail: common.mustCall((reason) => { failReason = reason; }),
140+
};
141+
assert.throws(
142+
() => pipeToSync(fromSync('data'), writer, { failOnIncompleteClose: true }),
143+
(error) => {
144+
assert.strictEqual(error.code, 'ERR_INVALID_STATE');
145+
assert.strictEqual(error, failReason);
146+
return true;
147+
});
148+
149+
// preventFail takes precedence.
150+
assert.throws(
151+
() => pipeToSync(fromSync('data'), {
152+
writeSync() { return true; },
153+
endSync: common.mustCall(() => -1),
154+
fail: common.mustNotCall(),
155+
}, { failOnIncompleteClose: true, preventFail: true }),
156+
{ code: 'ERR_INVALID_STATE' });
157+
158+
// It has no effect when the writer closes synchronously.
159+
assert.strictEqual(pipeToSync(fromSync('data'), {
160+
writeSync() { return true; },
161+
endSync: common.mustCall(() => 4),
162+
fail: common.mustNotCall(),
163+
}, { failOnIncompleteClose: true }), 4);
164+
165+
// A push() writer is failed with the error, so its consumer sees it.
166+
const { writer: pushWriter, readable } = push();
167+
let thrown;
168+
assert.throws(() => {
169+
try {
170+
pipeToSync(fromSync('abc'), pushWriter, { failOnIncompleteClose: true });
171+
} catch (error) {
172+
thrown = error;
173+
throw error;
174+
}
175+
}, { code: 'ERR_INVALID_STATE' });
176+
await assert.rejects(text(readable), (error) => error === thrown);
177+
}
178+
132179
Promise.all([
133180
testPipeToSyncEndSyncFailure(),
134181
testPipeToSyncEndSyncFailureDoesNotFailWriter(),
182+
testPipeToSyncFailOnIncompleteClose(),
135183
testPipeToSyncNoEndSync(),
136184
testPipeToSyncPreventFail(),
137185
testPipeToSyncPreventClose(),

0 commit comments

Comments
 (0)