diff --git a/CHANGELOG.md b/CHANGELOG.md index f8bcaea..66f9d5a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -253,6 +253,15 @@ policy in [docs/versioning.md](docs/versioning.md). ### Fixed +- **A dialed WebSocket released on its own io thread no longer terminates + the process.** When the last handle to a `BeastWebSocketClient::Dial` + socket dropped inside one of its own completions (typically a `Detached` + loop that owned the socket and finished there), the destructor joined the + thread it was running on, and `std::terminate` fired with "Resource + deadlock avoided". That thread now owns a share of the connection and + session: a destructor running on it detaches, and the thread frees them + once `run()` returns. Found as an intermittent failure in the consumer + module's async acceptance test. - **A modeled `@httpHeader("Accept")` input member reaches the wire.** An operation with a response `@httpPayload` emits that payload's content type as the request's Accept header — but it did so with an unconditional `Set`, diff --git a/runtime/src/http/beast_transport.cc b/runtime/src/http/beast_transport.cc index 8123c7a..5f165c4 100644 --- a/runtime/src/http/beast_transport.cc +++ b/runtime/src/http/beast_transport.cc @@ -2417,54 +2417,69 @@ struct DialedConnection { } }; +// The io machinery a dialed socket shares with its io thread. Declaration +// order is load-bearing: members destroy in reverse, so the session (whose +// websocket stream references the io_context) is released BEFORE the +// connection destroys that context. +struct DialedIo { + std::unique_ptr connection; + std::shared_ptr session; +}; + // The handle the application owns: delegates to the session and keeps the -// io thread alive for the connection's lifetime. Destruction (an app -// thread — pump callbacks never own this object) aborts the session, stops -// the io, and joins. +// io thread alive for the connection's lifetime. Destruction aborts the +// session, stops the io, and joins, unless it runs on the io thread itself. +// That happens when a completion owns the last handle (a Detached loop that +// owns its dialed socket finishes on this thread): a thread cannot join +// itself, so it is detached, and its own reference to DialedIo frees the io +// machinery once run() returns, which the stop makes immediate. class DialedWebSocket final : public WebSocket { public: DialedWebSocket(std::shared_ptr session, std::unique_ptr connection) - : connection_(std::move(connection)), - session_(std::move(session)), - work_guard_(asio::make_work_guard(connection_->io)), - runner_([this] { RunIoThread(connection_->io, "beast client websocket"); }) {} + : io_(std::make_shared( + DialedIo{.connection = std::move(connection), .session = std::move(session)})), + work_guard_(asio::make_work_guard(io_->connection->io)), + runner_([io = io_] { RunIoThread(io->connection->io, "beast client websocket"); }) {} ~DialedWebSocket() override { - session_->Abort("client released the session"); + io_->session->Abort("client released the session"); work_guard_.reset(); - connection_->io.stop(); - if (runner_.joinable()) { - runner_.join(); + io_->connection->io.stop(); + if (!runner_.joinable()) { + return; } + if (runner_.get_id() == std::this_thread::get_id()) { + runner_.detach(); + return; + } + runner_.join(); } - Outcome> Receive() override { return session_->Receive(); } + Outcome> Receive() override { + return io_->session->Receive(); + } Outcome> Receive(std::chrono::milliseconds timeout) override { - return session_->Receive(timeout); + return io_->session->Receive(timeout); } Outcome Send(const eventstream::Message& message) override { - return session_->Send(message); + return io_->session->Send(message); } - void Close() override { session_->Close(); } + void Close() override { io_->session->Close(); } void ReceiveAsync(WebSocket::ReceiveCallback callback) override { - session_->ReceiveAsync(std::move(callback)); + io_->session->ReceiveAsync(std::move(callback)); } void ReceiveAsync(std::chrono::milliseconds timeout, WebSocket::ReceiveCallback callback) override { - session_->ReceiveAsync(timeout, std::move(callback)); + io_->session->ReceiveAsync(timeout, std::move(callback)); } void SendAsync(const eventstream::Message& message, WebSocket::SendCallback callback) override { - session_->SendAsync(message, std::move(callback)); + io_->session->SendAsync(message, std::move(callback)); } - bool SupportsAsync() const override { return session_->SupportsAsync(); } + bool SupportsAsync() const override { return io_->session->SupportsAsync(); } private: - // Declaration order is load-bearing: members destroy in reverse, so the - // session (whose websocket stream references the io_context) must be - // released BEFORE connection_ destroys that context. - std::unique_ptr connection_; - std::shared_ptr session_; + std::shared_ptr io_; // shared with runner_, which may outlive this asio::executor_work_guard work_guard_; std::thread runner_; }; diff --git a/runtime/tests/http/beast_websocket_test.cc b/runtime/tests/http/beast_websocket_test.cc index a8f7cd1..67a3247 100644 --- a/runtime/tests/http/beast_websocket_test.cc +++ b/runtime/tests/http/beast_websocket_test.cc @@ -106,6 +106,55 @@ TEST(BeastWebSocketTest, MessagesRoundTripBothWaysOverTheUpgrade) { server.Stop(); } +TEST(BeastWebSocketTest, ADetachedLoopMayReleaseTheLastHandleOnTheSocketsOwnIoThread) { + // A Detached client loop that owns its dialed socket resumes, finishes and + // is destroyed on that socket's io thread, so the last reference drops + // there. Teardown must not join the thread it runs on (std::terminate, + // "Resource deadlock avoided") or free the io_context under it. + std::promise parked; + std::shared_future go = parked.get_future().share(); + BeastServerTransport::Options options; + options.on_websocket = [go](const HttpRequest&, WebSocket& socket) { + go.wait(); // push only once the client's receive is parked + (void)socket.Send(Text("news", "unsolicited")); + while (true) { + auto message = socket.Receive(); + if (!message.ok() || !message->has_value()) return; + } + }; + BeastServerTransport server(options); + ASSERT_TRUE(server.Start(NotFoundHandler()).ok()); + auto dialed = BeastWebSocketClient::Dial({.host = "127.0.0.1", .port = server.port()}); + ASSERT_TRUE(dialed.ok()) << dialed.error().message(); + + std::promise delivered; + std::promise frame_gone; + // The frame owns the only reference once `dialed` is released below. + struct SignalOnExit { + std::promise* gone; + ~SignalOnExit() { gone->set_value(); } + }; + [](std::shared_ptr socket, std::promise* delivered, + std::promise* gone) -> eventstream::Detached { + const SignalOnExit signal{gone}; // destroyed after `stream`, so after the release + eventstream::AsyncEventStream stream( + std::move(socket), [](const Message& m) -> Outcome { return m; }, + [](const Message& m) -> Outcome { return m; }); + auto message = co_await stream.Receive(); + delivered->set_value(message.ok() && message->has_value() ? (**message).payload.ToString() + : ""); + }(std::exchange(*dialed, nullptr), &delivered, &frame_gone); + parked.set_value(); + + auto body = delivered.get_future(); + ASSERT_EQ(body.wait_for(std::chrono::seconds(5)), std::future_status::ready); + EXPECT_EQ(body.get(), "unsolicited"); + // The frame (and with it ~DialedWebSocket) finished on the io thread + // without taking the process down. + EXPECT_EQ(frame_gone.get_future().wait_for(std::chrono::seconds(5)), std::future_status::ready); + server.Stop(); +} + TEST(BeastWebSocketTest, AReceiveDeadlineExpiresOnAQuietWireAndSparesTheSession) { // The echo server says nothing unsolicited, so a receive with no deadline // would park here forever — the hang this overload exists to prevent.