Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`,
Expand Down
63 changes: 39 additions & 24 deletions runtime/src/http/beast_transport.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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<DialedConnection> connection;
std::shared_ptr<WebSocketSessionBase> 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<WebSocketSessionBase> session,
std::unique_ptr<DialedConnection> 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>(
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<std::optional<eventstream::Message>> Receive() override { return session_->Receive(); }
Outcome<std::optional<eventstream::Message>> Receive() override {
return io_->session->Receive();
}
Outcome<std::optional<eventstream::Message>> Receive(std::chrono::milliseconds timeout) override {
return session_->Receive(timeout);
return io_->session->Receive(timeout);
}
Outcome<Unit> 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<DialedConnection> connection_;
std::shared_ptr<WebSocketSessionBase> session_;
std::shared_ptr<DialedIo> io_; // shared with runner_, which may outlive this
asio::executor_work_guard<asio::io_context::executor_type> work_guard_;
std::thread runner_;
};
Expand Down
49 changes: 49 additions & 0 deletions runtime/tests/http/beast_websocket_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> parked;
std::shared_future<void> 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<std::string> delivered;
std::promise<void> frame_gone;
// The frame owns the only reference once `dialed` is released below.
struct SignalOnExit {
std::promise<void>* gone;
~SignalOnExit() { gone->set_value(); }
};
[](std::shared_ptr<WebSocket> socket, std::promise<std::string>* delivered,
std::promise<void>* gone) -> eventstream::Detached {
const SignalOnExit signal{gone}; // destroyed after `stream`, so after the release
eventstream::AsyncEventStream<Message, Message> stream(
std::move(socket), [](const Message& m) -> Outcome<Message> { return m; },
[](const Message& m) -> Outcome<Message> { return m; });
auto message = co_await stream.Receive();
delivered->set_value(message.ok() && message->has_value() ? (**message).payload.ToString()
: "<none>");
}(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.
Expand Down
Loading