// Copyright (c) 2017-2022 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #pragma once #include "basics.h" #include "events.h" #include #include #include #include #include #include #include #include #include namespace workerd { class ActorObserver; } namespace workerd::api { class Blob; template struct DeferredProxy; class CloseEvent: public Event { public: CloseEvent(uint code, kj::String reason, bool clean) : Event("close"), code(code), reason(kj::mv(reason)), clean(clean) {} CloseEvent(kj::String type, int code, kj::String reason, bool clean) : Event(kj::mv(type)), code(code), reason(kj::mv(reason)), clean(clean) {} struct Initializer { jsg::Optional code; jsg::Optional reason; jsg::Optional wasClean; JSG_STRUCT(code, reason, wasClean); JSG_STRUCT_TS_OVERRIDE(CloseEventInit); }; static jsg::Ref constructor( jsg::Lock& js, kj::String type, jsg::Optional initializer) { Initializer init = kj::mv(initializer).orDefault({}); return js.alloc(kj::mv(type), init.code.orDefault(0), kj::mv(init.reason).orDefault(jsg::USVString(kj::str())), init.wasClean.orDefault(false)); } int getCode() { return code; } kj::StringPtr getReason() { return reason; } bool getWasClean() { return clean; } JSG_RESOURCE_TYPE(CloseEvent) { JSG_INHERIT(Event); JSG_READONLY_INSTANCE_PROPERTY(code, getCode); JSG_READONLY_INSTANCE_PROPERTY(reason, getReason); JSG_READONLY_INSTANCE_PROPERTY(wasClean, getWasClean); JSG_TS_ROOT(); // CloseEvent will be referenced from the `WebSocketEventMap` define } void visitForMemoryInfo(jsg::MemoryTracker& tracker) const { tracker.trackField("reason", reason); } private: int code; kj::String reason; bool clean; }; WD_STRONG_BOOL(AllowHalfOpen); // The forward declaration is necessary so we can make some // WebSocket methods accessible to WebSocketPair via friend declaration. class WebSocket; class WebSocketPair: public jsg::Object { private: struct IteratorState final { jsg::Ref pair; size_t index = 0; void visitForGc(jsg::GcVisitor& visitor) { visitor.visit(pair); } JSG_MEMORY_INFO(IteratorState) { tracker.trackField("pair", pair); } }; public: WebSocketPair(jsg::Ref first, jsg::Ref second) : sockets{kj::mv(first), kj::mv(second)} {} static jsg::Ref constructor(jsg::Lock& js); jsg::Ref getFirst() { return sockets[0].addRef(); } jsg::Ref getSecond() { return sockets[1].addRef(); } JSG_ITERATOR(PairIterator, entries, jsg::Ref, IteratorState, iteratorNext); JSG_RESOURCE_TYPE(WebSocketPair) { // TODO(soon): These really should be using an indexed property handler rather // than named instance properties but jsg does not yet have support for that. JSG_READONLY_INSTANCE_PROPERTY(0, getFirst); JSG_READONLY_INSTANCE_PROPERTY(1, getSecond); JSG_ITERABLE(entries); JSG_TS_OVERRIDE(const WebSocketPair: { new (): { 0: WebSocket; 1: WebSocket }; }); // Ensure correct typing with `Object.values()`. // Without this override, the generated definition will look like: // // ```ts // declare class WebSocketPair { // constructor(); // readonly 0: WebSocket; // readonly 1: WebSocket; // } // ``` // // Trying to call `Object.values(new WebSocketPair())` will result // in the following `any` typed values: // // ```ts // const [one, two] = Object.values(new WebSocketPair()); // // ^? const one: any // ``` // // With this override in place, `one` and `two` will be typed `WebSocket`. } void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; private: jsg::Ref sockets[2]; static kj::Maybe> iteratorNext(jsg::Lock& js, IteratorState& state) { if (state.index >= 2) { return kj::none; } return state.pair->sockets[state.index++].addRef(); } void visitForGc(jsg::GcVisitor& visitor) { visitor.visit(sockets[0]); visitor.visit(sockets[1]); } }; class WebSocket: public EventTarget { private: // Forward declarations. struct PackedWebSocket; struct Native; public: // WebSocket ready states. static constexpr int READY_STATE_CONNECTING = 0; static constexpr int READY_STATE_OPEN = 1; static constexpr int READY_STATE_CLOSING = 2; static constexpr int READY_STATE_CLOSED = 3; // Creates the Native object when we recreate the WebSocket when waking from hibernation. IoOwn initNative(IoContext& ioContext, kj::WebSocket& ws, kj::Array tags, bool closedOutgoingConn); // Some properties of the `api::WebSocket` that need to survive hibernation. When we initiate // the hibernation process, we want to move these properties out of the `api::WebSocket`. // When we recreate the websocket due to activity, we move the properties back in. struct HibernationPackage { kj::Maybe url; kj::Maybe protocol; kj::Maybe extensions; kj::Maybe> serializedAttachment; // `maybeTags` is only non-empty when we're recreating the api::WebSocket. // We don't need to populate it when hibernating because the tags are already // stored in the HibernationManager. kj::Maybe> maybeTags; // True forever once the JS WebSocket calls `close()`. bool closedOutgoingConnection = false; // Whether the WebSocket allows half-open close state. AllowHalfOpen allowHalfOpen = AllowHalfOpen::YES; }; ~WebSocket() noexcept(false) { weakRef->invalidate(); } // This WebSocket constructor is only used when WebSockets wake up from hibernation. // It will immediately set the `state` to `Accepted`, but it limits the behavior by specifying it // as `Hibernatable` -- thereby making most api::WebSocket methods inaccessible. WebSocket(jsg::Lock& js, IoContext& ioContext, kj::WebSocket& ws, HibernationPackage package); // Similar to how the JS `constructor()` creates a WebSocket, when waking from hibernation // we want to be able to recreate WebSockets from C++ that will be delivered to JS code. static jsg::Ref hibernatableFromNative( jsg::Lock& js, kj::WebSocket& ws, HibernationPackage package); // The JS WebSocket constructor needs to initiate a connection, but we need to return the // WebSocket object to the caller in Javascript immediately. We will defer the connection logic // to the `initConnection` method. WebSocket(jsg::Lock& js, kj::Own native); // The JS WebSocket constructor needs to initiate a connection, but we need to return the // WebSocket object to the caller in Javascript immediately. We will defer the connection logic // to the `initConnection` method. WebSocket(jsg::Lock& js, kj::String url); // We initiate a `new WebSocket()` connection and set up a continuation that handles the // response once it's available. This includes assigning the native websocket and dispatching the // relevant `open`/`error` events. void initConnection(jsg::Lock& js, kj::Promise); // Pumps messages from this WebSocket to `other`, and from `other` to this, making sure to // register pending events as appropriate. Used to connect a websocket to a client via an HTTP // response. // // Only one of this or accept() is allowed to be invoked. // // As an exception to the usual KJ convention, it is not necessary for the JavaScript `WebSocket` // object to be kept live while waiting for the promise returned by couple() to complete. Instead, // the promise takes direct ownership of the underlying KJ-native WebSocket (as well as `other`). kj::Promise> couple(kj::Own other, RequestObserver& request); // Extract the kj::WebSocket from this api::WebSocket (if applicable). The kj::WebSocket will be // owned elsewhere, but the api::WebSocket will retain a reference. kj::Own acceptAsHibernatable(kj::Array tags); void tryReleaseNative(jsg::Lock& js); // Accesses the tags of the hibernatable websocket. kj::Array getHibernatableTags(); enum class HibernatableReleaseState { // The way we release Hibernatable WebSockets slightly differs from regular WebSockets. // We can't access the isolate after the event runs. `NONE` indicates we are not releasing. NONE, CLOSE, ERROR }; // Called when a Hibernatable WebSocket wants to dispatch a close/error event, this modifies // our `Accepted` state to prepare the state to transition to `Released`. void initiateHibernatableRelease(jsg::Lock& js, kj::Own ws, kj::Array tags, HibernatableReleaseState releaseState); bool awaitingHibernatableError(); bool awaitingHibernatableRelease(); // Should only be called on one end of a WebSocketPair. // Relevant for WebSocket Hibernation: the end we return in the Response must be in the // AwaitingAcceptanceOrCoupling state. bool peerIsAwaitingCoupling(); HibernationPackage buildPackageForHibernation(); // --------------------------------------------------------------------------- // JS API. struct AcceptOptions { jsg::Optional allowHalfOpen; JSG_STRUCT(allowHalfOpen); }; // Creates a new outbound WebSocket. static jsg::Ref constructor(jsg::Lock& js, kj::String url, jsg::Optional, kj::String>> protocols); // Begin delivering events locally. void accept(jsg::Lock& js, jsg::Optional options); // Same as accept(), but websockets that are created with `new WebSocket()` in JS cannot call // accept(). Instead, we only permit the C++ constructor to call this "internal" version of accept() // so that the websocket can start processing messages once the connection has been established. void internalAccept(jsg::Lock& js, kj::Maybe> cs); // We defer the actual logic of accept() and internalAccept() to this method, since they largely // share code. void startReadLoop(jsg::Lock& js, kj::Maybe> cs); void send(jsg::Lock& js, kj::OneOf, kj::String> message); void close(jsg::Lock& js, jsg::Optional code, jsg::Optional reason); // Used to get/set the attachment for hibernation. // If the object isn't serialized, it will not survive hibernation. void serializeAttachment(jsg::Lock& js, jsg::JsValue attachment); // Used to get/set the attachment for hibernation. // If the object isn't serialized, it will not survive hibernation. kj::Maybe deserializeAttachment(jsg::Lock& js); // Used to get/store the last auto request/response timestamp for this WebSocket. // These methods are c++ only and are not exposed to our js interface. // Also used to track hibernatable websockets auto-response sends. void setAutoResponseStatus(kj::Maybe time, kj::Promise autoResponsePromise); // Used to get/store the last auto request/response timestamp for this WebSocket. // These methods are c++ only and are not exposed to our js interface. kj::Maybe getAutoResponseTimestamp(); kj::Promise sendAutoResponse(kj::String message, kj::WebSocket& ws); int getReadyState(); bool isAccepted(); bool isReleased(); // For internal use only. // We need to access the underlying KJ WebSocket so we can determine the compression configuration // it uses (if any). kj::Maybe getPreferredExtensions(kj::WebSocket::ExtensionsContext ctx); kj::Maybe getUrl(); kj::Maybe getProtocol(); kj::Maybe getExtensions(); kj::StringPtr getBinaryType(); void setBinaryType(kj::String value); JSG_RESOURCE_TYPE(WebSocket, CompatibilityFlags::Reader flags) { JSG_INHERIT(EventTarget); JSG_METHOD(accept); JSG_METHOD(send); JSG_METHOD(close); JSG_METHOD(serializeAttachment); JSG_METHOD(deserializeAttachment); JSG_STATIC_CONSTANT(READY_STATE_CONNECTING); JSG_STATIC_CONSTANT_NAMED(CONNECTING, WebSocket::READY_STATE_CONNECTING); JSG_STATIC_CONSTANT(READY_STATE_OPEN); JSG_STATIC_CONSTANT_NAMED(OPEN, WebSocket::READY_STATE_OPEN); JSG_STATIC_CONSTANT(READY_STATE_CLOSING); JSG_STATIC_CONSTANT_NAMED(CLOSING, WebSocket::READY_STATE_CLOSING); JSG_STATIC_CONSTANT(READY_STATE_CLOSED); JSG_STATIC_CONSTANT_NAMED(CLOSED, WebSocket::READY_STATE_CLOSED); // Previously, we were setting all properties as instance properties, // which broke the ability to subclass the Event object. With the // compatibility flag set, we instead attach the properties to the // prototype. if (flags.getJsgPropertyOnPrototypeTemplate()) { JSG_READONLY_PROTOTYPE_PROPERTY(readyState, getReadyState); JSG_READONLY_PROTOTYPE_PROPERTY(url, getUrl); JSG_READONLY_PROTOTYPE_PROPERTY(protocol, getProtocol); JSG_READONLY_PROTOTYPE_PROPERTY(extensions, getExtensions); JSG_PROTOTYPE_PROPERTY(binaryType, getBinaryType, setBinaryType); } else { JSG_READONLY_INSTANCE_PROPERTY(readyState, getReadyState); JSG_READONLY_INSTANCE_PROPERTY(url, getUrl); JSG_READONLY_INSTANCE_PROPERTY(protocol, getProtocol); JSG_READONLY_INSTANCE_PROPERTY(extensions, getExtensions); JSG_INSTANCE_PROPERTY(binaryType, getBinaryType, setBinaryType); } JSG_TS_DEFINE(type WebSocketEventMap = { close: CloseEvent; message: MessageEvent; open: Event; error: ErrorEvent; }); JSG_TS_OVERRIDE(extends EventTarget { get binaryType(): "blob" | "arraybuffer"; set binaryType(value: "blob" | "arraybuffer"); }); } void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; kj::Own> addWeakRef() { return weakRef->addRef(); } private: kj::Own> weakRef; kj::Maybe url; kj::Maybe protocol = kj::String(); kj::Maybe extensions = kj::String(); // The binaryType attribute per the WHATWG WebSocket spec. Defaults to "blob" when the // websocket_standard_binary_type compat flag is enabled, "arraybuffer" otherwise. enum class BinaryType { BLOB, ARRAYBUFFER }; BinaryType binaryType_ = BinaryType::ARRAYBUFFER; kj::Maybe autoResponseTimestamp; // All WebSockets have this property. It starts out null but can // be assigned to any serializable value. The property will survive hibernation. // We have to serialize each time we call the setter so we can determine if the size limit // has been breached. kj::Maybe> serializedAttachment; // Tracks farNative->closedOutgoing, but we need to access it when we trigger Hibernation so it // cannot be `IoOwn`ed as `farNative` is. This informs the HibernatableWebSocket if we called // `close()`, thereby preventing calls to `send()` even after we wake from hibernation. bool closedOutgoingForHib = false; // When YES, a server-initiated close does NOT automatically send a reciprocal close frame, // leaving readyState as CLOSING (2) when the close event fires. The application is then // responsible for calling close() explicitly. When NO (spec-compliant default with the // web_socket_auto_reply_to_close compat flag), a close reply is sent automatically and // readyState is CLOSED (3) when the close event fires. // Default is YES (legacy behavior); overridden from the compat flag at construction time. AllowHalfOpen allowHalfOpen = AllowHalfOpen::YES; // Maximum allowed size for WebSocket messages inline static const size_t SUGGESTED_MAX_MESSAGE_SIZE = 1u << 20; // Maximum size of a WebSocket attachment. inline static const size_t MAX_ATTACHMENT_SIZE = 1024 * 16; struct AwaitingConnection { // A canceler associated with the pending websocket connection for `new Websocket()`. kj::Canceler canceler; }; struct AwaitingAcceptanceOrCoupling { explicit AwaitingAcceptanceOrCoupling(kj::Own ws): ws(kj::mv(ws)) {} kj::Own ws; }; struct Accepted { // A `Hibernatable` WebSocket shares a sub-set of behavior that's already implemented for an // `Accepted` WebSocket, so we can think of it a sub-state. struct Hibernatable { kj::WebSocket& ws; // If we have initiated a hibernatable error/close event, we need to take back ownership of // the kj::WebSocket so any final queued messages will deliver. We store this owned websocket // in `attachedForClose`. Since the `ws` reference is still valid, we prevent usage of // `attachedForClose` directly in favor of using continuing to use `ws` directly. kj::Maybe> attachedForClose; // We can't move the state to Released after the Hibernatable Close/Error event runs, since // we don't have a request on the thread by the time the event completes. // // If we are "releasing", we may prevent the websocket from doing certain things like calling // send/close. We're more restrictive if we're delivering an Error than delivering a Close. HibernatableReleaseState releaseState = HibernatableReleaseState::NONE; // There are two possible states for tagsRef: // 1. kj::Array // - Tags are owned by the HibernationManager, we just reference them to save memory. // 2. kj::Array // - We're going to be dispatching a Close or an Error event, i.e. the // HibernatableWebSocket is free to go away. We can no longer rely on tags stored in // the HibernationManager, so instead we copy the data into the api::WebSocket. // // We could just copy all tags into api::WebSocet every time we reactivate/wake from // hibernation, but it could add up to 2.56KB of memory for each websocket. // With a maximum of 32k websockets, that could put a lot of memory pressure on the DO. kj::OneOf, kj::Array> tagsRef; }; explicit Accepted(kj::Own ws, Native& native, IoContext& context); explicit Accepted(Hibernatable ws, Native& native, IoContext& context); ~Accepted() noexcept(false); // A simple wrapper to make it easier to access the underlying kj::WebSocket. class WrappedWebSocket { public: explicit WrappedWebSocket(Hibernatable ws); explicit WrappedWebSocket(kj::Own ws); kj::WebSocket* operator->(); kj::WebSocket& operator*(); kj::Maybe&> getIfNotHibernatable(); kj::Maybe getIfHibernatable(); kj::Array getHibernatableTags(); // Transitions our Hibernatable websocket to a "Releasing" state. // The websocket will transition to `Released` when convenient. void initiateHibernatableRelease(jsg::Lock& js, kj::Own ws, kj::Array tags, HibernatableReleaseState state); bool isAwaitingRelease(); bool isAwaitingError(); private: kj::OneOf, Hibernatable> inner; }; WrappedWebSocket ws; bool isHibernatable(); kj::Promise createAbortTask(Native& native, IoContext& context); // Listens for ws->whenAborted() and possibly triggers a proactive shutdown. kj::Promise whenAbortedTask = nullptr; kj::Maybe> actorMetrics; // This canceler wraps the pump loop as a precaution to make sure we can't exit the Accepted // state with a pump task still happening asynchronously. In practice the canceler should usually // be empty when destroyed because we do not leave the Accepted state if we're still pumping. // Even in the case of IoContext premature cancellation, the pump task should be canceled // by the IoContext before the Canceler is destroyed. kj::Canceler canceler; }; struct Released {}; using NativeState = kj::OneOf; friend kj::StringPtr KJ_STRINGIFY(const NativeState&); struct Native { NativeState state; // Is there currently a task running to pump outgoing messages? bool isPumping = false; // Has a Close message been enqueued for send? (It may still be in outgoingMessages. Check // closedOutgoing && !isPumping to check if it has gone out.) bool closedOutgoing = false; // Has a Close message been received, or has a premature disconnection occurred? bool closedIncoming = false; // Have we detected that the peer has stopped accepting messages? We may want to clean up more // proactively in this case. bool outgoingAborted = false; }; // The underlying native WebSocket (or a promise that will emplace one). // // The state transitions look like so: // - Starts as `AwaitingConnection` if the `WebSocket(url ...)` ctor is used. // - Starts as `AwaitingAcceptanceOrCoupling` if the `WebSocket(native)` ctor is used. // - Transitions from `AwaitingConnection` to `AwaitingAcceptanceOrCoupling` when the native // connection is established and to `Accepted` once the read loop starts. // - Transitions from `AwaitingConnection` to `Released` when connection establishment fails. // - Transitions from `AwaitingAcceptanceOrCoupling` to `Accepted` when it is accepted. // - Transitions from `AwaitingAcceptanceOrCoupling` to `Released` when it is coupled to another // web socket. // - Transitions from `Accepted` to `Released` when outgoing pump is done and either both // directions have seen "close" messages or an error has occurred. IoOwn farNative; // If any error has occurred. kj::Maybe> error; struct GatedMessage { kj::Maybe> outputLock; // must wait for this before actually sending kj::WebSocket::Message message; size_t pendingAutoResponses = 0; }; using OutgoingMessagesMap = kj::Table; // Queue of messages to be sent. This is wrapped in an IoOwn so that the pump loop can safely // access the map without locking the isolate. IoOwn outgoingMessages; // Keep track of current hibernatable websockets auto-response status to avoid racing // between regular websocket messages, and auto-responses. struct AutoResponse { // The ongoing auto-response promise, used for pump() synchronization. // Wrapped in IoOwn when an IoContext is available (for GC-safe destruction that avoids // violating DISALLOW_KJ_IO_DESTRUCTORS_SCOPE), or plain kj::Own when no IoContext is // available (e.g. when called from the hibernation manager's readLoop). using OwnedAutoResponsePromise = kj::OneOf>, kj::Own>>; kj::Maybe ongoingAutoResponse; workerd::util::Queue pendingAutoResponseDeque; size_t queuedAutoResponses = 0; bool isPumping = false; bool isClosed = false; JSG_MEMORY_INFO(AutoResponse) { tracker.trackFieldWithSize("ongoingAutoResponse", sizeof(kj::Promise)); pendingAutoResponseDeque.forEach( [&](const kj::String& message) { tracker.trackField(nullptr, message); }); } }; AutoResponse autoResponseStatus; kj::Maybe> observer; // Contains a websocket and possibly some data from the WebSocketResponse headers. struct PackedWebSocket { kj::Own ws; kj::Maybe proto; kj::Maybe extensions; }; // So that each end of a WebSocketPair can keep track of its peer. // We use a weak ref to track the peer to avoid having a strong ref cycle // between the two WebSocket instances that would cause them to leak. This // can mean, however, that it's possible for one of the peers to be garbage // collected while the other still exists. This should be fairly unusual tho. kj::Maybe>> peer; void visitForGc(jsg::GcVisitor& visitor) { visitor.visit(error); } void setPeer(kj::Own> peer); friend jsg::Ref WebSocketPair::constructor(jsg::Lock&); void dispatchOpen(jsg::Lock& js); void ensurePumping(jsg::Lock& js); // Returns the number of pending auto-responses that should be sent before the next outgoing // message, and advances the queuedAutoResponses counter. Called each time a GatedMessage is // inserted into outgoingMessages to guarantee auto-response ordering. size_t getPendingAutoResponseCount(); // Write messages from `outgoingMessages` into `ws`. // // These are not necessarily called under isolate lock, but they are called on the given // context's thread. They are declared `static` to prove they don't access the JavaScript // object's members in a thread-unsafe way. `outgoingMessages` and `ws` are both `IoOwn`ed // objects so are safe to access from the thread without the isolate lock. The whole task is // owned by the `IoContext` so it'll be canceled if the `IoContext` is destroyed. static kj::Promise pump(IoContext& context, OutgoingMessagesMap& outgoingMessages, kj::WebSocket& ws, Native& native, AutoResponse& autoResponse, kj::Maybe>& observer); kj::Promise> readLoop( kj::Maybe> cs, size_t maxMessageSize); void reportError(jsg::Lock& js, kj::Exception&& e); void reportError(jsg::Lock& js, jsg::JsRef err); void assertNoError(jsg::Lock& js); }; #define EW_WEBSOCKET_ISOLATE_TYPES \ api::CloseEvent, api::CloseEvent::Initializer, api::WebSocket, api::WebSocket::AcceptOptions, \ api::WebSocketPair, api::WebSocketPair::PairIterator, \ api::WebSocketPair::PairIterator:: \ Next // The list of websocket.h types that are added to worker.c++'s JSG_DECLARE_ISOLATE_TYPE } // namespace workerd::api