From 3d844c6f213ea4f75c64bb2619ec7b7b53b4c8d7 Mon Sep 17 00:00:00 2001 From: Mike Aizatsky Date: Fri, 21 Aug 2026 11:53:32 -0700 Subject: [PATCH] always copy websocket tags --- src/workerd/api/hibernatable-adapter.c++ | 5 +-- src/workerd/api/hibernatable-adapter.h | 4 +- src/workerd/api/hibernatable-web-socket.c++ | 2 +- src/workerd/api/web-socket.c++ | 37 +++++++------------ src/workerd/api/web-socket.h | 22 ++++------- src/workerd/io/hibernation-manager-test.c++ | 32 ++++++++++++++++ src/workerd/io/legacy-hibernation-manager.c++ | 10 +---- src/workerd/io/legacy-hibernation-manager.h | 12 +----- 8 files changed, 61 insertions(+), 63 deletions(-) diff --git a/src/workerd/api/hibernatable-adapter.c++ b/src/workerd/api/hibernatable-adapter.c++ index 850ad04d749..4889be24da6 100644 --- a/src/workerd/api/hibernatable-adapter.c++ +++ b/src/workerd/api/hibernatable-adapter.c++ @@ -18,7 +18,7 @@ HibernatableWebSocketAdapter::HibernatableWebSocketAdapter(jsg::Lock& js, } HibernatableWebSocketAdapter::HibernatableWebSocketAdapter( - WebSocket& shellParam, kj::WebSocket& ws, kj::Array tags) + WebSocket& shellParam, kj::WebSocket& ws, kj::Array tags) : shell(shellParam) { KJ_UNIMPLEMENTED("EW-10817: HibernatableWebSocketAdapter transition ctor not yet implemented"); } @@ -109,8 +109,7 @@ bool HibernatableWebSocketAdapter::isAwaitingCoupling() { return false; } -kj::Own HibernatableWebSocketAdapter::acceptAsHibernatable( - kj::Array) { +kj::Own HibernatableWebSocketAdapter::acceptAsHibernatable(kj::Array) { unimplemented(); } diff --git a/src/workerd/api/hibernatable-adapter.h b/src/workerd/api/hibernatable-adapter.h index c11f7a7c29c..712107cd12f 100644 --- a/src/workerd/api/hibernatable-adapter.h +++ b/src/workerd/api/hibernatable-adapter.h @@ -35,7 +35,7 @@ class HibernatableWebSocketAdapter final: public WebSocketAdapter { // available (`HibernationManager::acceptWebSocket` does not currently plumb one through). // When the implementation actually needs JS state, it can grab the lock from the ambient // IoContext. - HibernatableWebSocketAdapter(WebSocket& shell, kj::WebSocket& ws, kj::Array tags); + HibernatableWebSocketAdapter(WebSocket& shell, kj::WebSocket& ws, kj::Array tags); ~HibernatableWebSocketAdapter() noexcept(false); @@ -67,7 +67,7 @@ class HibernatableWebSocketAdapter final: public WebSocketAdapter { bool isHibernatable() override; void setObserver(kj::Own observer) override; - kj::Own acceptAsHibernatable(kj::Array tags) override; + kj::Own acceptAsHibernatable(kj::Array tags) override; void initiateHibernatableRelease(jsg::Lock& js, kj::Own ws, kj::Array tags, diff --git a/src/workerd/api/hibernatable-web-socket.c++ b/src/workerd/api/hibernatable-web-socket.c++ index bb3cdbdc213..abbd869d683 100644 --- a/src/workerd/api/hibernatable-web-socket.c++ +++ b/src/workerd/api/hibernatable-web-socket.c++ @@ -41,7 +41,7 @@ HibernatableWebSocketEvent::ItemsForRelease HibernatableWebSocketEvent::prepareF // the HibernatableWebSocket (it removes it from `webSocketsForEventHandler`). auto websocketRef = hibernatableWebSocket.value->getActiveOrUnhibernate(lock); auto ownedWebSocket = kj::mv(KJ_REQUIRE_NONNULL(hibernatableWebSocket.value->ws)); - auto tags = hibernatableWebSocket.value->cloneTags(); + auto tags = hibernatableWebSocket.value->getTags(); // Now that we've obtained the websocket for the event, let's free up the slots we had allocated. manager.webSocketsForEventHandler.erase(hibernatableWebSocket); diff --git a/src/workerd/api/web-socket.c++ b/src/workerd/api/web-socket.c++ index 098fce96948..a4826fd9adf 100644 --- a/src/workerd/api/web-socket.c++ +++ b/src/workerd/api/web-socket.c++ @@ -173,7 +173,7 @@ void WebSocket::setObserver(kj::Own observer) { impl->setObserver(kj::mv(observer)); } -kj::Own WebSocket::acceptAsHibernatable(kj::Array tags) { +kj::Own WebSocket::acceptAsHibernatable(kj::Array tags) { // TODO(EW-10817): When the HibernatableWebSocketAdapter path is functional, re-enable the // autogate-driven swap below — extract the kj::WebSocket from the legacy adapter (via // `LegacyWebSocketAdapter::extractForHibernatableTransition`, also commented out today) @@ -248,13 +248,11 @@ void WebSocket::visitForMemoryInfo(jsg::MemoryTracker& tracker) const { // LegacyWebSocketAdapter — operational implementation. // ============================================================================= -IoOwn LegacyWebSocketAdapter::initNative(IoContext& ioContext, - kj::WebSocket& ws, - kj::Array tags, - bool closedOutgoingConn) { +IoOwn LegacyWebSocketAdapter::initNative( + IoContext& ioContext, kj::WebSocket& ws, kj::Array tags, bool closedOutgoingConn) { auto nativeObj = kj::heap(); nativeObj->state.init( - Accepted::Hibernatable{.ws = ws, .tagsRef = kj::mv(tags)}, *nativeObj, ioContext); + Accepted::Hibernatable{.ws = ws, .tags = kj::mv(tags)}, *nativeObj, ioContext); // We might have called `close()` when this WebSocket was previously active. // If so, we want to prevent any future calls to `send()`. nativeObj->closedOutgoing = closedOutgoingConn; @@ -1470,14 +1468,14 @@ void LegacyWebSocketAdapter::setPeer(jsg::WeakRef other) { peer = kj::mv(other); } -kj::Own LegacyWebSocketAdapter::acceptAsHibernatable(kj::Array tags) { +kj::Own LegacyWebSocketAdapter::acceptAsHibernatable(kj::Array tags) { KJ_IF_SOME(hibernatable, farNative->state.tryGet()) { // We can only request hibernation if we have not called accept. auto ws = kj::mv(hibernatable.ws); // We pass a reference to the kj::WebSocket for the api::WebSocket to refer to when calling // `send()` or `close()`. - farNative->state.init(Accepted::Hibernatable{.ws = *ws, .tagsRef = kj::mv(tags)}, - *farNative, IoContext::current()); + farNative->state.init( + Accepted::Hibernatable{.ws = *ws, .tags = kj::mv(tags)}, *farNative, IoContext::current()); return kj::mv(ws); } JSG_FAIL_REQUIRE(TypeError, @@ -1597,21 +1595,12 @@ kj::Maybe LegacyWebSocketAdapte } kj::Array LegacyWebSocketAdapter::Accepted::WrappedWebSocket::getHibernatableTags() { - KJ_SWITCH_ONEOF(KJ_REQUIRE_NONNULL(inner.tryGet()).tagsRef) { - KJ_CASE_ONEOF(ref, kj::Array) { - // Tags are still owned by the HibernationManager - return kj::heapArray(ref); - } - KJ_CASE_ONEOF(arr, kj::Array) { - // We have the array already, let's copy it and return. - auto cpy = kj::heapArray(arr.size()); - for (auto& i: kj::indices(arr)) { - cpy[i] = arr[i].asPtr(); - } - return cpy; - } + auto& tags = KJ_REQUIRE_NONNULL(inner.tryGet()).tags; + auto result = kj::heapArray(tags.size()); + for (auto i: kj::indices(tags)) { + result[i] = tags[i]; } - KJ_UNREACHABLE; + return result; } void LegacyWebSocketAdapter::Accepted::WrappedWebSocket::initiateHibernatableRelease(jsg::Lock& js, @@ -1622,7 +1611,7 @@ void LegacyWebSocketAdapter::Accepted::WrappedWebSocket::initiateHibernatableRel hibernatable.releaseState = state; // Note that we move the owned kj::WebSocket here. hibernatable.attachedForClose = kj::mv(ws); - hibernatable.tagsRef.init>(kj::mv(tags)); + hibernatable.tags = kj::mv(tags); } bool LegacyWebSocketAdapter::Accepted::WrappedWebSocket::isAwaitingRelease() { diff --git a/src/workerd/api/web-socket.h b/src/workerd/api/web-socket.h index 63ee58a1404..fe3a8ac490a 100644 --- a/src/workerd/api/web-socket.h +++ b/src/workerd/api/web-socket.h @@ -196,7 +196,7 @@ class WebSocket: public EventTarget { // `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; + kj::Maybe> maybeTags; // True forever once the JS WebSocket calls `close()`. bool closedOutgoingConnection = false; @@ -241,7 +241,7 @@ class WebSocket: public EventTarget { // 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); + kj::Own acceptAsHibernatable(kj::Array tags); // Accesses the tags of the hibernatable websocket. kj::Array getHibernatableTags(); @@ -505,7 +505,7 @@ class WebSocketAdapter { // Hibernation transitions: extracts the kj::WebSocket so the HibernationManager can // assume ownership; the adapter retains a bare reference for sends. - virtual kj::Own acceptAsHibernatable(kj::Array tags) = 0; + virtual kj::Own acceptAsHibernatable(kj::Array tags) = 0; // HibernationManager coordination: signals that the underlying connection is winding // down and the adapter should arrange to deliver its terminal close/error event. @@ -603,7 +603,7 @@ class LegacyWebSocketAdapter final: public WebSocketAdapter { bool isHibernatable() override; void setObserver(kj::Own observer) override; - kj::Own acceptAsHibernatable(kj::Array tags) override; + kj::Own acceptAsHibernatable(kj::Array tags) override; void initiateHibernatableRelease(jsg::Lock& js, kj::Own ws, kj::Array tags, @@ -668,12 +668,8 @@ class LegacyWebSocketAdapter final: public WebSocketAdapter { // Close. WebSocket::HibernatableReleaseState releaseState = WebSocket::HibernatableReleaseState::NONE; - // There are two possible states for tagsRef: - // 1. kj::Array — tags 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 - // and the HibernatableWebSocket is free to go away mid-dispatch; copy locally. - kj::OneOf, kj::Array> tagsRef; + // The native state can outlive the HibernationManager, so it owns its tag strings. + kj::Array tags; }; explicit Accepted(kj::Own ws, Native& native, IoContext& context); @@ -781,10 +777,8 @@ class LegacyWebSocketAdapter final: public WebSocketAdapter { // Creates a fresh `Native` set up in the Accepted-Hibernatable state. Used by the // hibernation-revival constructor. - IoOwn initNative(IoContext& ioContext, - kj::WebSocket& ws, - kj::Array tags, - bool closedOutgoingConn); + IoOwn initNative( + IoContext& ioContext, kj::WebSocket& ws, kj::Array tags, bool closedOutgoingConn); void dispatchOpen(jsg::Lock& js); void ensurePumping(jsg::Lock& js); diff --git a/src/workerd/io/hibernation-manager-test.c++ b/src/workerd/io/hibernation-manager-test.c++ index 70364bb1068..b7313eef2a5 100644 --- a/src/workerd/io/hibernation-manager-test.c++ +++ b/src/workerd/io/hibernation-manager-test.c++ @@ -309,6 +309,38 @@ KJ_TEST("HibernationManager: failed event dispatches remove WebSocket") { fixture.drainAndDestroy(kj::mv(request)); } +KJ_TEST("HibernationManager: retained native WebSocket tags outlive manager teardown") { + DispatchStats stats; + TestFixture fixture(stubLoopbackParams(stats, kj::str("retained-tags"))); + auto hm = makeTestHm(fixture); + auto request = fixture.newIncomingRequest(); + constexpr kj::StringPtr tag = + "a-long-hibernatable-websocket-tag-that-must-remain-valid-after-manager-teardown"_kj; + auto end1 KJ_UNUSED = acceptNewWebSocket(fixture, *request, *hm, tag); + + kj::Maybe> retained; + fixture.enterContext(*request, [&](const TestFixture::Environment& env) { + auto websockets = hm->getWebSockets(env.js, tag); + KJ_ASSERT(websockets.size() == 1); + auto tags = websockets[0]->getHibernatableTags(); + KJ_ASSERT(tags.size() == 1); + KJ_ASSERT(tags[0] == tag, tags[0]); + retained = websockets[0].addRef(); + }); + + // The native WebSocket can be retained independently of the manager. Its tags must remain + // readable after the manager removes the final socket and destroys the corresponding bucket. + fixture.enterContext(*request, [&](const TestFixture::Environment&) { + hm = nullptr; + + auto tags = KJ_REQUIRE_NONNULL(retained)->getHibernatableTags(); + KJ_ASSERT(tags.size() == 1); + KJ_ASSERT(tags[0] == tag, tags[0]); + }); + + fixture.drainAndDestroy(kj::mv(request)); +} + KJ_TEST("HibernationManager: DO sends binary message to eyeball") { DispatchStats stats; TestFixture fixture(stubLoopbackParams(stats, kj::str("do-send-bin"))); diff --git a/src/workerd/io/legacy-hibernation-manager.c++ b/src/workerd/io/legacy-hibernation-manager.c++ index 0af14e77cba..10b73876798 100644 --- a/src/workerd/io/legacy-hibernation-manager.c++ +++ b/src/workerd/io/legacy-hibernation-manager.c++ @@ -44,15 +44,7 @@ LegacyHibernationManagerImpl::HibernatableWebSocket::~HibernatableWebSocket() no } } -kj::Array LegacyHibernationManagerImpl::HibernatableWebSocket::getTags() { - auto tags = kj::heapArray(tagItems.size()); - for (auto i: kj::indices(tagItems)) { - tags[i] = tagItems[i].tag; - } - return tags; -} - -kj::Array LegacyHibernationManagerImpl::HibernatableWebSocket::cloneTags() { +kj::Array LegacyHibernationManagerImpl::HibernatableWebSocket::getTags() { auto tags = kj::heapArray(tagItems.size()); for (auto i: kj::indices(tagItems)) { tags[i] = kj::str(tagItems[i].tag); diff --git a/src/workerd/io/legacy-hibernation-manager.h b/src/workerd/io/legacy-hibernation-manager.h index 82e25ad54af..96da0656a4a 100644 --- a/src/workerd/io/legacy-hibernation-manager.h +++ b/src/workerd/io/legacy-hibernation-manager.h @@ -85,16 +85,8 @@ class LegacyHibernationManagerImpl final: public Worker::Actor::HibernationManag ~HibernatableWebSocket() noexcept(false); KJ_DISALLOW_COPY_AND_MOVE(HibernatableWebSocket); - // Returns the tags associated with this HibernatableWebSocket. - kj::Array getTags(); - - // Returns the tags associated with this HibernatableWebSocket. - // Note that this returns an array of Strings, unlike `getTags()`. - // Copying the strings each time tags are requested would be expensive, - // so we only do it when we're delivering a close/error event because - // we will be destroying the HibernatableWebSocket object, - // which the tags need to outlive. - kj::Array cloneTags(); + // Returns an owned copy of the tags associated with this HibernatableWebSocket. + kj::Array getTags(); // Returns a reference to the active websocket. If the websocket is currently hibernating, // we have to unhibernate it first. The process moves values from the HibernatableWebSocket