Skip to content
File

Blob: src/workerd/api/web-socket.h

cpp679 lines
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#pragma once
6 
7#include "basics.h"
8#include "events.h"
9 
10#include <workerd/io/io-gate.h>
11#include <workerd/io/observer.h>
12#include <workerd/jsg/jsg.h>
13#include <workerd/util/checked-queue.h>
14#include <workerd/util/strong-bool.h>
15#include <workerd/util/weak-refs.h>
16 
17#include <kj/compat/http.h>
18 
19#include <cstdlib>
20#include <list>
21 
22namespace workerd {
23class ActorObserver;
24}
25 
26namespace workerd::api {
27 
28class Blob;
29 
30template <typename T>
31struct DeferredProxy;
32 
33class CloseEvent: public Event {
34 public:
35 CloseEvent(uint code, kj::String reason, bool clean)
36 : Event("close"),
37 code(code),
38 reason(kj::mv(reason)),
39 clean(clean) {}
40 CloseEvent(kj::String type, int code, kj::String reason, bool clean)
41 : Event(kj::mv(type)),
42 code(code),
43 reason(kj::mv(reason)),
44 clean(clean) {}
45 
46 struct Initializer {
47 jsg::Optional<int> code;
48 jsg::Optional<jsg::USVString> reason;
49 jsg::Optional<bool> wasClean;
50 
51 JSG_STRUCT(code, reason, wasClean);
52 JSG_STRUCT_TS_OVERRIDE(CloseEventInit);
53 };
54 static jsg::Ref<CloseEvent> constructor(
55 jsg::Lock& js, kj::String type, jsg::Optional<Initializer> initializer) {
56 Initializer init = kj::mv(initializer).orDefault({});
57 return js.alloc<CloseEvent>(kj::mv(type), init.code.orDefault(0),
58 kj::mv(init.reason).orDefault(jsg::USVString(kj::str())), init.wasClean.orDefault(false));
59 }
60 
61 int getCode() {
62 return code;
63 }
64 kj::StringPtr getReason() {
65 return reason;
66 }
67 bool getWasClean() {
68 return clean;
69 }
70 
71 JSG_RESOURCE_TYPE(CloseEvent) {
72 JSG_INHERIT(Event);
73 
74 JSG_READONLY_INSTANCE_PROPERTY(code, getCode);
75 JSG_READONLY_INSTANCE_PROPERTY(reason, getReason);
76 JSG_READONLY_INSTANCE_PROPERTY(wasClean, getWasClean);
77 
78 JSG_TS_ROOT();
79 // CloseEvent will be referenced from the `WebSocketEventMap` define
80 }
81 
82 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
83 tracker.trackField("reason", reason);
84 }
85 
86 private:
87 int code;
88 kj::String reason;
89 bool clean;
90};
91 
92WD_STRONG_BOOL(AllowHalfOpen);
93 
94// The forward declaration is necessary so we can make some
95// WebSocket methods accessible to WebSocketPair via friend declaration.
96class WebSocket;
97 
98class WebSocketPair: public jsg::Object {
99 private:
100 struct IteratorState final {
101 jsg::Ref<WebSocketPair> pair;
102 size_t index = 0;
103 
104 void visitForGc(jsg::GcVisitor& visitor) {
105 visitor.visit(pair);
106 }
107 
108 JSG_MEMORY_INFO(IteratorState) {
109 tracker.trackField("pair", pair);
110 }
111 };
112 
113 public:
114 WebSocketPair(jsg::Ref<WebSocket> first, jsg::Ref<WebSocket> second)
115 : sockets{kj::mv(first), kj::mv(second)} {}
116 
117 static jsg::Ref<WebSocketPair> constructor(jsg::Lock& js);
118 
119 jsg::Ref<WebSocket> getFirst() {
120 return sockets[0].addRef();
121 }
122 jsg::Ref<WebSocket> getSecond() {
123 return sockets[1].addRef();
124 }
125 
126 JSG_ITERATOR(PairIterator, entries, jsg::Ref<WebSocket>, IteratorState, iteratorNext);
127 
128 JSG_RESOURCE_TYPE(WebSocketPair) {
129 // TODO(soon): These really should be using an indexed property handler rather
130 // than named instance properties but jsg does not yet have support for that.
131 JSG_READONLY_INSTANCE_PROPERTY(0, getFirst);
132 JSG_READONLY_INSTANCE_PROPERTY(1, getSecond);
133 JSG_ITERABLE(entries);
134 
135 JSG_TS_OVERRIDE(const WebSocketPair: {
136 new (): { 0: WebSocket; 1: WebSocket };
137 });
138 // Ensure correct typing with `Object.values()`.
139 // Without this override, the generated definition will look like:
140 //
141 // ```ts
142 // declare class WebSocketPair {
143 // constructor();
144 // readonly 0: WebSocket;
145 // readonly 1: WebSocket;
146 // }
147 // ```
148 //
149 // Trying to call `Object.values(new WebSocketPair())` will result
150 // in the following `any` typed values:
151 //
152 // ```ts
153 // const [one, two] = Object.values(new WebSocketPair());
154 // // ^? const one: any
155 // ```
156 //
157 // With this override in place, `one` and `two` will be typed `WebSocket`.
158 }
159 
160 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const;
161 
162 private:
163 jsg::Ref<WebSocket> sockets[2];
164 
165 static kj::Maybe<jsg::Ref<WebSocket>> iteratorNext(jsg::Lock& js, IteratorState& state) {
166 if (state.index >= 2) {
167 return kj::none;
168 }
169 return state.pair->sockets[state.index++].addRef();
170 }
171 
172 void visitForGc(jsg::GcVisitor& visitor) {
173 visitor.visit(sockets[0]);
174 visitor.visit(sockets[1]);
175 }
176};
177 
178class WebSocket: public EventTarget {
179 private:
180 // Forward declarations.
181 struct PackedWebSocket;
182 struct Native;
183 
184 public:
185 // WebSocket ready states.
186 static constexpr int READY_STATE_CONNECTING = 0;
187 static constexpr int READY_STATE_OPEN = 1;
188 static constexpr int READY_STATE_CLOSING = 2;
189 static constexpr int READY_STATE_CLOSED = 3;
190 
191 // Creates the Native object when we recreate the WebSocket when waking from hibernation.
192 IoOwn<Native> initNative(IoContext& ioContext,
193 kj::WebSocket& ws,
194 kj::Array<kj::StringPtr> tags,
195 bool closedOutgoingConn);
196 
197 // Some properties of the `api::WebSocket` that need to survive hibernation. When we initiate
198 // the hibernation process, we want to move these properties out of the `api::WebSocket`.
199 // When we recreate the websocket due to activity, we move the properties back in.
200 struct HibernationPackage {
201 kj::Maybe<kj::String> url;
202 kj::Maybe<kj::String> protocol;
203 kj::Maybe<kj::String> extensions;
204 kj::Maybe<kj::Array<byte>> serializedAttachment;
205 
206 // `maybeTags` is only non-empty when we're recreating the api::WebSocket.
207 // We don't need to populate it when hibernating because the tags are already
208 // stored in the HibernationManager.
209 kj::Maybe<kj::Array<kj::StringPtr>> maybeTags;
210 
211 // True forever once the JS WebSocket calls `close()`.
212 bool closedOutgoingConnection = false;
213 
214 // Whether the WebSocket allows half-open close state.
215 AllowHalfOpen allowHalfOpen = AllowHalfOpen::YES;
216 };
217 
218 ~WebSocket() noexcept(false) {
219 weakRef->invalidate();
220 }
221 
222 // This WebSocket constructor is only used when WebSockets wake up from hibernation.
223 // It will immediately set the `state` to `Accepted`, but it limits the behavior by specifying it
224 // as `Hibernatable` -- thereby making most api::WebSocket methods inaccessible.
225 WebSocket(jsg::Lock& js, IoContext& ioContext, kj::WebSocket& ws, HibernationPackage package);
226 
227 // Similar to how the JS `constructor()` creates a WebSocket, when waking from hibernation
228 // we want to be able to recreate WebSockets from C++ that will be delivered to JS code.
229 static jsg::Ref<WebSocket> hibernatableFromNative(
230 jsg::Lock& js, kj::WebSocket& ws, HibernationPackage package);
231 
232 // The JS WebSocket constructor needs to initiate a connection, but we need to return the
233 // WebSocket object to the caller in Javascript immediately. We will defer the connection logic
234 // to the `initConnection` method.
235 WebSocket(jsg::Lock& js, kj::Own<kj::WebSocket> native);
236 
237 // The JS WebSocket constructor needs to initiate a connection, but we need to return the
238 // WebSocket object to the caller in Javascript immediately. We will defer the connection logic
239 // to the `initConnection` method.
240 WebSocket(jsg::Lock& js, kj::String url);
241 
242 // We initiate a `new WebSocket()` connection and set up a continuation that handles the
243 // response once it's available. This includes assigning the native websocket and dispatching the
244 // relevant `open`/`error` events.
245 void initConnection(jsg::Lock& js, kj::Promise<PackedWebSocket>);
246 
247 // Pumps messages from this WebSocket to `other`, and from `other` to this, making sure to
248 // register pending events as appropriate. Used to connect a websocket to a client via an HTTP
249 // response.
250 //
251 // Only one of this or accept() is allowed to be invoked.
252 //
253 // As an exception to the usual KJ convention, it is not necessary for the JavaScript `WebSocket`
254 // object to be kept live while waiting for the promise returned by couple() to complete. Instead,
255 // the promise takes direct ownership of the underlying KJ-native WebSocket (as well as `other`).
256 kj::Promise<DeferredProxy<void>> couple(kj::Own<kj::WebSocket> other, RequestObserver& request);
257 
258 // Extract the kj::WebSocket from this api::WebSocket (if applicable). The kj::WebSocket will be
259 // owned elsewhere, but the api::WebSocket will retain a reference.
260 kj::Own<kj::WebSocket> acceptAsHibernatable(kj::Array<kj::StringPtr> tags);
261 
262 void tryReleaseNative(jsg::Lock& js);
263 
264 // Accesses the tags of the hibernatable websocket.
265 kj::Array<kj::StringPtr> getHibernatableTags();
266 
267 enum class HibernatableReleaseState {
268 // The way we release Hibernatable WebSockets slightly differs from regular WebSockets.
269 // We can't access the isolate after the event runs. `NONE` indicates we are not releasing.
270 NONE,
271 CLOSE,
272 ERROR
273 };
274 
275 // Called when a Hibernatable WebSocket wants to dispatch a close/error event, this modifies
276 // our `Accepted` state to prepare the state to transition to `Released`.
277 void initiateHibernatableRelease(jsg::Lock& js,
278 kj::Own<kj::WebSocket> ws,
279 kj::Array<kj::String> tags,
280 HibernatableReleaseState releaseState);
281 
282 bool awaitingHibernatableError();
283 
284 bool awaitingHibernatableRelease();
285 
286 // Should only be called on one end of a WebSocketPair.
287 // Relevant for WebSocket Hibernation: the end we return in the Response must be in the
288 // AwaitingAcceptanceOrCoupling state.
289 bool peerIsAwaitingCoupling();
290 
291 HibernationPackage buildPackageForHibernation();
292 
293 // ---------------------------------------------------------------------------
294 // JS API.
295 
296 struct AcceptOptions {
297 jsg::Optional<bool> allowHalfOpen;
298 
299 JSG_STRUCT(allowHalfOpen);
300 };
301 
302 // Creates a new outbound WebSocket.
303 static jsg::Ref<WebSocket> constructor(jsg::Lock& js,
304 kj::String url,
305 jsg::Optional<kj::OneOf<kj::Array<kj::String>, kj::String>> protocols);
306 
307 // Begin delivering events locally.
308 void accept(jsg::Lock& js, jsg::Optional<AcceptOptions> options);
309 
310 // Same as accept(), but websockets that are created with `new WebSocket()` in JS cannot call
311 // accept(). Instead, we only permit the C++ constructor to call this "internal" version of accept()
312 // so that the websocket can start processing messages once the connection has been established.
313 void internalAccept(jsg::Lock& js, kj::Maybe<kj::Own<InputGate::CriticalSection>> cs);
314 
315 // We defer the actual logic of accept() and internalAccept() to this method, since they largely
316 // share code.
317 void startReadLoop(jsg::Lock& js, kj::Maybe<kj::Own<InputGate::CriticalSection>> cs);
318 
319 void send(jsg::Lock& js, kj::OneOf<kj::Array<byte>, kj::String> message);
320 void close(jsg::Lock& js, jsg::Optional<int> code, jsg::Optional<jsg::USVString> reason);
321 
322 // Used to get/set the attachment for hibernation.
323 // If the object isn't serialized, it will not survive hibernation.
324 void serializeAttachment(jsg::Lock& js, jsg::JsValue attachment);
325 
326 // Used to get/set the attachment for hibernation.
327 // If the object isn't serialized, it will not survive hibernation.
328 kj::Maybe<jsg::JsValue> deserializeAttachment(jsg::Lock& js);
329 
330 // Used to get/store the last auto request/response timestamp for this WebSocket.
331 // These methods are c++ only and are not exposed to our js interface.
332 // Also used to track hibernatable websockets auto-response sends.
333 void setAutoResponseStatus(kj::Maybe<kj::Date> time, kj::Promise<void> autoResponsePromise);
334 
335 // Used to get/store the last auto request/response timestamp for this WebSocket.
336 // These methods are c++ only and are not exposed to our js interface.
337 kj::Maybe<kj::Date> getAutoResponseTimestamp();
338 
339 kj::Promise<void> sendAutoResponse(kj::String message, kj::WebSocket& ws);
340 
341 int getReadyState();
342 
343 bool isAccepted();
344 bool isReleased();
345 
346 // For internal use only.
347 // We need to access the underlying KJ WebSocket so we can determine the compression configuration
348 // it uses (if any).
349 kj::Maybe<kj::String> getPreferredExtensions(kj::WebSocket::ExtensionsContext ctx);
350 
351 kj::Maybe<kj::StringPtr> getUrl();
352 kj::Maybe<kj::StringPtr> getProtocol();
353 kj::Maybe<kj::StringPtr> getExtensions();
354 
355 kj::StringPtr getBinaryType();
356 void setBinaryType(kj::String value);
357 
358 JSG_RESOURCE_TYPE(WebSocket, CompatibilityFlags::Reader flags) {
359 JSG_INHERIT(EventTarget);
360 JSG_METHOD(accept);
361 JSG_METHOD(send);
362 JSG_METHOD(close);
363 JSG_METHOD(serializeAttachment);
364 JSG_METHOD(deserializeAttachment);
365 
366 JSG_STATIC_CONSTANT(READY_STATE_CONNECTING);
367 JSG_STATIC_CONSTANT_NAMED(CONNECTING, WebSocket::READY_STATE_CONNECTING);
368 
369 JSG_STATIC_CONSTANT(READY_STATE_OPEN);
370 JSG_STATIC_CONSTANT_NAMED(OPEN, WebSocket::READY_STATE_OPEN);
371 
372 JSG_STATIC_CONSTANT(READY_STATE_CLOSING);
373 JSG_STATIC_CONSTANT_NAMED(CLOSING, WebSocket::READY_STATE_CLOSING);
374 
375 JSG_STATIC_CONSTANT(READY_STATE_CLOSED);
376 JSG_STATIC_CONSTANT_NAMED(CLOSED, WebSocket::READY_STATE_CLOSED);
377 
378 // Previously, we were setting all properties as instance properties,
379 // which broke the ability to subclass the Event object. With the
380 // compatibility flag set, we instead attach the properties to the
381 // prototype.
382 if (flags.getJsgPropertyOnPrototypeTemplate()) {
383 JSG_READONLY_PROTOTYPE_PROPERTY(readyState, getReadyState);
384 JSG_READONLY_PROTOTYPE_PROPERTY(url, getUrl);
385 JSG_READONLY_PROTOTYPE_PROPERTY(protocol, getProtocol);
386 JSG_READONLY_PROTOTYPE_PROPERTY(extensions, getExtensions);
387 JSG_PROTOTYPE_PROPERTY(binaryType, getBinaryType, setBinaryType);
388 } else {
389 JSG_READONLY_INSTANCE_PROPERTY(readyState, getReadyState);
390 JSG_READONLY_INSTANCE_PROPERTY(url, getUrl);
391 JSG_READONLY_INSTANCE_PROPERTY(protocol, getProtocol);
392 JSG_READONLY_INSTANCE_PROPERTY(extensions, getExtensions);
393 JSG_INSTANCE_PROPERTY(binaryType, getBinaryType, setBinaryType);
394 }
395 
396 JSG_TS_DEFINE(type WebSocketEventMap = {
397 close: CloseEvent;
398 message: MessageEvent;
399 open: Event;
400 error: ErrorEvent;
401 });
402 JSG_TS_OVERRIDE(extends EventTarget<WebSocketEventMap> {
403 get binaryType(): "blob" | "arraybuffer";
404 set binaryType(value: "blob" | "arraybuffer");
405 });
406 }
407 
408 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const;
409 
410 kj::Own<WeakRef<WebSocket>> addWeakRef() {
411 return weakRef->addRef();
412 }
413 
414 private:
415 kj::Own<WeakRef<WebSocket>> weakRef;
416 kj::Maybe<kj::String> url;
417 kj::Maybe<kj::String> protocol = kj::String();
418 kj::Maybe<kj::String> extensions = kj::String();
419 // The binaryType attribute per the WHATWG WebSocket spec. Defaults to "blob" when the
420 // websocket_standard_binary_type compat flag is enabled, "arraybuffer" otherwise.
421 enum class BinaryType { BLOB, ARRAYBUFFER };
422 BinaryType binaryType_ = BinaryType::ARRAYBUFFER;
423 
424 kj::Maybe<kj::Date> autoResponseTimestamp;
425 // All WebSockets have this property. It starts out null but can
426 // be assigned to any serializable value. The property will survive hibernation.
427 // We have to serialize each time we call the setter so we can determine if the size limit
428 // has been breached.
429 kj::Maybe<kj::Array<byte>> serializedAttachment;
430 
431 // Tracks farNative->closedOutgoing, but we need to access it when we trigger Hibernation so it
432 // cannot be `IoOwn`ed as `farNative` is. This informs the HibernatableWebSocket if we called
433 // `close()`, thereby preventing calls to `send()` even after we wake from hibernation.
434 bool closedOutgoingForHib = false;
435 
436 // When YES, a server-initiated close does NOT automatically send a reciprocal close frame,
437 // leaving readyState as CLOSING (2) when the close event fires. The application is then
438 // responsible for calling close() explicitly. When NO (spec-compliant default with the
439 // web_socket_auto_reply_to_close compat flag), a close reply is sent automatically and
440 // readyState is CLOSED (3) when the close event fires.
441 // Default is YES (legacy behavior); overridden from the compat flag at construction time.
442 AllowHalfOpen allowHalfOpen = AllowHalfOpen::YES;
443 
444 // Maximum allowed size for WebSocket messages
445 inline static const size_t SUGGESTED_MAX_MESSAGE_SIZE = 1u << 20;
446 
447 // Maximum size of a WebSocket attachment.
448 inline static const size_t MAX_ATTACHMENT_SIZE = 1024 * 16;
449 
450 struct AwaitingConnection {
451 // A canceler associated with the pending websocket connection for `new Websocket()`.
452 kj::Canceler canceler;
453 };
454 struct AwaitingAcceptanceOrCoupling {
455 explicit AwaitingAcceptanceOrCoupling(kj::Own<kj::WebSocket> ws): ws(kj::mv(ws)) {}
456 kj::Own<kj::WebSocket> ws;
457 };
458 struct Accepted {
459 // A `Hibernatable` WebSocket shares a sub-set of behavior that's already implemented for an
460 // `Accepted` WebSocket, so we can think of it a sub-state.
461 struct Hibernatable {
462 kj::WebSocket& ws;
463 // If we have initiated a hibernatable error/close event, we need to take back ownership of
464 // the kj::WebSocket so any final queued messages will deliver. We store this owned websocket
465 // in `attachedForClose`. Since the `ws` reference is still valid, we prevent usage of
466 // `attachedForClose` directly in favor of using continuing to use `ws` directly.
467 kj::Maybe<kj::Own<void>> attachedForClose;
468 
469 // We can't move the state to Released after the Hibernatable Close/Error event runs, since
470 // we don't have a request on the thread by the time the event completes.
471 //
472 // If we are "releasing", we may prevent the websocket from doing certain things like calling
473 // send/close. We're more restrictive if we're delivering an Error than delivering a Close.
474 HibernatableReleaseState releaseState = HibernatableReleaseState::NONE;
475 
476 // There are two possible states for tagsRef:
477 // 1. kj::Array<kj::StringPtr>
478 // - Tags are owned by the HibernationManager, we just reference them to save memory.
479 // 2. kj::Array<kj::String>
480 // - We're going to be dispatching a Close or an Error event, i.e. the
481 // HibernatableWebSocket is free to go away. We can no longer rely on tags stored in
482 // the HibernationManager, so instead we copy the data into the api::WebSocket.
483 //
484 // We could just copy all tags into api::WebSocet every time we reactivate/wake from
485 // hibernation, but it could add up to 2.56KB of memory for each websocket.
486 // With a maximum of 32k websockets, that could put a lot of memory pressure on the DO.
487 kj::OneOf<kj::Array<kj::StringPtr>, kj::Array<kj::String>> tagsRef;
488 };
489 
490 explicit Accepted(kj::Own<kj::WebSocket> ws, Native& native, IoContext& context);
491 explicit Accepted(Hibernatable ws, Native& native, IoContext& context);
492 
493 ~Accepted() noexcept(false);
494 
495 // A simple wrapper to make it easier to access the underlying kj::WebSocket.
496 class WrappedWebSocket {
497 public:
498 explicit WrappedWebSocket(Hibernatable ws);
499 explicit WrappedWebSocket(kj::Own<kj::WebSocket> ws);
500 
501 kj::WebSocket* operator->();
502 
503 kj::WebSocket& operator*();
504 
505 kj::Maybe<kj::Own<kj::WebSocket>&> getIfNotHibernatable();
506 kj::Maybe<Hibernatable&> getIfHibernatable();
507 kj::Array<kj::StringPtr> getHibernatableTags();
508 
509 // Transitions our Hibernatable websocket to a "Releasing" state.
510 // The websocket will transition to `Released` when convenient.
511 void initiateHibernatableRelease(jsg::Lock& js,
512 kj::Own<kj::WebSocket> ws,
513 kj::Array<kj::String> tags,
514 HibernatableReleaseState state);
515 
516 bool isAwaitingRelease();
517 bool isAwaitingError();
518 
519 private:
520 kj::OneOf<kj::Own<kj::WebSocket>, Hibernatable> inner;
521 };
522 
523 WrappedWebSocket ws;
524 
525 bool isHibernatable();
526 
527 kj::Promise<void> createAbortTask(Native& native, IoContext& context);
528 // Listens for ws->whenAborted() and possibly triggers a proactive shutdown.
529 kj::Promise<void> whenAbortedTask = nullptr;
530 
531 kj::Maybe<kj::Own<ActorObserver>> actorMetrics;
532 
533 // This canceler wraps the pump loop as a precaution to make sure we can't exit the Accepted
534 // state with a pump task still happening asynchronously. In practice the canceler should usually
535 // be empty when destroyed because we do not leave the Accepted state if we're still pumping.
536 // Even in the case of IoContext premature cancellation, the pump task should be canceled
537 // by the IoContext before the Canceler is destroyed.
538 kj::Canceler canceler;
539 };
540 
541 struct Released {};
542 using NativeState =
543 kj::OneOf<AwaitingConnection, AwaitingAcceptanceOrCoupling, Accepted, Released>;
544 friend kj::StringPtr KJ_STRINGIFY(const NativeState&);
545 
546 struct Native {
547 NativeState state;
548 
549 // Is there currently a task running to pump outgoing messages?
550 bool isPumping = false;
551 
552 // Has a Close message been enqueued for send? (It may still be in outgoingMessages. Check
553 // closedOutgoing && !isPumping to check if it has gone out.)
554 bool closedOutgoing = false;
555 
556 // Has a Close message been received, or has a premature disconnection occurred?
557 bool closedIncoming = false;
558 
559 // Have we detected that the peer has stopped accepting messages? We may want to clean up more
560 // proactively in this case.
561 bool outgoingAborted = false;
562 };
563 
564 // The underlying native WebSocket (or a promise that will emplace one).
565 //
566 // The state transitions look like so:
567 // - Starts as `AwaitingConnection` if the `WebSocket(url ...)` ctor is used.
568 // - Starts as `AwaitingAcceptanceOrCoupling` if the `WebSocket(native)` ctor is used.
569 // - Transitions from `AwaitingConnection` to `AwaitingAcceptanceOrCoupling` when the native
570 // connection is established and to `Accepted` once the read loop starts.
571 // - Transitions from `AwaitingConnection` to `Released` when connection establishment fails.
572 // - Transitions from `AwaitingAcceptanceOrCoupling` to `Accepted` when it is accepted.
573 // - Transitions from `AwaitingAcceptanceOrCoupling` to `Released` when it is coupled to another
574 // web socket.
575 // - Transitions from `Accepted` to `Released` when outgoing pump is done and either both
576 // directions have seen "close" messages or an error has occurred.
577 IoOwn<Native> farNative;
578 
579 // If any error has occurred.
580 kj::Maybe<jsg::JsRef<jsg::JsValue>> error;
581 
582 struct GatedMessage {
583 kj::Maybe<kj::Promise<void>> outputLock; // must wait for this before actually sending
584 kj::WebSocket::Message message;
585 size_t pendingAutoResponses = 0;
586 };
587 using OutgoingMessagesMap = kj::Table<GatedMessage, kj::InsertionOrderIndex>;
588 // Queue of messages to be sent. This is wrapped in an IoOwn so that the pump loop can safely
589 // access the map without locking the isolate.
590 IoOwn<OutgoingMessagesMap> outgoingMessages;
591 
592 // Keep track of current hibernatable websockets auto-response status to avoid racing
593 // between regular websocket messages, and auto-responses.
594 struct AutoResponse {
595 // The ongoing auto-response promise, used for pump() synchronization.
596 // Wrapped in IoOwn when an IoContext is available (for GC-safe destruction that avoids
597 // violating DISALLOW_KJ_IO_DESTRUCTORS_SCOPE), or plain kj::Own when no IoContext is
598 // available (e.g. when called from the hibernation manager's readLoop).
599 using OwnedAutoResponsePromise =
600 kj::OneOf<IoOwn<kj::Promise<void>>, kj::Own<kj::Promise<void>>>;
601 kj::Maybe<OwnedAutoResponsePromise> ongoingAutoResponse;
602 workerd::util::Queue<kj::String> pendingAutoResponseDeque;
603 size_t queuedAutoResponses = 0;
604 bool isPumping = false;
605 bool isClosed = false;
606 
607 JSG_MEMORY_INFO(AutoResponse) {
608 tracker.trackFieldWithSize("ongoingAutoResponse", sizeof(kj::Promise<void>));
609 pendingAutoResponseDeque.forEach(
610 [&](const kj::String& message) { tracker.trackField(nullptr, message); });
611 }
612 };
613 
614 AutoResponse autoResponseStatus;
615 
616 kj::Maybe<kj::Own<WebSocketObserver>> observer;
617 
618 // Contains a websocket and possibly some data from the WebSocketResponse headers.
619 struct PackedWebSocket {
620 kj::Own<kj::WebSocket> ws;
621 kj::Maybe<kj::String> proto;
622 kj::Maybe<kj::String> extensions;
623 };
624 
625 // So that each end of a WebSocketPair can keep track of its peer.
626 // We use a weak ref to track the peer to avoid having a strong ref cycle
627 // between the two WebSocket instances that would cause them to leak. This
628 // can mean, however, that it's possible for one of the peers to be garbage
629 // collected while the other still exists. This should be fairly unusual tho.
630 kj::Maybe<kj::Own<WeakRef<WebSocket>>> peer;
631 
632 void visitForGc(jsg::GcVisitor& visitor) {
633 visitor.visit(error);
634 }
635 
636 void setPeer(kj::Own<WeakRef<WebSocket>> peer);
637 
638 friend jsg::Ref<WebSocketPair> WebSocketPair::constructor(jsg::Lock&);
639 
640 void dispatchOpen(jsg::Lock& js);
641 
642 void ensurePumping(jsg::Lock& js);
643 
644 // Returns the number of pending auto-responses that should be sent before the next outgoing
645 // message, and advances the queuedAutoResponses counter. Called each time a GatedMessage is
646 // inserted into outgoingMessages to guarantee auto-response ordering.
647 size_t getPendingAutoResponseCount();
648 
649 // Write messages from `outgoingMessages` into `ws`.
650 //
651 // These are not necessarily called under isolate lock, but they are called on the given
652 // context's thread. They are declared `static` to prove they don't access the JavaScript
653 // object's members in a thread-unsafe way. `outgoingMessages` and `ws` are both `IoOwn`ed
654 // objects so are safe to access from the thread without the isolate lock. The whole task is
655 // owned by the `IoContext` so it'll be canceled if the `IoContext` is destroyed.
656 static kj::Promise<void> pump(IoContext& context,
657 OutgoingMessagesMap& outgoingMessages,
658 kj::WebSocket& ws,
659 Native& native,
660 AutoResponse& autoResponse,
661 kj::Maybe<kj::Own<WebSocketObserver>>& observer);
662 
663 kj::Promise<kj::Maybe<kj::Exception>> readLoop(
664 kj::Maybe<kj::Own<InputGate::CriticalSection>> cs, size_t maxMessageSize);
665 
666 void reportError(jsg::Lock& js, kj::Exception&& e);
667 void reportError(jsg::Lock& js, jsg::JsRef<jsg::JsValue> err);
668 
669 void assertNoError(jsg::Lock& js);
670};
671 
672#define EW_WEBSOCKET_ISOLATE_TYPES \
673 api::CloseEvent, api::CloseEvent::Initializer, api::WebSocket, api::WebSocket::AcceptOptions, \
674 api::WebSocketPair, api::WebSocketPair::PairIterator, \
675 api::WebSocketPair::PairIterator:: \
676 Next // The list of websocket.h types that are added to worker.c++'s JSG_DECLARE_ISOLATE_TYPE
677 
678} // namespace workerd::api