File
Blob: src/workerd/api/messagechannel.c++
| 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 | |
| 9 | namespace workerd::api { |
| 10 | MessagePort::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 | |
| 44 | void 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. |
| 62 | void 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. |
| 85 | void 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. |
| 91 | void 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 | |
| 132 | void 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 | |
| 141 | void 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. |
| 152 | void 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 | |
| 174 | kj::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 | |
| 179 | void 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 | |
| 194 | jsg::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 |