fix(messaging): preserve MessagePort document lifecycle

Create detached branded ports from retained constructors, retire stale owners
before transfer, and reject manual dispatch on retired receiver realms after
validating Event state. Capture postMessage destinations after argument
conversion and before serialization so reentry cannot revive old connections
or cancel an already selected live destination.

Cover constructor realms, frame reinsertion, endpoint cleanup, borrowed
dispatch, transfer and serialization reentry.
This commit is contained in:
ldm0
2026-09-27 06:43:25 +08:00
parent 4839869050
commit 56b77a2a3b
12 changed files with 428 additions and 39 deletions
+16
View File
@@ -134,6 +134,22 @@ where
})
}
/// Capture the destination before an embedding runs author code.
pub fn message_port_peer_id(&self, port_id: MessagePortId) -> Option<MessagePortId> {
self.ports.lock().get(&port_id)?.peer_id
}
/// Enqueue to a previously captured destination, even if its source has
/// since been disentangled. Port ids survive transfer to another owner.
pub fn enqueue_message_to_endpoint(&self, port_id: MessagePortId, payload: P) -> bool {
let mut ports = self.ports.lock();
let Some(state) = ports.get_mut(&port_id) else {
return false;
};
state.pending_messages.push_back(payload);
true
}
/// Enqueue a message to the peer endpoint and return the target peer id.
pub fn enqueue_message_to_message_port(
&self,
@@ -417,7 +417,6 @@ pub(crate) use self::shared::{
runtime_message_allowed_for_current_target, structured_clone_value,
structured_clone_value_with_options, structured_deserialize_value_for_message_event,
structured_serialize_value_for_post_message,
structured_serialize_value_for_post_message_with_source_port,
structured_serialize_value_for_window_post_message,
structured_serialize_value_for_window_post_message_options,
wasm_module_message_allowed_for_target, wasm_module_message_allowed_for_target_origin,
@@ -31,7 +31,8 @@ pub(in crate::context_bootstrap) use state::{
};
pub(in crate::context_bootstrap::message_ports) use state::{
forget_message_port_wrapper, message_port_is_closed, message_port_is_started,
new_message_port_object, set_message_port_peer, set_message_port_started,
new_detached_message_port_object, new_message_port_object,
retire_message_port_if_owner_is_stale, set_message_port_peer, set_message_port_started,
};
const MESSAGE_PORT_EVENT_LISTENERS_SLOT: &str = "__moliMessagePortEventListeners";
@@ -25,6 +25,24 @@ pub(in crate::context_bootstrap) fn message_channel_constructor_callback<'s>(
}
let Some(realm) = MessagePortRealmBinding::current(scope) else {
let global = scope.get_current_context().global(scope);
if !crate::context_bootstrap::event_target_dispatch::target_execution_context_is_live(
scope, global,
) {
// A retained constructor still creates branded objects after its
// document is destroyed, but those ports have no live endpoint.
let Some(port1) = new_detached_message_port_object(scope) else {
return;
};
let Some(port2) = new_detached_message_port_object(scope) else {
return;
};
MessageChannelObjectDeclaration::new(port1, port2)
.initialize(scope, args.this())
.expect("MessageChannel declaration should initialize detached ports");
rv.set(args.this().into());
return;
}
throw_type_error(
scope,
"Failed to construct 'MessageChannel': Execution context is unavailable.",
@@ -15,19 +15,29 @@ pub(in crate::context_bootstrap) fn message_port_post_message_callback<'s>(
return;
};
let transfer_arg = (args.length() > 1).then(|| args.get(1));
let Some(data) =
crate::context_bootstrap::structured_serialize_value_for_post_message_with_source_port(
scope,
args.get(0),
transfer_arg,
"MessagePort",
Some(port_id),
)
let Some(transfers) =
parse_post_message_transfer_list(scope, transfer_arg, "MessagePort", Some(port_id))
else {
return;
};
if let Some(peer_id) = current_message_port_registry(scope)
.and_then(|registry| registry.enqueue_message_to_message_port(port_id, data))
// Web IDL conversion precedes destination selection. Serialization follows
// it, so closing or retiring the source in a data getter does not redirect
// or cancel the message already being sent to a live destination.
retire_message_port_if_owner_is_stale(scope, port_id);
let registry = current_message_port_registry(scope);
let peer_id = registry
.as_ref()
.and_then(|registry| registry.message_port_peer_id(port_id));
let data = transfers.serialize(scope, args.get(0));
// Serialization can run getters that destroy the source document.
retire_message_port_if_owner_is_stale(scope, port_id);
let Some(data) = data else {
return;
};
if let Some(peer_id) = peer_id
&& !data.transferred_message_ports().contains(&peer_id)
&& let Some(registry) = registry
&& registry.enqueue_message_to_endpoint(peer_id, data)
{
// Publish directly to the peer's stable owner route. The Page scheduler
// cannot execute the task until the current PageVm turn is restored;
@@ -122,7 +122,7 @@ struct MessagePortObjectDeclaration<'scope> {
prototype: v8::Local<'scope, v8::Object>,
#[webapi(slot = MESSAGE_PORT_ID_SLOT)]
port_id: v8::Local<'scope, v8::BigInt>,
port_id: v8::Local<'scope, v8::Value>,
#[webapi(slot = MESSAGE_PORT_PEER_SLOT, init = "undefined")]
peer: (),
@@ -237,6 +237,25 @@ pub(in crate::context_bootstrap::message_ports) fn new_message_port_object<'s>(
scope: &mut v8::PinScope<'s, '_>,
port_id: MessagePortId,
realm: &MessagePortRealmBinding,
) -> Option<v8::Local<'s, v8::Object>> {
let port = new_message_port_wrapper(scope, Some(port_id))?;
if !realm.register_wrapper(scope, port_id, port) {
return None;
}
Some(port)
}
pub(in crate::context_bootstrap::message_ports) fn new_detached_message_port_object<'s>(
scope: &mut v8::PinScope<'s, '_>,
) -> Option<v8::Local<'s, v8::Object>> {
let port = new_message_port_wrapper(scope, None)?;
set_message_port_bool_slot(scope, port, MESSAGE_PORT_CLOSED_SLOT, true);
Some(port)
}
fn new_message_port_wrapper<'s>(
scope: &mut v8::PinScope<'s, '_>,
port_id: Option<MessagePortId>,
) -> Option<v8::Local<'s, v8::Object>> {
let port = message_port_object_declaration(scope, port_id)?
.bind(scope)
@@ -245,22 +264,22 @@ pub(in crate::context_bootstrap::message_ports) fn new_message_port_object<'s>(
// transferable communication endpoint and its queue attachment.
mark_simple_event_target_slot(scope, port, MESSAGE_PORT_EVENT_LISTENERS_SLOT);
install_simple_event_target_ordered_handlers(scope, port);
if !realm.register_wrapper(scope, port_id, port) {
return None;
}
Some(port)
}
fn message_port_object_declaration<'s>(
scope: &mut v8::PinScope<'s, '_>,
port_id: MessagePortId,
port_id: Option<MessagePortId>,
) -> Option<MessagePortObjectDeclaration<'s>> {
let prototype = super::super::exposed_interfaces::ensure_intrinsic_interface_prototype(
scope,
"MessagePort",
)
.ok()?;
let port_id = v8::BigInt::new_from_u64(scope, port_id);
let port_id = match port_id {
Some(port_id) => v8::BigInt::new_from_u64(scope, port_id).into(),
None => v8::undefined(scope).into(),
};
Some(MessagePortObjectDeclaration::new(prototype, port_id))
}
@@ -352,11 +371,21 @@ pub(crate) fn detach_message_port_owner_for_transfer(
scope: &mut v8::PinScope<'_, '_>,
port_id: MessagePortId,
) {
retire_message_port_if_owner_is_stale(scope, port_id);
if let Some(registry) = current_message_port_registry(scope) {
registry.detach_message_port_owner_for_transfer(port_id);
}
}
pub(in crate::context_bootstrap::message_ports) fn retire_message_port_if_owner_is_stale(
scope: &mut v8::PinScope<'_, '_>,
port_id: MessagePortId,
) {
if let Some(host_ptr) = context_host_ptr_from_global_bridge(scope) {
unsafe { &mut *host_ptr }.retire_message_port_if_owner_is_stale(port_id);
}
}
pub(in crate::context_bootstrap::message_ports) fn message_port_wrapper_for_id<'s>(
scope: &mut v8::PinScope<'s, '_>,
port_id: MessagePortId,
@@ -14,7 +14,7 @@ use crate::{
webidl,
};
struct PostMessageTransferList<'s> {
pub(in crate::context_bootstrap) struct PostMessageTransferList<'s> {
array_buffers: Vec<v8::Local<'s, v8::ArrayBuffer>>,
message_ports: Vec<v8::Local<'s, v8::Object>>,
readable_streams: Vec<v8::Local<'s, v8::Object>>,
@@ -22,6 +22,26 @@ struct PostMessageTransferList<'s> {
transform_streams: Vec<v8::Local<'s, v8::Object>>,
}
impl<'s> PostMessageTransferList<'s> {
pub(in crate::context_bootstrap) fn serialize(
self,
scope: &mut v8::PinScope<'s, '_>,
value: v8::Local<'s, v8::Value>,
) -> Option<V8StructuredClonePayload> {
let mut payload = serialize_for_wire_for_runtime_message(
scope,
value,
&self.array_buffers,
&self.message_ports,
&self.readable_streams,
&self.writable_streams,
&self.transform_streams,
)?;
attach_runtime_message_source(&mut payload, RuntimeMessageSourceSecurity::current(scope));
Some(payload)
}
}
#[derive(Clone, Copy)]
enum TransferListOperation<'a> {
PostMessage(&'a str),
@@ -122,17 +142,7 @@ pub(crate) fn structured_serialize_value_for_post_message_with_source_port<'s>(
) -> Option<V8StructuredClonePayload> {
let transfers =
parse_post_message_transfer_list(scope, transfer_arg, interface_name, source_port_id)?;
let mut payload = serialize_for_wire_for_runtime_message(
scope,
value,
&transfers.array_buffers,
&transfers.message_ports,
&transfers.readable_streams,
&transfers.writable_streams,
&transfers.transform_streams,
)?;
attach_runtime_message_source(&mut payload, RuntimeMessageSourceSecurity::current(scope));
Some(payload)
transfers.serialize(scope, value)
}
pub(crate) fn structured_serialize_value_for_window_post_message<'s>(
@@ -330,7 +340,7 @@ pub(crate) fn structured_clone_value_for_storage<'s>(
structured_deserialize_value(scope, &bytes)
}
fn parse_post_message_transfer_list<'s>(
pub(in crate::context_bootstrap) fn parse_post_message_transfer_list<'s>(
scope: &mut v8::PinScope<'s, '_>,
transfer_arg: Option<v8::Local<'s, v8::Value>>,
interface_name: &str,
@@ -82,6 +82,18 @@ impl RendererMessagePortRegistry {
self.inner.enqueue_message_to_message_port(port_id, payload)
}
pub(crate) fn message_port_peer_id(&self, port_id: MessagePortId) -> Option<MessagePortId> {
self.inner.message_port_peer_id(port_id)
}
pub(crate) fn enqueue_message_to_endpoint(
&self,
port_id: MessagePortId,
payload: V8StructuredClonePayload,
) -> bool {
self.inner.enqueue_message_to_endpoint(port_id, payload)
}
pub(crate) fn take_pending_message_port_message(
&self,
port_id: MessagePortId,
@@ -69,6 +69,19 @@ impl JsContextHost {
v8::Local<'s, v8::Context>,
v8::Local<'s, v8::Object>,
)> {
if self.retire_message_port_if_owner_is_stale(port_id) {
return None;
}
let entry = self.message_port_wrappers.get(&port_id)?;
Some((
entry.identity.dispatch_scope(),
entry.identity.realm_token(),
v8::Local::new(scope, &entry.context),
v8::Local::new(scope, &entry.wrapper),
))
}
pub(crate) fn retire_message_port_if_owner_is_stale(&mut self, port_id: MessagePortId) -> bool {
let stale_owner = self
.message_port_wrappers
.get(&port_id)
@@ -81,15 +94,9 @@ impl JsContextHost {
?identity,
"closed MessagePort for retired execution context"
);
return None;
return true;
}
let entry = self.message_port_wrappers.get(&port_id)?;
Some((
entry.identity.dispatch_scope(),
entry.identity.realm_token(),
v8::Local::new(scope, &entry.context),
v8::Local::new(scope, &entry.wrapper),
))
false
}
pub(crate) fn message_port_execution_context_identity(
@@ -0,0 +1,154 @@
function retiredChannelProbe() {
const failures = [], rows = [];
const check = (value, label) => { if (!value) failures.push(label); };
for (const reinsert of [false, true]) {
const frame = document.createElement('iframe');
document.body.appendChild(frame);
const realm = frame.contentWindow;
const C = realm.MessageChannel, P = realm.MessagePort, T = realm.TypeError;
const getters = ['port1', 'port2'].map(name =>
Object.getOwnPropertyDescriptor(MessageChannel.prototype, name).get);
frame.remove();
if (reinsert) document.body.appendChild(frame);
class Subchannel extends C {}
for (const construct of [() => new C(), () => new Subchannel(),
() => Reflect.construct(C, [], function Custom() {})]) {
try {
const channel = construct(), ports = getters.map(get => get.call(channel));
check(ports[0] !== ports[1], 'distinct retired ports');
Object.setPrototypeOf(channel, null);
Object.freeze(channel);
for (const [index, port] of ports.entries()) {
check(Object.getPrototypeOf(port) === P.prototype, 'retired port realm');
check(getters[index].call(channel) === port, 'retired channel brand');
let calls = 0;
port.addEventListener('probe', () => calls++);
const event = new Event('probe');
check(port.dispatchEvent(event) === false && calls === 0, 'retired dispatch');
check(event.target === null && event.eventPhase === 0, 'retired event untouched');
try { structuredClone(port, {transfer:[port]}); failures.push('transferred born-detached port'); }
catch (error) { check(error.name === 'DataCloneError', 'born-detached transfer error'); }
port.start(); port.postMessage('ignored'); port.close(); port.close();
}
rows.push([reinsert, true]);
} catch (error) { failures.push('retained constructor: ' + error.name); }
}
try { C(); failures.push('constructor called without new'); }
catch (error) { check(error instanceof T, 'retired constructor error realm'); }
if (reinsert) {
const fresh = new frame.contentWindow.MessageChannel();
let calls = 0;
fresh.port1.addEventListener('probe', () => calls++);
check(fresh.port1.dispatchEvent(new Event('probe')) === true && calls === 1,
'reinserted frame has a fresh active realm');
check(Object.getPrototypeOf(fresh.port1) !== P.prototype, 'fresh port realm');
fresh.port1.close(); fresh.port2.close();
}
frame.remove();
}
return {failures, rows};
}
function retiredTargetProbe() {
const failures = [], rows = [];
const check = (value, label) => { if (!value) failures.push(label); };
for (const kind of ['EventTarget', 'MessagePort', 'FileReader', 'AbortSignal']) {
for (const state of ['live', 'removed', 'reinserted']) {
const frame = document.createElement('iframe');
document.body.appendChild(frame);
const w = frame.contentWindow;
const T = w.TypeError, childDispatch = w.EventTarget.prototype.dispatchEvent;
let channel;
const target = kind === 'EventTarget' ? new w.EventTarget() :
kind === 'MessagePort' ? (channel = new w.MessageChannel()).port1 :
kind === 'FileReader' ? new w.FileReader() : new w.AbortController().signal;
const methods = [target.dispatchEvent];
if (kind !== 'AbortSignal') methods.push(EventTarget.prototype.dispatchEvent);
let calls = 0;
target.addEventListener('probe', () => calls++);
if (state !== 'live') frame.remove();
if (state === 'reinserted') document.body.appendChild(frame);
for (const [index, dispatch] of methods.entries()) {
const label = kind + ':' + state + ':' + index;
const before = calls, event = new Event('probe');
try {
const returned = dispatch.call(target, event);
check(returned === (state === 'live'), label + ' return');
check(calls - before === (state === 'live' ? 1 : 0), label + ' listeners');
check(event.target === (state === 'live' ? target : null), label + ' event target');
check(event.currentTarget === null && event.eventPhase === 0, label + ' event state');
for (const value of [null, {}, new Proxy(event, {})]) {
try { dispatch.call(target, value); failures.push(label + ' accepted invalid Event'); }
catch (error) { check(error instanceof (index === 0 ? T : TypeError), label + ' error realm'); }
}
try { dispatch.call(target, document.createEvent('Event')); failures.push(label + ' uninitialized'); }
catch (error) { check(error.name === 'InvalidStateError', label + ' initialization before liveness'); }
if (state !== 'live') {
const active = new EventTarget(), reused = new Event('probe');
active.dispatchEvent(reused);
check(dispatch.call(target, reused) === false && reused.target === active,
label + ' preserve previous target');
active.addEventListener('probe', event => {
try { dispatch.call(target, event); failures.push(label + ' accepted dispatching Event'); }
catch (error) { check(error.name === 'InvalidStateError', label + ' dispatch flag before liveness'); }
});
active.dispatchEvent(new Event('probe'));
if (kind !== 'AbortSignal') {
let parentCalls = 0;
active.addEventListener('parent', () => parentCalls++);
check(childDispatch.call(active, new Event('parent')) === true && parentCalls === 1,
label + ' receiver realm governs borrowed dispatch');
}
}
rows.push(label);
} catch (error) { failures.push(label + ': ' + error.name); }
}
channel?.port1.close(); channel?.port2.close(); frame.remove();
}
}
return {failures, rows};
}
async function retiredPortTransferProbe() {
const rows = [];
for (const mode of ['transfer-before-removal', 'transfer-after-removal',
'remove-during-transfer', 'post-after-removal', 'remove-during-options',
'remove-during-post', 'close-during-post', 'remove-and-transfer-during-post']) {
const frame = document.createElement('iframe'); document.body.appendChild(frame);
const channel = new frame.contentWindow.MessageChannel();
const ports = [channel.port1, channel.port2];
let clones = [], messages = [], extra, error = null;
try {
if (mode === 'transfer-before-removal') {
clones = structuredClone(ports, {transfer:ports}); frame.remove();
} else if (mode === 'transfer-after-removal') {
frame.remove(); clones = structuredClone(ports, {transfer:ports});
} else if (mode === 'remove-during-transfer') {
clones = structuredClone({get ports() {frame.remove(); return ports;}}, {transfer:ports}).ports;
} else {
clones = [ports[0], structuredClone(ports[1], {transfer:[ports[1]]})];
if (mode === 'post-after-removal') frame.remove();
}
await new Promise(resolve => {
setTimeout(resolve, 200);
clones[1].onmessage = event => {messages.push(event.data.value ?? event.data);};
const reentrant = ['remove-during-post', 'close-during-post', 'remove-and-transfer-during-post'].includes(mode);
const data = reentrant ? {get value() {
if (mode === 'close-during-post') clones[0].close();
else frame.remove();
if (mode === 'remove-and-transfer-during-post') {
extra = structuredClone(clones[0], {transfer:[clones[0]]});
}
return 'first';
}} : 'first';
if (mode === 'remove-during-options') {
clones[0].postMessage(data, {get transfer() {frame.remove(); return [];}});
} else clones[0].postMessage(data);
if (reentrant) clones[0].postMessage('second');
});
} catch (caught) {error = caught.name;}
finally {for (const port of [...ports, ...clones]) port.close(); extra?.close(); frame.remove();}
rows.push({mode, messages, error});
}
return rows;
}
@@ -0,0 +1,132 @@
use super::*;
const LIFECYCLE_PROBE: &str = include_str!("message_port_lifecycle.js");
#[test]
fn message_channel_retained_constructor_creates_detached_ports_without_endpoints() {
let mut vm = new_parsed_test_vm(
"https://retired-message-channel.test/",
"<!doctype html><body>",
);
let registry = vm._context_host.borrow().message_port_registry();
assert_eq!(registry.endpoint_count(), 0);
let value = vm
.eval(&format!(
"{LIFECYCLE_PROBE}\nJSON.stringify(retiredChannelProbe())"
))
.unwrap();
let result: serde_json::Value = serde_json::from_str(&value).unwrap();
assert_eq!(result["failures"], serde_json::json!([]), "{result}");
assert_eq!(result["rows"].as_array().unwrap().len(), 6, "{result}");
assert_eq!(
registry.endpoint_count(),
0,
"retired constructors must not allocate endpoints"
);
}
#[test]
fn simple_event_targets_validate_events_before_rejecting_retired_receiver_realms() {
let mut vm = new_parsed_test_vm(
"https://retired-event-target.test/",
"<!doctype html><body>",
);
let value = vm
.eval(&format!(
"{LIFECYCLE_PROBE}\nJSON.stringify(retiredTargetProbe())"
))
.unwrap();
let result: serde_json::Value = serde_json::from_str(&value).unwrap();
assert_eq!(result["failures"], serde_json::json!([]), "{result}");
assert_eq!(result["rows"].as_array().unwrap().len(), 21, "{result}");
}
#[test]
fn message_port_transfer_disentangles_retired_owners_before_rebinding_endpoints() {
for (mode, remaining) in [("before", 2), ("after", 0), ("getter", 0)] {
let mut vm = new_parsed_test_vm(
"https://retired-port-transfer.test/",
"<!doctype html><body>",
);
let registry = vm._context_host.borrow().message_port_registry();
let value = vm
.eval(&format!(
r#"
const frame = document.createElement('iframe'); document.body.appendChild(frame);
const channel = new frame.contentWindow.MessageChannel();
const ports = [channel.port1, channel.port2];
if ('{mode}' === 'after') frame.remove();
const clones = structuredClone({{get ports() {{
if ('{mode}' === 'getter') frame.remove();
return ports;
}}}}, {{transfer:ports}}).ports;
frame.remove();
String(clones.length === 2 && clones.every(port => port instanceof MessagePort))
"#
))
.unwrap();
assert_eq!(value, "true", "{mode}");
assert_eq!(registry.endpoint_count(), remaining, "{mode}");
vm.eval("for (const port of clones) port.close()").unwrap();
assert_eq!(registry.endpoint_count(), 0, "{mode}");
}
}
#[tokio::test]
async fn message_port_retirement_preserves_the_destination_selected_before_serialization() {
let loader = ResourceRequestClient::new(&moli_fetch::FetchConfig::default()).unwrap();
let mut vm = new_storage_page_task_executor_test_vm_with_loader(
"https://retired-port-delivery.test/",
&loader,
);
vm.eval(&format!(
r#"{LIFECYCLE_PROBE}
retiredPortTransferProbe().then(value => globalThis.retiredPortResult = value,
error => globalThis.retiredPortResult = String(error));
"#
))
.unwrap();
tokio::time::timeout(std::time::Duration::from_secs(10), async {
while vm
.eval("globalThis.retiredPortResult !== undefined")
.unwrap()
!= "true"
{
vm.run_one_oldest_ready_page_task_executor_turn(&loader)
.await
.unwrap();
}
})
.await
.expect("MessagePort lifecycle probe should settle");
let value = vm.eval("JSON.stringify(retiredPortResult)").unwrap();
let rows: serde_json::Value = serde_json::from_str(&value).unwrap();
let rows = rows.as_array().expect("port lifecycle rows");
assert_eq!(rows.len(), 8);
for row in rows {
assert_eq!(row["error"], serde_json::Value::Null, "{row}");
let accepted = matches!(
row["mode"].as_str().unwrap(),
"transfer-before-removal"
| "remove-during-post"
| "close-during-post"
| "remove-and-transfer-during-post"
);
assert_eq!(
row["messages"],
if accepted {
serde_json::json!(["first"])
} else {
serde_json::json!([])
},
"{row}"
);
}
assert_eq!(
vm._context_host
.borrow()
.message_port_registry()
.endpoint_count(),
0
);
}
@@ -34,6 +34,7 @@ mod media_devices;
mod media_query_list_events;
mod message_channel;
mod message_port_events;
mod message_port_lifecycle;
mod misc;
mod navigation;
mod performance;