From f9e1af95c2e91e76444f4f226846f4359fff38fd Mon Sep 17 00:00:00 2001 From: Basilisk-Dev Date: Fri, 8 May 2026 22:30:41 -0400 Subject: [PATCH] Fix readable stream transfer queueing --- js/src/builtin/Stream.cpp | 121 ++++++++++++++++++++++++++++++++++---- 1 file changed, 110 insertions(+), 11 deletions(-) diff --git a/js/src/builtin/Stream.cpp b/js/src/builtin/Stream.cpp index 7234ed596a..31f2238085 100644 --- a/js/src/builtin/Stream.cpp +++ b/js/src/builtin/Stream.cpp @@ -5058,7 +5058,15 @@ js::ByteLengthQueuingStrategy::constructor(JSContext* cx, unsigned argc, Value* { CallArgs args = CallArgsFromVp(argc, vp); - RootedObject strategy(cx, NewBuiltinClassInstance(cx)); + if (!ThrowIfNotConstructing(cx, args, "ByteLengthQueuingStrategy")) + return false; + + RootedObject proto(cx); + RootedObject newTarget(cx, &args.newTarget().toObject()); + if (!GetPrototypeFromConstructor(cx, newTarget, &proto)) + return false; + + RootedObject strategy(cx, NewObjectWithClassProto(cx, proto)); if (!strategy) return false; @@ -5104,7 +5112,16 @@ js::CountQueuingStrategy::constructor(JSContext* cx, unsigned argc, Value* vp) { CallArgs args = CallArgsFromVp(argc, vp); - Rooted strategy(cx, NewBuiltinClassInstance(cx)); + if (!ThrowIfNotConstructing(cx, args, "CountQueuingStrategy")) + return false; + + RootedObject proto(cx); + RootedObject newTarget(cx, &args.newTarget().toObject()); + if (!GetPrototypeFromConstructor(cx, newTarget, &proto)) + return false; + + Rooted strategy(cx); + strategy = NewObjectWithClassProto(cx, proto); if (!strategy) return false; @@ -6808,25 +6825,107 @@ js::InitStreamExtras(JSContext* cx, HandleObject global) if (stream.locked) return null; try { - return { reader: stream.getReader() }; + var record = { + reader: stream.getReader(), + pendingRead: undefined, + pumping: undefined, + closed: false, + errored: false, + canceled: false + }; + ensureReadableTransferRead(record); + return record; } catch (e) { return null; } } + function ensureReadableTransferRead(record) { + if (!record.pendingRead && !record.closed && !record.errored && + !record.canceled) { + try { + record.pendingRead = Promise.resolve(record.reader.read()); + } catch (e) { + record.pendingRead = Promise.reject(e); + } + silenceRejection(record.pendingRead); + } + return record.pendingRead; + } + + function errorReadableTransfer(record, controller, error) { + if (record.canceled || record.closed || record.errored) + return; + record.errored = true; + try { controller.error(error); } catch (e) {} + } + + function pumpReadableTransfer(record, controller) { + if (record.pumping || record.closed || record.errored || record.canceled) + return; + + var pendingRead = ensureReadableTransferRead(record); + if (!pendingRead) + return; + + record.pumping = pendingRead.then(function(result) { + record.pendingRead = undefined; + record.pumping = undefined; + + if (record.canceled || record.errored || record.closed) + return; + + if (result.done) { + record.closed = true; + try { controller.close(); } catch (e) {} + return; + } + + if (isUncloneableStreamChunk(result.value)) { + var cloneError = makeDataCloneError(); + silenceRejection(record.reader.cancel(cloneError)); + errorReadableTransfer(record, controller, cloneError); + return; + } + + try { + controller.enqueue(result.value); + } catch (e) { + silenceRejection(record.reader.cancel(e)); + errorReadableTransfer(record, controller, e); + return; + } + + if (controller.desiredSize > 0) + pumpReadableTransfer(record, controller); + }, function(error) { + record.pendingRead = undefined; + record.pumping = undefined; + errorReadableTransfer(record, controller, error); + }); + silenceRejection(record.pumping); + } + + function cancelReadableTransfer(record, reason) { + record.canceled = true; + Promise.resolve().then(function() { + try { + silenceRejection(record.reader.cancel(reason)); + } catch (e) {} + }); + return undefined; + } + function __uxpReceiveReadableStreamTransfer(record) { return new ReadableStream({ + start: function(controller) { + pumpReadableTransfer(record, controller); + }, pull: function(controller) { - return record.reader.read().then(function(result) { - if (result.done) { - controller.close(); - return; - } - controller.enqueue(result.value); - }); + pumpReadableTransfer(record, controller); }, cancel: function(reason) { - return record.reader.cancel(reason); + return cancelReadableTransfer(record, reason); } }); }