From 7f01c2943d9741c95b715f4d599a5676c820e4fd Mon Sep 17 00:00:00 2001 From: Karl Seguin Date: Thu, 16 Jul 2026 13:12:33 +0800 Subject: [PATCH] webapi: Add SharedWorker Shared on the Session (by url+name), but owned by a Page. When the created page goes, so does the SharedWorker. --- src/browser/Frame.zig | 9 + src/browser/Page.zig | 11 + src/browser/Session.zig | 6 + src/browser/js/Env.zig | 31 +- src/browser/js/Execution.zig | 8 + src/browser/js/Snapshot.zig | 12 +- src/browser/js/bridge.zig | 15 +- .../tests/shared_worker/counter-worker.js | 8 + .../tests/shared_worker/echo-worker.js | 8 + .../tests/shared_worker/listener-worker.js | 10 + .../tests/shared_worker/shared_worker.html | 137 ++++++ src/browser/webapi/EventTarget.zig | 3 + src/browser/webapi/MessagePort.zig | 109 +++-- src/browser/webapi/SharedWorker.zig | 146 +++++++ .../webapi/SharedWorkerGlobalScope.zig | 407 ++++++++++++++++++ src/browser/webapi/Window.zig | 2 +- src/browser/webapi/Worker.zig | 25 +- src/browser/webapi/WorkerGlobalScope.zig | 14 + src/browser/webapi/event/MessageEvent.zig | 39 +- 19 files changed, 935 insertions(+), 65 deletions(-) create mode 100644 src/browser/tests/shared_worker/counter-worker.js create mode 100644 src/browser/tests/shared_worker/echo-worker.js create mode 100644 src/browser/tests/shared_worker/listener-worker.js create mode 100644 src/browser/tests/shared_worker/shared_worker.html create mode 100644 src/browser/webapi/SharedWorker.zig create mode 100644 src/browser/webapi/SharedWorkerGlobalScope.zig diff --git a/src/browser/Frame.zig b/src/browser/Frame.zig index 23b7122cc..1381d1086 100644 --- a/src/browser/Frame.zig +++ b/src/browser/Frame.zig @@ -52,6 +52,7 @@ const Screen = @import("webapi/Screen.zig"); const VisualViewport = @import("webapi/VisualViewport.zig"); const AbstractRange = @import("webapi/AbstractRange.zig"); const Worker = @import("webapi/Worker.zig"); +const MessagePort = @import("webapi/MessagePort.zig"); const CSSStyleSheet = @import("webapi/css/CSSStyleSheet.zig"); const CustomElementDefinition = @import("webapi/CustomElementDefinition.zig"); const PageTransitionEvent = @import("webapi/event/PageTransitionEvent.zig"); @@ -194,6 +195,9 @@ _live_ranges: std.DoublyLinkedList = .{}, // channels in this frame's origin _broadcast_channels: std.DoublyLinkedList = .{}, +// List of MessagePorts living in this frame's context. +_message_ports: std.DoublyLinkedList = .{}, + // MutationObserver / IntersectionObserver bookkeeping. See frame/observers.zig. _mutation: observers.Mutation = .{}, _intersection: observers.Intersection = .{}, @@ -469,6 +473,11 @@ pub fn deinit(self: *Frame) void { page.releaseArena(qn.arena); } + while (self._message_ports.first) |node| { + const port: *MessagePort = @alignCast(@fieldParentPtr("_node", node)); + port.close(); // removes from self._message_ports + } + { // Release all objects we're referencing { diff --git a/src/browser/Page.zig b/src/browser/Page.zig index 413468bda..d63c467ef 100644 --- a/src/browser/Page.zig +++ b/src/browser/Page.zig @@ -27,6 +27,8 @@ const Session = @import("Session.zig"); const Factory = @import("Factory.zig"); const Viewport = @import("Viewport.zig"); +const SharedWorkerGlobalScope = @import("webapi/SharedWorkerGlobalScope.zig"); + const v8 = js.v8; const Allocator = std.mem.Allocator; const IS_DEBUG = builtin.mode == .Debug; @@ -116,6 +118,10 @@ popups: std.ArrayList(*Frame) = .empty, // ideal from a memory point of view). closed_frames: std.ArrayList(*Frame) = .empty, +// SharedWorkerGlobalScopes created by this Page's frames (also registered in +// session.shared_workers so other pages can connect). +shared_workers: std.ArrayList(*SharedWorkerGlobalScope) = .empty, + // In-flight navigation for a root page. When not null, this page will "replace" // the referenced page once the response header arrives. This is necessary // because, during navigation, both the "old" and "new" pages remain addressable @@ -166,6 +172,11 @@ pub fn deinit(self: *Page) void { self.frame.deinit(); + for (self.shared_workers.items) |scope| { + scope.deinit(); + } + self.shared_workers = .empty; + const session = self.session; lp.metrics.js_heap_size_bytes.observe(session.browser.env.isolate.getHeapStatistics().total_physical_size); defer session.browser.env.memoryPressureNotification(.moderate); diff --git a/src/browser/Session.zig b/src/browser/Session.zig index a85f0fa91..f0027fab0 100644 --- a/src/browser/Session.zig +++ b/src/browser/Session.zig @@ -33,6 +33,7 @@ const Browser = @import("Browser.zig"); pub const Runner = @import("Runner.zig"); const Notification = @import("../Notification.zig"); const QueuedNavigation = Frame.QueuedNavigation; +const SharedWorkerGlobalScope = @import("webapi/SharedWorkerGlobalScope.zig"); const log = lp.log; const ArenaPool = App.ArenaPool; @@ -68,6 +69,11 @@ arena_pool: *ArenaPool, // both the live page and its in-flight replacement. pages: std.ArrayList(*Page) = .{}, +// Live SharedWorkerGlobalScopes, keyed by "url\x00name", so every +// `new SharedWorker(url, name)` in the session connects to the same instance. +// Owned by the Page that creates it. +shared_workers: std.StringHashMapUnmanaged(*SharedWorkerGlobalScope) = .empty, + _page_destruction_queue: std.ArrayList(*Page) = .{}, // Round-robin cursor for fair page iteration (processQueuedNavigation) diff --git a/src/browser/js/Env.zig b/src/browser/js/Env.zig index b0a07709d..e0dbc5cde 100644 --- a/src/browser/js/Env.zig +++ b/src/browser/js/Env.zig @@ -31,6 +31,7 @@ const App = @import("../../App.zig"); const Frame = @import("../Frame.zig"); const Window = @import("../webapi/Window.zig"); const WorkerGlobalScope = @import("../webapi/WorkerGlobalScope.zig"); +const SharedWorkerGlobalScope = @import("../webapi/SharedWorkerGlobalScope.zig"); const DedicatedWorkerGlobalScope = @import("../webapi/DedicatedWorkerGlobalScope.zig"); const v8 = js.v8; @@ -274,8 +275,12 @@ fn _createContext(self: *Env, global: anytype, params: ContextParams) !*Context break :blk @ptrCast(v8.v8__Global__Get(existing, isolate.handle)); }; - // Restore the context from the snapshot (0 = Page, 1 = Worker) - const snapshot_index: u32 = if (comptime is_frame) 0 else 1; + // Restore the context from the snapshot + // (0 = Page, 1 = DedicatedWorker, 2 = SharedWorker) + const snapshot_index: u32 = if (comptime is_frame) 0 else switch (global._type) { + .dedicated => 1, + .shared => 2, + }; const v8_context = v8.v8__Context__FromSnapshot__Config(isolate.handle, snapshot_index, &.{ .global_template = null, .global_object = reuse_global_object, @@ -297,11 +302,19 @@ fn _createContext(self: *Env, global: anytype, params: ContextParams) !*Context .prototype_chain = (&Window.JsApi.Meta.prototype_chain).ptr, .prototype_len = @intCast(Window.JsApi.Meta.prototype_chain.len), .subtype = .node, - } else .{ - .value = @ptrCast(global._type.dedicated), - .prototype_chain = (&DedicatedWorkerGlobalScope.JsApi.Meta.prototype_chain).ptr, - .prototype_len = @intCast(DedicatedWorkerGlobalScope.JsApi.Meta.prototype_chain.len), - .subtype = null, + } else switch (global._type) { + .dedicated => |scope| .{ + .value = @ptrCast(scope), + .prototype_chain = (&DedicatedWorkerGlobalScope.JsApi.Meta.prototype_chain).ptr, + .prototype_len = @intCast(DedicatedWorkerGlobalScope.JsApi.Meta.prototype_chain.len), + .subtype = null, + }, + .shared => |scope| .{ + .value = @ptrCast(scope), + .prototype_chain = (&SharedWorkerGlobalScope.JsApi.Meta.prototype_chain).ptr, + .prototype_len = @intCast(SharedWorkerGlobalScope.JsApi.Meta.prototype_chain.len), + .subtype = null, + }, }; v8.v8__Object__SetAlignedPointerInInternalField(global_obj, 0, tao); @@ -349,7 +362,9 @@ fn _createContext(self: *Env, global: anytype, params: ContextParams) !*Context // Register in the identity map. Multiple contexts can be created for the // same global (via CDP), so we only register the first one. - const identity_ptr = if (comptime is_frame) @intFromPtr(global.window) else @intFromPtr(global._type.dedicated); + const identity_ptr = if (comptime is_frame) @intFromPtr(global.window) else switch (global._type) { + inline else => |scope| @intFromPtr(scope), + }; const gop = try params.identity.identity_map.getOrPut(params.identity_arena, identity_ptr); if (gop.found_existing == false) { var global_global: v8.Global = undefined; diff --git a/src/browser/js/Execution.zig b/src/browser/js/Execution.zig index 5efd164c6..020175adc 100644 --- a/src/browser/js/Execution.zig +++ b/src/browser/js/Execution.zig @@ -115,6 +115,14 @@ pub fn getBroadcastChannels(self: *const Execution) *std.DoublyLinkedList { }; } +// The owning global's (Frame or WGS) list of live MessagePorts, walked at +// that global's teardown to sever cross-context entanglement. +pub fn messagePorts(self: *const Execution) *std.DoublyLinkedList { + return switch (self.js.global) { + inline else => |g| &g._message_ports, + }; +} + // The global's serialized origin (e.g. "https://example.com"), or null for an // opaque origin. pub fn origin(self: *const Execution) ?[]const u8 { diff --git a/src/browser/js/Snapshot.zig b/src/browser/js/Snapshot.zig index 7441051ea..bebd89a65 100644 --- a/src/browser/js/Snapshot.zig +++ b/src/browser/js/Snapshot.zig @@ -31,7 +31,9 @@ const v8 = js.v8; const log = lp.log; const JsApis = bridge.JsApis; const PageJsApis = bridge.PageJsApis; -const WorkerJsApis = bridge.WorkerJsApis; +const SharedWorkerJsApis = bridge.SharedWorkerJsApis; +const DedicatedWorkerJsApis = bridge.DedicatedWorkerJsApis; + const IS_DEBUG = @import("builtin").mode == .Debug; const Snapshot = @This(); @@ -209,9 +211,15 @@ pub fn create() !Snapshot { { const DedicatedWorkerGlobalScope = @import("../webapi/DedicatedWorkerGlobalScope.zig"); - const index = try createSnapshotContext(.worker, &WorkerJsApis, DedicatedWorkerGlobalScope.JsApi, isolate, snapshot_creator.?, &templates); + const index = try createSnapshotContext(.worker, &DedicatedWorkerJsApis, DedicatedWorkerGlobalScope.JsApi, isolate, snapshot_creator.?, &templates); std.debug.assert(index == 1); } + + { + const SharedWorkerGlobalScope = @import("../webapi/SharedWorkerGlobalScope.zig"); + const index = try createSnapshotContext(.worker, &SharedWorkerJsApis, SharedWorkerGlobalScope.JsApi, isolate, snapshot_creator.?, &templates); + std.debug.assert(index == 2); + } } const blob = v8.v8__SnapshotCreator__createBlob(snapshot_creator, v8.kKeep); diff --git a/src/browser/js/bridge.zig b/src/browser/js/bridge.zig index 7bbd62a48..8e4a9eb1f 100644 --- a/src/browser/js/bridge.zig +++ b/src/browser/js/bridge.zig @@ -1113,6 +1113,7 @@ pub const PageJsApis = flattenTypes(&.{ @import("../webapi/MessagePort.zig"), @import("../webapi/BroadcastChannel.zig"), @import("../webapi/Worker.zig"), + @import("../webapi/SharedWorker.zig"), @import("../webapi/media/MediaError.zig"), @import("../webapi/media/TextTrackCue.zig"), @import("../webapi/media/VTTCue.zig"), @@ -1181,11 +1182,13 @@ pub const PageJsApis = flattenTypes(&.{ @import("../webapi/collections/DOMStringList.zig"), }); -// APIs available on Worker context globals (constructors like URL, Headers, etc.) +// APIs available on every Worker context global (constructors like URL, +// Headers, etc.), regardless of worker kind. Each kind's snapshot context +// adds its own global-scope type on top (DedicatedWorkerJsApis, +// SharedWorkerJsApis below). // This is a subset of PageJsApis plus WorkerGlobalScope. // TODO: Expand this list to include all worker-appropriate APIs. -pub const WorkerJsApis = flattenTypes(&.{ - @import("../webapi/DedicatedWorkerGlobalScope.zig"), +const worker_common_apis = [_]type{ @import("../webapi/WorkerGlobalScope.zig"), @import("../webapi/WorkerLocation.zig"), @import("../webapi/Navigator.zig"), @@ -1252,7 +1255,10 @@ pub const WorkerJsApis = flattenTypes(&.{ @import("../webapi/MessageChannel.zig"), @import("../webapi/MessagePort.zig"), @import("../webapi/collections/DOMStringList.zig"), -}); +}; + +pub const DedicatedWorkerJsApis = flattenTypes(&([_]type{@import("../webapi/DedicatedWorkerGlobalScope.zig")} ++ worker_common_apis)); +pub const SharedWorkerJsApis = flattenTypes(&([_]type{@import("../webapi/SharedWorkerGlobalScope.zig")} ++ worker_common_apis)); // Master list of ALL JS APIs across all contexts. // Used by Env (class IDs, templates), JsApiLookup, and anywhere that needs @@ -1261,6 +1267,7 @@ pub const WorkerJsApis = flattenTypes(&.{ pub const JsApis = blk: { const base = PageJsApis ++ [_]type{ @import("../webapi/DedicatedWorkerGlobalScope.zig").JsApi, + @import("../webapi/SharedWorkerGlobalScope.zig").JsApi, @import("../webapi/WorkerGlobalScope.zig").JsApi, @import("../webapi/WorkerLocation.zig").JsApi, }; diff --git a/src/browser/tests/shared_worker/counter-worker.js b/src/browser/tests/shared_worker/counter-worker.js new file mode 100644 index 000000000..675d68982 --- /dev/null +++ b/src/browser/tests/shared_worker/counter-worker.js @@ -0,0 +1,8 @@ +let counter = 0; +onconnect = (e) => { + const port = e.ports[0]; + port.onmessage = () => { + counter++; + port.postMessage({ count: counter, name: self.name }); + }; +}; diff --git a/src/browser/tests/shared_worker/echo-worker.js b/src/browser/tests/shared_worker/echo-worker.js new file mode 100644 index 000000000..544e66664 --- /dev/null +++ b/src/browser/tests/shared_worker/echo-worker.js @@ -0,0 +1,8 @@ +onconnect = (e) => { + // Per spec the connect event's source is the port itself (testharness.js + // relies on this rather than ports[0]). + const port = e.source; + port.onmessage = (event) => { + port.postMessage({ echo: event.data, from: 'shared-worker', sourceIsPort: e.source === e.ports[0] }); + }; +}; diff --git a/src/browser/tests/shared_worker/listener-worker.js b/src/browser/tests/shared_worker/listener-worker.js new file mode 100644 index 000000000..128d4cf9f --- /dev/null +++ b/src/browser/tests/shared_worker/listener-worker.js @@ -0,0 +1,10 @@ +// Exercises the addEventListener flavor on both the global ('connect') and +// the port ('message' + explicit start()), instead of the auto-starting +// onconnect/onmessage attribute setters. +self.addEventListener('connect', (e) => { + const port = e.ports[0]; + port.addEventListener('message', (event) => { + port.postMessage('listener:' + event.data); + }); + port.start(); +}); diff --git a/src/browser/tests/shared_worker/shared_worker.html b/src/browser/tests/shared_worker/shared_worker.html new file mode 100644 index 000000000..30757854c --- /dev/null +++ b/src/browser/tests/shared_worker/shared_worker.html @@ -0,0 +1,137 @@ + + + + + + + + + + + + + + + diff --git a/src/browser/webapi/EventTarget.zig b/src/browser/webapi/EventTarget.zig index b083969be..e1686b3f5 100644 --- a/src/browser/webapi/EventTarget.zig +++ b/src/browser/webapi/EventTarget.zig @@ -37,6 +37,7 @@ pub const Type = union(enum) { node: *@import("Node.zig"), window: *@import("Window.zig"), worker: *@import("Worker.zig"), + shared_worker: *@import("SharedWorker.zig"), worker_global_scope: *@import("WorkerGlobalScope.zig"), xhr: *@import("net/XMLHttpRequestEventTarget.zig"), abort_signal: *@import("AbortSignal.zig"), @@ -214,6 +215,7 @@ pub fn format(self: *EventTarget, writer: *std.Io.Writer) !void { .generic => writer.writeAll(""), .window => writer.writeAll(""), .worker => writer.writeAll(""), + .shared_worker => writer.writeAll(""), .worker_global_scope => writer.writeAll(""), .xhr => writer.writeAll(""), .abort_signal => writer.writeAll(""), @@ -243,6 +245,7 @@ pub fn toString(self: *EventTarget) []const u8 { .generic => return "[object EventTarget]", .window => return "[object Window]", .worker => return "[object Worker]", + .shared_worker => return "[object SharedWorker]", .worker_global_scope => return "[object WorkerGlobalScope]", .xhr => return "[object XMLHttpRequestEventTarget]", .abort_signal => return "[object AbortSignal]", diff --git a/src/browser/webapi/MessagePort.zig b/src/browser/webapi/MessagePort.zig index 54e775cde..773f869f2 100644 --- a/src/browser/webapi/MessagePort.zig +++ b/src/browser/webapi/MessagePort.zig @@ -1,4 +1,4 @@ -// Copyright (C) 2023-2025 Lightpanda (Selecy SAS) +// Copyright (C) 2023-2026 Lightpanda (Selecy SAS) // // Francis Bouvier // Pierre Tachoire @@ -16,6 +16,7 @@ // You should have received a copy of the GNU Affero General Public License // along with this program. If not, see . +const std = @import("std"); const lp = @import("lightpanda"); const js = @import("../js/js.zig"); @@ -29,16 +30,30 @@ const Execution = js.Execution; const MessagePort = @This(); _proto: *EventTarget, + +// The context this port lives in. The two ends of an entangled pair can be in +// different contexts. Delivery always goes through the *receiving* port's context. +_exec: *Execution, _enabled: bool = false, _closed: bool = false, _on_message: ?js.Function.Global = null, _on_message_error: ?js.Function.Global = null, _entangled_port: ?*MessagePort = null, +// queued message received before the port is started. +_pending: std.ArrayList(js.Value.Global) = .empty, + +// Link list into the owning Frame or WorkerGlobalScope. When the frame/WGS is +// shutdown, the port will be closed. +_node: std.DoublyLinkedList.Node = .{}, + pub fn init(exec: *Execution) !*MessagePort { - return exec._factory.eventTarget(MessagePort{ + const self = try exec._factory.eventTarget(MessagePort{ ._proto = undefined, + ._exec = exec, }); + exec.messagePorts().append(&self._node); + return self; } pub fn asEventTarget(self: *MessagePort) *EventTarget { @@ -50,7 +65,7 @@ pub fn entangle(port1: *MessagePort, port2: *MessagePort) void { port2._entangled_port = port1; } -pub fn postMessage(self: *MessagePort, message: js.Value, exec: *Execution) !void { +pub fn postMessage(self: *MessagePort, message: js.Value) !void { if (self._closed) { return; } @@ -60,13 +75,14 @@ pub fn postMessage(self: *MessagePort, message: js.Value, exec: *Execution) !voi return; } - // StructuredSerialize runs synchronously (per spec): clone the message into a - // fresh, self-owned temp now so the receiver gets an independent copy (a - // mutation on one side isn't visible to the other) and an unserializable - // value throws a DataCloneError to the caller. Mirrors Worker.postMessage. + // StructuredSerialize runs synchronously (per spec): clone the message now + // so the receiver gets an independent copy (a mutation on one side isn't + // visible to the other) and an unserializable value throws a DataCloneError + // to the caller. The clone is made directly into the receiving port's + // context, which may not be the sender's. const cloned = blk: { var ls: js.Local.Scope = undefined; - exec.js.localScope(&ls); + other._exec.js.localScope(&ls); defer ls.deinit(); // Contain any V8 exception from a failed serialization so it surfaces as @@ -82,29 +98,40 @@ pub fn postMessage(self: *MessagePort, message: js.Value, exec: *Execution) !voi }; errdefer cloned.release(); - // Create callback to deliver message - const callback = try exec._factory.create(PostMessageCallback{ - .exec = exec, - .port = other, - .message = cloned, - }); + if (!other._enabled) { + try other._pending.append(other._exec.arena, cloned); + return; + } - try exec.js.scheduler.add(callback, PostMessageCallback.run, 0, .{ - .name = "MessagePort.postMessage", - .low_priority = false, - .finalizer = PostMessageCallback.cancelled, - }); + try other.scheduleDelivery(cloned); } pub fn start(self: *MessagePort) void { - if (self._closed) { + if (self._closed or self._enabled) { return; } self._enabled = true; + + for (self._pending.items) |message| { + self.scheduleDelivery(message) catch |err| { + log.warn(.dom, "MessagePort.start drain", .{ .err = err }); + message.release(); + }; + } + self._pending.clearRetainingCapacity(); } pub fn close(self: *MessagePort) void { + if (self._closed) { + return; + } self._closed = true; + self._exec.messagePorts().remove(&self._node); + + for (self._pending.items) |message| { + message.release(); + } + self._pending.clearRetainingCapacity(); // Break entanglement if (self._entangled_port) |other| { @@ -119,6 +146,9 @@ pub fn getOnMessage(self: *const MessagePort) ?js.Function.Global { pub fn setOnMessage(self: *MessagePort, cb: ?js.Function.Global) !void { self._on_message = cb; + if (cb != null) { + self.start(); + } } pub fn getOnMessageError(self: *const MessagePort) ?js.Function.Global { @@ -129,38 +159,55 @@ pub fn setOnMessageError(self: *MessagePort, cb: ?js.Function.Global) !void { self._on_message_error = cb; } -const PostMessageCallback = struct { +// Queues delivery of `message` (a clone already living in this port's +// context) on this port's scheduler. +fn scheduleDelivery(self: *MessagePort, message: js.Value.Global) !void { + const exec = self._exec; + const callback = try exec._factory.create(DeliverCallback{ + .port = self, + .message = message, + }); + errdefer exec._factory.destroy(callback); + + try exec.js.scheduler.add(callback, DeliverCallback.run, 0, .{ + .name = "MessagePort.postMessage", + .low_priority = false, + .finalizer = DeliverCallback.cancelled, + }); +} + +const DeliverCallback = struct { port: *MessagePort, message: js.Value.Global, - exec: *Execution, // Called by the scheduler if the task is dropped before it runs. `run` and // `cancelled` are mutually exclusive, so the temp is released exactly once. fn cancelled(ctx: *anyopaque) void { - const self: *PostMessageCallback = @ptrCast(@alignCast(ctx)); + const self: *DeliverCallback = @ptrCast(@alignCast(ctx)); self.message.release(); self.deinit(); } - fn deinit(self: *PostMessageCallback) void { - self.exec._factory.destroy(self); + fn deinit(self: *DeliverCallback) void { + self.port._exec._factory.destroy(self); } fn run(ctx: *anyopaque) !?u32 { - const self: *PostMessageCallback = @ptrCast(@alignCast(ctx)); + const self: *DeliverCallback = @ptrCast(@alignCast(ctx)); defer self.deinit(); - const exec = self.exec; + const port = self.port; + const exec = port._exec; // The MessageEvent takes ownership of the cloned temp and releases it on // teardown; on any path where we don't hand it over, release it here so // it doesn't leak. - if (self.port._closed) { + if (port._closed) { self.message.release(); return null; } - const target = self.port.asEventTarget(); - if (!exec.hasDirectListeners(target, "message", self.port._on_message)) { + const target = port.asEventTarget(); + if (!exec.hasDirectListeners(target, "message", port._on_message)) { self.message.release(); return null; } @@ -175,7 +222,7 @@ const PostMessageCallback = struct { return null; }).asEvent(); - exec.dispatch(target, event, self.port._on_message, .{ .context = "MessagePort message" }) catch |err| { + exec.dispatch(target, event, port._on_message, .{ .context = "MessagePort message" }) catch |err| { log.err(.dom, "MessagePort.postMessage", .{ .err = err }); }; diff --git a/src/browser/webapi/SharedWorker.zig b/src/browser/webapi/SharedWorker.zig new file mode 100644 index 000000000..991175480 --- /dev/null +++ b/src/browser/webapi/SharedWorker.zig @@ -0,0 +1,146 @@ +// Copyright (C) 2023-2026 Lightpanda (Selecy SAS) +// +// Francis Bouvier +// Pierre Tachoire +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as +// published by the Free Software Foundation, either version 3 of the +// License, or (at your option) any later version. +// +// This program is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +// The page-side handle for a shared worker. Unlike Worker, this owns nothing: +// the worker itself (SharedWorkerGlobalScope) is owned by the page that first +// created it and found through the Session's registry, keyed by +// (resolved url, name) — so every `new SharedWorker(url, name)` in the +// session, from any frame or page, talks to the same instance. All this +// handle keeps is its end of the connection's MessagePort pair. +// +// Divergence from the spec: the worker's lifetime is tied to its creating +// page, not to the set of connected clients. If the creator dies while +// another page is still connected, the worker dies and the surviving port +// goes quiet (entanglement is severed; nothing dangles). + +const std = @import("std"); +const lp = @import("lightpanda"); + +const js = @import("../js/js.zig"); + +const URL = @import("../URL.zig"); +const Frame = @import("../Frame.zig"); + +const Worker = @import("Worker.zig"); +const EventTarget = @import("EventTarget.zig"); +const MessagePort = @import("MessagePort.zig"); +const SharedWorkerGlobalScope = @import("SharedWorkerGlobalScope.zig"); + +const log = lp.log; + +const SharedWorker = @This(); + +_proto: *EventTarget, +_port: *MessagePort, +_on_error: ?js.Function.Global = null, + +const NameOrOpts = union(enum) { + name: []const u8, + options: Opts, + + const Opts = struct { + name: []const u8 = "", + type: Worker.WorkerType = .classic, + }; +}; + +pub fn init(url: []const u8, name_or_options: ?NameOrOpts, frame: *Frame) !*SharedWorker { + const options: NameOrOpts.Opts = if (name_or_options) |noo| switch (noo) { + .name => |n| .{ .name = n }, + .options => |o| o, + } else .{}; + + const resolved_url = try URL.resolve(frame.call_arena, frame.base(), url, .{ .encoding = frame.charset }); + + const scope = blk: { + const session = frame._session; + // \x00 can appear in neither a URL nor a name + const lookup_key = try std.fmt.allocPrint(frame.call_arena, "{s}\x00{s}", .{ resolved_url, options.name }); + if (session.shared_workers.get(lookup_key)) |existing| { + break :blk existing; + } + + const s = try SharedWorkerGlobalScope.init(frame, resolved_url, options.name, options.type); + errdefer s.deinit(); + + const page = frame._page; + try page.shared_workers.append(page.frame_arena, s); + errdefer _ = page.shared_workers.pop(); + + try s.register(lookup_key); + + break :blk s; + }; + + const port = try scope.connect(&frame.js.execution); + return frame._page.factory.eventTarget(SharedWorker{ + ._proto = undefined, + ._port = port, + }); +} + +pub fn asEventTarget(self: *SharedWorker) *EventTarget { + return self._proto; +} + +pub fn getPort(self: *const SharedWorker) *MessagePort { + return self._port; +} + +pub fn getOnError(self: *const SharedWorker) ?js.Function.Global { + return self._on_error; +} + +pub fn setOnError(self: *SharedWorker, setter: ?FunctionSetter) void { + self._on_error = getFunctionFromSetter(setter); +} + +const FunctionSetter = union(enum) { + func: js.Function.Global, + anything: js.Value, +}; + +fn getFunctionFromSetter(setter_: ?FunctionSetter) ?js.Function.Global { + const setter = setter_ orelse return null; + return switch (setter) { + .func => |func| func, + .anything => null, + }; +} + +pub const JsApi = struct { + pub const bridge = js.Bridge(SharedWorker); + + pub const Meta = struct { + pub const name = "SharedWorker"; + pub const prototype_chain = bridge.prototypeChain(); + pub var class_id: bridge.ClassId = undefined; + }; + + pub const constructor = bridge.constructor(SharedWorker.init, .{}); + + pub const port = bridge.accessor(SharedWorker.getPort, null, .{}); + pub const onerror = bridge.accessor(SharedWorker.getOnError, SharedWorker.setOnError, .{}); +}; + +const testing = @import("../../testing.zig"); +test "WebApi: SharedWorker" { + const filter: testing.LogFilter = .init(&.{.http}); + defer filter.deinit(); + try testing.htmlRunner("shared_worker", .{ .timeout_ms = 8000 }); +} diff --git a/src/browser/webapi/SharedWorkerGlobalScope.zig b/src/browser/webapi/SharedWorkerGlobalScope.zig new file mode 100644 index 000000000..bf3a6719f --- /dev/null +++ b/src/browser/webapi/SharedWorkerGlobalScope.zig @@ -0,0 +1,407 @@ +// Copyright (C) 2023-2026 Lightpanda (Selecy SAS) +// +// Francis Bouvier +// Pierre Tachoire +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as +// published by the Free Software Foundation, either version 3 of the +// License, or (at your option) any later version. +// +// This program is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +const std = @import("std"); +const lp = @import("lightpanda"); + +const js = @import("../js/js.zig"); + +const URL = @import("../URL.zig"); +const Frame = @import("../Frame.zig"); +const Transfer = @import("../../network/HttpClient.zig").Transfer; + +const Worker = @import("Worker.zig"); +const MessagePort = @import("MessagePort.zig"); +const WorkerGlobalScope = @import("WorkerGlobalScope.zig"); +const MessageEvent = @import("event/MessageEvent.zig"); + +const log = lp.log; +const Allocator = std.mem.Allocator; +const IS_DEBUG = @import("builtin").mode == .Debug; + +const SharedWorkerGlobalScope = @This(); + +_proto: *WorkerGlobalScope, +_arena: Allocator, + +_url: [:0]const u8, +_name: []const u8, +_type: Worker.WorkerType, + +// Key under which this scope is registered in session.shared_workers. Empty +// until registered; deinit/close use it to unregister. +_registry_key: []const u8 = "", + +// used by HttpClient when generating notification +// Ultimately used by CDP to generate request/loader ids. +_frame_id: u32, +_loader_id: u32, + +_closed: bool = false, +_script_loaded: bool = false, +_script_arena: ?Allocator = null, +_script_buffer: std.ArrayList(u8) = .empty, +_http_transfer: ?*Transfer = null, + +_on_connect: ?js.Function.Global = null, + +// Ports pending connnection, created before the initial script finished +_pending_connects: std.ArrayList(*MessagePort) = .empty, + +pub fn init(frame: *Frame, url: [:0]const u8, name: []const u8, worker_type: Worker.WorkerType) !*SharedWorkerGlobalScope { + const session = frame._session; + + const arena = try session.getArena(.small, "SharedWorker"); + errdefer session.releaseArena(arena); + + const owned_url = try arena.dupeZ(u8, url); + const self = try arena.create(SharedWorkerGlobalScope); + const proto = try WorkerGlobalScope.init( + arena, + owned_url, + .{ .shared = self }, + worker_type == .module, + session.nextFrameId(), + session.nextLoaderId(), + frame, + ); + self.* = .{ + ._proto = proto, + ._arena = arena, + ._url = owned_url, + ._name = try arena.dupe(u8, name), + ._type = worker_type, + ._frame_id = proto._frame_id, + ._loader_id = proto._loader_id, + }; + errdefer proto.deinit(); + + if (!session.worker_loading_enabled) { + log.debug(.browser, "shared worker disabled", .{ .url = owned_url }); + return self; + } + + self._script_arena = try session.getArena(.large, "SharedWorker.script"); + errdefer self.releaseScriptArena(); + + const transfer = proto.newRequest(.{ + .ctx = self, + .method = .GET, + .url = owned_url, + .frame_id = self._frame_id, + .loader_id = self._loader_id, + .resource_type = .script, + .cookie_jar = &session.cookie_jar, + .cookie_origin = owned_url, + .notification = session.notification, + .header_callback = httpHeaderCallback, + .data_callback = httpDataCallback, + .done_callback = httpDoneCallback, + .error_callback = httpErrorCallback, + .shutdown_callback = httpShutdownCallback, + }) catch |err| { + log.err(.browser, "SharedWorker request", .{ .url = owned_url, .err = err }); + return err; + }; + + // Held for deinit's abort; the done, error and shutdown callbacks clear + // it. The shutdown one matters: WorkerGlobalScope.deinit's abortOwner + // kills the transfer, and deinit must not abort a freed transfer. + self._http_transfer = transfer; + + transfer.submit() catch |err| { + log.err(.browser, "SharedWorker request", .{ .url = owned_url, .err = err }); + return err; + }; + return self; +} + +// Called from Page.deinit of the owning (creating) page. +pub fn deinit(self: *SharedWorkerGlobalScope) void { + if (self._http_transfer) |transfer| { + transfer.abort(error.Abort); + self._http_transfer = null; + } + self.releaseScriptArena(); + self.unregister(); + self._proto.deinit(); + self._proto._session.releaseArena(self._arena); +} + +pub fn register(self: *SharedWorkerGlobalScope, lookup_key: []const u8) !void { + // The key lives in our arena: every removal path (deinit, close) unregisters + // before the arena is released. + const key = try self._arena.dupe(u8, lookup_key); + + const session = self._proto._session; + try session.shared_workers.put(session.arena, key, self); + self._registry_key = key; +} + +fn unregister(self: *SharedWorkerGlobalScope) void { + if (self._registry_key.len == 0) { + return; + } + _ = self._proto._session.shared_workers.remove(self._registry_key); + self._registry_key = ""; +} + +// Establishes a new connection from a client context: creates the entangled +// port pair and queues the connect event. Returns the client's end. +pub fn connect(self: *SharedWorkerGlobalScope, client_exec: *js.Execution) !*MessagePort { + const client_port = try MessagePort.init(client_exec); + const worker_port = try MessagePort.init(&self._proto.js.execution); + MessagePort.entangle(client_port, worker_port); + + if (self._script_loaded == false) { + try self._pending_connects.append(self._arena, worker_port); + return client_port; + } + + try self.scheduleConnect(worker_port); + return client_port; +} + +pub fn getName(self: *const SharedWorkerGlobalScope) []const u8 { + return self._name; +} + +pub fn close(self: *SharedWorkerGlobalScope) void { + // Once closed, new SharedWorker(url, name) must create a fresh instance. + self.unregister(); + // TODO: we should also stop new tasks from being scheduled + self._proto._session.idb.detachContext(self._proto.js); + self._proto.js.scheduler.reset(); + self._closed = true; +} + +pub fn getOnConnect(self: *const SharedWorkerGlobalScope) ?js.Function.Global { + return self._on_connect; +} + +pub fn setOnConnect(self: *SharedWorkerGlobalScope, setter: ?WorkerGlobalScope.FunctionSetter) void { + self._on_connect = WorkerGlobalScope.getFunctionFromSetter(setter); +} + +fn httpHeaderCallback(transfer: *Transfer) !Transfer.HeaderResult { + const self: *SharedWorkerGlobalScope = @ptrCast(@alignCast(transfer.req.ctx)); + + const status = transfer.responseStatus() orelse return .abort; + if (status < 200 or status >= 300) { + log.warn(.browser, "SharedWorker status", .{ + .url = self._url, + .status = status, + }); + return .abort; + } + + if (transfer.getContentLength()) |cl| { + try self._script_buffer.ensureTotalCapacity(self._script_arena.?, cl); + } + + return .proceed; +} + +fn httpDataCallback(transfer: *Transfer, data: []const u8) !void { + const self: *SharedWorkerGlobalScope = @ptrCast(@alignCast(transfer.req.ctx)); + try self._script_buffer.appendSlice(self._script_arena.?, data); +} + +fn httpDoneCallback(ctx: *anyopaque) !void { + const self: *SharedWorkerGlobalScope = @ptrCast(@alignCast(ctx)); + self._http_transfer = null; + defer self.releaseScriptArena(); + + const url = self._url; + const script = self._script_buffer.items; + + if (comptime IS_DEBUG) { + log.info(.browser, "shared worker fetch done", .{ + .url = url, + .len = script.len, + }); + } + + try self.loadInitialScript(script); +} + +fn httpShutdownCallback(ctx: *anyopaque) void { + const self: *SharedWorkerGlobalScope = @ptrCast(@alignCast(ctx)); + self._http_transfer = null; + self.releaseScriptArena(); +} + +fn httpErrorCallback(ctx: *anyopaque, err: anyerror) void { + const self: *SharedWorkerGlobalScope = @ptrCast(@alignCast(ctx)); + self._http_transfer = null; + self.releaseScriptArena(); + + log.err(.browser, "shared worker fetch error", .{ + .url = self._url, + .err = err, + }); + + // The worker will never load and onconnect will never be registered. + // Drain the buffered connects so they get dispatched (and dropped at the + // "no listener" check) rather than accumulating until teardown. + self._script_loaded = true; + self.drainPendingConnects(); +} + +fn loadInitialScript(self: *SharedWorkerGlobalScope, script: []const u8) !void { + const js_context = self._proto.js; + + if (js_context.env.terminatePending()) { + return; + } + + // The flip-and-drain runs after eval (and the trailing runMacrotasks) — + // by which point the outer script has had its only chance to register + // onconnect. On eval-throw the defer still fires; the connect events get + // scheduled and then drop at the "no listener" check. + defer { + self._script_loaded = true; + self.drainPendingConnects(); + } + + var ls: js.Local.Scope = undefined; + js_context.localScope(&ls); + defer ls.deinit(); + + var try_catch: js.TryCatch = undefined; + try_catch.init(&ls.local); + defer try_catch.deinit(); + + switch (self._type) { + .classic => _ = ls.local.eval(script, self._url) catch |err| { + if (js_context.env.terminatePending()) { + return; + } + + const caught = try_catch.caughtOrError(self._script_arena.?, err); + log.err(.browser, "shared worker script error", .{ .url = self._url, .caught = caught }); + return; + }, + .module => js_context.module(false, &ls.local, script, self._url, true) catch |err| { + if (js_context.env.terminatePending()) { + return; + } + + const caught = try_catch.caughtOrError(self._script_arena.?, err); + log.err(.browser, "shared worker module error", .{ .url = self._url, .caught = caught }); + return; + }, + } + + ls.local.runMacrotasks(); +} + +// Idempotent: reached from the done, error and shutdown callbacks and from +// deinit (the abort there re-enters via the error callback). +fn releaseScriptArena(self: *SharedWorkerGlobalScope) void { + const arena = self._script_arena orelse return; + self._script_arena = null; + self._script_buffer = .empty; + self._proto._session.releaseArena(arena); +} + +fn drainPendingConnects(self: *SharedWorkerGlobalScope) void { + for (self._pending_connects.items) |port| { + self.scheduleConnect(port) catch |err| { + log.warn(.browser, "shared worker drain connect failed", .{ .err = err }); + }; + } + self._pending_connects.clearRetainingCapacity(); +} + +fn scheduleConnect(self: *SharedWorkerGlobalScope, port: *MessagePort) !void { + const wgs = self._proto; + const session = wgs._session; + + const connect_arena = try session.getArena(.tiny, "SharedWorkerGlobalScope.connect"); + errdefer session.releaseArena(connect_arena); + + const callback = try connect_arena.create(ConnectCallback); + callback.* = .{ + .port = port, + .worker_scope = self, + .arena = connect_arena, + }; + + try wgs.js.scheduler.add(callback, ConnectCallback.run, 0, .{ + .name = "SharedWorkerGlobalScope.connect", + .low_priority = false, + .finalizer = ConnectCallback.cancelled, + }); +} + +const ConnectCallback = struct { + port: *MessagePort, + arena: Allocator, + worker_scope: *SharedWorkerGlobalScope, + + fn cancelled(ctx: *anyopaque) void { + const self: *ConnectCallback = @ptrCast(@alignCast(ctx)); + self.deinit(); + } + + fn deinit(self: *ConnectCallback) void { + self.worker_scope._proto._session.releaseArena(self.arena); + } + + fn run(ctx: *anyopaque) !?u32 { + const self: *ConnectCallback = @ptrCast(@alignCast(ctx)); + defer self.deinit(); + + const worker_scope = self.worker_scope; + const wgs = worker_scope._proto; + const target = wgs.asEventTarget(); + + const on_connect = worker_scope._on_connect; + if (!wgs._event_manager.hasDirectListeners(target, "connect", on_connect)) { + return null; + } + + // Per spec the connect event's data is the empty string and the new + // port is both the sole entry of `ports` and the `source`. + const event = (try MessageEvent.initTrusted(comptime .wrap("connect"), .{ + .data = .{ .string = "" }, + .source = .{ .port = self.port }, + .ports = &.{self.port}, + .bubbles = false, + .cancelable = false, + }, wgs._page)).asEvent(); + + try wgs.dispatch(target, event, on_connect, .{ .context = "SharedWorkerGlobalScope.connect" }); + return null; + } +}; + +pub const JsApi = struct { + pub const bridge = js.Bridge(SharedWorkerGlobalScope); + + pub const Meta = struct { + pub const name = "SharedWorkerGlobalScope"; + pub const prototype_chain = bridge.prototypeChain(); + pub var class_id: bridge.ClassId = undefined; + }; + + pub const name = bridge.accessor(SharedWorkerGlobalScope.getName, null, .{}); + pub const close = bridge.function(SharedWorkerGlobalScope.close, .{}); + pub const onconnect = bridge.accessor(SharedWorkerGlobalScope.getOnConnect, SharedWorkerGlobalScope.setOnConnect, .{}); +}; diff --git a/src/browser/webapi/Window.zig b/src/browser/webapi/Window.zig index c42f1da6a..c44648e4e 100644 --- a/src/browser/webapi/Window.zig +++ b/src/browser/webapi/Window.zig @@ -1093,7 +1093,7 @@ const PostMessageCallback = struct { const event = (try MessageEvent.initTrusted(comptime .wrap("message"), .{ .data = .{ .value = self.message }, .origin = self.origin, - .source = self.source, + .source = .{ .window = self.source }, .ports = self.ports, .bubbles = false, .cancelable = false, diff --git a/src/browser/webapi/Worker.zig b/src/browser/webapi/Worker.zig index bbcbe600a..61fe9a1b4 100644 --- a/src/browser/webapi/Worker.zig +++ b/src/browser/webapi/Worker.zig @@ -55,6 +55,7 @@ _worker_scope: *DedicatedWorkerGlobalScope, _url: [:0]const u8, _type: WorkerType = .classic, _script_loaded: bool = false, +_script_arena: ?Allocator = null, _script_buffer: std.ArrayList(u8) = .empty, _http_transfer: ?*Transfer = null, @@ -70,7 +71,7 @@ const WorkerOptions = struct { pub fn init(url: []const u8, options: ?WorkerOptions, frame: *Frame) !*Worker { const session = frame._session; - const arena = try session.getArena(.large, "Worker"); + const arena = try session.getArena(.small, "Worker"); errdefer session.releaseArena(arena); const resolved_url = try URL.resolve(arena, frame.base(), url, .{ .encoding = frame.charset }); @@ -101,6 +102,9 @@ pub fn init(url: []const u8, options: ?WorkerOptions, frame: *Frame) !*Worker { return self; } + self._script_arena = try session.getArena(.large, "Worker.script"); + errdefer self.releaseScriptArena(); + const transfer = frame.newRequest(.{ .ctx = self, .method = .GET, @@ -144,6 +148,7 @@ pub fn deinit(self: *Worker) void { res.abort(error.Abort); self._http_transfer = null; } + self.releaseScriptArena(); self._worker_scope.deinit(); self._frame._session.releaseArena(self._arena); } @@ -165,7 +170,7 @@ fn httpHeaderCallback(transfer: *Transfer) !Transfer.HeaderResult { } if (transfer.getContentLength()) |cl| { - try self._script_buffer.ensureTotalCapacity(self._arena, cl); + try self._script_buffer.ensureTotalCapacity(self._script_arena.?, cl); } return .proceed; @@ -173,12 +178,13 @@ fn httpHeaderCallback(transfer: *Transfer) !Transfer.HeaderResult { fn httpDataCallback(transfer: *Transfer, data: []const u8) !void { const self: *Worker = @ptrCast(@alignCast(transfer.req.ctx)); - try self._script_buffer.appendSlice(self._arena, data); + try self._script_buffer.appendSlice(self._script_arena.?, data); } fn httpDoneCallback(ctx: *anyopaque) !void { const self: *Worker = @ptrCast(@alignCast(ctx)); self._http_transfer = null; + defer self.releaseScriptArena(); const url = self._url; const script = self._script_buffer.items; @@ -235,7 +241,7 @@ fn loadInitialScript(self: *Worker, script: []const u8) !void { return; } - const caught = try_catch.caughtOrError(self._arena, err); + const caught = try_catch.caughtOrError(self._script_arena.?, err); log.err(.browser, "worker script error", .{ .url = self._url, .caught = caught }); self.fireErrorEvent(caught.exception orelse @errorName(err), null); return; @@ -245,7 +251,7 @@ fn loadInitialScript(self: *Worker, script: []const u8) !void { return; } - const caught = try_catch.caughtOrError(self._arena, err); + const caught = try_catch.caughtOrError(self._script_arena.?, err); log.err(.browser, "worker module error", .{ .url = self._url, .caught = caught }); self.fireErrorEvent(caught.exception orelse @errorName(err), null); return; @@ -258,11 +264,13 @@ fn loadInitialScript(self: *Worker, script: []const u8) !void { fn httpShutdownCallback(ctx: *anyopaque) void { const self: *Worker = @ptrCast(@alignCast(ctx)); self._http_transfer = null; + self.releaseScriptArena(); } fn httpErrorCallback(ctx: *anyopaque, err: anyerror) void { const self: *Worker = @ptrCast(@alignCast(ctx)); self._http_transfer = null; + self.releaseScriptArena(); log.err(.browser, "worker fetch error", .{ .url = self._url, @@ -279,6 +287,13 @@ fn httpErrorCallback(ctx: *anyopaque, err: anyerror) void { self.fireErrorEvent(@errorName(err), null); } +fn releaseScriptArena(self: *Worker) void { + const arena = self._script_arena orelse return; + self._script_arena = null; + self._script_buffer = .empty; + self._frame._session.releaseArena(arena); +} + // Fire an error event on the Worker object (parent context) fn fireErrorEvent(self: *Worker, message: []const u8, error_value: ?js.Value.Global) void { self._fireErrorEvent(message, error_value) catch |err| { diff --git a/src/browser/webapi/WorkerGlobalScope.zig b/src/browser/webapi/WorkerGlobalScope.zig index 2bd8c352e..85c707ed3 100644 --- a/src/browser/webapi/WorkerGlobalScope.zig +++ b/src/browser/webapi/WorkerGlobalScope.zig @@ -46,6 +46,8 @@ const ErrorEvent = @import("event/ErrorEvent.zig"); const Fetch = @import("net/Fetch.zig"); const idb = @import("storage/idb/idb.zig"); const CookieStore = @import("storage/CookieStore.zig"); +const MessagePort = @import("MessagePort.zig"); +const SharedWorkerGlobalScope = @import("SharedWorkerGlobalScope.zig"); const DedicatedWorkerGlobalScope = @import("DedicatedWorkerGlobalScope.zig"); const builtin = @import("builtin"); @@ -98,6 +100,9 @@ _script_manager: ScriptManagerBase, // channels in this worker's origin _broadcast_channels: std.DoublyLinkedList = .{}, +// List of MessagePorts living in this worker's context. +_message_ports: std.DoublyLinkedList = .{}, + // These fields represent the "Window"-like component of the WGS _proto: *EventTarget, _console: Console = .init, @@ -115,6 +120,7 @@ _location: WorkerLocation, _timers: Timers = .{}, pub const Type = union(enum) { + shared: *SharedWorkerGlobalScope, dedicated: *DedicatedWorkerGlobalScope, }; @@ -190,6 +196,14 @@ pub fn deinit(self: *WorkerGlobalScope) void { browser.http_client.abortOwner(&self._http_owner); + // Close this worker's MessagePorts before the context dies: this severs + // entanglement with page-side ports (which may outlive us) and releases + // any still-queued messages. + while (self._message_ports.first) |node| { + const port: *MessagePort = @alignCast(@fieldParentPtr("_node", node)); + port.close(); // removes from self._message_ports + } + self._identity.deinit(); self._script_manager.deinit(); diff --git a/src/browser/webapi/event/MessageEvent.zig b/src/browser/webapi/event/MessageEvent.zig index 02fb1fe83..4631384b0 100644 --- a/src/browser/webapi/event/MessageEvent.zig +++ b/src/browser/webapi/event/MessageEvent.zig @@ -36,14 +36,19 @@ _proto: *Event, _data: ?Data = null, _origin: []const u8 = "", _last_event_id: []const u8 = "", -_source: ?*Window = null, +_source: ?Source = null, _ports: []const *MessagePort = &.{}, +pub const Source = union(enum) { + window: *Window, + port: *MessagePort, +}; + const MessageEventOptions = struct { data: ?Data = null, origin: ?[]const u8 = null, lastEventId: ?[]const u8 = null, - source: ?*Window = null, + source: ?Source = null, ports: []const *MessagePort = &.{}, }; @@ -124,18 +129,24 @@ pub fn getLastEventId(self: *const MessageEvent) []const u8 { return self._last_event_id; } -pub fn getSource(self: *const MessageEvent, exec: *js.Execution) ?Window.Access { - switch (exec.js.global) { - .frame => |frame| { - const source = self._source orelse return null; - return Window.Access.init(frame.window, source); - }, - .worker => { - // source for worker should always be null - if (comptime IS_DEBUG) { - std.debug.assert(self._source == null); - } - return null; +const SourceAccess = union(enum) { + window: Window.Access, + port: *MessagePort, +}; + +pub fn getSource(self: *const MessageEvent, exec: *js.Execution) ?SourceAccess { + const source = self._source orelse return null; + switch (source) { + .port => |port| return .{ .port = port }, + .window => |window| switch (exec.js.global) { + .frame => |frame| return .{ .window = Window.Access.init(frame.window, window) }, + .worker => { + // a window source should never reach a worker context + if (comptime IS_DEBUG) { + std.debug.assert(false); + } + return null; + }, }, } }