From 81dc65c0be7f709fab710210aa2d56974eeead9f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Wies=C5=82aw=20=C5=A0olt=C3=A9s?= Date: Wed, 16 Sep 2026 22:03:27 +0200 Subject: [PATCH] Implement Worker MessagePort and iframe transport --- .../native/webscene_v8_runtime.cpp | 6 + .../webscene_v8_runtime_browser_apis.inc | 54 +- .../webscene_v8_runtime_cache_and_frames.inc | 1 + .../native/webscene_v8_runtime_clone.inc | 758 +++++++++++++++++- .../native/webscene_v8_runtime_document.inc | 4 + .../native/webscene_v8_runtime_dom_core.inc | 12 +- .../native/webscene_v8_runtime_lifecycle.inc | 3 + .../native/webscene_v8_runtime_navigation.inc | 2 + .../native/webscene_v8_runtime_state.inc | 10 + .../webscene_v8_runtime_state_types.inc | 55 +- .../native/webscene_v8_runtime_tasks.inc | 39 +- .../native/webscene_v8_runtime_workers.inc | 143 +++- .../tests/hybrid_v8_runtime_tests.cpp | 353 +++++++- .../iframe-window-messageport-origin.html | 54 ++ .../resources/workers/message-port-helper.js | 1 + .../resources/workers/message-port.js | 19 + .../worker-messageport-structured-clone.html | 44 + ...worker-termination-queue-messageerror.html | 74 ++ .../webscene-aureon-runtime-profile.json | 34 + 19 files changed, 1589 insertions(+), 77 deletions(-) create mode 100644 tests/WebPlatformSubset/contracts/iframe-window-messageport-origin.html create mode 100644 tests/WebPlatformSubset/contracts/resources/workers/message-port-helper.js create mode 100644 tests/WebPlatformSubset/contracts/resources/workers/message-port.js create mode 100644 tests/WebPlatformSubset/contracts/worker-messageport-structured-clone.html create mode 100644 tests/WebPlatformSubset/contracts/worker-termination-queue-messageerror.html diff --git a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime.cpp b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime.cpp index c4f5d2d76..7df4a7f0f 100644 --- a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime.cpp +++ b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime.cpp @@ -305,6 +305,10 @@ struct v8_dom_runtime::implementation final { { prune_persistent_compilation_cache(); initialize_v8_process(); + { + std::lock_guard lock(message_port_wake->mutex); + message_port_wake->notify = runtime_work_available; + } if (!force_dedicated_isolate && std::getenv("WEBSCENE_V8_SHARED_ISOLATE") != nullptr) { try { shared_isolate = acquire_shared_isolate(); @@ -3752,6 +3756,7 @@ struct v8_dom_runtime::implementation final { install_clipboard_api(local_context); install_websocket_globals(local_context); install_editor_web_platform_globals(local_context); + install_message_channel(local_context); install_tree_walker_platform(local_context); install_custom_elements_platform(local_context); local_context->Global()->Set(local_context,js_string(isolate,"__webSceneRevokeObjectUrl"),v8::Function::New(local_context,revoke_object_url).ToLocalChecked()).Check(); @@ -5137,6 +5142,7 @@ bool v8_dom_runtime::has_pending_tasks() const noexcept || impl_->websocket_transport.has_pending_events() || !impl_->pending_window_messages.empty() || impl_->has_worker_messages() + || impl_->has_message_port_messages() || impl_->has_ready_fetch_task() || !impl_->pending_dialog_close_events.empty() || !impl_->pending_programmatic_scroll_events.empty() diff --git a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_browser_apis.inc b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_browser_apis.inc index e0f91adcd..779400d95 100644 --- a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_browser_apis.inc +++ b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_browser_apis.inc @@ -221,6 +221,53 @@ if (target_context.IsEmpty()) return; if (source_context.IsEmpty()) source_context = target_context; + if (self->pending_window_messages.size() + >= maximum_pending_window_messages + || self->pending_window_message_bytes + >= maximum_pending_window_message_bytes) { + throw_dom_exception( + info, + "Window message queue capacity exhausted", + "QuotaExceededError"); + return; + } + + auto target_origin = std::string("*"); + v8::Local transfer = v8::Undefined(isolate); + if (info.Length() > 1 && info[1]->IsObject() && !info[1]->IsArray()) { + v8::Local target_origin_value; + if (!info[1].As()->Get( + isolate->GetCurrentContext(), + js_string(isolate, "targetOrigin")).ToLocal(&target_origin_value)) { + return; + } + if (!target_origin_value->IsUndefined()) { + target_origin = to_utf8(isolate, target_origin_value); + } + if (!info[1].As()->Get( + isolate->GetCurrentContext(), + js_string(isolate, "transfer")).ToLocal(&transfer)) { + return; + } + } else { + if (info.Length() > 1 && !info[1]->IsUndefined()) { + target_origin = to_utf8(isolate, info[1]); + } + if (info.Length() > 2) transfer = info[2]; + } + + clone_packet packet; + if (!serialize_clone( + info, + info[0], + transfer, + packet, + maximum_pending_window_message_bytes + - self->pending_window_message_bytes)) { + return; + } + const auto packet_bytes = clone_packet_size(packet); + const auto read_origin = [&](v8::Local local_context) { v8::Context::Scope scope(local_context); v8::Local location_value; @@ -241,11 +288,10 @@ self->pending_window_messages.push_back(pending_window_message{ v8::Global(isolate, target_context), v8::Global(isolate, source_context), - v8::Global(isolate, info[0]), - info.Length() > 1 && !info[1]->IsUndefined() - ? to_utf8(isolate, info[1]) - : std::string("*"), + std::move(packet), + std::move(target_origin), read_origin(source_context)}); + self->pending_window_message_bytes += packet_bytes; } static void get_window_frames( diff --git a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_cache_and_frames.inc b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_cache_and_frames.inc index 97dcd0e3f..4488d7a12 100644 --- a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_cache_and_frames.inc +++ b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_cache_and_frames.inc @@ -1899,6 +1899,7 @@ install_clipboard_api(local_context); install_websocket_globals(local_context); install_editor_web_platform_globals(local_context); + install_message_channel(local_context); install_tree_walker_platform(local_context); install_custom_elements_platform(local_context); install_fetch_globals(local_context); diff --git a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_clone.inc b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_clone.inc index 0188ae2cf..72e1b5031 100644 --- a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_clone.inc +++ b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_clone.inc @@ -1,80 +1,754 @@ - struct clone_packet { - std::vector bytes; - std::vector> transfers; - }; - struct clone_delegate : v8::ValueSerializer::Delegate { + static constexpr size_t maximum_clone_packet_bytes = 16U * 1024U * 1024U; + + static size_t clone_packet_size(const clone_packet& packet) + { + auto size = packet.bytes.size(); + for (const auto& transfer : packet.transfers) { + if (transfer == nullptr + || transfer->ByteLength() > std::numeric_limits::max() - size) { + return std::numeric_limits::max(); + } + size += transfer->ByteLength(); + } + return size; + } + + static v8::Local message_port_private(v8::Isolate* isolate) + { + return v8::Private::ForApi(isolate, js_string(isolate, "WebScene.MessagePort")); + } + + message_port_binding* message_port_binding_from_object( + v8::Local realm, + v8::Local object) + { + v8::Local value; + if (!object->GetPrivate(realm, message_port_private(isolate)).ToLocal(&value) + || !value->IsExternal()) { + return nullptr; + } + return static_cast( + value.As()->Value(v8::kExternalPointerTypeTagDefault)); + } + + static void clear_message_port_wake( + const std::shared_ptr& endpoint) + { + if (endpoint == nullptr || endpoint->channel == nullptr) return; + std::lock_guard lock(endpoint->channel->mutex); + endpoint->channel->wake[endpoint->side].reset(); + } + + static bool message_port_endpoint_closed( + const std::shared_ptr& endpoint) + { + if (endpoint == nullptr || endpoint->channel == nullptr) return true; + std::lock_guard lock(endpoint->channel->mutex); + return endpoint->channel->closed[endpoint->side]; + } + + static void message_port_wrapper_collected( + const v8::WeakCallbackInfo& info) + { + auto* binding = info.GetParameter(); + if (binding == nullptr) return; + if (binding->endpoint != nullptr + && binding->endpoint->channel != nullptr) { + std::lock_guard lock(binding->endpoint->channel->mutex); + const auto side = binding->endpoint->side; + binding->endpoint->channel->closed[side] = true; + binding->endpoint->channel->incoming[side].clear(); + binding->endpoint->channel->incoming_bytes[side] = 0U; + binding->endpoint->channel->wake[side].reset(); + } + binding->endpoint.reset(); + binding->context.Reset(); + binding->wrapper.Reset(); + binding->transferred = false; + binding->reclaimed = true; + } + + v8::MaybeLocal create_message_port_wrapper( + v8::Local realm, + std::shared_ptr endpoint) + { + if (endpoint == nullptr || endpoint->channel == nullptr) return {}; + const auto reusable = std::find_if( + message_port_bindings.begin(), + message_port_bindings.end(), + [](const auto& binding) { + return binding != nullptr && binding->reclaimed; + }); + if (reusable == message_port_bindings.end() + && message_port_bindings.size() >= maximum_message_port_bindings) { + isolate->ThrowException(v8::Exception::RangeError( + js_string(isolate, "MessagePort binding capacity exhausted"))); + return {}; + } + + auto global = realm->Global(); + v8::Local event_target_value; + if (!global->Get(realm, js_string(isolate, "EventTarget")) + .ToLocal(&event_target_value) + || !event_target_value->IsFunction()) { + return {}; + } + v8::Local wrapper; + if (!event_target_value.As()->NewInstance(realm, 0, nullptr) + .ToLocal(&wrapper)) { + return {}; + } + v8::Local constructor_value; + v8::Local prototype_value; + if (global->Get(realm, js_string(isolate, "MessagePort")) + .ToLocal(&constructor_value) + && constructor_value->IsFunction() + && constructor_value.As()->Get( + realm, + js_string(isolate, "prototype")).ToLocal(&prototype_value) + && prototype_value->IsObject()) { + wrapper->SetPrototype(realm, prototype_value).Check(); + } + + message_port_binding* raw = nullptr; + if (reusable != message_port_bindings.end()) { + raw = reusable->get(); + } else { + auto binding = std::make_unique(); + raw = binding.get(); + message_port_bindings.push_back(std::move(binding)); + } + raw->endpoint = std::move(endpoint); + raw->transferred = false; + raw->reclaimed = false; + raw->context.Reset(isolate, realm); + raw->wrapper.Reset(isolate, wrapper); + raw->wrapper.SetWeak( + raw, + message_port_wrapper_collected, + v8::WeakCallbackType::kParameter); + + auto external = v8::External::New( + isolate, + raw, + v8::kExternalPointerTypeTagDefault); + wrapper->SetPrivate(realm, message_port_private(isolate), external).Check(); + wrapper->Set( + realm, + js_string(isolate, "postMessage"), + v8::Function::New(realm, message_port_post_message, external, 1) + .ToLocalChecked()).Check(); + wrapper->Set( + realm, + js_string(isolate, "start"), + v8::Function::New(realm, message_port_start, external) + .ToLocalChecked()).Check(); + wrapper->Set( + realm, + js_string(isolate, "close"), + v8::Function::New(realm, message_port_close, external) + .ToLocalChecked()).Check(); + wrapper->Set(realm, js_string(isolate, "onmessage"), v8::Null(isolate)).Check(); + wrapper->Set( + realm, + js_string(isolate, "onmessageerror"), + v8::Null(isolate)).Check(); + + { + std::lock_guard lock(raw->endpoint->channel->mutex); + raw->endpoint->channel->wake[raw->endpoint->side] = message_port_wake; + } + return wrapper; + } + + struct clone_delegate final : v8::ValueSerializer::Delegate { + implementation* runtime; const v8::FunctionCallbackInfo& info; - explicit clone_delegate(const v8::FunctionCallbackInfo& value) : info(value) {} - void ThrowDataCloneError(v8::Local message) override { + const std::vector& port_transfers; + v8::ValueSerializer* serializer{}; + + clone_delegate( + implementation* runtime_value, + const v8::FunctionCallbackInfo& info_value, + const std::vector& port_transfer_values) + : runtime(runtime_value) + , info(info_value) + , port_transfers(port_transfer_values) + { + } + + void ThrowDataCloneError(v8::Local message) override + { const auto text = to_utf8(info.GetIsolate(), message); throw_dom_exception(info, text.c_str(), "DataCloneError"); } + + bool HasCustomHostObject(v8::Isolate*) override { return true; } + + v8::Maybe IsHostObject( + v8::Isolate* isolate, + v8::Local object) override + { + return v8::Just( + runtime->message_port_binding_from_object( + isolate->GetCurrentContext(), + object) != nullptr); + } + + v8::Maybe WriteHostObject( + v8::Isolate* isolate, + v8::Local object) override + { + const auto* binding = runtime->message_port_binding_from_object( + isolate->GetCurrentContext(), + object); + const auto found = std::find(port_transfers.begin(), port_transfers.end(), binding); + if (binding == nullptr || found == port_transfers.end()) { + ThrowDataCloneError(js_string( + isolate, + "MessagePort must be included in the transfer list")); + return v8::Nothing(); + } + serializer->WriteUint32(static_cast( + std::distance(port_transfers.begin(), found))); + return v8::Just(true); + } }; - static bool serialize_clone_impl(const v8::FunctionCallbackInfo& info, - v8::Local value, v8::Local transfer, clone_packet& packet) { - auto* isolate = info.GetIsolate(); auto realm = isolate->GetCurrentContext(); - clone_delegate delegate(info); - v8::ValueSerializer serializer(isolate, &delegate); + + struct clone_deserialize_delegate final : v8::ValueDeserializer::Delegate { + implementation* runtime; + v8::Local realm; + const clone_packet& packet; + v8::ValueDeserializer* deserializer{}; + std::vector> wrappers; + + clone_deserialize_delegate( + implementation* runtime_value, + v8::Local realm_value, + const clone_packet& packet_value) + : runtime(runtime_value) + , realm(realm_value) + , packet(packet_value) + , wrappers(packet.port_transfers.size()) + { + } + + v8::MaybeLocal wrapper(uint32_t index) + { + if (index >= packet.port_transfers.size()) return {}; + if (!wrappers[index].IsEmpty()) { + return wrappers[index].Get(runtime->isolate); + } + v8::Local value; + if (!runtime->create_message_port_wrapper( + realm, + packet.port_transfers[index]).ToLocal(&value)) { + return {}; + } + wrappers[index].Reset(runtime->isolate, value); + return value; + } + + v8::MaybeLocal ReadHostObject(v8::Isolate*) override + { + uint32_t index = 0U; + if (!deserializer->ReadUint32(&index)) return {}; + return wrapper(index); + } + }; + + static bool serialize_clone_impl( + const v8::FunctionCallbackInfo& info, + v8::Local value, + v8::Local transfer, + clone_packet& packet, + size_t maximum_bytes = maximum_clone_packet_bytes) + { + auto* self = current(info.GetIsolate()); + auto* isolate = info.GetIsolate(); + auto realm = isolate->GetCurrentContext(); std::vector> buffers; + std::vector ports; if (!transfer->IsUndefined()) { if (!transfer->IsArray()) { - isolate->ThrowException(v8::Exception::TypeError(js_string(isolate, "transfer must be an array"))); return false; + isolate->ThrowException(v8::Exception::TypeError( + js_string(isolate, "transfer must be an array"))); + return false; } auto list = transfer.As(); - for (uint32_t i = 0; i < list->Length(); ++i) { + for (uint32_t index = 0; index < list->Length(); ++index) { v8::Local candidate; - if (!list->Get(realm, i).ToLocal(&candidate)) return false; - if (!candidate->IsArrayBuffer()) { - throw_dom_exception(info, "Unsupported transferable", "DataCloneError"); return false; + if (!list->Get(realm, index).ToLocal(&candidate)) return false; + if (candidate->IsArrayBuffer()) { + auto buffer = candidate.As(); + if (!buffer->IsDetachable() || buffer->WasDetached() + || std::find(buffers.begin(), buffers.end(), buffer) != buffers.end()) { + throw_dom_exception( + info, + "Invalid or duplicate transferable", + "DataCloneError"); + return false; + } + buffers.push_back(buffer); + continue; } - auto buffer = candidate.As(); - if (!buffer->IsDetachable() || buffer->WasDetached() - || std::find(buffers.begin(), buffers.end(), buffer) != buffers.end()) { - throw_dom_exception(info, "Invalid or duplicate transferable", "DataCloneError"); return false; + auto* port = candidate->IsObject() && self != nullptr + ? self->message_port_binding_from_object( + realm, + candidate.As()) + : nullptr; + if (port == nullptr || port->transferred + || message_port_endpoint_closed(port->endpoint) + || std::find(ports.begin(), ports.end(), port) != ports.end()) { + throw_dom_exception( + info, + "Invalid, duplicate, or unsupported transferable", + "DataCloneError"); + return false; } - buffers.push_back(buffer); - serializer.TransferArrayBuffer(i, buffer); + ports.push_back(port); } } + + clone_delegate delegate(self, info, ports); + v8::ValueSerializer serializer(isolate, &delegate); + delegate.serializer = &serializer; + for (uint32_t index = 0; index < buffers.size(); ++index) { + serializer.TransferArrayBuffer(index, buffers[index]); + } serializer.WriteHeader(); if (!serializer.WriteValue(realm, value).FromMaybe(false)) return false; auto [data, size] = serializer.Release(); std::unique_ptr memory(data, std::free); packet.bytes.assign(data, data + size); for (auto buffer : buffers) packet.transfers.push_back(buffer->GetBackingStore()); - // Serialization must succeed before ownership changes on the sender. - for (auto buffer : buffers) if (!buffer->Detach({}).FromMaybe(false)) return false; + for (auto* port : ports) packet.port_transfers.push_back(port->endpoint); + if (clone_packet_size(packet) > maximum_bytes) { + packet = {}; + throw_dom_exception( + info, + "Structured clone exceeds the bounded message size", + "QuotaExceededError"); + return false; + } + + // Serialization and quota checks must succeed before sender ownership + // changes. This preserves all transferables when cloning is rejected. + for (auto buffer : buffers) { + if (!buffer->Detach({}).FromMaybe(false)) return false; + } + for (auto* port : ports) { + clear_message_port_wake(port->endpoint); + port->transferred = true; + port->endpoint.reset(); + port->context.Reset(); + } return true; } - static bool serialize_clone(const v8::FunctionCallbackInfo& info, - v8::Local value, v8::Local transfer, clone_packet& packet) { - try {return serialize_clone_impl(info,value,transfer,packet);} - catch(const std::exception& e){ - info.GetIsolate()->ThrowException(v8::Exception::RangeError(js_string(info.GetIsolate(),e.what()))); + + static bool serialize_clone( + const v8::FunctionCallbackInfo& info, + v8::Local value, + v8::Local transfer, + clone_packet& packet, + size_t maximum_bytes = maximum_clone_packet_bytes) + { + try { + return serialize_clone_impl(info, value, transfer, packet, maximum_bytes); + } catch (const std::exception& exception) { + info.GetIsolate()->ThrowException(v8::Exception::RangeError( + js_string(info.GetIsolate(), exception.what()))); return false; } } - static v8::MaybeLocal deserialize_clone(v8::Isolate* isolate, - v8::Local realm, const clone_packet& packet) { - v8::ValueDeserializer deserializer(isolate, packet.bytes.data(), packet.bytes.size()); - for (uint32_t i = 0; i < packet.transfers.size(); ++i) - deserializer.TransferArrayBuffer(i, v8::ArrayBuffer::New(isolate, packet.transfers[i])); + + v8::MaybeLocal deserialize_clone( + v8::Isolate* isolate_value, + v8::Local realm, + const clone_packet& packet, + std::vector>* transferred_ports = nullptr) + { + clone_deserialize_delegate delegate(this, realm, packet); + v8::ValueDeserializer deserializer( + isolate_value, + packet.bytes.data(), + packet.bytes.size(), + &delegate); + delegate.deserializer = &deserializer; + for (uint32_t index = 0; index < packet.transfers.size(); ++index) { + deserializer.TransferArrayBuffer( + index, + v8::ArrayBuffer::New(isolate_value, packet.transfers[index])); + } if (!deserializer.ReadHeader(realm).FromMaybe(false)) return {}; - return deserializer.ReadValue(realm); + v8::Local value; + if (!deserializer.ReadValue(realm).ToLocal(&value)) return {}; + if (transferred_ports != nullptr) { + transferred_ports->reserve(packet.port_transfers.size()); + for (uint32_t index = 0; index < packet.port_transfers.size(); ++index) { + v8::Local port; + if (!delegate.wrapper(index).ToLocal(&port)) return {}; + transferred_ports->push_back(port); + } + } + return value; } - static void structured_clone(const v8::FunctionCallbackInfo& info) { + + static void structured_clone(const v8::FunctionCallbackInfo& info) + { if (info.Length() == 0) { - info.GetIsolate()->ThrowException(v8::Exception::TypeError(js_string(info.GetIsolate(), "structuredClone requires a value"))); return; + info.GetIsolate()->ThrowException(v8::Exception::TypeError( + js_string(info.GetIsolate(), "structuredClone requires a value"))); + return; } auto realm = info.GetIsolate()->GetCurrentContext(); - if(info.Length()>1&&!info[1]->IsNullOrUndefined()&&!info[1]->IsObject()){ - info.GetIsolate()->ThrowException(v8::Exception::TypeError(js_string(info.GetIsolate(),"StructuredSerializeOptions must be a dictionary")));return; + if (info.Length() > 1 && !info[1]->IsNullOrUndefined() + && !info[1]->IsObject()) { + info.GetIsolate()->ThrowException(v8::Exception::TypeError( + js_string(info.GetIsolate(), "StructuredSerializeOptions must be a dictionary"))); + return; } v8::Local transfer = v8::Undefined(info.GetIsolate()); if (info.Length() > 1 && info[1]->IsObject() - && !info[1].As()->Get(realm, js_string(info.GetIsolate(), "transfer")).ToLocal(&transfer)) return; + && !info[1].As()->Get( + realm, + js_string(info.GetIsolate(), "transfer")).ToLocal(&transfer)) { + return; + } clone_packet packet; if (!serialize_clone(info, info[0], transfer, packet)) return; v8::Local result; - if (deserialize_clone(info.GetIsolate(), realm, packet).ToLocal(&result)) info.GetReturnValue().Set(result); + if (current(info.GetIsolate())->deserialize_clone( + info.GetIsolate(), realm, packet).ToLocal(&result)) { + info.GetReturnValue().Set(result); + } + } + + static message_port_binding* message_port_binding_from_callback( + const v8::FunctionCallbackInfo& info) + { + if (!info.Data()->IsExternal()) return nullptr; + return static_cast( + info.Data().As()->Value( + v8::kExternalPointerTypeTagDefault)); + } + + static bool transfer_list_contains( + v8::Local realm, + v8::Local transfer, + v8::Local value) + { + if (!transfer->IsArray()) return false; + auto list = transfer.As(); + for (uint32_t index = 0; index < list->Length(); ++index) { + v8::Local candidate; + if (list->Get(realm, index).ToLocal(&candidate) + && candidate->IsObject() + && candidate.As()->StrictEquals(value)) { + return true; + } + } + return false; + } + + static void message_port_post_message( + const v8::FunctionCallbackInfo& info) + { + auto* self = current(info.GetIsolate()); + auto* binding = message_port_binding_from_callback(info); + if (self == nullptr || binding == nullptr || binding->transferred + || binding->endpoint == nullptr || info.Length() < 1) { + return; + } + + auto realm = info.GetIsolate()->GetCurrentContext(); + v8::Local transfer = v8::Undefined(info.GetIsolate()); + if (info.Length() > 1) { + if (info[1]->IsObject() && !info[1]->IsArray()) { + if (!info[1].As()->Get( + realm, + js_string(info.GetIsolate(), "transfer")).ToLocal(&transfer)) { + return; + } + } else { + transfer = info[1]; + } + } + if (transfer_list_contains(realm, transfer, info.This())) { + throw_dom_exception( + info, + "A MessagePort cannot transfer itself", + "DataCloneError"); + return; + } + + auto endpoint = binding->endpoint; + const auto destination = 1U - endpoint->side; + size_t remaining_bytes = 0U; + { + std::lock_guard lock(endpoint->channel->mutex); + if (endpoint->channel->closed[endpoint->side] + || endpoint->channel->closed[destination]) { + return; + } + if (endpoint->channel->incoming[destination].size() + >= maximum_message_port_queue_messages + || endpoint->channel->incoming_bytes[destination] + >= maximum_message_port_queue_bytes) { + throw_dom_exception( + info, + "MessagePort queue capacity exhausted", + "QuotaExceededError"); + return; + } + remaining_bytes = maximum_message_port_queue_bytes + - endpoint->channel->incoming_bytes[destination]; + } + + clone_packet packet; + if (!serialize_clone(info, info[0], transfer, packet, remaining_bytes)) return; + const auto packet_bytes = clone_packet_size(packet); + std::shared_ptr wake; + { + std::lock_guard lock(endpoint->channel->mutex); + if (endpoint->channel->closed[endpoint->side] + || endpoint->channel->closed[destination]) { + return; + } + if (endpoint->channel->incoming[destination].size() + >= maximum_message_port_queue_messages + || packet_bytes > maximum_message_port_queue_bytes + - endpoint->channel->incoming_bytes[destination]) { + throw_dom_exception( + info, + "MessagePort queue capacity exhausted", + "QuotaExceededError"); + return; + } + endpoint->channel->incoming_bytes[destination] += packet_bytes; + endpoint->channel->incoming[destination].push_back(std::move(packet)); + wake = endpoint->channel->wake[destination].lock(); + } + if (wake != nullptr) wake->signal(); + } + + static void message_port_start( + const v8::FunctionCallbackInfo& info) + { + auto* self = current(info.GetIsolate()); + auto* binding = message_port_binding_from_callback(info); + if (self != nullptr && binding != nullptr && !binding->transferred + && binding->endpoint != nullptr) { + self->message_port_wake->signal(); + } + } + + static void close_message_port_binding(message_port_binding& binding) + { + if (binding.transferred || binding.endpoint == nullptr + || binding.endpoint->channel == nullptr) { + return; + } + const auto& endpoint = binding.endpoint; + std::lock_guard lock(endpoint->channel->mutex); + endpoint->channel->closed[endpoint->side] = true; + endpoint->channel->wake[endpoint->side].reset(); + } + + static void message_port_close( + const v8::FunctionCallbackInfo& info) + { + auto* binding = message_port_binding_from_callback(info); + if (binding != nullptr) close_message_port_binding(*binding); + } + + static void message_port_construct( + const v8::FunctionCallbackInfo& info) + { + info.GetIsolate()->ThrowException(v8::Exception::TypeError( + js_string(info.GetIsolate(), "Illegal constructor"))); + } + + static void message_channel_construct( + const v8::FunctionCallbackInfo& info) + { + auto* self = current(info.GetIsolate()); + if (self == nullptr) return; + if (!info.IsConstructCall()) { + info.GetIsolate()->ThrowException(v8::Exception::TypeError( + js_string(info.GetIsolate(), "MessageChannel requires new"))); + return; + } + const auto available_bindings = static_cast(std::count_if( + self->message_port_bindings.begin(), + self->message_port_bindings.end(), + [](const auto& binding) { + return binding != nullptr && binding->reclaimed; + })); + if (self->message_port_bindings.size() - available_bindings + 2U + > maximum_message_port_bindings) { + info.GetIsolate()->ThrowException(v8::Exception::RangeError( + js_string(info.GetIsolate(), "MessagePort binding capacity exhausted"))); + return; + } + auto channel = std::make_shared(); + auto first = std::make_shared(); + first->channel = channel; + first->side = 0U; + auto second = std::make_shared(); + second->channel = std::move(channel); + second->side = 1U; + auto realm = info.GetIsolate()->GetCurrentContext(); + v8::Local port1; + v8::Local port2; + if (!self->create_message_port_wrapper(realm, std::move(first)).ToLocal(&port1) + || !self->create_message_port_wrapper(realm, std::move(second)).ToLocal(&port2)) { + return; + } + info.This()->Set(realm, js_string(info.GetIsolate(), "port1"), port1).Check(); + info.This()->Set(realm, js_string(info.GetIsolate(), "port2"), port2).Check(); + info.GetReturnValue().Set(info.This()); + } + + void install_message_channel(v8::Local realm) + { + auto global = realm->Global(); + auto port_constructor = v8::Function::New( + realm, + message_port_construct).ToLocalChecked(); + auto channel_constructor = v8::Function::New( + realm, + message_channel_construct).ToLocalChecked(); + + v8::Local event_target_value; + v8::Local event_target_prototype; + v8::Local port_prototype; + if (global->Get(realm, js_string(isolate, "EventTarget")) + .ToLocal(&event_target_value) + && event_target_value->IsFunction() + && event_target_value.As()->Get( + realm, + js_string(isolate, "prototype")).ToLocal(&event_target_prototype) + && port_constructor->Get( + realm, + js_string(isolate, "prototype")).ToLocal(&port_prototype) + && event_target_prototype->IsObject() + && port_prototype->IsObject()) { + port_prototype.As()->SetPrototype( + realm, + event_target_prototype).Check(); + } + global->DefineOwnProperty( + realm, + js_string(isolate, "MessagePort"), + port_constructor, + v8::None).Check(); + global->DefineOwnProperty( + realm, + js_string(isolate, "MessageChannel"), + channel_constructor, + v8::None).Check(); + } + + bool has_message_port_messages() const + { + for (const auto& binding : message_port_bindings) { + if (binding == nullptr || binding->transferred + || binding->endpoint == nullptr) { + continue; + } + std::lock_guard lock(binding->endpoint->channel->mutex); + if (!binding->endpoint->channel->incoming[binding->endpoint->side].empty()) { + return true; + } + } + return false; + } + + bool drain_message_port_message() + { + for (const auto& binding : message_port_bindings) { + if (binding == nullptr || binding->transferred + || binding->endpoint == nullptr) { + continue; + } + clone_packet packet; + { + std::lock_guard lock(binding->endpoint->channel->mutex); + auto& queue = binding->endpoint->channel->incoming[ + binding->endpoint->side]; + if (queue.empty()) continue; + packet = std::move(queue.front()); + queue.pop_front(); + const auto bytes = clone_packet_size(packet); + auto& queued_bytes = binding->endpoint->channel->incoming_bytes[ + binding->endpoint->side]; + queued_bytes = bytes > queued_bytes ? 0U : queued_bytes - bytes; + } + + auto realm = binding->context.Get(isolate); + auto target = binding->wrapper.Get(isolate); + if (realm.IsEmpty() || target.IsEmpty()) return true; + v8::Context::Scope scope(realm); + v8::TryCatch caught(isolate); + std::vector> ports; + v8::Local value; + const auto decoded = deserialize_clone( + isolate, + realm, + packet, + &ports).ToLocal(&value); + if (!decoded) caught.Reset(); + auto event = v8::Object::New(isolate); + event->Set( + realm, + js_string(isolate, "type"), + js_string(isolate, decoded ? "message" : "messageerror")).Check(); + if (decoded) { + event->Set(realm, js_string(isolate, "data"), value).Check(); + } + auto port_array = v8::Array::New(isolate, static_cast(ports.size())); + for (uint32_t index = 0; index < ports.size(); ++index) { + port_array->Set(realm, index, ports[index]).Check(); + } + event->Set(realm, js_string(isolate, "ports"), port_array).Check(); + event->Set(realm, js_string(isolate, "origin"), js_string(isolate, "")).Check(); + event->Set(realm, js_string(isolate, "lastEventId"), js_string(isolate, "")).Check(); + + v8::Local dispatch; + if (!target->Get(realm, js_string(isolate, "dispatchEvent")).ToLocal(&dispatch) + || !dispatch->IsFunction()) { + return true; + } + v8::Local arguments[]{event}; + if (dispatch.As()->Call( + realm, + target, + 1, + arguments).IsEmpty()) { + last_error = describe_reported_exception(caught, realm); + return false; + } + perform_microtask_checkpoint(); + return true; + } + return true; + } + + void stop_message_ports() + { + if (message_port_wake != nullptr) message_port_wake->stop(); + for (auto& binding : message_port_bindings) { + if (binding != nullptr) close_message_port_binding(*binding); + } + } + + void reset_message_ports_for_navigation() + { + for (auto& binding : message_port_bindings) { + if (binding != nullptr) close_message_port_binding(*binding); + } + message_port_bindings.clear(); + pending_window_messages.clear(); + pending_window_message_bytes = 0U; } diff --git a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_document.inc b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_document.inc index 0acf71182..c89d68569 100644 --- a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_document.inc +++ b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_document.inc @@ -1251,6 +1251,10 @@ global->DefineOwnProperty(local_context, js_string(isolate, "isSecureContext"), v8::Boolean::New(isolate, secure), v8::ReadOnly).Check(); global->Set(local_context, js_string(isolate, "location"), location).Check(); + global->Set( + local_context, + js_string(isolate, "origin"), + js_string(isolate, origin.c_str())).Check(); v8::Local document_value; if (global->Get( local_context, diff --git a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_dom_core.inc b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_dom_core.inc index 3df274b4d..1bda85c96 100644 --- a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_dom_core.inc +++ b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_dom_core.inc @@ -800,8 +800,18 @@ std::erase_if( pending_window_messages, [&](const auto& task) { - return task.target_context.Get(isolate) == local_context; + if (task.target_context.Get(isolate) != local_context) return false; + const auto bytes = clone_packet_size(task.packet); + pending_window_message_bytes = bytes > pending_window_message_bytes + ? 0U : pending_window_message_bytes - bytes; + return true; }); + for (auto& binding : message_port_bindings) { + if (binding != nullptr + && binding->context.Get(isolate) == local_context) { + close_message_port_binding(*binding); + } + } std::erase_if( pending_dialog_close_events, [&](const auto& task) { diff --git a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_lifecycle.inc b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_lifecycle.inc index 35cfb1741..127e247c4 100644 --- a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_lifecycle.inc +++ b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_lifecycle.inc @@ -49,6 +49,7 @@ ~implementation() { + stop_message_ports(); stop_workers(); #if defined(WEBSCENE_NATIVE_ENGINE_ENABLE_MEDIA) clear_media_bindings(); @@ -142,6 +143,8 @@ media_query_lists.clear(); timers.clear(); pending_window_messages.clear(); + pending_window_message_bytes = 0U; + message_port_bindings.clear(); pending_dialog_close_events.clear(); pending_programmatic_scroll_events.clear(); pending_interop_promises.clear(); diff --git a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_navigation.inc b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_navigation.inc index e819180e1..736dc7336 100644 --- a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_navigation.inc +++ b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_navigation.inc @@ -115,6 +115,8 @@ retire_document_graphics(); #endif stop_workers(); + dedicated_workers.clear(); + reset_message_ports_for_navigation(); file_targets.clear(); { std::lock_guard lock(file_requests_mutex); file_requests.clear(); } object_url_file_data.clear(); diff --git a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_state.inc b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_state.inc index dbfd5d2dc..b7d04f366 100644 --- a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_state.inc +++ b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_state.inc @@ -125,7 +125,17 @@ std::vector resize_listeners; std::vector> media_query_lists; std::vector timers; + static constexpr size_t maximum_pending_window_messages = 256U; + static constexpr size_t maximum_pending_window_message_bytes = + 16U * 1024U * 1024U; std::deque pending_window_messages; + size_t pending_window_message_bytes{}; + static constexpr size_t maximum_message_port_bindings = 4096U; + static constexpr size_t maximum_message_port_queue_messages = 256U; + static constexpr size_t maximum_message_port_queue_bytes = 16U * 1024U * 1024U; + std::shared_ptr message_port_wake = + std::make_shared(); + std::vector> message_port_bindings; std::vector pending_dialog_close_events; std::vector pending_programmatic_scroll_events; diff --git a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_state_types.inc b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_state_types.inc index d5d66ddbb..81bad743b 100644 --- a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_state_types.inc +++ b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_state_types.inc @@ -18,10 +18,63 @@ double animation_timestamp_ms{-1}; }; + struct message_port_endpoint; + + struct clone_packet { + std::vector bytes; + std::vector> transfers; + std::vector> port_transfers; + }; + + struct message_port_runtime_wake final { + std::mutex mutex; + bool stopped{false}; + std::function notify; + + void signal() + { + std::function callback; + { + std::lock_guard lock(mutex); + if (stopped) return; + callback = notify; + } + if (callback) callback(); + } + + void stop() + { + std::lock_guard lock(mutex); + stopped = true; + notify = {}; + } + }; + + struct message_port_channel final { + std::mutex mutex; + std::array, 2U> incoming; + std::array incoming_bytes{}; + std::array closed{}; + std::array, 2U> wake; + }; + + struct message_port_endpoint final { + std::shared_ptr channel; + size_t side{}; + }; + + struct message_port_binding final { + std::shared_ptr endpoint; + v8::Global context; + v8::Global wrapper; + bool transferred{false}; + bool reclaimed{false}; + }; + struct pending_window_message final { v8::Global target_context; v8::Global source_context; - v8::Global data; + clone_packet packet; std::string target_origin; std::string source_origin; }; diff --git a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_tasks.inc b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_tasks.inc index 509fd86a0..75b651d10 100644 --- a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_tasks.inc +++ b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_tasks.inc @@ -349,11 +349,13 @@ if (pending_window_messages.empty()) return true; auto task = std::move(pending_window_messages.front()); pending_window_messages.pop_front(); + const auto packet_bytes = clone_packet_size(task.packet); + pending_window_message_bytes = packet_bytes > pending_window_message_bytes + ? 0U : pending_window_message_bytes - packet_bytes; auto target_context = task.target_context.Get(isolate); auto source_context = task.source_context.Get(isolate); - auto data = task.data.Get(isolate); - if (target_context.IsEmpty() || source_context.IsEmpty() || data.IsEmpty()) { + if (target_context.IsEmpty() || source_context.IsEmpty()) { return true; } @@ -384,15 +386,26 @@ return true; } + std::vector> transferred_ports; + v8::Local data; + v8::TryCatch decode_error(isolate); + const auto decoded = deserialize_clone( + isolate, + target_context, + task.packet, + &transferred_ports).ToLocal(&data); + if (!decoded) decode_error.Reset(); auto event = v8::Object::New(isolate); event->Set( target_context, js_string(isolate, "type"), - js_string(isolate, "message")).Check(); - event->Set( - target_context, - js_string(isolate, "data"), - data).Check(); + js_string(isolate, decoded ? "message" : "messageerror")).Check(); + if (decoded) { + event->Set( + target_context, + js_string(isolate, "data"), + data).Check(); + } event->Set( target_context, js_string(isolate, "origin"), @@ -405,10 +418,13 @@ target_context, js_string(isolate, "lastEventId"), js_string(isolate, "")).Check(); - event->Set( - target_context, - js_string(isolate, "ports"), - v8::Array::New(isolate)).Check(); + auto ports = v8::Array::New( + isolate, + static_cast(transferred_ports.size())); + for (uint32_t index = 0; index < transferred_ports.size(); ++index) { + ports->Set(target_context, index, transferred_ports[index]).Check(); + } + event->Set(target_context, js_string(isolate, "ports"), ports).Check(); event->Set( target_context, js_string(isolate, "bubbles"), @@ -570,6 +586,7 @@ if(drain_media())return true; #endif if (has_worker_messages()) return drain_worker_message(); + if (has_message_port_messages()) return drain_message_port_message(); if (websocket_transport.has_pending_events()) { return drain_websocket_event(); } diff --git a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_workers.inc b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_workers.inc index d6de5e7d2..5b230a3f4 100644 --- a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_workers.inc +++ b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime_workers.inc @@ -2,6 +2,8 @@ std::mutex mutex; std::condition_variable ready; std::deque incoming, outgoing; + size_t incoming_bytes{}; + size_t outgoing_bytes{}; std::deque errors; bool stopped{false}; bool wake_requested{false}; @@ -16,6 +18,7 @@ stopped = true; if (running_isolate) running_isolate->TerminateExecution(); incoming.clear(); outgoing.clear(); errors.clear(); + incoming_bytes = 0U; outgoing_bytes = 0U; } ready.notify_all(); if (thread.joinable()) thread.join(); @@ -25,6 +28,12 @@ std::vector> dedicated_workers; worker_state* worker_owner{}; bool force_dedicated_isolate{false}; + bool worker_module_realm{false}; + static constexpr size_t maximum_dedicated_workers = 64U; + static constexpr size_t maximum_worker_queue_messages = 256U; + static constexpr size_t maximum_worker_queue_bytes = 16U * 1024U * 1024U; + static constexpr size_t maximum_worker_errors = 16U; + static constexpr size_t maximum_worker_error_length = 4096U; static v8::MaybeLocal worker_transfer( const v8::FunctionCallbackInfo& info) { @@ -37,9 +46,31 @@ auto* state = static_cast(info.Data().As()->Value(v8::kExternalPointerTypeTagDefault)); v8::Local transfer; if (!worker_transfer(info).ToLocal(&transfer)) return; + size_t remaining_bytes = 0U; + { + std::lock_guard lock(state->mutex); + if (state->stopped) return; + if (state->incoming.size() >= maximum_worker_queue_messages + || state->incoming_bytes >= maximum_worker_queue_bytes) { + throw_dom_exception(info,"Worker message queue capacity exhausted","QuotaExceededError"); + return; + } + remaining_bytes = maximum_worker_queue_bytes - state->incoming_bytes; + } clone_packet packet; - if (!serialize_clone(info, info[0], transfer, packet)) return; - { std::lock_guard lock(state->mutex); if (state->stopped) return; state->incoming.push_back(std::move(packet)); } + if (!serialize_clone(info, info[0], transfer, packet, remaining_bytes)) return; + const auto bytes = clone_packet_size(packet); + { + std::lock_guard lock(state->mutex); + if (state->stopped) return; + if (state->incoming.size() >= maximum_worker_queue_messages + || bytes > maximum_worker_queue_bytes - state->incoming_bytes) { + throw_dom_exception(info,"Worker message queue capacity exhausted","QuotaExceededError"); + return; + } + state->incoming_bytes += bytes; + state->incoming.push_back(std::move(packet)); + } state->ready.notify_one(); } static void worker_post_parent_impl(const v8::FunctionCallbackInfo& info) { @@ -48,9 +79,31 @@ auto* state = self->worker_owner; v8::Local transfer; if (!worker_transfer(info).ToLocal(&transfer)) return; + size_t remaining_bytes = 0U; + { + std::lock_guard lock(state->mutex); + if (state->stopped) return; + if (state->outgoing.size() >= maximum_worker_queue_messages + || state->outgoing_bytes >= maximum_worker_queue_bytes) { + throw_dom_exception(info,"Worker message queue capacity exhausted","QuotaExceededError"); + return; + } + remaining_bytes = maximum_worker_queue_bytes - state->outgoing_bytes; + } clone_packet packet; - if (!serialize_clone(info, info[0], transfer, packet)) return; - { std::lock_guard lock(state->mutex); if (state->stopped) return; state->outgoing.push_back(std::move(packet)); } + if (!serialize_clone(info, info[0], transfer, packet, remaining_bytes)) return; + const auto bytes = clone_packet_size(packet); + { + std::lock_guard lock(state->mutex); + if (state->stopped) return; + if (state->outgoing.size() >= maximum_worker_queue_messages + || bytes > maximum_worker_queue_bytes - state->outgoing_bytes) { + throw_dom_exception(info,"Worker message queue capacity exhausted","QuotaExceededError"); + return; + } + state->outgoing_bytes += bytes; + state->outgoing.push_back(std::move(packet)); + } if (state->notify_parent) state->notify_parent(); } static void worker_terminate(const v8::FunctionCallbackInfo& info) { @@ -63,6 +116,34 @@ std::lock_guard lock(self->worker_owner->mutex); self->worker_owner->stopped = true; } } + + static void worker_import_scripts(const v8::FunctionCallbackInfo& info) { + auto* self = current(info.GetIsolate()); + if (!self || !self->worker_owner) return; + if (self->worker_module_realm) { + info.GetIsolate()->ThrowException(v8::Exception::TypeError( + js_string(info.GetIsolate(),"importScripts is unavailable in a module worker"))); + return; + } + auto realm = info.GetIsolate()->GetCurrentContext(); + for (int index = 0; index < info.Length(); ++index) { + const auto url = resolve_resource_url( + to_utf8(info.GetIsolate(),info[index]), + self->current_base_address()); + if (resource_origin(url) != resource_origin(self->current_base_address())) { + throw_dom_exception(info,"Worker script must be same origin","SecurityError"); + return; + } + std::string source,resolved,error; + if (!self->load_text_resource(url,{},WEBSCENE_RESOURCE_SCRIPT,source,resolved) + || !self->execute_in_context(realm,source,resolved,error)) { + info.GetIsolate()->ThrowException(v8::Exception::Error( + js_dom_string(info.GetIsolate(), + error.empty() ? "Unable to load worker script: "+url : error))); + return; + } + } + } bool dispatch_worker_message(v8::Local realm, v8::Local target, const clone_packet* packet, const std::string& error = {}) { v8::Context::Scope scope(realm); @@ -71,8 +152,17 @@ event->Set(realm, js_string(isolate,"type"), js_string(isolate, packet ? "message" : "error")).Check(); if (packet) { v8::Local value; - if (!deserialize_clone(isolate, realm, *packet).ToLocal(&value)) return false; - event->Set(realm, js_string(isolate,"data"), value).Check(); + std::vector> ports; + if (!deserialize_clone(isolate, realm, *packet, &ports).ToLocal(&value)) { + caught.Reset(); + event->Set(realm,js_string(isolate,"type"),js_string(isolate,"messageerror")).Check(); + } else { + event->Set(realm, js_string(isolate,"data"), value).Check(); + auto port_array=v8::Array::New(isolate,static_cast(ports.size())); + for(uint32_t index=0;indexSet(realm,index,ports[index]).Check(); + event->Set(realm,js_string(isolate,"ports"),port_array).Check(); + } } else event->Set(realm, js_string(isolate,"message"), js_dom_string(isolate,error)).Check(); v8::Local dispatch; if (!target->Get(realm, js_string(isolate,"dispatchEvent")).ToLocal(&dispatch) || !dispatch->IsFunction()) return false; @@ -96,7 +186,11 @@ { std::lock_guard lock(state->mutex); if (!state->errors.empty()) { error = std::move(state->errors.front()); state->errors.pop_front(); } - else if (!state->outgoing.empty()) { packet = std::move(state->outgoing.front()); state->outgoing.pop_front(); } + else if (!state->outgoing.empty()) { + packet = std::move(state->outgoing.front()); state->outgoing.pop_front(); + const auto bytes=clone_packet_size(*packet); + state->outgoing_bytes=bytes>state->outgoing_bytes?0U:state->outgoing_bytes-bytes; + } else continue; } return dispatch_worker_message(state->realm.Get(isolate),state->wrapper.Get(isolate), @@ -114,10 +208,6 @@ } const auto url = resolve_resource_url(to_utf8(isolate,info[0]),self->current_base_address()); // Until worker CORS/credentials are qualified, enforce a same-origin script. - auto origin = [](const std::string& value) { - auto scheme = value.find("://"); - return scheme == std::string::npos ? std::string{} : value.substr(0,value.find('/',scheme+3)); - }; // Capture a same-origin object URL before entering another isolate. The // worker owns its source even if the parent revokes the URL immediately. std::optional embedded_source; @@ -129,7 +219,8 @@ } embedded_source.emplace(found->second.bytes.begin(), found->second.bytes.end()); } - if (!embedded_source && origin(url) != origin(self->current_base_address())) { + if (!embedded_source + && resource_origin(url) != resource_origin(self->current_base_address())) { throw_dom_exception(info,"Worker script must be same origin","SecurityError"); return; } bool module = false; @@ -148,7 +239,7 @@ if(!info[1].As()->Get(realm,js_string(isolate,"name")).ToLocal(&name))return; if(!name->IsUndefined())worker_name=to_utf8(isolate,name); } - if (self->dedicated_workers.size() >= 1024) { + if (self->dedicated_workers.size() >= maximum_dedicated_workers) { isolate->ThrowException(v8::Exception::RangeError(js_string(isolate,"Worker capacity exhausted"))); return; } v8::Local target_constructor; @@ -167,10 +258,16 @@ target->Set(realm,js_string(isolate,"terminate"),v8::Function::New(realm,worker_terminate,data).ToLocalChecked()).Check(); auto loader=self->load_resource_callback; auto root=self->resource_root; + auto worker_origin=resource_origin(self->current_base_address()); self->dedicated_workers.push_back(std::move(state)); - raw->thread=std::thread([raw,url,module,embedded_source=std::move(embedded_source),worker_name=std::move(worker_name),loader=std::move(loader),root=std::move(root)] { + raw->thread=std::thread([raw,url,module,worker_origin=std::move(worker_origin),embedded_source=std::move(embedded_source),worker_name=std::move(worker_name),loader=std::move(loader),root=std::move(root)] { auto report=[&](const std::string& text) { - {std::lock_guard lock(raw->mutex);if(raw->stopped)return;raw->errors.push_back(text);} + { + std::lock_guard lock(raw->mutex); + if(raw->stopped)return; + if(raw->errors.size()>=maximum_worker_errors)raw->errors.pop_front(); + raw->errors.push_back(text.substr(0U,maximum_worker_error_length)); + } if(raw->notify_parent)raw->notify_parent(); }; try { @@ -184,7 +281,7 @@ if(child->isolate)child->isolate->CancelTerminateExecution(); } } guard{raw,child}; - child->force_dedicated_isolate=true;child->worker_owner=raw; + child->force_dedicated_isolate=true;child->worker_owner=raw;child->worker_module_realm=module; child->runtime_work_available=[raw]{ {std::lock_guard lock(raw->mutex);raw->wake_requested=true;} raw->ready.notify_one(); @@ -202,8 +299,16 @@ auto global=realm->Global(); child->document_base_address=url; child->set_context_location(realm,global,url); + if(embedded_source) { + v8::Local location; + if(global->Get(realm,js_string(child->isolate,"location")).ToLocal(&location) + &&location->IsObject()) + location.As()->Set(realm,js_string(child->isolate,"origin"),js_dom_string(child->isolate,worker_origin)).Check(); + global->Set(realm,js_string(child->isolate,"origin"),js_dom_string(child->isolate,worker_origin)).Check(); + } global->Set(realm,js_string(child->isolate,"postMessage"),v8::Function::New(realm,worker_post_parent,{},1).ToLocalChecked()).Check(); global->Set(realm,js_string(child->isolate,"close"),v8::Function::New(realm,worker_close).ToLocalChecked()).Check(); + global->Set(realm,js_string(child->isolate,"importScripts"),v8::Function::New(realm,worker_import_scripts).ToLocalChecked()).Check(); global->Set(realm,js_string(child->isolate,"self"),global).Check(); global->Set(realm,js_string(child->isolate,"name"),js_dom_string(child->isolate,worker_name)).Check(); global->Delete(realm,js_string(child->isolate,"document")).FromMaybe(false); @@ -222,7 +327,11 @@ [&]{return raw->stopped||raw->wake_requested||!raw->incoming.empty();}); raw->wake_requested=false; if(raw->stopped)break; - if(!raw->incoming.empty()){packet=std::move(raw->incoming.front());raw->incoming.pop_front();} + if(!raw->incoming.empty()){ + packet=std::move(raw->incoming.front());raw->incoming.pop_front(); + const auto bytes=clone_packet_size(*packet); + raw->incoming_bytes=bytes>raw->incoming_bytes?0U:raw->incoming_bytes-bytes; + } } v8::Isolate::Scope entered(child->isolate);v8::HandleScope handles(child->isolate); auto realm=child->context.Get(child->isolate);v8::Context::Scope scope(realm); diff --git a/experiments/WebScene.NativeEngine.Probe/tests/hybrid_v8_runtime_tests.cpp b/experiments/WebScene.NativeEngine.Probe/tests/hybrid_v8_runtime_tests.cpp index aa2c96d00..cbb453367 100644 --- a/experiments/WebScene.NativeEngine.Probe/tests/hybrid_v8_runtime_tests.cpp +++ b/experiments/WebScene.NativeEngine.Probe/tests/hybrid_v8_runtime_tests.cpp @@ -3,6 +3,8 @@ #include #include #include +#include +#include #include void require(bool value, const char* message) { if (!value) throw std::runtime_error(message); } void test_compiled_template_shared_document() { @@ -97,6 +99,331 @@ void test_blob_worker_source_lifetime() { require(runtime.execute("blobWorker.terminate()", "worker-cleanup"), runtime.last_error().c_str()); } +void test_worker_message_port_transfer_and_throughput() { + webscene_native::native_document document; + webscene_native::v8_dom_runtime runtime(document, + []{return webscene_native::v8_dom_runtime::viewport_metrics{640,480,1,0};}); + bool completed = false; + runtime.register_compiled_template("worker-port-result", [&](auto& dom, const std::string& value) -> auto& { + require(value == "10000", "Worker MessagePort result changed"); + completed = true; + return dom.create_element("span"); + }); + require(runtime.initialize(), "Worker MessagePort runtime failed"); + require(runtime.execute(R"JS( + const source = ` + onmessage = event => { + const port = event.data; + if (!(port instanceof MessagePort)) throw Error('Transferred value is not a MessagePort'); + if (typeof document !== 'undefined' || typeof importScripts !== 'function') + throw Error('Worker global shape changed'); + port.onmessage = message => port.postMessage(message.data + 1); + port.start(); + }`; + const url = URL.createObjectURL(new Blob([source], {type:'text/javascript'})); + const worker = new Worker(url, {name:'message-port-contract'}); + URL.revokeObjectURL(url); + const channel = new MessageChannel(); + if (!(channel.port1 instanceof MessagePort) || !(channel.port1 instanceof EventTarget)) + throw Error('MessagePort interface identity changed'); + let count = 0; + channel.port1.onmessage = event => { + count = event.data; + if (count < 10000) channel.port1.postMessage(count); + else { + worker.terminate(); + channel.port1.close(); + document.createCompiledTemplate('worker-port-result', count); + } + }; + worker.postMessage(channel.port2, [channel.port2]); + channel.port2.postMessage('detached endpoint must be inert'); + channel.port1.postMessage(0); + )JS", "worker-message-port-throughput"), runtime.last_error().c_str()); + const auto started = std::chrono::steady_clock::now(); + while (!completed + && std::chrono::steady_clock::now() - started < std::chrono::seconds(12)) { + require(runtime.pump_task(), runtime.last_error().c_str()); + std::this_thread::yield(); + } + const auto elapsed = std::chrono::duration_cast( + std::chrono::steady_clock::now() - started); + require(completed, "Worker MessagePort throughput did not complete"); + require(elapsed < std::chrono::seconds(10), "10,000 MessagePort round trips exceeded 10 seconds"); +} + +void test_message_port_clone_and_queue_bounds() { + webscene_native::native_document document; + webscene_native::v8_dom_runtime runtime(document, + []{return webscene_native::v8_dom_runtime::viewport_metrics{640,480,1,0};}); + require(runtime.initialize(), "MessagePort bound runtime failed"); + require(runtime.execute(R"JS( + const clonedChannel = new MessageChannel(); + const clone = structuredClone( + {port: clonedChannel.port2}, + {transfer: [clonedChannel.port2]}); + if (!(clone.port instanceof MessagePort)) throw Error('structuredClone lost MessagePort identity'); + let duplicateRejected = false; + try { structuredClone(clone.port, {transfer:[clone.port, clone.port]}); } + catch (error) { duplicateRejected = error.name === 'DataCloneError'; } + if (!duplicateRejected) throw Error('Duplicate MessagePort transfer was accepted'); + clone.port.close(); + clonedChannel.port1.close(); + + const bounded = new MessageChannel(); + const oversized = new ArrayBuffer(16 * 1024 * 1024); + let byteSaturated = false; + try { bounded.port1.postMessage(oversized, [oversized]); } + catch (error) { byteSaturated = error.name === 'QuotaExceededError'; } + if (!byteSaturated || oversized.byteLength === 0) + throw Error('MessagePort byte bound detached or accepted an oversized packet'); + for (let index = 0; index < 256; ++index) bounded.port1.postMessage(index); + let saturated = false; + try { bounded.port1.postMessage(256); } + catch (error) { saturated = error.name === 'QuotaExceededError'; } + if (!saturated) throw Error('MessagePort queue accepted an unbounded backlog'); + bounded.port1.close(); + bounded.port2.close(); + + globalThis.closedPortDelivered = false; + const closing = new MessageChannel(); + closing.port2.onmessage = () => { closedPortDelivered = true; }; + closing.port2.close(); + closing.port2.close(); + closing.port1.postMessage('must not be delivered'); + closing.port1.close(); + )JS", "message-port-clone-and-bounds"), runtime.last_error().c_str()); + require(runtime.execute( + "if (closedPortDelivered) throw Error('Closed MessagePort received a queued message')", + "message-port-close"), runtime.last_error().c_str()); +} + +void test_message_port_binding_memory_is_bounded() { + webscene_native::native_document document; + webscene_native::v8_dom_runtime runtime(document, + []{return webscene_native::v8_dom_runtime::viewport_metrics{640,480,1,0};}); + require(runtime.initialize(), "MessagePort memory runtime failed"); + + auto create_and_release_batch = [&] { + require(runtime.execute(R"JS( + (() => { + for (let index = 0; index < 2048; ++index) { + const channel = new MessageChannel(); + channel.port1.close(); + channel.port2.close(); + } + })() + )JS", "message-port-memory-batch"), runtime.last_error().c_str()); + runtime.notify_low_memory(); + require(runtime.pump_task(), runtime.last_error().c_str()); + }; + + create_and_release_batch(); + const auto baseline = runtime.read_memory_metrics().used_heap_bytes; + create_and_release_batch(); + create_and_release_batch(); + const auto retained = runtime.read_memory_metrics().used_heap_bytes; + require( + retained <= baseline + 16U * 1024U * 1024U, + "Released MessagePort batches retained more than 16 MiB"); +} + +void test_iframe_worker_extension_host_port_bootstrap() { + webscene_native::native_document document; + const std::string root_url = "https://worker.test/index.html"; + const std::string root_html = R"HTML( + + )HTML"; + bool completed = false; + webscene_native::v8_dom_runtime runtime(document, + []{return webscene_native::v8_dom_runtime::viewport_metrics{640,480,1,0};}, {}, + [&](uint32_t kind, const std::string& url, const auto&, const std::string&, int64_t, + webscene_native::v8_dom_runtime::resource_response& response) { + if (kind != WEBSCENE_RESOURCE_DOCUMENT || url != root_url) return false; + response.content = root_html; + return true; + }); + runtime.register_compiled_template("iframe-port-result", [&](auto& dom, const std::string& value) -> auto& { + require(value == "42", "iframe MessagePort response changed"); + completed = true; + return dom.create_element("span"); + }); + require(runtime.initialize(), "iframe MessagePort runtime failed"); + require(runtime.load_url(root_url), runtime.last_error().c_str()); + for (unsigned i = 0; i < 5000 && !completed; ++i) { + require(runtime.pump_task(), runtime.last_error().c_str()); + if (!runtime.has_pending_tasks()) std::this_thread::sleep_for(std::chrono::milliseconds(1)); + } + require(completed, "Code OSS-shaped iframe MessagePort bootstrap did not complete"); +} + +void test_worker_termination_race() { + webscene_native::native_document document; + webscene_native::v8_dom_runtime runtime(document, + []{return webscene_native::v8_dom_runtime::viewport_metrics{640,480,1,0};}); + require(runtime.initialize(), "Worker termination runtime failed"); + require(runtime.execute(R"JS( + const url = URL.createObjectURL(new Blob(['while(true) {}'], {type:'text/javascript'})); + globalThis.terminatingWorker = new Worker(url); + URL.revokeObjectURL(url); + )JS", "worker-termination-start"), runtime.last_error().c_str()); + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + const auto started = std::chrono::steady_clock::now(); + require(runtime.execute(R"JS( + for (let index = 0; index < 256; ++index) terminatingWorker.postMessage(index); + let saturated = false; + try { terminatingWorker.postMessage(256); } + catch (error) { saturated = error.name === 'QuotaExceededError'; } + if (!saturated) throw Error('Worker queue accepted an unbounded backlog'); + terminatingWorker.terminate(); + terminatingWorker.terminate(); + )JS", + "worker-termination-stop"), runtime.last_error().c_str()); + require(std::chrono::steady_clock::now() - started < std::chrono::seconds(2), + "Worker termination did not cancel active execution promptly"); +} + +void test_worker_error_delivery() { + webscene_native::native_document document; + webscene_native::v8_dom_runtime runtime(document, + []{return webscene_native::v8_dom_runtime::viewport_metrics{640,480,1,0};}); + unsigned delivered = 0; + runtime.register_compiled_template("worker-error-result", [&](auto& dom, const std::string& value) -> auto& { + require(value == "1", "worker error event shape changed"); + ++delivered; + return dom.create_element("span"); + }); + require(runtime.initialize(), "worker error runtime failed"); + require(runtime.execute(R"JS( + const url = URL.createObjectURL(new Blob( + ['throw new Error("expected-worker-failure")'], + {type:'text/javascript'})); + const failingWorker = new Worker(url); + URL.revokeObjectURL(url); + failingWorker.onerror = event => { + document.createCompiledTemplate( + 'worker-error-result', + event.type === 'error' && event.message.includes('expected-worker-failure') ? 1 : 0); + failingWorker.terminate(); + }; + + const moduleUrl = URL.createObjectURL(new Blob( + ['export const broken = ;'], + {type:'text/javascript'})); + const failingModuleWorker = new Worker(moduleUrl, {type:'module'}); + URL.revokeObjectURL(moduleUrl); + failingModuleWorker.onerror = event => { + document.createCompiledTemplate( + 'worker-error-result', event.type === 'error' ? 1 : 0); + failingModuleWorker.terminate(); + }; + )JS", "worker-error-delivery"), runtime.last_error().c_str()); + for (unsigned index = 0; index < 2000 && delivered < 2U; ++index) { + require(runtime.pump_task(), runtime.last_error().c_str()); + if (!runtime.has_pending_tasks()) std::this_thread::sleep_for(std::chrono::milliseconds(1)); + } + require(delivered == 2U, "Classic and module worker failures did not dispatch error events"); +} + +void test_worker_and_port_navigation_shutdown() { + webscene_native::native_document document; + const std::string navigation_url = "https://worker.test/after-navigation.html"; + webscene_native::v8_dom_runtime runtime(document, + []{return webscene_native::v8_dom_runtime::viewport_metrics{640,480,1,0};}, {}, + [&](uint32_t kind, const std::string& url, const auto&, const std::string&, int64_t, + webscene_native::v8_dom_runtime::resource_response& response) { + if (kind != WEBSCENE_RESOURCE_DOCUMENT || url != navigation_url) return false; + response.content = "after navigation"; + return true; + }); + require(runtime.initialize(), "worker navigation runtime failed"); + require(runtime.execute(R"JS( + const url = URL.createObjectURL(new Blob(['while(true) {}'], {type:'text/javascript'})); + globalThis.navigationWorker = new Worker(url); + URL.revokeObjectURL(url); + globalThis.navigationChannel = new MessageChannel(); + navigationChannel.port1.postMessage('queued before navigation'); + )JS", "worker-navigation-start"), runtime.last_error().c_str()); + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + const auto started = std::chrono::steady_clock::now(); + require(runtime.load_url(navigation_url), runtime.last_error().c_str()); + require(std::chrono::steady_clock::now() - started < std::chrono::seconds(2), + "Navigation did not cancel worker execution promptly"); + require(runtime.execute(R"JS( + const afterNavigation = new MessageChannel(); + afterNavigation.port1.close(); + afterNavigation.port2.close(); + )JS", "worker-navigation-resources-reset"), runtime.last_error().c_str()); +} + +void test_window_messageerror_on_receiver_resource_exhaustion() { + webscene_native::native_document document; + webscene_native::v8_dom_runtime runtime(document, + []{return webscene_native::v8_dom_runtime::viewport_metrics{640,480,1,0};}); + bool completed = false; + runtime.register_compiled_template("messageerror-result", [&](auto& dom, const std::string& value) -> auto& { + require(value == "1", "iframe received the wrong structured-clone failure event"); + completed = true; + return dom.create_element("span"); + }); + require(runtime.initialize(), "window messageerror runtime failed"); + require(runtime.execute(R"JS( + addEventListener('message', event => { + if (event.data !== 'messageerror-ready') return; + const channels = []; + for (let index = 0; index < 2048; ++index) channels.push(new MessageChannel()); + globalThis.__webSceneMessageErrorChannels = channels; + const transferred = channels[channels.length - 1].port2; + event.source.postMessage({transferred}, '*', [transferred]); + }); + const frame = document.createElement('iframe'); + document.body.appendChild(frame); + const child = frame.contentDocument; + child.open(); + child.write(` diff --git a/tests/WebPlatformSubset/contracts/resources/workers/message-port-helper.js b/tests/WebPlatformSubset/contracts/resources/workers/message-port-helper.js new file mode 100644 index 000000000..e5d0507c0 --- /dev/null +++ b/tests/WebPlatformSubset/contracts/resources/workers/message-port-helper.js @@ -0,0 +1 @@ +self.workerHelperValue = 42; diff --git a/tests/WebPlatformSubset/contracts/resources/workers/message-port.js b/tests/WebPlatformSubset/contracts/resources/workers/message-port.js new file mode 100644 index 000000000..768662d57 --- /dev/null +++ b/tests/WebPlatformSubset/contracts/resources/workers/message-port.js @@ -0,0 +1,19 @@ +self.importScripts('./message-port-helper.js'); + +self.onmessage = event => { + const { payload, port, bytes } = event.data; + const valid = event.ports.length === 1 + && event.ports[0] === port + && port instanceof MessagePort + && port instanceof EventTarget + && payload.self === payload + && payload.map.get(7) === 'seven' + && payload.date.getTime() === 456 + && bytes[0] === 41 + && self.workerHelperValue === 42 + && typeof document === 'undefined' + && origin === location.origin; + bytes[0] += 1; + port.postMessage({ valid, bytes }, [bytes.buffer]); + port.close(); +}; diff --git a/tests/WebPlatformSubset/contracts/worker-messageport-structured-clone.html b/tests/WebPlatformSubset/contracts/worker-messageport-structured-clone.html new file mode 100644 index 000000000..16b68b7a2 --- /dev/null +++ b/tests/WebPlatformSubset/contracts/worker-messageport-structured-clone.html @@ -0,0 +1,44 @@ + +Worker structured clone and transferable MessagePort contract + diff --git a/tests/WebPlatformSubset/contracts/worker-termination-queue-messageerror.html b/tests/WebPlatformSubset/contracts/worker-termination-queue-messageerror.html new file mode 100644 index 000000000..33ca9d85e --- /dev/null +++ b/tests/WebPlatformSubset/contracts/worker-termination-queue-messageerror.html @@ -0,0 +1,74 @@ + +Worker termination and bounded messaging contract + diff --git a/tests/WebPlatformSubset/webscene-aureon-runtime-profile.json b/tests/WebPlatformSubset/webscene-aureon-runtime-profile.json index 169748c0a..37a6744b2 100644 --- a/tests/WebPlatformSubset/webscene-aureon-runtime-profile.json +++ b/tests/WebPlatformSubset/webscene-aureon-runtime-profile.json @@ -21,6 +21,40 @@ "reason": "Project-owned regression contract; full upstream WPT conformance remains separate.", "nativeNavigation": true }, + { + "path": "contracts/worker-messageport-structured-clone.html", + "type": "contract", + "capabilities": [ + "dedicated-workers", + "structured-clone", + "messageport-transfer", + "arraybuffer-transfer" + ], + "reason": "Focused Code OSS worker-host transport candidate with cycles, built-in clone values, transferable ownership, and MessageEvent shape.", + "nativeNavigation": true + }, + { + "path": "contracts/iframe-window-messageport-origin.html", + "type": "contract", + "capabilities": [ + "iframe-window-postmessage", + "same-origin-event-origin", + "messageport-transfer" + ], + "reason": "Focused Code OSS iframe bootstrap candidate matching the extension-host MessagePort handoff.", + "nativeNavigation": true + }, + { + "path": "contracts/worker-termination-queue-messageerror.html", + "type": "contract", + "capabilities": [ + "worker-termination", + "bounded-message-queues", + "messageerror" + ], + "reason": "WebScene resource-bound candidate covering prompt termination, deterministic saturation, and receiver-side clone failure delivery.", + "nativeNavigation": true + }, { "path": "contracts/svg-inline-style.html", "type": "contract",