diff --git a/dom/base/StructuredCloneHolder.cpp b/dom/base/StructuredCloneHolder.cpp index a7afd310cb..ed82cf9540 100644 --- a/dom/base/StructuredCloneHolder.cpp +++ b/dom/base/StructuredCloneHolder.cpp @@ -138,6 +138,120 @@ StructuredCloneCallbacksError(JSContext* aCx, NS_WARNING("Failed to clone data."); } +struct SameThreadStreamTransferData +{ + SameThreadStreamTransferData(JSContext* aCx, JS::HandleObject aRecord) + : mRecord(aCx, aRecord) + {} + + ~SameThreadStreamTransferData() + { + mRecord.reset(); + } + + JS::PersistentRootedObject mRecord; +}; + +bool +CallStreamTransferHelper(JSContext* aCx, + const char* aHelperName, + JS::HandleValue aArgument, + JS::MutableHandleValue aResult) +{ + JS::Rooted global(aCx, JS::CurrentGlobalOrNull(aCx)); + if (!global) { + return false; + } + + // Resolve the stream extras lazily before looking up the internal helper. + JS::Rooted ignored(aCx); + if (!JS_GetProperty(aCx, global, "WritableStream", &ignored)) { + return false; + } + + JS::Rooted helper(aCx); + if (!JS_GetProperty(aCx, global, aHelperName, &helper)) { + return false; + } + + if (!helper.isObject() || !JS::IsCallable(&helper.toObject())) { + aResult.setUndefined(); + return true; + } + + JS::Rooted argument(aCx, aArgument); + if (!JS_WrapValue(aCx, &argument)) { + return false; + } + + return JS::Call(aCx, JS::UndefinedHandleValue, helper, + JS::HandleValueArray(argument), aResult); +} + +bool +TryWriteSameThreadStreamTransfer(JSContext* aCx, + JS::Handle aObj, + const char* aHelperName, + uint32_t aTransferTag, + uint32_t* aTag, + JS::TransferableOwnership* aOwnership, + void** aContent, + uint64_t* aExtraData, + bool* aHandled) +{ + *aHandled = false; + + JS::Rooted argument(aCx, JS::ObjectValue(*aObj)); + JS::Rooted transferRecord(aCx); + if (!CallStreamTransferHelper(aCx, aHelperName, argument, &transferRecord)) { + return false; + } + + if (transferRecord.isUndefined()) { + return true; + } + + if (!transferRecord.isObject()) { + return false; + } + + JS::Rooted record(aCx, &transferRecord.toObject()); + SameThreadStreamTransferData* data = + new SameThreadStreamTransferData(aCx, record); + + *aTag = aTransferTag; + *aOwnership = JS::SCTAG_TMO_CUSTOM; + *aContent = data; + *aExtraData = 0; + *aHandled = true; + return true; +} + +bool +ReadSameThreadStreamTransfer(JSContext* aCx, + void* aContent, + const char* aHelperName, + JS::MutableHandleObject aReturnObject) +{ + MOZ_ASSERT(aContent); + SameThreadStreamTransferData* data = + static_cast(aContent); + + JS::Rooted record(aCx, JS::ObjectValue(*data->mRecord)); + JS::Rooted result(aCx); + if (!CallStreamTransferHelper(aCx, aHelperName, record, &result)) { + return false; + } + + if (!result.isObject()) { + return false; + } + + aReturnObject.set(&result.toObject()); + delete data; + return true; +} + } // anonymous namespace const JSStructuredCloneCallbacks StructuredCloneHolder::sCallbacks = { @@ -1375,6 +1489,20 @@ StructuredCloneHolder::CustomReadTransferHandler(JSContext* aCx, return true; } + if (mStructuredCloneScope == StructuredCloneScope::SameProcessSameThread) { + if (aTag == SCTAG_DOM_TRANSFERRED_WRITABLESTREAM) { + return ReadSameThreadStreamTransfer(aCx, aContent, + "__uxpReceiveWritableStreamTransfer", + aReturnObject); + } + + if (aTag == SCTAG_DOM_TRANSFERRED_READABLESTREAM) { + return ReadSameThreadStreamTransfer(aCx, aContent, + "__uxpReceiveReadableStreamTransfer", + aReturnObject); + } + } + return false; } @@ -1443,6 +1571,31 @@ StructuredCloneHolder::CustomWriteTransferHandler(JSContext* aCx, } } + if (mStructuredCloneScope == StructuredCloneScope::SameProcessSameThread) { + bool handled = false; + if (!TryWriteSameThreadStreamTransfer(aCx, obj, + "__uxpTransferWritableStream", + SCTAG_DOM_TRANSFERRED_WRITABLESTREAM, + aTag, aOwnership, aContent, + aExtraData, &handled)) { + return false; + } + if (handled) { + return true; + } + + if (!TryWriteSameThreadStreamTransfer(aCx, obj, + "__uxpTransferReadableStream", + SCTAG_DOM_TRANSFERRED_READABLESTREAM, + aTag, aOwnership, aContent, + aExtraData, &handled)) { + return false; + } + if (handled) { + return true; + } + } + return false; } @@ -1480,6 +1633,16 @@ StructuredCloneHolder::CustomFreeTransferHandler(uint32_t aTag, delete data; return; } + + if ((aTag == SCTAG_DOM_TRANSFERRED_WRITABLESTREAM || + aTag == SCTAG_DOM_TRANSFERRED_READABLESTREAM) && + mStructuredCloneScope == StructuredCloneScope::SameProcessSameThread) { + MOZ_ASSERT(aContent); + SameThreadStreamTransferData* data = + static_cast(aContent); + delete data; + return; + } } bool diff --git a/dom/base/StructuredCloneTags.h b/dom/base/StructuredCloneTags.h index 09b91f7afb..86ea88ba26 100644 --- a/dom/base/StructuredCloneTags.h +++ b/dom/base/StructuredCloneTags.h @@ -68,6 +68,10 @@ enum StructuredCloneTags { // This tag is used by both main thread and workers. SCTAG_DOM_URLSEARCHPARAMS, + // Same-thread transferable stream records. These are not supported by IDB. + SCTAG_DOM_TRANSFERRED_READABLESTREAM, + SCTAG_DOM_TRANSFERRED_WRITABLESTREAM, + // When adding a new tag for IDB, please don't add it to the end of the list! // Tags that are supported by IDB must not ever change. See the static assert // in IDBObjectStore.cpp, method CommonStructuredCloneReadCallback. diff --git a/js/src/builtin/Stream.cpp b/js/src/builtin/Stream.cpp index 144d222dd6..7234ed596a 100644 --- a/js/src/builtin/Stream.cpp +++ b/js/src/builtin/Stream.cpp @@ -6697,6 +6697,140 @@ js::InitStreamExtras(JSContext* cx, HandleObject global) Object.defineProperty(TransformStream.prototype, Symbol.toStringTag, { value: "TransformStream", configurable: true }); + function isReadableStreamObject(value) { + if (!isObject(value) || !global.ReadableStream) + return false; + try { + return value instanceof global.ReadableStream; + } catch (e) { + return false; + } + } + + function isUncloneableStreamChunk(value) { + return isObject(value) && + (isWritableStream(value) || + isReadableStreamObject(value) || + transformState.has(value)); + } + + function makeDataCloneError() { + return makeDOMException("The object could not be cloned.", "DataCloneError"); + } + + function abortTransferredWritable(record, reason) { + try { + return record.writer.abort(reason); + } catch (e) { + return Promise.reject(e); + } + } + + function rejectTransferredWritableWrite(record, error) { + var abortPromise = abortTransferredWritable(record, error); + silenceRejection(abortPromise); + return Promise.reject(error); + } + + function forwardTransferredWritableWrite(record, chunk) { + if (isUncloneableStreamChunk(chunk)) + return rejectTransferredWritableWrite(record, makeDataCloneError()); + + var prior = record.tail || Promise.resolve(undefined); + var firstWrite = !record.started; + record.started = true; + + var forwarding = prior.then(function() { + var writePromise; + try { + writePromise = Promise.resolve(record.writer.write(chunk)); + } catch (e) { + writePromise = Promise.reject(e); + } + record.inFlight = writePromise; + silenceRejection(writePromise); + return writePromise; + }); + + record.tail = forwarding; + silenceRejection(forwarding); + + if (firstWrite) + return undefined; + + return prior.then(function() { return undefined; }); + } + + function closeTransferredWritable(record) { + var tail = record.tail || Promise.resolve(undefined); + var closePromise = tail.then(function() { + return record.writer.close(); + }); + record.tail = closePromise; + silenceRejection(closePromise); + return closePromise; + } + + function __uxpTransferWritableStream(stream) { + if (!isWritableStream(stream)) + return undefined; + if (stream.locked) + return null; + try { + return { + writer: stream.getWriter(), + started: false, + inFlight: undefined, + tail: undefined + }; + } catch (e) { + return null; + } + } + + function __uxpReceiveWritableStreamTransfer(record) { + return new WritableStream({ + write: function(chunk) { + return forwardTransferredWritableWrite(record, chunk); + }, + close: function() { + return closeTransferredWritable(record); + }, + abort: function(reason) { + return abortTransferredWritable(record, reason); + } + }); + } + + function __uxpTransferReadableStream(stream) { + if (!isReadableStreamObject(stream)) + return undefined; + if (stream.locked) + return null; + try { + return { reader: stream.getReader() }; + } catch (e) { + return null; + } + } + + function __uxpReceiveReadableStreamTransfer(record) { + return new ReadableStream({ + pull: function(controller) { + return record.reader.read().then(function(result) { + if (result.done) { + controller.close(); + return; + } + controller.enqueue(result.value); + }); + }, + cancel: function(reason) { + return record.reader.cancel(reason); + } + }); + } + function pipeTo(readable, destination, options) { options = options === undefined ? {} : requiredObject(options, "options"); var signal = options.signal; @@ -6842,6 +6976,18 @@ js::InitStreamExtras(JSContext* cx, HandleObject global) { value: pipeTo, configurable: true }); Object.defineProperty(global, "__uxpStreamsPipeThrough", { value: pipeThrough, configurable: true }); + Object.defineProperty(global, "__uxpTransferWritableStream", + { value: __uxpTransferWritableStream, + configurable: true }); + Object.defineProperty(global, "__uxpReceiveWritableStreamTransfer", + { value: __uxpReceiveWritableStreamTransfer, + configurable: true }); + Object.defineProperty(global, "__uxpTransferReadableStream", + { value: __uxpTransferReadableStream, + configurable: true }); + Object.defineProperty(global, "__uxpReceiveReadableStreamTransfer", + { value: __uxpReceiveReadableStreamTransfer, + configurable: true }); if (global.ReadableStream && !global.ReadableStream.from) { Object.defineProperty(global.ReadableStream, "from", {