Add same-thread stream transfer support

This commit is contained in:
Basilisk-Dev 2026-05-08 21:43:14 -04:00 committed by wuggy
commit 58e0e6a192
3 changed files with 313 additions and 0 deletions

View file

@ -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<JSObject*> global(aCx, JS::CurrentGlobalOrNull(aCx));
if (!global) {
return false;
}
// Resolve the stream extras lazily before looking up the internal helper.
JS::Rooted<JS::Value> ignored(aCx);
if (!JS_GetProperty(aCx, global, "WritableStream", &ignored)) {
return false;
}
JS::Rooted<JS::Value> helper(aCx);
if (!JS_GetProperty(aCx, global, aHelperName, &helper)) {
return false;
}
if (!helper.isObject() || !JS::IsCallable(&helper.toObject())) {
aResult.setUndefined();
return true;
}
JS::Rooted<JS::Value> 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<JSObject*> aObj,
const char* aHelperName,
uint32_t aTransferTag,
uint32_t* aTag,
JS::TransferableOwnership* aOwnership,
void** aContent,
uint64_t* aExtraData,
bool* aHandled)
{
*aHandled = false;
JS::Rooted<JS::Value> argument(aCx, JS::ObjectValue(*aObj));
JS::Rooted<JS::Value> transferRecord(aCx);
if (!CallStreamTransferHelper(aCx, aHelperName, argument, &transferRecord)) {
return false;
}
if (transferRecord.isUndefined()) {
return true;
}
if (!transferRecord.isObject()) {
return false;
}
JS::Rooted<JSObject*> 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<SameThreadStreamTransferData*>(aContent);
JS::Rooted<JS::Value> record(aCx, JS::ObjectValue(*data->mRecord));
JS::Rooted<JS::Value> 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<SameThreadStreamTransferData*>(aContent);
delete data;
return;
}
}
bool

View file

@ -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.

View file

@ -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", {