Support readable stream transfer to workers

This commit is contained in:
Basilisk-Dev 2026-05-08 22:58:09 -04:00 • committed by wuggy
commit 7d8bf841ed
2 changed files with 354 additions and 0 deletions

View file

@ -152,6 +152,11 @@ struct SameThreadStreamTransferData
JS::PersistentRootedObject mRecord; JS::PersistentRootedObject mRecord;
}; };
struct MessagePortStreamTransferData
{
MessagePortIdentifier mIdentifier;
};
bool bool
CallStreamTransferHelper(JSContext* aCx, CallStreamTransferHelper(JSContext* aCx,
const char* aHelperName, const char* aHelperName,
@ -227,6 +232,50 @@ TryWriteSameThreadStreamTransfer(JSContext* aCx,
return true; return true;
} }
bool
TryWriteReadableStreamPortTransfer(JSContext* aCx,
JS::Handle<JSObject*> aObj,
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> transferPortValue(aCx);
if (!CallStreamTransferHelper(aCx, "__uxpTransferReadableStreamPort",
argument, &transferPortValue)) {
return false;
}
if (transferPortValue.isUndefined()) {
return true;
}
if (!transferPortValue.isObject()) {
return false;
}
JS::Rooted<JSObject*> transferPortObj(aCx, &transferPortValue.toObject());
MessagePort* port = nullptr;
nsresult rv = UNWRAP_OBJECT(MessagePort, &transferPortObj, port);
if (NS_FAILED(rv) || !port) {
return false;
}
MessagePortStreamTransferData* data = new MessagePortStreamTransferData();
port->CloneAndDisentangle(data->mIdentifier);
*aTag = SCTAG_DOM_TRANSFERRED_READABLESTREAM;
*aOwnership = JS::SCTAG_TMO_CUSTOM;
*aContent = data;
*aExtraData = 0;
*aHandled = true;
return true;
}
bool bool
ReadSameThreadStreamTransfer(JSContext* aCx, ReadSameThreadStreamTransfer(JSContext* aCx,
void* aContent, void* aContent,
@ -252,6 +301,48 @@ ReadSameThreadStreamTransfer(JSContext* aCx,
return true; return true;
} }
bool
ReadReadableStreamPortTransfer(JSContext* aCx,
nsISupports* aParent,
void* aContent,
JS::MutableHandleObject aReturnObject)
{
MOZ_ASSERT(aContent);
MessagePortStreamTransferData* data =
static_cast<MessagePortStreamTransferData*>(aContent);
nsCOMPtr<nsIGlobalObject> global = do_QueryInterface(aParent);
ErrorResult rv;
RefPtr<MessagePort> port =
MessagePort::Create(global, data->mIdentifier, rv);
if (NS_WARN_IF(rv.Failed())) {
rv.SuppressException();
return false;
}
JS::Rooted<JS::Value> portValue(aCx);
if (!GetOrCreateDOMReflector(aCx, port, &portValue)) {
JS_ClearPendingException(aCx);
return false;
}
JS::Rooted<JS::Value> result(aCx);
if (!CallStreamTransferHelper(aCx,
"__uxpReceiveReadableStreamTransferFromPort",
portValue, &result)) {
return false;
}
if (!result.isObject()) {
return false;
}
aReturnObject.set(&result.toObject());
delete data;
return true;
}
} // anonymous namespace } // anonymous namespace
const JSStructuredCloneCallbacks StructuredCloneHolder::sCallbacks = { const JSStructuredCloneCallbacks StructuredCloneHolder::sCallbacks = {
@ -1503,6 +1594,12 @@ StructuredCloneHolder::CustomReadTransferHandler(JSContext* aCx,
} }
} }
if (mStructuredCloneScope == StructuredCloneScope::SameProcessDifferentThread &&
aTag == SCTAG_DOM_TRANSFERRED_READABLESTREAM) {
return ReadReadableStreamPortTransfer(aCx, mParent, aContent,
aReturnObject);
}
return false; return false;
} }
@ -1596,6 +1693,17 @@ StructuredCloneHolder::CustomWriteTransferHandler(JSContext* aCx,
} }
} }
if (mStructuredCloneScope == StructuredCloneScope::SameProcessDifferentThread) {
bool handled = false;
if (!TryWriteReadableStreamPortTransfer(aCx, obj, aTag, aOwnership,
aContent, aExtraData, &handled)) {
return false;
}
if (handled) {
return true;
}
}
return false; return false;
} }
@ -1643,6 +1751,16 @@ StructuredCloneHolder::CustomFreeTransferHandler(uint32_t aTag,
delete data; delete data;
return; return;
} }
if (aTag == SCTAG_DOM_TRANSFERRED_READABLESTREAM &&
mStructuredCloneScope == StructuredCloneScope::SameProcessDifferentThread) {
MOZ_ASSERT(aContent);
MessagePortStreamTransferData* data =
static_cast<MessagePortStreamTransferData*>(aContent);
MessagePort::ForceClose(data->mIdentifier);
delete data;
return;
}
} }
bool bool

View file

@ -6916,6 +6916,236 @@ js::InitStreamExtras(JSContext* cx, HandleObject global)
return undefined; return undefined;
} }
function closeReadableTransferPort(port) {
try { port.close(); } catch (e) {}
}
function startReadableTransferPort(port) {
try {
if (port.start)
port.start();
} catch (e) {}
}
function tryPostReadableTransferPort(port, message) {
try {
port.postMessage(message);
return true;
} catch (e) {
return false;
}
}
function postReadableTransferPortError(record, error) {
var posted = tryPostReadableTransferPort(record.port,
{ type: "error", value: error });
if (!posted) {
tryPostReadableTransferPort(record.port, {
type: "error",
name: error && error.name,
message: error && error.message
});
}
record.closed = true;
closeReadableTransferPort(record.port);
}
function requestReadableTransferPortChunk(state) {
if (state.closed || state.errored || state.canceled)
return;
if (!tryPostReadableTransferPort(state.port, { type: "pull" })) {
state.errored = true;
if (state.controller)
try { state.controller.error(makeDataCloneError()); } catch (e) {}
closeReadableTransferPort(state.port);
}
}
function errorReadableTransferPortState(state, error) {
if (state.errored)
return;
state.errored = true;
try {
if (state.controller)
state.controller.error(error);
} catch (e) {}
closeReadableTransferPort(state.port);
}
function pumpReadableTransferPort(record) {
if (record.reading || record.closed || record.canceled ||
record.credits <= 0)
return;
--record.credits;
record.reading = true;
var readPromise;
try {
readPromise = Promise.resolve(record.reader.read());
} catch (e) {
readPromise = Promise.reject(e);
}
silenceRejection(readPromise);
readPromise.then(function(result) {
record.reading = false;
if (record.canceled || record.closed)
return;
if (result.done) {
record.closed = true;
tryPostReadableTransferPort(record.port, { type: "close" });
closeReadableTransferPort(record.port);
return;
}
if (isUncloneableStreamChunk(result.value)) {
var cloneError = makeDataCloneError();
silenceRejection(record.reader.cancel(cloneError));
postReadableTransferPortError(record, cloneError);
return;
}
if (!tryPostReadableTransferPort(record.port,
{ type: "chunk", value: result.value })) {
record.closed = true;
closeReadableTransferPort(record.port);
return;
}
pumpReadableTransferPort(record);
}, function(error) {
record.reading = false;
if (record.canceled || record.closed)
return;
postReadableTransferPortError(record, error);
});
}
function __uxpTransferReadableStreamPort(stream) {
if (!isReadableStreamObject(stream))
return undefined;
if (stream.locked || typeof global.MessageChannel !== "function")
return null;
var channel;
try {
channel = new global.MessageChannel();
var record = {
reader: stream.getReader(),
port: channel.port1,
reading: false,
credits: 1,
closed: false,
canceled: false
};
record.port.onmessage = function(event) {
var message = event.data;
if (!message || record.closed || record.canceled)
return;
if (message.type === "pull") {
++record.credits;
pumpReadableTransferPort(record);
return;
}
if (message.type === "cancel") {
record.canceled = true;
Promise.resolve().then(function() {
try {
silenceRejection(record.reader.cancel(message.reason));
} catch (e) {}
});
closeReadableTransferPort(record.port);
}
};
startReadableTransferPort(record.port);
pumpReadableTransferPort(record);
return channel.port2;
} catch (e) {
if (channel) {
closeReadableTransferPort(channel.port1);
closeReadableTransferPort(channel.port2);
}
return null;
}
}
function errorFromReadableTransferPortMessage(message) {
if ("value" in message)
return message.value;
return makeDOMException(message.message || "ReadableStream transfer failed",
message.name || "DataCloneError");
}
function __uxpReceiveReadableStreamTransferFromPort(port) {
var state = {
port: port,
controller: undefined,
closed: false,
errored: false,
canceled: false
};
port.onmessage = function(event) {
var message = event.data;
if (!message || state.closed || state.errored || state.canceled)
return;
if (message.type === "chunk") {
if (isUncloneableStreamChunk(message.value)) {
errorReadableTransferPortState(state, makeDataCloneError());
return;
}
try {
state.controller.enqueue(message.value);
} catch (e) {
errorReadableTransferPortState(state, e);
return;
}
if (state.controller.desiredSize > 0)
requestReadableTransferPortChunk(state);
return;
}
if (message.type === "close") {
state.closed = true;
try { state.controller.close(); } catch (e) {}
closeReadableTransferPort(port);
return;
}
if (message.type === "error") {
errorReadableTransferPortState(
state, errorFromReadableTransferPortMessage(message));
}
};
return new ReadableStream({
start: function(controller) {
state.controller = controller;
startReadableTransferPort(port);
requestReadableTransferPortChunk(state);
},
pull: function() {
requestReadableTransferPortChunk(state);
},
cancel: function(reason) {
state.canceled = true;
tryPostReadableTransferPort(port, { type: "cancel", reason: reason });
closeReadableTransferPort(port);
return undefined;
}
});
}
function __uxpReceiveReadableStreamTransfer(record) { function __uxpReceiveReadableStreamTransfer(record) {
return new ReadableStream({ return new ReadableStream({
start: function(controller) { start: function(controller) {
@ -7087,6 +7317,12 @@ js::InitStreamExtras(JSContext* cx, HandleObject global)
Object.defineProperty(global, "__uxpReceiveReadableStreamTransfer", Object.defineProperty(global, "__uxpReceiveReadableStreamTransfer",
{ value: __uxpReceiveReadableStreamTransfer, { value: __uxpReceiveReadableStreamTransfer,
configurable: true }); configurable: true });
Object.defineProperty(global, "__uxpTransferReadableStreamPort",
{ value: __uxpTransferReadableStreamPort,
configurable: true });
Object.defineProperty(global, "__uxpReceiveReadableStreamTransferFromPort",
{ value: __uxpReceiveReadableStreamTransferFromPort,
configurable: true });
if (global.ReadableStream && !global.ReadableStream.from) { if (global.ReadableStream && !global.ReadableStream.from) {
Object.defineProperty(global.ReadableStream, "from", { Object.defineProperty(global.ReadableStream, "from", {