Skip to content

Commit 5d1832a

Browse files
authored
Release the pipe source when validation fails after locking (#7498)
1 parent 911aada commit 5d1832a

4 files changed

Lines changed: 167 additions & 94 deletions

File tree

‎src/per_isolate/webstreams/readable.ts‎

Lines changed: 99 additions & 91 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ import type {
1818
UnderlyingDefaultSource,
1919
UnderlyingSource,
2020
WritableStream as WritableStreamType,
21+
WritableStreamDefaultWriter as WritableStreamDefaultWriterType,
2122
} from './types';
2223
import type {
2324
ByteQueueEntry,
@@ -227,7 +228,7 @@ let readableStreamCancel: <R>(
227228
let readableStreamPipeThroughTo: <R>(
228229
source: ReadableStream<R>,
229230
destination: WritableStreamType<R>,
230-
options?: StreamPipeOptions
231+
options: ConvertedPipeOptions
231232
) => Promise<void>;
232233
let readableStreamPipeTo: <R>(
233234
source: ReadableStream<R>,
@@ -2911,35 +2912,36 @@ class ReadableStreamDrainingReader<R> {
29112912
}
29122913
}
29132914

2914-
// The pipe (spec ReadableStreamPipeTo). Internal operations only on both
2915-
// ends — locks are held for the duration. Chunks are read only while the
2916-
// destination desires them: the pump moves buffered chunks until
2917-
// desiredSize runs out, and reads the next chunk when none is buffered. It
2918-
// resumes from the destination's ready hook when a write completes, and
2919-
// when a read delivers, so a buffered backlog moves without a promise per
2920-
// wake-up. Writes are not awaited individually (the close path queues
2921-
// behind them, and write failures surface through the destination's closed
2922-
// promise). A shutdown waits for the writes made before running its
2923-
// action, including the write of a chunk that a pending read delivers
2924-
// during the wait.
2925-
function pipeToInternal<R>(
2926-
source: ReadableStream<R>,
2927-
destination: WritableStreamType<R>,
2928-
options: StreamPipeOptions = kEmptyDictionary as StreamPipeOptions
2929-
): Promise<void> {
2930-
// Spec-mandated read order (§4.9.1): preventAbort, preventCancel,
2931-
// preventClose, signal. WPT piping/throwing-options.any.js verifies
2932-
// that getter side-effects occur in exactly this sequence.
2933-
const preventAbort = !!options.preventAbort;
2934-
const preventCancel = !!options.preventCancel;
2935-
const preventClose = !!options.preventClose;
2936-
const signal = options.signal;
2915+
// StreamPipeOptions after WebIDL conversion: plain data, no getters.
2916+
interface ConvertedPipeOptions {
2917+
readonly preventAbort: boolean;
2918+
readonly preventCancel: boolean;
2919+
readonly preventClose: boolean;
2920+
readonly signal: AbortSignal | undefined;
2921+
}
2922+
2923+
// WebIDL conversion of StreamPipeOptions. The pipe methods run it before
2924+
// their locked checks, so a getter cannot change a lock after it is checked.
2925+
// Members are read once each in spec order (WPT
2926+
// piping/throwing-options.any.js).
2927+
function convertPipeOptions(options: unknown): ConvertedPipeOptions {
2928+
// WebIDL: null and undefined become {}.
2929+
if (options == null) options = kEmptyDictionary;
2930+
if (!isActualObject(options)) {
2931+
throw new TypeError('Pipe options must be an object');
2932+
}
2933+
const dict = options as StreamPipeOptions;
2934+
const preventAbort = !!dict.preventAbort;
2935+
const preventCancel = !!dict.preventCancel;
2936+
const preventClose = !!dict.preventClose;
2937+
const signal = dict.signal;
29372938
if (signal !== undefined) {
29382939
// Brand check. Under the modern JSG layout the captured `aborted` getter
29392940
// throws for non-AbortSignal receivers; under the instance-property
29402941
// layout (old compat dates) the capture is a plain read that cannot
29412942
// brand-check, so additionally require the boolean a genuine signal's
2942-
// own data property carries.
2943+
// own data property carries. (A forged {aborted: boolean} slips through
2944+
// under old dates only; the native fast path's C++ unwrap rejects it.)
29432945
let aborted: unknown;
29442946
try {
29452947
aborted = AbortSignalAbortedGet(signal);
@@ -2950,10 +2952,45 @@ function pipeToInternal<R>(
29502952
throw new TypeError('options.signal must be an AbortSignal');
29512953
}
29522954
}
2955+
return {
2956+
__proto__: null,
2957+
preventAbort,
2958+
preventCancel,
2959+
preventClose,
2960+
signal,
2961+
} as ConvertedPipeOptions;
2962+
}
2963+
2964+
// The pipe (spec ReadableStreamPipeTo). Internal operations only on both
2965+
// ends — locks are held for the duration. Chunks are read only while the
2966+
// destination desires them: the pump moves buffered chunks until
2967+
// desiredSize runs out, and reads the next chunk when none is buffered. It
2968+
// resumes from the destination's ready hook when a write completes, and
2969+
// when a read delivers, so a buffered backlog moves without a promise per
2970+
// wake-up. Writes are not awaited individually (the close path queues
2971+
// behind them, and write failures surface through the destination's closed
2972+
// promise). A shutdown waits for the writes made before running its
2973+
// action, including the write of a chunk that a pending read delivers
2974+
// during the wait.
2975+
function pipeToInternal<R>(
2976+
source: ReadableStream<R>,
2977+
destination: WritableStreamType<R>,
2978+
options: ConvertedPipeOptions
2979+
): Promise<void> {
2980+
const { preventAbort, preventCancel, preventClose, signal } = options;
29532981

2954-
// Lock both ends.
2982+
// Lock both ends. The callers have checked both locks; release the reader
2983+
// if the writer still cannot be acquired.
29552984
const reader = new ReadableStreamDefaultReader<R>(source);
2956-
const writer = writableInternals.acquireWriter(destination);
2985+
let writer: WritableStreamDefaultWriterType<R>;
2986+
try {
2987+
writer = writableInternals.acquireWriter(destination);
2988+
} catch (e) {
2989+
// The catch here is purely defensive. The acquireWriter
2990+
// should not actually throw.
2991+
readableStreamReaderGenericRelease(reader);
2992+
throw e;
2993+
}
29572994
setReadableStreamDisturbed(source);
29582995

29592996
const { promise, resolve, reject } =
@@ -3529,7 +3566,7 @@ class ReadableStream<R> {
35293566
readableStreamPipeThroughTo = <R>(
35303567
source: ReadableStream<R>,
35313568
destination: WritableStreamType<R>,
3532-
options?: StreamPipeOptions
3569+
options: ConvertedPipeOptions
35333570
) => {
35343571
// The pending-closure gate (see #pendingClosure). The prototype
35353572
// pipeThrough reaches pipeToInternal through here WITHOUT passing
@@ -3558,9 +3595,20 @@ class ReadableStream<R> {
35583595
readableStreamPipeTo = <R>(
35593596
source: ReadableStream<R>,
35603597
destination: WritableStreamType<R>,
3561-
options: StreamPipeOptions = kEmptyDictionary as StreamPipeOptions
3598+
options: unknown = kEmptyDictionary
35623599
): Promise<void> => {
35633600
try {
3601+
// WebIDL argument conversion first: the destination brand check,
3602+
// then the options, all before the locked checks and before the
3603+
// fast path below permanently consumes both endpoints. Both paths
3604+
// receive the converted values and never re-read the user's object.
3605+
if (!writableInternals.isWritableStream(destination)) {
3606+
throw new TypeError(
3607+
"Failed to execute 'pipeTo': destination is not a WritableStream"
3608+
);
3609+
}
3610+
const converted = convertPipeOptions(options);
3611+
const { preventAbort, preventCancel, preventClose } = converted;
35643612
if (isReadableStreamLocked(source)) {
35653613
throw new TypeError('Cannot pipe a stream that is locked');
35663614
}
@@ -3570,47 +3618,6 @@ class ReadableStream<R> {
35703618
if (source.#pendingClosure) {
35713619
throw pendingClosureError();
35723620
}
3573-
// WebIDL: null coerces to {} for optional dictionaries.
3574-
if (options === null) options = kEmptyDictionary as StreamPipeOptions;
3575-
if (!isActualObject(options)) {
3576-
throw new TypeError('Pipe options must be an object');
3577-
}
3578-
// Convert the options ONCE, before any extraction: the fast path
3579-
// below permanently consumes both endpoints, so validation (and
3580-
// the user-observable getter side effects, in the spec-mandated
3581-
// §4.9.1 order) must happen while both streams are still
3582-
// untouched. Both paths below receive the converted values as
3583-
// plain data properties; neither re-reads the user's options
3584-
// object.
3585-
const preventAbort = !!options.preventAbort;
3586-
const preventCancel = !!options.preventCancel;
3587-
const preventClose = !!options.preventClose;
3588-
const signal = options.signal;
3589-
if (signal !== undefined) {
3590-
// Brand check (same check and error text as the JS pump). Under
3591-
// the modern JSG layout the captured `aborted` getter throws for
3592-
// non-AbortSignal receivers; under the instance-property layout
3593-
// (old compat dates) the capture is a plain read that cannot
3594-
// brand-check, so additionally require the boolean a genuine
3595-
// signal's own data property carries. (A deliberately forged
3596-
// {aborted: boolean} object can slip through under old dates
3597-
// only; the fast path's C++ typed unwrap still rejects it.)
3598-
let aborted: unknown;
3599-
try {
3600-
aborted = AbortSignalAbortedGet(signal);
3601-
} catch {
3602-
throw new TypeError('options.signal must be an AbortSignal');
3603-
}
3604-
if (typeof aborted !== 'boolean') {
3605-
throw new TypeError('options.signal must be an AbortSignal');
3606-
}
3607-
}
3608-
const converted = {
3609-
preventAbort,
3610-
preventCancel,
3611-
preventClose,
3612-
signal,
3613-
} as StreamPipeOptions;
36143621

36153622
// PIPE DISPATCH: if both source and dest are native-backed, take
36163623
// the fast path -- extract both and let the sink's pipeFrom hook
@@ -3637,16 +3644,16 @@ class ReadableStream<R> {
36373644
// TODO(streams-ts): revisit extending the fast path to the
36383645
// prevent* options (e.g. reversible extraction) so those pipes can
36393646
// also run entirely at the C++ layer.
3640-
// Brand-check the destination before probing for native markers so
3641-
// a Proxy's getOwnPropertyDescriptor trap cannot observe the symbol.
3647+
// The destination is brand-checked above, so a Proxy's
3648+
// getOwnPropertyDescriptor trap cannot observe the symbol.
36423649
const sourceExtractor = ObjectGetOwnPropertyDescriptor(
36433650
source,
36443651
kExtractNativeSource
36453652
)?.value as ((this: ReadableStream<R>) => object) | undefined;
3646-
const sinkExtractor = writableInternals.isWritableStream(destination)
3647-
? (ObjectGetOwnPropertyDescriptor(destination, kExtractNativeSink)
3648-
?.value as ((this: object) => object) | undefined)
3649-
: undefined;
3653+
const sinkExtractor = ObjectGetOwnPropertyDescriptor(
3654+
destination,
3655+
kExtractNativeSink
3656+
)?.value as ((this: object) => object) | undefined;
36503657
if (
36513658
sourceExtractor !== undefined &&
36523659
sinkExtractor !== undefined &&
@@ -3674,7 +3681,7 @@ class ReadableStream<R> {
36743681
unknown
36753682
>;
36763683
const pipeFrom = nativeSink.pipeFrom as
3677-
| ((source: object, opts: StreamPipeOptions) => Promise<void>)
3684+
| ((source: object, opts: ConvertedPipeOptions) => Promise<void>)
36783685
| undefined;
36793686
if (pipeFrom === undefined) {
36803687
throw new TypeError(
@@ -4356,29 +4363,30 @@ class ReadableStream<R> {
43564363
options: StreamPipeOptions = kEmptyDictionary as StreamPipeOptions
43574364
): ReadableStreamType<T> {
43584365
assertIsReadableStream(this);
4359-
if (isReadableStreamLocked(this)) {
4360-
throw new TypeError('Cannot pipe a stream that is locked');
4361-
}
4362-
// WebIDL dictionary conversion: alphabetical member order means
4363-
// "readable" is read (and brand-checked) BEFORE "writable". The
4364-
// WPT pipe-through.any.js tests verify that a bad readable stops
4365-
// the writable getter from ever being called.
4366+
// WebIDL argument conversion precedes the locked checks, so no user
4367+
// code runs between those checks and pipeToInternal taking the locks.
4368+
// Dictionary members are read in alphabetical order: "readable" is read
4369+
// and brand-checked before "writable" (WPT pipe-through.any.js).
43664370
const readable = transform.readable;
43674371
if (!isReadableStream(readable)) {
43684372
throw new TypeError(
43694373
"Failed to execute 'pipeThrough': readable is not a ReadableStream"
43704374
);
43714375
}
43724376
const writable = transform.writable;
4373-
if (writable.locked) {
4374-
throw new TypeError('Cannot pipe to a locked writable stream');
4377+
if (!writableInternals.isWritableStream(writable)) {
4378+
throw new TypeError(
4379+
"Failed to execute 'pipeThrough': writable is not a WritableStream"
4380+
);
43754381
}
4376-
// WebIDL: null coerces to {} for optional dictionaries.
4377-
if (options === null) options = kEmptyDictionary as StreamPipeOptions;
4378-
if (!isActualObject(options)) {
4379-
throw new TypeError('Pipe options must be an object');
4382+
const converted = convertPipeOptions(options);
4383+
if (isReadableStreamLocked(this)) {
4384+
throw new TypeError('Cannot pipe a stream that is locked');
4385+
}
4386+
if (writableInternals.isWritableStreamLocked(writable)) {
4387+
throw new TypeError('Cannot pipe to a locked writable stream');
43804388
}
4381-
const promise = readableStreamPipeThroughTo(this, writable, options);
4389+
const promise = readableStreamPipeThroughTo(this, writable, converted);
43824390
markPromiseHandled(promise);
43834391
return readable;
43844392
}

‎src/tests/streams/piping/AGENTS.md‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -40,8 +40,9 @@ ends, preventAbort/preventCancel incl. TRUTHY coercion, dest stays
4040
usable under preventAbort); option plumbing (getter order
4141
[preventAbort, preventCancel, preventClose, signal], throwing-getter
4242
identity with no locks taken, bad-signal TypeError); pipeThrough
43-
locked-endpoint sync throws; custom error type/instance preservation;
44-
cancel-propagation through native identity AND JS transforms (source
43+
locked-endpoint sync throws; a bad destination, an option getter that
44+
locks it, or a shadowed `locked` fails without locking the source;
45+
custom error type/instance preservation; cancel-propagation through native identity AND JS transforms (source
4546
ends CLOSED, all locks release); external close()/abort() on a piped
4647
(locked) destination rejects while the pipe proceeds; backward write-
4748
error propagation with hook identity; backpressure through
@@ -75,7 +76,7 @@ the source FIRST, then releasing the write (`pipeStopsPullingWhenDestStalls`).
7576
| Module | Coverage |
7677
| --- | --- |
7778
| `pipe-matrix.js` | migrated pipe-streams-test.js wholesale (35): pipeThrough + pipeTo across JS↔native in all directions, prevent* combos, pre-aborted and mid-read AbortSignals, tee'd pipes, queued-destination close (ledger #1-#4, #13) |
78-
| `api-surface.js` | brand checks (ledger #5), option getter order, throwing getters, invalid signal, locked pipeThrough endpoints |
79+
| `api-surface.js` | brand checks (ledger #5), option getter order, throwing getters, invalid signal, locked pipeThrough endpoints, lock safety when validation fails |
7980
| `abort-signal.js` | the pipe's AbortSignal is an abort algorithm: a synthetic 'abort' event is ignored, a listener's stopImmediatePropagation() cannot block the abort, an abort after the pipe settles does nothing |
8081
| `error-propagation.js` | forward matrix (starts-errored × prevent* × truthy), hwm-0 dest (ledger #6), custom-error preservation (migrated from streams-error-edge-cases-test.js) |
8182
| `close-propagation.js` | the WPT-disabled backward territory, bounded: external close/abort on piped dest, write-throw backward propagation, idle dest-controller error and one with a write in flight (ledger #7) |

‎src/tests/streams/piping/api-surface.js‎

Lines changed: 61 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -118,3 +118,64 @@ export const pipeThroughLockedEndpoints = {
118118
strictEqual(rs2.locked, false);
119119
},
120120
};
121+
122+
// A TypeError without internal names (`#writer`), and no lock taken.
123+
const cleanTypeError = (e) => e.name === 'TypeError' && !/#/.test(e.message);
124+
125+
// A non-WritableStream destination fails before the source is locked.
126+
export const badDestinationLeavesSourceUnlocked = {
127+
async test() {
128+
const rs = new ReadableStream({});
129+
throws(
130+
() => rs.pipeThrough({ readable: new ReadableStream(), writable: {} }),
131+
cleanTypeError
132+
);
133+
strictEqual(rs.locked, false);
134+
await rejects(rs.pipeTo({}), cleanTypeError);
135+
strictEqual(rs.locked, false);
136+
},
137+
};
138+
139+
// Options are read before the locked checks, so an option getter that
140+
// locks the destination fails the pipe without locking the source.
141+
export const optionGetterLocksDestination = {
142+
async test() {
143+
const rs = new ReadableStream({});
144+
const t = new TransformStream();
145+
throws(
146+
() =>
147+
rs.pipeThrough(t, {
148+
get preventAbort() {
149+
t.writable.getWriter();
150+
return false;
151+
},
152+
}),
153+
cleanTypeError
154+
);
155+
strictEqual(rs.locked, false);
156+
const ws = new WritableStream({});
157+
await rejects(
158+
rs.pipeTo(ws, {
159+
get preventAbort() {
160+
ws.getWriter();
161+
return false;
162+
},
163+
}),
164+
cleanTypeError
165+
);
166+
strictEqual(rs.locked, false);
167+
},
168+
};
169+
170+
// pipeThrough checks the writable's lock internally, not through its
171+
// user-visible `locked`.
172+
export const shadowedWritableLocked = {
173+
test() {
174+
const rs = new ReadableStream({});
175+
const t = new TransformStream();
176+
t.writable.getWriter();
177+
Object.defineProperty(t.writable, 'locked', { value: false });
178+
throws(() => rs.pipeThrough(t), cleanTypeError);
179+
strictEqual(rs.locked, false);
180+
},
181+
};

‎src/tests/streams/piping/main.js‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,9 @@ export {
4848
invalidSignalRejected,
4949
brandChecks,
5050
pipeThroughLockedEndpoints,
51+
badDestinationLeavesSourceUnlocked,
52+
optionGetterLocksDestination,
53+
shadowedWritableLocked,
5154
} from 'api-surface';
5255

5356
export {

0 commit comments

Comments
 (0)