File
Blob: src/workerd/api/messagechannel.h
| 1 | #pragma once |
| 2 | |
| 3 | #include <workerd/api/basics.h> |
| 4 | #include <workerd/io/io-context.h> |
| 5 | #include <workerd/jsg/jsg.h> |
| 6 | #include <workerd/jsg/modules-new.h> |
| 7 | #include <workerd/jsg/ser.h> |
| 8 | #include <workerd/jsg/url.h> |
| 9 | #include <workerd/util/weak-refs.h> |
| 10 | |
| 11 | namespace workerd::api { |
| 12 | |
| 13 | // A closely approximate implementation of the Web platform standard MessagePort. |
| 14 | // MessagePorts always come in pairs. When a message is posted to |
| 15 | // one it is delivered to the other, and vice versa. When one port |
| 16 | // is closed both ports are closed. |
| 17 | // |
| 18 | // This intentionally does not implement the full MessagePort spec and we know |
| 19 | // that it varies from the standard definition in a number of ways: |
| 20 | // |
| 21 | // - It does not support transfer lists. We do not implement the transfer |
| 22 | // list semantics, but we do validate the transfer list input to an extent. |
| 23 | // - It does not support serialization/deserialization. It's not possible to |
| 24 | // send a MessagePort anywhere currently. |
| 25 | // - The `messageerror` event is only partially implemented. Currently, if a |
| 26 | // message data cannot be serialized/deserialized it will throw an error |
| 27 | // synchronously when posted rather than dispatching the `messageerror` event |
| 28 | // on the receiving port, this is just easiest to implement for now and makes |
| 29 | // the most sense for our current use case since the MessagePort only ever |
| 30 | // passes messages around within the same isolate (that is, we're not sending |
| 31 | // the serialized data off anywhere, we're just cloning it and dispatching it.) |
| 32 | // - We intentionally do not implement the "port message queue" semantics exactly |
| 33 | // as they are described in the spec. When a MessagePort has an onmessage listener, |
| 34 | // the message delivery is flowing, when there is no onmessage listener, the |
| 35 | // messages are queued up until the port is started. Because we are storing |
| 36 | // these as JS values, we don't worry about extra memory accounting for the queue. |
| 37 | // - We do not emit the close event on entangled ports when one of them is GC'd. |
| 38 | // - We do not check to see if a MessagePort is entangled with another when we |
| 39 | // call entangle because there's only one way to entangle them currently and |
| 40 | // it's impossible for them to be already entangled. |
| 41 | // - We do not implement disentangle steps other than to invalidate the weak |
| 42 | // ref to the other port when one of them is closed. |
| 43 | // - We do not prevent a MessagePort from being garbage collected while it has |
| 44 | // messages queued up. Eventually when we implement ser/deser this might change. |
| 45 | // - Unlike the implementation in Node.js, not closing a MessagePort does not |
| 46 | // prevent anything from exiting. It's best to close MessagePorts manually |
| 47 | // but the current implementation does not require it. |
| 48 | // |
| 49 | // Because of these differences we do not currently run the full suite of web |
| 50 | // platform tests against our implementation -- we know most of them will fail |
| 51 | // since most of them depend on the ability to transfer MessagePorts or depend |
| 52 | // on the mechanisms we do not implement. And yes, we know that this means that |
| 53 | // if we need stricter compliance with the spec in the future we will likely |
| 54 | // need to introduce a compat flag. |
| 55 | class MessagePort final: public EventTarget { |
| 56 | public: |
| 57 | // While we do not support transfer lists in the implementation |
| 58 | // currently, we do want to validate those inputs. |
| 59 | using TransferList = kj::Array<jsg::JsRef<jsg::JsValue>>; |
| 60 | struct PostMessageOptions { |
| 61 | jsg::Optional<TransferList> transfer; |
| 62 | JSG_STRUCT(transfer); |
| 63 | }; |
| 64 | using TransferListOrOptions = kj::OneOf<TransferList, PostMessageOptions>; |
| 65 | |
| 66 | MessagePort(); |
| 67 | ~MessagePort() noexcept(false) { |
| 68 | closeImpl(); |
| 69 | } |
| 70 | |
| 71 | // MessagePort instances cannot be created directly. |
| 72 | // Use `new MessageChannel()` |
| 73 | static jsg::Ref<MessagePort> constructor() = delete; |
| 74 | |
| 75 | void postMessage(jsg::Lock& js, |
| 76 | jsg::Optional<jsg::JsRef<jsg::JsValue>> data = kj::none, |
| 77 | jsg::Optional<TransferListOrOptions> options = kj::none); |
| 78 | void closeImpl(); |
| 79 | void close(jsg::Lock& js); |
| 80 | void start(jsg::Lock& js); |
| 81 | |
| 82 | // Support the onmessage getter and setter. Per the spec, when |
| 83 | // onmessage is set, the MessagePort is automatically started, |
| 84 | // but when addEventListener is set, start must be called |
| 85 | // separately. That's a kind of a weird rule but ok. To support |
| 86 | // that we need to define an onmessage getter/setter pair. |
| 87 | kj::Maybe<jsg::JsValue> getOnMessage(jsg::Lock& js); |
| 88 | void setOnMessage(jsg::Lock& js, jsg::JsValue value); |
| 89 | |
| 90 | JSG_RESOURCE_TYPE(MessagePort) { |
| 91 | JSG_INHERIT(EventTarget); |
| 92 | JSG_METHOD(postMessage); |
| 93 | JSG_METHOD(close); |
| 94 | JSG_METHOD(start); |
| 95 | JSG_PROTOTYPE_PROPERTY(onmessage, getOnMessage, setOnMessage); |
| 96 | } |
| 97 | |
| 98 | jsg::Ref<MessagePort> addRef() { |
| 99 | return JSG_THIS; |
| 100 | } |
| 101 | bool isClosed() const { |
| 102 | return state.is<Closed>(); |
| 103 | } |
| 104 | |
| 105 | void deliver(jsg::Lock& js, const jsg::JsValue& data); |
| 106 | |
| 107 | // Bind two message ports together such that messages posted to |
| 108 | // one are delivered to the other. |
| 109 | static void entangle(MessagePort& port1, MessagePort& port2); |
| 110 | |
| 111 | kj::Maybe<MessagePort&> getOther() { |
| 112 | return other->tryGet().map([](MessagePort& o) -> MessagePort& { return o; }); |
| 113 | } |
| 114 | |
| 115 | // TODO(soon): Support serialization/deserialization to use MessagePort |
| 116 | // with JSRPC. We'll need to implement a rpc mechanism for passing the |
| 117 | // messages across the rpc boundary. |
| 118 | |
| 119 | private: |
| 120 | // When the MessagePort is in the pending state, messages posted to it |
| 121 | // will be buffered until the port is started. When the port is started, |
| 122 | // the buffered messages will be delivered immediately. |
| 123 | using Pending = kj::Vector<jsg::JsRef<jsg::JsValue>>; |
| 124 | struct Started {}; |
| 125 | struct Closed {}; |
| 126 | |
| 127 | void dispatchMessage(jsg::Lock& js, const jsg::JsValue& value); |
| 128 | |
| 129 | kj::Own<WeakRef<MessagePort>> addWeakRef() { |
| 130 | KJ_ASSERT(weakThis->isValid()); |
| 131 | return kj::addRef(*weakThis); |
| 132 | } |
| 133 | |
| 134 | kj::Own<WeakRef<MessagePort>> weakThis; |
| 135 | kj::OneOf<Pending, Started, Closed> state; |
| 136 | |
| 137 | // Two ports are entangled when they weakly reference each other. |
| 138 | // Keep in mind that this is a weak reference! So if one of the |
| 139 | // ports gets GC'd the other will will also end up being closed. |
| 140 | // To keep them both alive, maintain strong references to both |
| 141 | // ports! |
| 142 | kj::Own<WeakRef<MessagePort>> other; |
| 143 | kj::Maybe<jsg::JsRef<jsg::JsValue>> onmessageValue; |
| 144 | }; |
| 145 | |
| 146 | // MessageChannel is simple enough... create a couple of MessagePorts |
| 147 | // and entangle those so that they will exchange messages with each |
| 148 | // other. |
| 149 | class MessageChannel final: public jsg::Object { |
| 150 | public: |
| 151 | MessageChannel(jsg::Ref<MessagePort> port1, jsg::Ref<MessagePort> port2) |
| 152 | : port1(kj::mv(port1)), |
| 153 | port2(kj::mv(port2)) {} |
| 154 | |
| 155 | static jsg::Ref<MessageChannel> constructor(jsg::Lock& js); |
| 156 | |
| 157 | jsg::Ref<MessagePort> getPort1() { |
| 158 | return port1.addRef(); |
| 159 | } |
| 160 | jsg::Ref<MessagePort> getPort2() { |
| 161 | return port2.addRef(); |
| 162 | } |
| 163 | |
| 164 | JSG_RESOURCE_TYPE(MessageChannel) { |
| 165 | JSG_LAZY_READONLY_INSTANCE_PROPERTY(port1, getPort1); |
| 166 | JSG_LAZY_READONLY_INSTANCE_PROPERTY(port2, getPort2); |
| 167 | } |
| 168 | |
| 169 | private: |
| 170 | jsg::Ref<MessagePort> port1; |
| 171 | jsg::Ref<MessagePort> port2; |
| 172 | }; |
| 173 | |
| 174 | // Module that exposes MessageChannel and MessagePort for internal use by |
| 175 | // built-in modules like node:worker_threads without requiring the global |
| 176 | // expose_global_message_channel compat flag. |
| 177 | class MessageChannelModule final: public jsg::Object { |
| 178 | public: |
| 179 | MessageChannelModule() = default; |
| 180 | MessageChannelModule(jsg::Lock&, const jsg::Url&) {} |
| 181 | |
| 182 | JSG_RESOURCE_TYPE(MessageChannelModule) { |
| 183 | JSG_NESTED_TYPE(MessageChannel); |
| 184 | JSG_NESTED_TYPE(MessagePort); |
| 185 | } |
| 186 | }; |
| 187 | |
| 188 | template <class Registry> |
| 189 | void registerMessageChannelModule(Registry& registry, auto featureFlags) { |
| 190 | registry.template addBuiltinModule<MessageChannelModule>( |
| 191 | "cloudflare-internal:messagechannel", workerd::jsg::ModuleRegistry::Type::INTERNAL); |
| 192 | } |
| 193 | |
| 194 | template <typename TypeWrapper> |
| 195 | kj::Own<jsg::modules::ModuleBundle> getInternalMessageChannelModuleBundle(auto featureFlags) { |
| 196 | jsg::modules::ModuleBundle::BuiltinBuilder builder( |
| 197 | jsg::modules::ModuleBundle::BuiltinBuilder::Type::BUILTIN_ONLY); |
| 198 | static const auto kSpecifier = "cloudflare-internal:messagechannel"_url; |
| 199 | builder.addObject<MessageChannelModule, TypeWrapper>(kSpecifier); |
| 200 | return builder.finish(); |
| 201 | } |
| 202 | |
| 203 | } // namespace workerd::api |
| 204 | |
| 205 | #define EW_MESSAGECHANNEL_ISOLATE_TYPES \ |
| 206 | api::MessagePort, api::MessageChannel, api::MessagePort::PostMessageOptions, \ |
| 207 | api::MessageChannelModule |