Skip to content
File

Blob: src/workerd/api/messagechannel.c++

7.1 KB
1#include "messagechannel.h"
2 
3#include "events.h"
4 
5#include <workerd/io/worker.h>
6#include <workerd/jsg/ser.h>
7#include <workerd/util/weak-refs.h>
8 
9namespace workerd::api {
10MessagePort::MessagePort()
11 : weakThis(kj::refcounted<WeakRef<MessagePort>>(kj::Badge<MessagePort>{}, *this)),
12 state(Pending()) {
13 // We set a callback on the underlying EventTarget to be notified when
14 // a listener for the message event is added or removed. When there
15 // are no listeners, we move back to the Pending state, otherwise we
16 // will switch to the Started state if necessary.
17 setEventListenerCallback([&](jsg::Lock& js, kj::StringPtr name, size_t count) {
18 if (name == "message"_kj) {
19 KJ_SWITCH_ONEOF(state) {
20 KJ_CASE_ONEOF(pending, Pending) {
21 // If we are in the pending state, start the port if we have listeners.
22 // This is technically not spec compliant, but it is what Node.js
23 // supports. Specifically, adding a new message listener using the
24 // addEventListener method is *technically* not supposed to start
25 // the port but we're going to do what Node.js does.
26 if (count > 0 || onmessageValue != kj::none) {
27 start(js);
28 }
29 }
30 KJ_CASE_ONEOF(started, Started) {
31 // If we are in the started state, stop the port if there are no listeners.
32 if (count == 0 && onmessageValue == kj::none) {
33 state = Pending();
34 }
35 }
36 KJ_CASE_ONEOF(_, Closed) {
37 // Nothing to do. We're already closed so we don't care.
38 }
39 }
40 }
41 });
42}
43 
44void MessagePort::dispatchMessage(jsg::Lock& js, const jsg::JsValue& value) {
45 JSG_TRY(js) {
46 auto message = js.alloc<MessageEvent>(js, kj::str("message"), value, kj::String(), JSG_THIS);
47 dispatchEventImpl(js, kj::mv(message));
48 }
49 JSG_CATCH(exception) {
50 // There was an error dispatching the message event.
51 // We will dispatch a messageerror event instead.
52 auto message = js.alloc<MessageEvent>(
53 js, kj::str("message"), jsg::JsValue(exception.getHandle(js)), kj::String(), JSG_THIS);
54 dispatchEventImpl(js, kj::mv(message));
55 // Now, if this dispatchEventImpl throws, we just blow up. Don't try to catch it.
56 }
57}
58 
59// Deliver the message to this port, buffering if necessary if the port
60// has not been started. Buffered messages will be delivered when the
61// port is started later.
62void MessagePort::deliver(jsg::Lock& js, const jsg::JsValue& value) {
63 KJ_SWITCH_ONEOF(state) {
64 KJ_CASE_ONEOF(pending, Pending) {
65 // We have not yet started the port so buffer the message.
66 // It will be delivered when the port is started.
67 // We don't know how many messages will be buffered, if any,
68 // so we avoid reserving space in the array.
69 pending.add(jsg::JsRef(js, value));
70 }
71 KJ_CASE_ONEOF(started, Started) {
72 js.resolvedPromise().then(
73 js, [self = JSG_THIS, value = jsg::JsRef(js, value)](jsg::Lock& js) mutable {
74 self->dispatchMessage(js, value.getHandle(js));
75 });
76 }
77 KJ_CASE_ONEOF(_, Closed) {
78 // Nothing to do in this case. Drop the message on the floor.
79 }
80 }
81}
82 
83// Binds two ports to each other such that messages posted to one
84// are delivered on the other.
85void MessagePort::entangle(MessagePort& port1, MessagePort& port2) {
86 port1.other = port2.addWeakRef();
87 port2.other = port1.addWeakRef();
88}
89 
90// Post a message to the entangled port.
91void MessagePort::postMessage(jsg::Lock& js,
92 jsg::Optional<jsg::JsRef<jsg::JsValue>> data,
93 jsg::Optional<TransferListOrOptions> options) {
94 
95 // We don't currently support transfer lists, even for local
96 // same-isolate delivery.
97 // TODO(conform): Implement transfer later?
98 bool hasTransfer = false;
99 KJ_SWITCH_ONEOF(kj::mv(options).orDefault(PostMessageOptions{})) {
100 KJ_CASE_ONEOF(list, TransferList) {
101 hasTransfer = list.size() > 0;
102 }
103 KJ_CASE_ONEOF(opts, PostMessageOptions) {
104 KJ_IF_SOME(list, opts.transfer) {
105 hasTransfer = list.size() > 0;
106 }
107 }
108 }
109 JSG_REQUIRE(!hasTransfer, Error, "Transfer list is not supported");
110 
111 // If the port is closed, other will be kj::none and we will just drop the message.
112 other->runIfAlive([&](MessagePort& o) {
113 jsg::Serializer ser(js);
114 
115 KJ_IF_SOME(d, data) {
116 ser.write(js, d.getHandle(js));
117 } else {
118 ser.write(js, js.undefined());
119 }
120 
121 auto released = ser.release();
122 JSG_REQUIRE(released.sharedArrayBuffers.size() == 0, TypeError,
123 "SharedArrayBuffer is unsupported with MessagePort");
124 
125 // Now, deserialize the message into a JsValue
126 jsg::Deserializer deserializer(js, released);
127 auto clonedData = deserializer.readValue(js);
128 o.deliver(js, clonedData);
129 });
130}
131 
132void MessagePort::closeImpl() {
133 // Any pending messages will be dropped on the floor, except for those that were
134 // already scheduled for delivery in the `start()` or `deliver()` methods.
135 if (state.is<Closed>()) return;
136 state = Closed{};
137 weakThis->invalidate();
138 other->runIfAlive([&](MessagePort& o) { o.closeImpl(); });
139}
140 
141void MessagePort::close(jsg::Lock& js) {
142 if (state.is<Closed>()) return;
143 state = Closed{};
144 weakThis->invalidate();
145 other->runIfAlive([&](MessagePort& o) { o.close(js); });
146 auto closeEvent = js.alloc<Event>(kj::str("close"), Event::Init{}, true);
147 dispatchEventImpl(js, kj::mv(closeEvent));
148}
149 
150// Start delivering messages on this port. Any messages that are
151// buffered will be drained immediately.
152void MessagePort::start(jsg::Lock& js) {
153 KJ_SWITCH_ONEOF(state) {
154 KJ_CASE_ONEOF(pending, Pending) {
155 auto list = kj::mv(pending);
156 state = Started{};
157 // We're going to dispatch the messages using a microtask so that the actual
158 // delivery is deferred to match Node.js' behavior as close as possible.
159 js.resolvedPromise().then(js, [list = kj::mv(list), self = JSG_THIS](jsg::Lock& js) mutable {
160 for (auto& item: list) {
161 self->dispatchMessage(js, item.getHandle(js));
162 }
163 });
164 }
165 KJ_CASE_ONEOF(_, Started) {
166 // Nothing to do in this case. We are already started!
167 }
168 KJ_CASE_ONEOF(_, Closed) {
169 // Nothing to do in this case. Can't start after closing.
170 }
171 }
172}
173 
174kj::Maybe<jsg::JsValue> MessagePort::getOnMessage(jsg::Lock& js) {
175 return onmessageValue.map(
176 [&](jsg::JsRef<jsg::JsValue>& ref) -> jsg::JsValue { return ref.getHandle(js); });
177}
178 
179void MessagePort::setOnMessage(jsg::Lock& js, jsg::JsValue value) {
180 if (!value.isObject() && !value.isFunction()) {
181 onmessageValue = kj::none;
182 // If we have no handlers and no onmessage ...
183 if (getHandlerCount("message"_kj) == 0 && onmessageValue == kj::none) {
184 // ...Put the port back into a pending state where messages
185 // will be enqueued until another listener is attached.
186 state = Pending();
187 }
188 } else {
189 onmessageValue = jsg::JsRef<jsg::JsValue>(js, value);
190 start(js);
191 }
192}
193 
194jsg::Ref<MessageChannel> MessageChannel::constructor(jsg::Lock& js) {
195 auto port1 = js.alloc<MessagePort>();
196 auto port2 = js.alloc<MessagePort>();
197 MessagePort::entangle(*port1, *port2);
198 return js.alloc<MessageChannel>(kj::mv(port1), kj::mv(port2));
199}
200 
201} // namespace workerd::api