Skip to content
File

Blob: src/workerd/api/web-socket.c++

54.3 KB
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#include "web-socket.h"
6 
7#include "blob.h"
8#include "events.h"
9#include "messagechannel.h"
10#include "util.h"
11 
12#include <workerd/io/features.h>
13#include <workerd/io/io-context.h>
14#include <workerd/io/worker.h>
15#include <workerd/jsg/jsg.h>
16#include <workerd/jsg/ser.h>
17#include <workerd/util/autogate.h>
18#include <workerd/util/sentry.h>
19 
20#include <kj/compat/url.h>
21 
22namespace workerd::api {
23 
24kj::StringPtr KJ_STRINGIFY(const WebSocket::NativeState& state) {
25 // TODO(someday) We might care more about this `OneOf` than its which, that probably means
26 // returning a kj::String instead.
27 KJ_SWITCH_ONEOF(state) {
28 KJ_CASE_ONEOF(ac, WebSocket::AwaitingConnection) return "AwaitingConnection";
29 KJ_CASE_ONEOF(aaoc, WebSocket::AwaitingAcceptanceOrCoupling)
30 return "AwaitingAcceptanceOrCoupling";
31 KJ_CASE_ONEOF(a, WebSocket::Accepted) return "Accepted";
32 KJ_CASE_ONEOF(r, WebSocket::Released) return "Released";
33 }
34 KJ_UNREACHABLE;
35}
36 
37IoOwn<WebSocket::Native> WebSocket::initNative(IoContext& ioContext,
38 kj::WebSocket& ws,
39 kj::Array<kj::StringPtr> tags,
40 bool closedOutgoingConn) {
41 auto nativeObj = kj::heap<Native>();
42 nativeObj->state.init<Accepted>(
43 Accepted::Hibernatable{.ws = ws, .tagsRef = kj::mv(tags)}, *nativeObj, ioContext);
44 // We might have called `close()` when this WebSocket was previously active.
45 // If so, we want to prevent any future calls to `send()`.
46 nativeObj->closedOutgoing = closedOutgoingConn;
47 autoResponseStatus.isClosed = nativeObj->closedOutgoing;
48 return ioContext.addObject(kj::mv(nativeObj));
49}
50 
51WebSocket::WebSocket(
52 jsg::Lock& js, IoContext& ioContext, kj::WebSocket& ws, HibernationPackage package)
53 : weakRef(kj::refcounted<WeakRef<WebSocket>>(kj::Badge<WebSocket>{}, *this)),
54 url(kj::mv(package.url)),
55 protocol(kj::mv(package.protocol)),
56 extensions(kj::mv(package.extensions)),
57 binaryType_(FeatureFlags::get(js).getWebsocketBinaryTypeDefault() ? BinaryType::BLOB
58 : BinaryType::ARRAYBUFFER),
59 serializedAttachment(kj::mv(package.serializedAttachment)),
60 allowHalfOpen(package.allowHalfOpen),
61 farNative(initNative(ioContext,
62 ws,
63 kj::mv(KJ_REQUIRE_NONNULL(package.maybeTags)),
64 package.closedOutgoingConnection)),
65 outgoingMessages(IoContext::current().addObject(kj::heap<OutgoingMessagesMap>())) {}
66// This constructor is used when reinstantiating a websocket that had been hibernating, which is
67// why we can go straight to the Accepted state. However, note that we are actually in the
68// `Hibernatable` "sub-state"!
69 
70jsg::Ref<WebSocket> WebSocket::hibernatableFromNative(
71 jsg::Lock& js, kj::WebSocket& ws, HibernationPackage package) {
72 return js.alloc<WebSocket>(js, IoContext::current(), ws, kj::mv(package));
73}
74 
75WebSocket::WebSocket(jsg::Lock& js, kj::Own<kj::WebSocket> native)
76 : weakRef(kj::refcounted<WeakRef<WebSocket>>(kj::Badge<WebSocket>{}, *this)),
77 url(kj::none),
78 binaryType_(FeatureFlags::get(js).getWebsocketBinaryTypeDefault() ? BinaryType::BLOB
79 : BinaryType::ARRAYBUFFER),
80 allowHalfOpen(!FeatureFlags::get(js).getWebSocketAutoReplyToClose()),
81 farNative(nullptr),
82 outgoingMessages(IoContext::current().addObject(kj::heap<OutgoingMessagesMap>())) {
83 auto nativeObj = kj::heap<Native>();
84 nativeObj->state.init<AwaitingAcceptanceOrCoupling>(kj::mv(native));
85 farNative = IoContext::current().addObject(kj::mv(nativeObj));
86}
87 
88WebSocket::WebSocket(jsg::Lock& js, kj::String url)
89 : weakRef(kj::refcounted<WeakRef<WebSocket>>(kj::Badge<WebSocket>{}, *this)),
90 url(kj::mv(url)),
91 binaryType_(FeatureFlags::get(js).getWebsocketBinaryTypeDefault() ? BinaryType::BLOB
92 : BinaryType::ARRAYBUFFER),
93 allowHalfOpen(!FeatureFlags::get(js).getWebSocketAutoReplyToClose()),
94 farNative(nullptr),
95 outgoingMessages(IoContext::current().addObject(kj::heap<OutgoingMessagesMap>())) {
96 auto nativeObj = kj::heap<Native>();
97 nativeObj->state.init<AwaitingConnection>();
98 farNative = IoContext::current().addObject(kj::mv(nativeObj));
99}
100 
101void WebSocket::initConnection(jsg::Lock& js, kj::Promise<PackedWebSocket> prom) {
102 
103 auto& canceler = KJ_ASSERT_NONNULL(farNative->state.tryGet<AwaitingConnection>()).canceler;
104 
105 IoContext::current()
106 .awaitIo(js, canceler.wrap(kj::mv(prom)),
107 [this, self = JSG_THIS](jsg::Lock& js, PackedWebSocket packedSocket) mutable {
108 auto& native = *farNative;
109 KJ_IF_SOME(pending, native.state.tryGet<AwaitingConnection>()) {
110 // We've successfully established our web socket, we do not need to cancel anything.
111 pending.canceler.release();
112 }
113 
114 native.state.init<AwaitingAcceptanceOrCoupling>(
115 AwaitingAcceptanceOrCoupling{IoContext::current().addObject(kj::mv(packedSocket.ws))});
116 
117 // both `protocol` and `extensions` start off as empty strings.
118 // They become null if the connection is established and no protocol/extension was chosen.
119 // https://html.spec.whatwg.org/multipage/web-sockets.html#dom-websocket-protocol
120 KJ_IF_SOME(proto, packedSocket.proto) {
121 protocol = kj::mv(proto);
122 } else {
123 protocol = kj::none;
124 }
125 
126 KJ_IF_SOME(ext, packedSocket.extensions) {
127 extensions = kj::mv(ext);
128 } else {
129 extensions = kj::none;
130 }
131 
132 // Fire open event.
133 internalAccept(js, IoContext::current().getCriticalSection());
134 dispatchOpen(js);
135 }).catch_(js, [this, self = JSG_THIS](jsg::Lock& js, jsg::Value&& e) mutable {
136 // Fire error event.
137 // Sets readyState to CLOSING.
138 farNative->closedIncoming = true;
139 
140 // Sets readyState to CLOSED.
141 reportError(js, jsg::JsValue(e.getHandle(js)).addRef(js));
142 
143 dispatchEventImpl(
144 js, js.alloc<CloseEvent>(1006, kj::str("Failed to establish websocket connection"), false));
145 });
146 // Note that in this attach we pass a strong reference to the WebSocket. The reference will be
147 // dropped when either the connection promise completes or the IoContext is torn down,
148 // whichever comes first.
149}
150 
151namespace {
152 
153// See item 10 of https://datatracker.ietf.org/doc/html/rfc6455#section-4.1
154bool validProtoToken(const kj::StringPtr protocol) {
155 if (kj::size(protocol) == 0) {
156 return false;
157 }
158 
159 for (auto& c: protocol) {
160 // Note that this also includes separators 0x20 (SP) and 0x09 (HT), so we don't need to check
161 // for them below.
162 if (c < 0x21 || 0x7E < c) {
163 return false;
164 }
165 
166 switch (c) {
167 case '(':
168 case ')':
169 case '<':
170 case '>':
171 case '@':
172 case ',':
173 case ';':
174 case ':':
175 case '\\':
176 case '/':
177 case '[':
178 case ']':
179 case '?':
180 case '=':
181 case '{':
182 case '}':
183 return false;
184 default:
185 break;
186 }
187 }
188 return true;
189}
190 
191} // namespace
192 
193jsg::Ref<WebSocket> WebSocket::constructor(jsg::Lock& js,
194 kj::String url,
195 jsg::Optional<kj::OneOf<kj::Array<kj::String>, kj::String>> protocols) {
196 
197 auto& context = IoContext::current();
198 
199 // Check if we have a valid URL
200 constexpr auto urlOptions = kj::Url::Options{.percentDecode = false, .allowEmpty = true};
201 constexpr auto wsErr = "WebSocket Constructor: "_kj;
202 
203 // To be compatible with fetch() implementation:
204 // - First parse the URL with REMOTE_HREF which requires hostname.
205 // - Then stringify the URL with HTTP_PROXY_REQUEST which omits the userinfo.
206 kj::Url urlRecord = JSG_REQUIRE_NONNULL(kj::Url::tryParse(url, kj::Url::REMOTE_HREF, urlOptions),
207 DOMSyntaxError, wsErr, "The url is invalid.");
208 
209 JSG_REQUIRE(urlRecord.scheme == "ws" || urlRecord.scheme == "wss", DOMSyntaxError, wsErr,
210 "The url scheme must be ws or wss.");
211 // We want the caller to pass `ws/wss` as per the spec, but FL would treat these as http in
212 // `X-Forwarded-Proto`, so we want to ensure that `wss` results in `https`, not `http`.
213 if (urlRecord.scheme == "ws") {
214 urlRecord.scheme = kj::str("http");
215 } else if (urlRecord.scheme == "wss") {
216 urlRecord.scheme = kj::str("https");
217 }
218 
219 JSG_REQUIRE(
220 urlRecord.fragment == kj::none, DOMSyntaxError, wsErr, "The url fragment must be empty.");
221 
222 kj::HttpHeaders headers(context.getHeaderTable());
223 
224 // Set protocols header if necessary.
225 KJ_IF_SOME(variant, protocols) {
226 // String consisting of the protocol(s) we send to the server.
227 kj::Maybe<kj::String> maybeProtoString;
228 
229 KJ_SWITCH_ONEOF(variant) {
230 KJ_CASE_ONEOF(proto, kj::String) {
231 JSG_REQUIRE(
232 validProtoToken(proto), DOMSyntaxError, wsErr, "The protocol header token is invalid.");
233 maybeProtoString = kj::mv(proto);
234 }
235 KJ_CASE_ONEOF(protoArr, kj::Array<kj::String>) {
236 // Per the WebSocket spec, an empty protocols array is valid and equivalent to not
237 // specifying any protocols - we simply don't set the Sec-WebSocket-Protocol header.
238 if (protoArr.size() > 0) {
239 // Search for duplicates by checking for their presence in the set.
240 kj::HashSet<kj::String> present;
241 
242 for (const auto& proto: protoArr) {
243 JSG_REQUIRE(validProtoToken(proto), DOMSyntaxError, wsErr,
244 "One of the protocol header tokens is invalid.");
245 JSG_REQUIRE(!present.contains(proto), DOMSyntaxError, wsErr,
246 "The protocols header cannot have repeating values.");
247 
248 present.insert(kj::str(proto));
249 }
250 constexpr auto delim = ", "_kj;
251 maybeProtoString = kj::str(kj::delimited(protoArr, delim));
252 }
253 }
254 }
255 KJ_IF_SOME(protoString, maybeProtoString) {
256 auto protoHeaderId = context.getHeaderIds().secWebSocketProtocol;
257 headers.set(protoHeaderId, kj::mv(protoString));
258 }
259 }
260 
261 // Any userinfo, username and/or password, should be removed.
262 // Users should use Authorization header for this purpose.
263 kj::String connUrl =
264 uriEncodeControlChars(urlRecord.toString(kj::Url::HTTP_PROXY_REQUEST).asBytes());
265 auto ws = js.alloc<WebSocket>(js, kj::mv(url));
266 
267 headers.set(kj::HttpHeaderId::SEC_WEBSOCKET_EXTENSIONS, kj::str("permessage-deflate"));
268 // By default, browsers set the compression extension header for `new WebSocket()`.
269 
270 if (!FeatureFlags::get(js).getWebSocketCompression()) {
271 // If we haven't enabled the websocket compression compatibility flag, strip the header from the
272 // subrequest.
273 headers.unset(kj::HttpHeaderId::SEC_WEBSOCKET_EXTENSIONS);
274 }
275 
276 auto client = context.getHttpClient(0, false, kj::none, "websocket_open"_kjc);
277 auto prom =
278 ([](auto& context, auto connUrl, auto headers, auto client) -> kj::Promise<PackedWebSocket> {
279 auto response = co_await client->openWebSocket(connUrl, headers);
280 
281 JSG_REQUIRE(response.statusCode == 101, TypeError,
282 "Failed to establish the WebSocket connection: expected server to reply with HTTP "
283 "status code 101 (switching protocols), but received ",
284 response.statusCode, " instead.");
285 
286 KJ_SWITCH_ONEOF(response.webSocketOrBody) {
287 KJ_CASE_ONEOF(webSocket, kj::Own<kj::WebSocket>) {
288 auto maybeProtoPtr = response.headers->get(context.getHeaderIds().secWebSocketProtocol);
289 auto maybeExtensionsPtr = response.headers->get(kj::HttpHeaderId::SEC_WEBSOCKET_EXTENSIONS);
290 
291 kj::Maybe<kj::String> maybeProto;
292 kj::Maybe<kj::String> maybeExtensions;
293 
294 KJ_IF_SOME(proto, maybeProtoPtr) {
295 maybeProto = kj::str(proto);
296 }
297 
298 KJ_IF_SOME(extensions, maybeExtensionsPtr) {
299 maybeExtensions = kj::str(extensions);
300 }
301 
302 co_return PackedWebSocket{.ws = webSocket.attach(kj::mv(client)),
303 .proto = kj::mv(maybeProto),
304 .extensions = kj::mv(maybeExtensions)};
305 }
306 KJ_CASE_ONEOF(body, kj::Own<kj::AsyncInputStream>) {
307 JSG_FAIL_REQUIRE(
308 TypeError, "Worker received body in a response to a request for a WebSocket.");
309 }
310 }
311 KJ_UNREACHABLE
312 })(context, kj::mv(connUrl), kj::mv(headers), kj::mv(client));
313 
314 ws->initConnection(js, kj::mv(prom));
315 
316 return ws;
317}
318 
319kj::Promise<DeferredProxy<void>> WebSocket::couple(
320 kj::Own<kj::WebSocket> other, RequestObserver& request) {
321 auto& native = *farNative;
322 JSG_REQUIRE(!native.state.is<AwaitingConnection>(), TypeError,
323 "Can't return WebSocket in a Response if it was created with `new WebSocket()`");
324 JSG_REQUIRE(!native.state.is<Released>(), TypeError,
325 "Can't return WebSocket that was already used in a response.");
326 KJ_IF_SOME(state, native.state.tryGet<Accepted>()) {
327 if (state.isHibernatable()) {
328 JSG_FAIL_REQUIRE(
329 TypeError, "Can't return WebSocket in a Response after calling acceptWebSocket().");
330 } else {
331 JSG_FAIL_REQUIRE(TypeError, "Can't return WebSocket in a Response after calling accept().");
332 }
333 }
334 
335 // Tear down the IoOwn since we now need to extend the WebSocket to a `DeferredProxy` promise.
336 // This works because the `DeferredProxy` ends on the same event loop, but after the request
337 // context goes away.
338 kj::Own<kj::WebSocket> self =
339 kj::mv(KJ_ASSERT_NONNULL(native.state.tryGet<AwaitingAcceptanceOrCoupling>()).ws);
340 native.state.init<Released>();
341 
342 auto& context = IoContext::current();
343 
344 auto upstream = other->pumpTo(*self);
345 auto downstream = self->pumpTo(*other);
346 
347 auto tryGetPeer = [&]() -> kj::Maybe<WebSocket&> {
348 KJ_IF_SOME(p, peer) {
349 return p->tryGet();
350 }
351 return kj::none;
352 };
353 auto isHibernatable = [&](workerd::api::WebSocket& ws) {
354 KJ_IF_SOME(state, ws.farNative->state.tryGet<Accepted>()) {
355 return state.isHibernatable();
356 }
357 return false;
358 };
359 KJ_IF_SOME(p, tryGetPeer()) {
360 // We're terminating the WebSocket in this worker, so the upstream promise (which pumps
361 // messages from the client to this worker) counts as something the request is waiting for.
362 upstream = upstream.attach(context.registerPendingEvent());
363 
364 // We can observe websocket traffic in both directions by attaching an observer to the peer
365 // websocket which terminates in the worker.
366 KJ_IF_SOME(observer, request.tryCreateWebSocketObserver()) {
367 p.observer = kj::mv(observer);
368 }
369 }
370 
371 // We need to use `eagerlyEvaluate()` on both inputs to `joinPromises` to work around the awkward
372 // behavior of `joinPromises` lazily-evaluating tail continuations.
373 auto promise = kj::joinPromises(
374 kj::arr(upstream.eagerlyEvaluate(nullptr), downstream.eagerlyEvaluate(nullptr)))
375 .attach(kj::mv(self), kj::mv(other));
376 
377 KJ_IF_SOME(peer, tryGetPeer()) {
378 // Since the WebSocket is terminated locally, we generally want the request and associated
379 // IoContext to stay alive until the WebSocket connection has terminated.
380 //
381 // However, there is one exception to this: when the WebSocket is hibernatable, we don't want
382 // the existence of this connection to prevent the actor from being evicted, so we fall through
383 // to deferred proxying in this case.
384 if (!isHibernatable(peer)) {
385 co_await promise;
386 co_return;
387 }
388 }
389 
390 // Either:
391 // 1. This websocket is just proxying through, in which case we can allow the IoContext to go
392 // away while still being able to successfully pump the websocket connection.
393 // 2. This is a hibernatable websocket and we are falling through to deferred proxying to
394 // potentially allow for hibernation to occur.
395 
396 // To begin deferred proxying, we can use this magic `KJ_CO_MAGIC` expression, which fulfills
397 // our outer promise for a DeferredProxy<void>, which wraps a promise for the rest of this
398 // coroutine.
399 KJ_CO_MAGIC BEGIN_DEFERRED_PROXYING;
400 
401 co_return co_await promise;
402}
403 
404void WebSocket::accept(jsg::Lock& js, jsg::Optional<AcceptOptions> options) {
405 auto& native = *farNative;
406 JSG_REQUIRE(!native.state.is<AwaitingConnection>(), TypeError,
407 "Websockets obtained from the 'new WebSocket()' constructor cannot call accept");
408 JSG_REQUIRE(!native.state.is<Released>(), TypeError,
409 "Can't accept() WebSocket that was already used in a response.");
410 
411 KJ_IF_SOME(accepted, native.state.tryGet<Accepted>()) {
412 JSG_REQUIRE(!accepted.isHibernatable(), TypeError,
413 "Can't accept() WebSocket after enabling hibernation.");
414 // Technically, this means it's OK to invoke `accept()` once a `new WebSocket()` resolves to
415 // an established connection. This is probably okay? It might spare the worker devs a class of
416 // errors they do not care care about.
417 return;
418 }
419 
420 KJ_IF_SOME(opts, options) {
421 KJ_IF_SOME(value, opts.allowHalfOpen) {
422 allowHalfOpen = AllowHalfOpen(value);
423 }
424 }
425 
426 internalAccept(js, IoContext::current().getCriticalSection());
427}
428 
429void WebSocket::internalAccept(jsg::Lock& js, kj::Maybe<kj::Own<InputGate::CriticalSection>> cs) {
430 auto& native = *farNative;
431 auto nativeWs = kj::mv(KJ_ASSERT_NONNULL(native.state.tryGet<AwaitingAcceptanceOrCoupling>()).ws);
432 native.state.init<Accepted>(kj::mv(nativeWs), native, IoContext::current());
433 return startReadLoop(js, kj::mv(cs));
434}
435 
436WebSocket::Accepted::Accepted(kj::Own<kj::WebSocket> wsParam, Native& native, IoContext& context)
437 : ws(kj::mv(wsParam)),
438 whenAbortedTask(createAbortTask(native, context)) {
439 KJ_IF_SOME(a, context.getActor()) {
440 auto& metrics = a.getMetrics();
441 metrics.webSocketAccepted();
442 
443 // Save the metrics object for the destructor since the IoContext may not be accessible
444 // there.
445 actorMetrics = kj::addRef(metrics);
446 }
447}
448 
449WebSocket::Accepted::Accepted(Hibernatable wsParam, Native& native, IoContext& context)
450 : ws(kj::mv(wsParam)),
451 whenAbortedTask(createAbortTask(native, context)) {
452 KJ_IF_SOME(a, context.getActor()) {
453 auto& metrics = a.getMetrics();
454 metrics.webSocketAccepted();
455 
456 // Save the metrics object for the destructor since the IoContext may not be accessible
457 // there.
458 actorMetrics = kj::addRef(metrics);
459 }
460}
461 
462kj::Promise<void> WebSocket::Accepted::createAbortTask(Native& native, IoContext& context) {
463 try {
464 // whenAborted() is theoretically not supposed to throw, but some code paths, like
465 // AbortableWebSocket and Cap'n Proto disconnects, may end up throwing DISCONNECTED. Treat
466 // exceptions the same as if `whenAborted()` finished normally -- but log in catch catch
467 // block if it's not DISCONNECTED.
468 co_await ws->whenAborted();
469 
470 // Other end disconnected prematurely. We may be able to clean up our state.
471 native.outgoingAborted = true;
472 if (!native.isPumping && native.closedIncoming) {
473 // We can safely destroy the underlying WebSocket as it is no longer in use.
474 // HACK: Replacing the state will delete `whenAbortedTask`, which is the task that is
475 // currently executing, which will crash. We know we're at the end of the task here
476 // so detach it as a work-around.
477 whenAbortedTask.detach([](auto&&) {});
478 native.state.init<Released>();
479 } else {
480 // Either we haven't received the incoming disconnect yet, or there are writes
481 // in-flight. In either case, we need to wait for those to happen before we destroy the
482 // underlying object, or we might have a UAF situation. Those other operations should
483 // fail shortly and notice the `outgoingAborted` flag when they do.
484 }
485 } catch (...) {
486 auto ex = kj::getCaughtExceptionAsKj();
487 if (ex.getType() != kj::Exception::Type::DISCONNECTED) {
488 LOG_EXCEPTION("webSocketWhenAborted", ex);
489 }
490 }
491}
492 
493WebSocket::Accepted::~Accepted() noexcept(false) {
494 KJ_IF_SOME(a, actorMetrics) {
495 a.get()->webSocketClosed();
496 }
497}
498 
499// Default max WebSocket message size limit. Note that kj-http's own default is 1MiB
500// (`kj::WebSocket::SUGGESTED_MAX_MESSAGE_SIZE`). We've found this to be too small for many commmon
501// use cases, such as proxying Chrome Devtools Protocol messages.
502//
503// JS-RPC messages are size-limited to 32MiB, and it seems to be working well, so we're setting the
504// WebSocket default max message size to match that.
505static constexpr size_t WEBSOCKET_MAX_MESSAGE_SIZE = 32u << 20;
506 
507void WebSocket::startReadLoop(jsg::Lock& js, kj::Maybe<kj::Own<InputGate::CriticalSection>> cs) {
508 size_t maxMessageSize = WEBSOCKET_MAX_MESSAGE_SIZE;
509 if (FeatureFlags::get(js).getIncreaseWebsocketMessageSize()) {
510 maxMessageSize = 128u << 20;
511 }
512 
513 // If the kj::WebSocket happens to be an AbortableWebSocket (see util/abortable.h), then
514 // calling readLoop here could throw synchronously if the canceler has already been tripped.
515 // Using kj::evalNow() here let's us capture that and handle correctly.
516 //
517 // We catch exceptions and return Maybe<Exception> instead since we want to handle the exceptions
518 // in awaitIo() below, but we don't want the KJ exception converted to JavaScript before we can
519 // examine it.
520 kj::Promise<kj::Maybe<kj::Exception>> promise = readLoop(kj::mv(cs), maxMessageSize);
521 
522 auto& context = IoContext::current();
523 
524 auto hasLocalPeer = [&]() {
525 KJ_IF_SOME(p, peer) {
526 if (p->isValid()) {
527 return true;
528 }
529 }
530 return false;
531 };
532 if (!hasLocalPeer()) {
533 promise = promise.attach(context.registerPendingEvent());
534 }
535 
536 // We put the read loop in a `waitUntil`, since there would otherwise be a race condition between
537 // delivering the final close message and the request being canceled due to client disconnect.
538 // This `waitUntil` will not significantly extend the lifetime of the request in practice, as the
539 // request otherwise ends when the client disconnects, and the read loop will also end when the
540 // client disconnects -- we just want to ensure that they happen in the right order.
541 //
542 // TODO(bug): Using waitUntil() for this purpose is only correct for WebSockets originating from
543 // the eyeball. For an outgoing WebSocket, we should just do addTask(). Alternatively, perhaps
544 // we need to adjust the cancellation logic to wait for whenThreadIdle() before cancelling,
545 // which would then allow close messages to be delivered from eyeball connections without any
546 // use of waitUntil().
547 //
548 // TODO(cleanup): We have to use awaitIoLegacy() so that we can handle registerPendingEvent()
549 // manually. Ideally, we'd refactor things such that a WebSocketPair where both ends are
550 // accepted locally is implemented completely in JavaScript space, using jsg::Promise instead
551 // of kj::Promise, and then only use awaitIo() on truly remote WebSockets.
552 // TODO(cleanup): Should addWaitUntil() take jsg::Promise instead of kj::Promise?
553 context.addWaitUntil(context.awaitJs(js,
554 context.awaitIoLegacy(js, kj::mv(promise))
555 .then(js,
556 [this, thisHandle = JSG_THIS](
557 jsg::Lock& js, kj::Maybe<kj::Exception>&& maybeError) mutable {
558 auto& native = *farNative;
559 KJ_IF_SOME(e, maybeError) {
560 if (!native.closedIncoming && e.getType() == kj::Exception::Type::DISCONNECTED) {
561 // Report premature disconnect or cancel as a close event.
562 dispatchEventImpl(js,
563 js.alloc<CloseEvent>(
564 1006, kj::str("WebSocket disconnected without sending Close frame."), false));
565 native.closedIncoming = true;
566 // If there are no further messages to send, so we can discard the underlying connection.
567 tryReleaseNative(js);
568 } else {
569 native.closedIncoming = true;
570 reportError(js, e.clone());
571 }
572 }
573 })));
574}
575 
576void WebSocket::send(jsg::Lock& js, kj::OneOf<kj::Array<byte>, kj::String> message) {
577 auto& native = *farNative;
578 JSG_REQUIRE(!native.closedOutgoing, TypeError, "Can't call WebSocket send() after close().");
579 if (native.outgoingAborted || native.state.is<Released>()) {
580 // Per the spec, we should silently ignore send()s that happen after the connection is closed.
581 // NOTE: The spec claims send() should also silently ignore messages sent after a close message
582 // has been sent or received cleanly. We ignore this advice:
583 // * If close has been sent, i.e. close() has been called, then calling send() is clearly a
584 // bug, and we'd like to help people debug, so we throw an exception above. (This point is
585 // debatable, we could change it.)
586 // * It makes no sense that *receiving* a close message should prevent further calls to send().
587 // The spec seems broken here. What if you need to send a couple final messages for a clean
588 // shutdown?
589 return;
590 } else if (awaitingHibernatableError()) {
591 // Ready for the hibernatable error event state, after encountering an error, the websocket
592 // isn't able to send outbound messages; let's release it.
593 tryReleaseNative(js);
594 return;
595 }
596 
597 JSG_REQUIRE(native.state.is<Accepted>(), TypeError,
598 "You must call one of accept() or state.acceptWebSocket() on this WebSocket before sending "
599 "messages.");
600 
601 auto maybeOutputLock = IoContext::current().waitForOutputLocksIfNecessary();
602 auto msg = [&]() -> kj::WebSocket::Message {
603 KJ_SWITCH_ONEOF(message) {
604 KJ_CASE_ONEOF(text, kj::String) {
605 return kj::mv(text);
606 break;
607 }
608 KJ_CASE_ONEOF(data, kj::Array<byte>) {
609 return kj::mv(data);
610 break;
611 }
612 }
613 
614 KJ_UNREACHABLE;
615 }();
616 
617 outgoingMessages->insert(
618 GatedMessage{kj::mv(maybeOutputLock), kj::mv(msg), getPendingAutoResponseCount()});
619 
620 ensurePumping(js);
621}
622 
623void WebSocket::close(
624 jsg::Lock& js, jsg::Optional<int> code, jsg::Optional<jsg::USVString> reason) {
625 auto& native = *farNative;
626 
627 // Per the spec, close code and reason validation must happen before any readyState checks.
628 // See https://websockets.spec.whatwg.org/#dom-websocket-close step 1.
629 KJ_IF_SOME(c, code) {
630 if (FeatureFlags::get(js).getPedanticWpt()) {
631 // The WHATWG WebSocket spec allows only 1000 (Normal Closure) or the range 3000-4999
632 // (reserved for libraries/frameworks and private use, per RFC 6455 Section 7.4.2).
633 JSG_REQUIRE(c == 1000 || (c >= 3000 && c <= 4999), DOMInvalidAccessError,
634 "Invalid WebSocket close code: ", c, ".");
635 } else {
636 // Legacy behavior: accept the full 1000-4999 range, only rejecting four codes that
637 // RFC 6455 Section 7.4.1 reserves and forbids endpoints from sending in a Close frame:
638 // 1004 - Reserved (no defined meaning)
639 // 1005 - No Status Rcvd (only used internally to indicate no code was present)
640 // 1006 - Abnormal Closure (only used internally when connection drops without a Close frame)
641 // 1015 - TLS Handshake failure (only used internally, never sent over the wire)
642 JSG_REQUIRE(c >= 1000 && c < 5000 && c != 1004 && c != 1005 && c != 1006 && c != 1015,
643 DOMInvalidAccessError, "Invalid WebSocket close code: ", c, ".");
644 }
645 }
646 
647 // Per the WHATWG WebSocket spec and RFC 6455 Section 5.5, the close frame body must not exceed
648 // 125 bytes (2-byte status code + up to 123 bytes of reason). Throw a SyntaxError if the
649 // reason's UTF-8 encoding is longer than 123 bytes.
650 if (FeatureFlags::get(js).getWebsocketCloseReasonByteLimit()) {
651 KJ_IF_SOME(r, reason) {
652 JSG_REQUIRE(r.size() <= 123, DOMSyntaxError,
653 "WebSocket close reason must not be longer than 123 bytes when UTF-8 encoded.");
654 }
655 }
656 
657 // The default code of 1005 cannot have a reason, per the standard, so if a reason is specified
658 // then there must be a code, too.
659 JSG_REQUIRE(reason == kj::none || code != kj::none, DOMInvalidAccessError,
660 "If you specify a WebSocket close reason, you must also specify a code.");
661 
662 // Handle close before connection is established for websockets obtained through `new WebSocket()`.
663 KJ_IF_SOME(pending, native.state.tryGet<AwaitingConnection>()) {
664 pending.canceler.cancel(kj::str("Called close before connection was established."));
665 
666 // Strictly speaking, we might not be all the way released by now, but we definitely shouldn't
667 // worry about canceling again.
668 native.state.init<Released>();
669 return;
670 }
671 
672 if (native.closedOutgoing || native.outgoingAborted || native.state.is<Released>()) {
673 // See comments in send(), above, which also apply here. Note that we opt to ignore a
674 // double-close() per spec, whereas send()-after-close() throws (off-spec).
675 
676 return;
677 } else if (awaitingHibernatableError()) {
678 // Ready for the hibernatable error event state, after encountering an error, the websocket
679 // isn't able to send outbound messages; let's release it.
680 tryReleaseNative(js);
681 return;
682 }
683 JSG_REQUIRE(native.state.is<Accepted>(), TypeError,
684 "You must call one of accept() or state.acceptWebSocket() on this WebSocket before sending "
685 "messages.");
686 
687 assertNoError(js);
688 
689 outgoingMessages->insert(GatedMessage{IoContext::current().waitForOutputLocksIfNecessary(),
690 kj::WebSocket::Close{
691 // Code 1005 actually translates to sending a close message with no body on the wire.
692 static_cast<uint16_t>(code.orDefault(1005)),
693 kj::mv(reason).orDefault(jsg::USVString(kj::str())),
694 },
695 getPendingAutoResponseCount()});
696 
697 native.closedOutgoing = true;
698 closedOutgoingForHib = true;
699 ensurePumping(js);
700}
701 
702int WebSocket::getReadyState() {
703 auto& native = *farNative;
704 if ((native.closedIncoming && native.closedOutgoing) || error != kj::none) {
705 return READY_STATE_CLOSED;
706 } else if (native.closedIncoming || native.closedOutgoing) {
707 // Bizarrely, the spec uses the same state for a close message having been sent *or* received,
708 // even though these are very different states from the point of view of the application.
709 return READY_STATE_CLOSING;
710 } else if (native.state.is<AwaitingConnection>()) {
711 return READY_STATE_CONNECTING;
712 }
713 return READY_STATE_OPEN;
714}
715 
716bool WebSocket::isAccepted() {
717 return farNative->state.is<Accepted>();
718}
719 
720bool WebSocket::isReleased() {
721 return farNative->state.is<Released>();
722}
723 
724kj::Maybe<kj::String> WebSocket::getPreferredExtensions(kj::WebSocket::ExtensionsContext ctx) {
725 KJ_SWITCH_ONEOF(farNative->state) {
726 KJ_CASE_ONEOF(ws, AwaitingConnection) {
727 return kj::none;
728 }
729 KJ_CASE_ONEOF(container, AwaitingAcceptanceOrCoupling) {
730 return container.ws->getPreferredExtensions(ctx);
731 }
732 KJ_CASE_ONEOF(container, Accepted) {
733 return container.ws->getPreferredExtensions(ctx);
734 }
735 KJ_CASE_ONEOF(container, Released) {
736 return kj::none;
737 }
738 }
739 return kj::none;
740}
741 
742kj::Maybe<kj::StringPtr> WebSocket::getUrl() {
743 return url.map([](kj::StringPtr value) { return value; });
744}
745 
746kj::Maybe<kj::StringPtr> WebSocket::getProtocol() {
747 return protocol.map([](kj::StringPtr value) { return value; });
748}
749 
750kj::Maybe<kj::StringPtr> WebSocket::getExtensions() {
751 return extensions.map([](kj::StringPtr value) { return value; });
752}
753 
754kj::Maybe<jsg::JsValue> WebSocket::deserializeAttachment(jsg::Lock& js) {
755 return serializedAttachment.map([&](kj::ArrayPtr<byte> attachment) -> jsg::JsValue {
756 jsg::Deserializer deserializer(js, attachment, kj::none, kj::none,
757 jsg::Deserializer::Options{
758 .version = 15,
759 .readHeader = true,
760 });
761 
762 return deserializer.readValue(js);
763 });
764}
765 
766void WebSocket::serializeAttachment(jsg::Lock& js, jsg::JsValue attachment) {
767 jsg::Serializer serializer(js,
768 jsg::Serializer::Options{
769 .version = 15,
770 .omitHeader = false,
771 });
772 serializer.write(js, attachment);
773 auto released = serializer.release();
774 JSG_REQUIRE(released.data.size() <= MAX_ATTACHMENT_SIZE, Error,
775 "A WebSocket 'attachment' cannot be larger than ", MAX_ATTACHMENT_SIZE,
776 " bytes."
777 "'attachment' was ",
778 released.data.size(), " bytes.");
779 serializedAttachment = kj::mv(released.data);
780}
781 
782void WebSocket::setAutoResponseStatus(
783 kj::Maybe<kj::Date> time, kj::Promise<void> autoResponsePromise) {
784 autoResponseTimestamp = time;
785 KJ_IF_SOME(context, IoContext::tryCurrent()) {
786 autoResponseStatus.ongoingAutoResponse.emplace(
787 context.addObject(kj::heap(kj::mv(autoResponsePromise))));
788 } else {
789 // Called outside an IoContext (e.g. from the hibernation manager's readLoop).
790 // Use plain kj::Own; the caller manages the promise lifecycle.
791 autoResponseStatus.ongoingAutoResponse.emplace(kj::heap(kj::mv(autoResponsePromise)));
792 }
793}
794 
795kj::Maybe<kj::Date> WebSocket::getAutoResponseTimestamp() {
796 return autoResponseTimestamp;
797}
798 
799void WebSocket::dispatchOpen(jsg::Lock& js) {
800 dispatchEventImpl(js, js.alloc<Event>("open"));
801}
802 
803void WebSocket::ensurePumping(jsg::Lock& js) {
804 auto& native = *farNative;
805 if (!native.isPumping) {
806 auto& context = IoContext::current();
807 auto& accepted = KJ_ASSERT_NONNULL(native.state.tryGet<Accepted>());
808 auto promise = kj::evalNow([&]() {
809 return accepted.canceler.wrap(
810 pump(context, *outgoingMessages, *accepted.ws, native, autoResponseStatus, observer));
811 });
812 
813 // TODO(cleanup): We use awaitIoLegacy() here because we don't want this to count as a pending
814 // event if this is a WebSocketPair with the other end being handled in the same isolate.
815 // In that case, the pump can hang if accept() is never called on the other end. Ideally,
816 // this scenario would be handled in-isolate using jsg::Promise, but that would take some
817 // refactoring.
818 context.awaitIoLegacy(js, kj::mv(promise))
819 .then(js, [this, thisHandle = JSG_THIS](jsg::Lock& js) {
820 auto& native = *farNative;
821 if (native.outgoingAborted) {
822 if (awaitingHibernatableRelease()) {
823 // We have a hibernatable websocket -- we don't want to dispatch a regular error event.
824 tryReleaseNative(js);
825 } else {
826 // Apparently, the peer stopped accepting messages (probably, disconnected entirely), but
827 // this didn't cause our writes to fail, maybe due to timing. Let's set the error now.
828 reportError(js, KJ_EXCEPTION(DISCONNECTED, "WebSocket peer disconnected"));
829 }
830 } else if (native.closedIncoming && native.closedOutgoing) {
831 if (awaitingHibernatableRelease()) {
832 // TODO(someday): These async races can be pretty complicated, and while it's good to have
833 // tests to make sure we're not broken, it would be nice to refactor this code eventually.
834 
835 // Hibernatable WebSockets had a subtle race condition where one pump() promise would
836 // start right after a previous pump() completed, but before this continuation ran.
837 //
838 // This race prevented close messages from being sent from inside the webSocketClose()
839 // handler because prior to the CLOSE getting sent in the second pump(), the promise
840 // continuation following the first pump() would transition us from Accepted to Released,
841 // triggering the canceler and cancelling the outgoing CLOSE of the second pump() promise.
842 //
843 // For a more detailed explanation, see https://github.com/cloudflare/workerd/pull/1535.
844 tryReleaseNative(js);
845 } else if (native.state.is<Accepted>()) {
846 // Native WebSocket no longer needed; release.
847 native.state.init<Released>();
848 } else if (native.state.is<Released>()) {
849 // While we were awaiting the jsg::Promise, someone else released our state. That's fine.
850 } else {
851 KJ_FAIL_ASSERT("Unexpected native web socket state", native.state);
852 }
853 }
854 }, [this, thisHandle = JSG_THIS](jsg::Lock& js, jsg::Value&& exception) mutable {
855 if (awaitingHibernatableRelease()) {
856 // We have a hibernatable websocket -- we don't want to dispatch a regular error event.
857 tryReleaseNative(js);
858 } else {
859 reportError(js, jsg::JsValue(exception.getHandle(js)).addRef(js));
860 }
861 });
862 }
863}
864 
865kj::Promise<void> WebSocket::sendAutoResponse(kj::String message, kj::WebSocket& ws) {
866 if (autoResponseStatus.isPumping) {
867 autoResponseStatus.pendingAutoResponseDeque.push(kj::mv(message));
868 } else if (!autoResponseStatus.isClosed) {
869 auto p = ws.send(message).fork();
870 KJ_IF_SOME(context, IoContext::tryCurrent()) {
871 autoResponseStatus.ongoingAutoResponse.emplace(context.addObject(kj::heap(p.addBranch())));
872 } else {
873 // Called outside an IoContext (e.g. from the hibernation manager's readLoop).
874 autoResponseStatus.ongoingAutoResponse.emplace(kj::heap(p.addBranch()));
875 }
876 co_await p;
877 autoResponseStatus.ongoingAutoResponse = kj::none;
878 }
879}
880 
881namespace {
882 
883size_t countBytesFromMessage(const kj::WebSocket::Message& message) {
884 // This does not count the extra data of the RPC frame or the savings from any compression.
885 // We're incentivizing customers to use reasonably sized messages, not trying to get an exact
886 // count of how many bytes went over the wire.
887 
888 KJ_SWITCH_ONEOF(message) {
889 KJ_CASE_ONEOF(s, kj::String) {
890 return s.size();
891 }
892 KJ_CASE_ONEOF(a, kj::Array<byte>) {
893 return a.size();
894 }
895 KJ_CASE_ONEOF(c, kj::WebSocket::Close) {
896 // If we include the size of the close code, that could incentivize our customers to omit
897 // sending Close frames when appropriate. The same cannot be said for the close reason since
898 // someone could encapsulate their final message in it to save costs.
899 return c.reason.size();
900 }
901 }
902 
903 KJ_UNREACHABLE;
904}
905 
906} // namespace
907 
908kj::Promise<void> WebSocket::pump(IoContext& context,
909 OutgoingMessagesMap& outgoingMessages,
910 kj::WebSocket& ws,
911 Native& native,
912 AutoResponse& autoResponse,
913 kj::Maybe<kj::Own<WebSocketObserver>>& observer) {
914 KJ_ASSERT(!native.isPumping);
915 native.isPumping = true;
916 autoResponse.isPumping = true;
917 bool completed = false;
918 KJ_DEFER({
919 // We use a KJ_DEFER to set native.isPumping = false to ensure that it happens -- we had a bug
920 // in the past where this was handled by the caller of WebSocket::pump() and it allowed for
921 // messages to get stuck in `outgoingMessages` until the pump task was restarted.
922 native.isPumping = false;
923 
924 // Either we were already through all our outgoing messages or we experienced failure/
925 // cancellation and cannot send these anyway.
926 outgoingMessages.clear();
927 
928 autoResponse.isPumping = false;
929 
930 autoResponse.pendingAutoResponseDeque.clear();
931 
932 if (!completed) {
933 // We didn't make it to `completed = true` at the end of this function, so either an
934 // exception was thrown or the task was canceled. Either way, we cannot send any further
935 // messages, because the connection is in a broken state and will just throw more exceptions.
936 // Setting `outgoingAborted` stops us from even trying to queue any more messages.
937 native.outgoingAborted = true;
938 }
939 });
940 
941 // If we have a ongoingAutoResponse, we must co_await it here because there's a ws.send()
942 // in progress. Otherwise there can occur ws.send() race problems.
943 KJ_IF_SOME(promiseHolder, autoResponse.ongoingAutoResponse) {
944 KJ_SWITCH_ONEOF(promiseHolder) {
945 KJ_CASE_ONEOF(ioOwned, IoOwn<kj::Promise<void>>) {
946 co_await *ioOwned;
947 }
948 KJ_CASE_ONEOF(owned, kj::Own<kj::Promise<void>>) {
949 co_await *owned;
950 }
951 }
952 autoResponse.ongoingAutoResponse = kj::none;
953 }
954 
955 do {
956 while (outgoingMessages.size() > 0) {
957 GatedMessage gatedMessage = outgoingMessages.release(*outgoingMessages.ordered().begin());
958 KJ_IF_SOME(promise, gatedMessage.outputLock) {
959 co_await promise;
960 }
961 
962 auto size = countBytesFromMessage(gatedMessage.message);
963 
964 while (gatedMessage.pendingAutoResponses > 0) {
965 auto message = KJ_ASSERT_NONNULL(autoResponse.pendingAutoResponseDeque.pop());
966 gatedMessage.pendingAutoResponses--;
967 autoResponse.queuedAutoResponses--;
968 co_await ws.send(message);
969 }
970 
971 KJ_SWITCH_ONEOF(gatedMessage.message) {
972 KJ_CASE_ONEOF(text, kj::String) {
973 co_await ws.send(text);
974 break;
975 }
976 KJ_CASE_ONEOF(data, kj::Array<byte>) {
977 co_await ws.send(data);
978 break;
979 }
980 KJ_CASE_ONEOF(close, kj::WebSocket::Close) {
981 co_await ws.close(close.code, close.reason);
982 autoResponse.isClosed = true;
983 break;
984 }
985 }
986 
987 KJ_IF_SOME(o, observer) {
988 o->sentMessage(size);
989 }
990 
991 KJ_IF_SOME(a, context.getActor()) {
992 a.getMetrics().sentWebSocketMessage(size);
993 }
994 }
995 
996 // If there are any auto-responses left to process, we should do it now.
997 // We should also check if the last sent message was a close. Shouldn't happen.
998 while (!autoResponse.pendingAutoResponseDeque.empty() && !autoResponse.isClosed) {
999 auto message = KJ_ASSERT_NONNULL(autoResponse.pendingAutoResponseDeque.pop());
1000 co_await ws.send(message);
1001 }
1002 
1003 // While we were `co_await`ing the auto-response send, more messages could have been queued
1004 // into `outgoingMessages`. If so we'll need to start over, otherwise these messages would be
1005 // discarded in our KJ_DEFER block!
1006 } while (outgoingMessages.size() > 0);
1007 
1008 completed = true;
1009}
1010 
1011size_t WebSocket::getPendingAutoResponseCount() {
1012 auto count =
1013 autoResponseStatus.pendingAutoResponseDeque.size() - autoResponseStatus.queuedAutoResponses;
1014 autoResponseStatus.queuedAutoResponses = autoResponseStatus.pendingAutoResponseDeque.size();
1015 return count;
1016}
1017 
1018void WebSocket::tryReleaseNative(jsg::Lock& js) {
1019 // If the native WebSocket is no longer needed (the connection closed) and there are no more
1020 // messages to send, we can discard the underlying connection.
1021 auto& native = *farNative;
1022 if ((native.closedOutgoing || native.outgoingAborted) && !native.isPumping) {
1023 // Native WebSocket no longer needed; release.
1024 KJ_ASSERT(native.state.is<Accepted>());
1025 native.state.init<Released>();
1026 }
1027}
1028 
1029kj::Array<kj::StringPtr> WebSocket::getHibernatableTags() {
1030 auto& accepted = JSG_REQUIRE_NONNULL(farNative->state.tryGet<Accepted>(), Error,
1031 "you must call 'acceptWebSocket()' before attempting to access the tags of a WebSocket.");
1032 JSG_REQUIRE(accepted.isHibernatable(), Error, "only hibernatable websockets can have tags.");
1033 return accepted.ws.getHibernatableTags();
1034}
1035 
1036kj::Promise<kj::Maybe<kj::Exception>> WebSocket::readLoop(
1037 kj::Maybe<kj::Own<InputGate::CriticalSection>> cs, size_t maxMessageSize) {
1038 try {
1039 // Note that we'll throw if the websocket has enabled hibernation.
1040 auto& ws = *KJ_REQUIRE_NONNULL(
1041 KJ_ASSERT_NONNULL(farNative->state.tryGet<Accepted>()).ws.getIfNotHibernatable());
1042 auto& context = IoContext::current();
1043 while (true) {
1044 auto message = co_await ws.receive(maxMessageSize);
1045 
1046 auto size = countBytesFromMessage(message);
1047 KJ_IF_SOME(o, observer) {
1048 o->receivedMessage(size);
1049 }
1050 
1051 context.getLimitEnforcer().topUpActor();
1052 KJ_IF_SOME(a, context.getActor()) {
1053 a.getMetrics().receivedWebSocketMessage(size);
1054 }
1055 
1056 // Re-enter the context with context.run(). This is arguably a bit unusual compared to other
1057 // I/O which is delivered by return from context.awaitIo(), but the difference here is that we
1058 // have a long stream of events over time. It makes sense to use context.run() each time a new
1059 // event arrives.
1060 // TODO(cleanup): The way context.run is defined, a capturing lambda is required here, which
1061 // is a bit unfortunate. We could simply things somewhat with a variation that would allow
1062 // something like context.run(handleMessage, *this, kj::mv(message)) where the acquired lock,
1063 // and the additional arguments are passed into handleMessage, avoiding the need for the
1064 // lambda here entirely.
1065 auto result = co_await context.run([this, message = kj::mv(message)](auto& wLock) mutable {
1066 auto& native = *farNative;
1067 jsg::Lock& js = wLock;
1068 KJ_SWITCH_ONEOF(message) {
1069 KJ_CASE_ONEOF(text, kj::String) {
1070 dispatchEventImpl(js, js.alloc<MessageEvent>(js, js.str(text)));
1071 }
1072 KJ_CASE_ONEOF(data, kj::Array<byte>) {
1073 if (binaryType_ == BinaryType::BLOB) {
1074 // Per the WHATWG spec, deliver binary messages as Blob when binaryType is "blob".
1075 auto ab = jsg::JsArrayBuffer::create(js, data);
1076 auto blob = js.alloc<Blob>(js, jsg::JsBufferSource(ab), kj::str());
1077 dispatchEventImpl(js, js.alloc<MessageEvent>(js, kj::str("message"), kj::mv(blob)));
1078 } else {
1079 auto ab = js.arrayBuffer(kj::mv(data)).getHandle(js);
1080 dispatchEventImpl(js, js.alloc<MessageEvent>(js, jsg::JsValue(ab)));
1081 }
1082 }
1083 KJ_CASE_ONEOF(close, kj::WebSocket::Close) {
1084 native.closedIncoming = true;
1085 if (!allowHalfOpen.toBool() && !native.closedOutgoing && !native.outgoingAborted &&
1086 !native.state.is<Released>()) {
1087 // When allowHalfOpen is false (the spec-compliant default with the
1088 // web_socket_auto_reply_to_close compat flag), automatically send a reciprocal
1089 // Close frame through the outgoing message pump so that readyState is CLOSED (3)
1090 // when the close event fires. Skip if a close frame was already sent (e.g. the
1091 // application called close() before the server sent its Close), or if the outgoing
1092 // side is otherwise unusable.
1093 outgoingMessages->insert(
1094 GatedMessage{IoContext::current().waitForOutputLocksIfNecessary(),
1095 kj::WebSocket::Close{close.code, kj::str(close.reason)},
1096 getPendingAutoResponseCount()});
1097 
1098 native.closedOutgoing = true;
1099 closedOutgoingForHib = true;
1100 ensurePumping(js);
1101 }
1102 dispatchEventImpl(js, js.alloc<CloseEvent>(close.code, kj::mv(close.reason), true));
1103 // Native WebSocket no longer needed; release.
1104 tryReleaseNative(js);
1105 return false;
1106 }
1107 }
1108 
1109 return true;
1110 }, mapAddRef(cs));
1111 
1112 if (!result) co_return kj::none;
1113 }
1114 KJ_UNREACHABLE;
1115 } catch (...) {
1116 co_return kj::getCaughtExceptionAsKj();
1117 }
1118}
1119 
1120jsg::Ref<WebSocketPair> WebSocketPair::constructor(jsg::Lock& js) {
1121 auto pipe = kj::newWebSocketPipe();
1122 auto pair = js.alloc<WebSocketPair>(
1123 js.alloc<WebSocket>(js, kj::mv(pipe.ends[0])), js.alloc<WebSocket>(js, kj::mv(pipe.ends[1])));
1124 auto first = pair->getFirst();
1125 auto second = pair->getSecond();
1126 
1127 first->setPeer(second->addWeakRef());
1128 second->setPeer(first->addWeakRef());
1129 return kj::mv(pair);
1130}
1131 
1132jsg::Ref<WebSocketPair::PairIterator> WebSocketPair::entries(jsg::Lock& js) {
1133 return js.alloc<PairIterator>(IteratorState{
1134 .pair = JSG_THIS,
1135 .index = 0,
1136 });
1137}
1138 
1139void WebSocket::reportError(jsg::Lock& js, kj::Exception&& e) {
1140 reportError(js, js.exceptionToJsValue(e.clone()));
1141}
1142 
1143void WebSocket::reportError(jsg::Lock& js, jsg::JsRef<jsg::JsValue> err) {
1144 // If this is the first error, raise the error event.
1145 if (error == kj::none) {
1146 auto msg = kj::str(v8::Exception::CreateMessage(js.v8Isolate, err.getHandle(js))->Get());
1147 error = err.addRef(js);
1148 
1149 dispatchEventImpl(js,
1150 js.alloc<ErrorEvent>(
1151 ErrorEvent::ErrorEventInit{.message = kj::mv(msg), .error = kj::mv(err)}));
1152 
1153 // After an error we don't allow further send()s. If the receive loop has also ended then we
1154 // can destroy the connection. Note that we don't set closedOutgoing = true because that flag
1155 // is specifically to indicate that `close()` has been called, and it causes `send()` to throw
1156 // an exception complaining specifically that `close()` was called, which would be
1157 // inappropriate in this case.
1158 auto& native = *farNative;
1159 native.outgoingAborted = true;
1160 if (native.closedIncoming && !native.isPumping) {
1161 KJ_IF_SOME(pending, native.state.tryGet<AwaitingConnection>()) {
1162 // Nothing worth canceling if we're reporting an error from the connection establishment
1163 // continuations.
1164 pending.canceler.release();
1165 }
1166 
1167 // We're no longer pumping so let's make sure we release the native connection here.
1168 native.state.init<Released>();
1169 }
1170 }
1171}
1172 
1173void WebSocket::assertNoError(jsg::Lock& js) {
1174 KJ_IF_SOME(e, error) {
1175 js.throwException(e.addRef(js));
1176 }
1177}
1178 
1179void WebSocket::setPeer(kj::Own<WeakRef<WebSocket>> other) {
1180 peer = kj::mv(other);
1181}
1182 
1183kj::Own<kj::WebSocket> WebSocket::acceptAsHibernatable(kj::Array<kj::StringPtr> tags) {
1184 KJ_IF_SOME(hibernatable, farNative->state.tryGet<AwaitingAcceptanceOrCoupling>()) {
1185 // We can only request hibernation if we have not called accept.
1186 auto ws = kj::mv(hibernatable.ws);
1187 // We pass a reference to the kj::WebSocket for the api::WebSocket to refer to when calling
1188 // `send()` or `close()`.
1189 farNative->state.init<Accepted>(Accepted::Hibernatable{.ws = *ws, .tagsRef = kj::mv(tags)},
1190 *farNative, IoContext::current());
1191 return kj::mv(ws);
1192 }
1193 JSG_FAIL_REQUIRE(TypeError,
1194 "Tried to make an api::WebSocket hibernatable when it was in an incompatible state.");
1195}
1196 
1197void WebSocket::initiateHibernatableRelease(jsg::Lock& js,
1198 kj::Own<kj::WebSocket> ws,
1199 kj::Array<kj::String> tags,
1200 HibernatableReleaseState releaseState) {
1201 // TODO(soon): We probably want this to be an assert, since this is meant to be called once
1202 // at the end of a websocket connection>
1203 KJ_IF_SOME(state, farNative->state.tryGet<Accepted>()) {
1204 KJ_REQUIRE(state.isHibernatable(),
1205 "tried to initiate hibernatable release but websocket wasn't hibernatable");
1206 state.ws.initiateHibernatableRelease(js, kj::mv(ws), kj::mv(tags), releaseState);
1207 farNative->closedIncoming = true;
1208 } else {
1209 KJ_LOG(WARNING, "Unexpected Hibernatable WebSocket state on release", farNative->state);
1210 }
1211}
1212 
1213bool WebSocket::awaitingHibernatableError() {
1214 KJ_IF_SOME(accepted, farNative->state.tryGet<Accepted>()) {
1215 return (accepted.ws.isAwaitingError());
1216 }
1217 return false;
1218}
1219 
1220bool WebSocket::awaitingHibernatableRelease() {
1221 KJ_IF_SOME(accepted, farNative->state.tryGet<Accepted>()) {
1222 return (accepted.ws.isAwaitingRelease());
1223 }
1224 return false;
1225}
1226 
1227bool WebSocket::peerIsAwaitingCoupling() {
1228 bool answer = false;
1229 KJ_IF_SOME(p, peer) {
1230 p->runIfAlive([&answer](WebSocket& ws) {
1231 answer = ws.farNative->state.is<AwaitingAcceptanceOrCoupling>();
1232 });
1233 }
1234 return answer;
1235}
1236 
1237WebSocket::HibernationPackage WebSocket::buildPackageForHibernation() {
1238 // TODO(cleanup): It would be great if we could limit this so only the HibernationManager
1239 // (or a derived class) could call it.
1240 return HibernationPackage{
1241 .url = kj::mv(url),
1242 .protocol = kj::mv(protocol),
1243 .extensions = kj::mv(extensions),
1244 .serializedAttachment = kj::mv(serializedAttachment),
1245 .maybeTags = kj::none,
1246 .closedOutgoingConnection = closedOutgoingForHib,
1247 .allowHalfOpen = allowHalfOpen,
1248 };
1249}
1250 
1251WebSocket::Accepted::WrappedWebSocket::WrappedWebSocket(Hibernatable ws) {
1252 inner.init<Hibernatable>(kj::mv(ws));
1253}
1254 
1255WebSocket::Accepted::WrappedWebSocket::WrappedWebSocket(kj::Own<kj::WebSocket> ws) {
1256 inner.init<kj::Own<kj::WebSocket>>(kj::mv(ws));
1257}
1258 
1259kj::WebSocket* WebSocket::Accepted::WrappedWebSocket::operator->() {
1260 KJ_SWITCH_ONEOF(inner) {
1261 KJ_CASE_ONEOF(owned, kj::Own<kj::WebSocket>) {
1262 return owned.get();
1263 }
1264 KJ_CASE_ONEOF(hibernatable, Hibernatable) {
1265 return &hibernatable.ws;
1266 }
1267 }
1268 KJ_UNREACHABLE;
1269}
1270 
1271kj::WebSocket& WebSocket::Accepted::WrappedWebSocket::operator*() {
1272 KJ_SWITCH_ONEOF(inner) {
1273 KJ_CASE_ONEOF(owned, kj::Own<kj::WebSocket>) {
1274 return *owned;
1275 }
1276 KJ_CASE_ONEOF(hibernatable, Hibernatable) {
1277 return hibernatable.ws;
1278 }
1279 }
1280 KJ_UNREACHABLE;
1281}
1282 
1283kj::Maybe<kj::Own<kj::WebSocket>&> WebSocket::Accepted::WrappedWebSocket::getIfNotHibernatable() {
1284 // The implication of getting nullptr is that this websocket is hibernatable. This is useful
1285 // if the caller only ever expects to get a regular websocket, for example, if they are in
1286 // any method that should be inaccessible to hibernatable websockets (ex. readLoop).
1287 return inner.tryGet<kj::Own<kj::WebSocket>>();
1288}
1289 
1290kj::Maybe<WebSocket::Accepted::Hibernatable&> WebSocket::Accepted::WrappedWebSocket::
1291 getIfHibernatable() {
1292 return inner.tryGet<Hibernatable>();
1293}
1294 
1295kj::Array<kj::StringPtr> WebSocket::Accepted::WrappedWebSocket::getHibernatableTags() {
1296 KJ_SWITCH_ONEOF(KJ_REQUIRE_NONNULL(inner.tryGet<Hibernatable>()).tagsRef) {
1297 KJ_CASE_ONEOF(ref, kj::Array<kj::StringPtr>) {
1298 // Tags are still owned by the HibernationManager
1299 return kj::heapArray<kj::StringPtr>(ref);
1300 }
1301 KJ_CASE_ONEOF(arr, kj::Array<kj::String>) {
1302 // We have the array already, let's copy it and return.
1303 auto cpy = kj::heapArray<kj::StringPtr>(arr.size());
1304 for (auto& i: kj::indices(arr)) {
1305 cpy[i] = arr[i].asPtr();
1306 }
1307 return cpy;
1308 }
1309 }
1310 KJ_UNREACHABLE;
1311}
1312 
1313void WebSocket::Accepted::WrappedWebSocket::initiateHibernatableRelease(jsg::Lock& js,
1314 kj::Own<kj::WebSocket> ws,
1315 kj::Array<kj::String> tags,
1316 HibernatableReleaseState state) {
1317 auto& hibernatable = KJ_REQUIRE_NONNULL(getIfHibernatable());
1318 hibernatable.releaseState = state;
1319 // Note that we move the owned kj::WebSocket here.
1320 hibernatable.attachedForClose = kj::mv(ws);
1321 hibernatable.tagsRef.init<kj::Array<kj::String>>(kj::mv(tags));
1322}
1323 
1324bool WebSocket::Accepted::WrappedWebSocket::isAwaitingRelease() {
1325 KJ_IF_SOME(ws, getIfHibernatable()) {
1326 return (ws.releaseState != HibernatableReleaseState::NONE);
1327 }
1328 return false;
1329}
1330 
1331bool WebSocket::Accepted::WrappedWebSocket::isAwaitingError() {
1332 KJ_IF_SOME(ws, getIfHibernatable()) {
1333 return (ws.releaseState == HibernatableReleaseState::ERROR);
1334 }
1335 return false;
1336}
1337 
1338bool WebSocket::Accepted::isHibernatable() {
1339 return ws.getIfNotHibernatable() == kj::none;
1340}
1341 
1342void WebSocketPair::visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
1343 tracker.trackField(nullptr, sockets[0]);
1344 tracker.trackField(nullptr, sockets[1]);
1345}
1346 
1347kj::StringPtr WebSocket::getBinaryType() {
1348 return binaryType_ == BinaryType::BLOB ? "blob"_kj : "arraybuffer"_kj;
1349}
1350 
1351void WebSocket::setBinaryType(kj::String value) {
1352 if (value == "blob") {
1353 binaryType_ = BinaryType::BLOB;
1354 } else if (value == "arraybuffer") {
1355 binaryType_ = BinaryType::ARRAYBUFFER;
1356 }
1357 // Per the spec, invalid values are silently ignored — the existing binaryType is retained.
1358}
1359 
1360void WebSocket::visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
1361 tracker.trackField("url", url);
1362 tracker.trackField("protocol", protocol);
1363 tracker.trackField("extensions", extensions);
1364 KJ_IF_SOME(attachment, serializedAttachment) {
1365 tracker.trackFieldWithSize("attachment", attachment.size());
1366 }
1367 tracker.trackFieldWithSize("IoOwn<Native>", sizeof(IoOwn<Native>));
1368 tracker.trackField("error", error);
1369 tracker.trackFieldWithSize("IoOwn<OutgoingMessagesMap>", sizeof(IoOwn<OutgoingMessagesMap>));
1370 tracker.trackField("autoResponseStatus", autoResponseStatus);
1371}
1372 
1373} // namespace workerd::api