基于Boost.Beast的WebSocket客户端与服务端实现示例求助
Boost.Beast WebSocket 客户端与服务端完整实现
一、服务端实现(含消息响应、订阅推送)
服务端实现了消息响应、连接管理、主题订阅推送三个核心功能,修复了原代码中缺少订阅逻辑、错误处理不完善的问题:
#include <boost/asio.hpp> #include <boost/beast.hpp> #include <nlohmann/json.hpp> #include <iostream> #include <memory> #include <string> #include <unordered_map> #include <set> #include <mutex> #include <chrono> namespace asio = boost::asio; namespace beast = boost::beast; namespace websocket = beast::websocket; using tcp = asio::ip::tcp; using json = nlohmann::json; // 会话管理器:管理所有连接与订阅关系 class SessionManager { public: void add_session(std::shared_ptr<class WebSocketSession> session) { std::lock_guard<std::mutex> lock(mtx_); sessions_.insert(session); } void remove_session(std::shared_ptr<class WebSocketSession> session) { std::lock_guard<std::mutex> lock(mtx_); sessions_.erase(session); // 清理该会话的所有订阅 for (auto& [topic, subs] : subscriptions_) { subs.erase(session); } } void subscribe(std::shared_ptr<class WebSocketSession> session, const std::string& topic) { std::lock_guard<std::mutex> lock(mtx_); subscriptions_[topic].insert(session); } void unsubscribe(std::shared_ptr<class WebSocketSession> session, const std::string& topic) { std::lock_guard<std::mutex> lock(mtx_); if (subscriptions_.count(topic)) { subscriptions_[topic].erase(session); } } // 向指定主题推送消息 void publish(const std::string& topic, const std::string& message) { std::lock_guard<std::mutex> lock(mtx_); if (!subscriptions_.count(topic)) return; for (auto session : subscriptions_[topic]) { session->send(message); } } private: std::mutex mtx_; std::set<std::shared_ptr<class WebSocketSession>> sessions_; std::unordered_map<std::string, std::set<std::shared_ptr<class WebSocketSession>>> subscriptions_; }; class WebSocketSession : public std::enable_shared_from_this<WebSocketSession> { public: WebSocketSession(tcp::socket socket, SessionManager& manager) : ws_(std::move(socket)), manager_(manager) { manager_.add_session(shared_from_this()); } ~WebSocketSession() { manager_.remove_session(shared_from_this()); } void run() { ws_.async_accept( beast::bind_front_handler( &WebSocketSession::on_accept, shared_from_this() ) ); } // 异步发送消息,保证线程安全 void send(const std::string& message) { auto self = shared_from_this(); asio::post(ws_.get_executor(), [self, message]() { bool write_in_progress = !write_queue_.empty(); write_queue_.push_back(message); if (!write_in_progress) { do_write(); } }); } private: void on_accept(beast::error_code ec) { if (ec) { std::cerr << "连接接收失败: " << ec.message() << std::endl; return; } do_read(); } void do_read() { ws_.async_read( buffer_, beast::bind_front_handler( &WebSocketSession::on_read, shared_from_this() ) ); } void on_read(beast::error_code ec, std::size_t bytes_transferred) { boost::ignore_unused(bytes_transferred); if (ec) { if (ec == websocket::error::closed) return; std::cerr << "读取消息失败: " << ec.message() << std::endl; return; } try { auto msg_str = beast::buffers_to_string(buffer_.data()); json msg = json::parse(msg_str); buffer_.consume(buffer_.size()); // 处理不同类型请求 if (msg.contains("echo")) { json response = { {"type", "echo_response"}, {"original", msg["echo"]} }; send(response.dump()); } else if (msg.contains("subscribe")) { std::string topic = msg["subscribe"]; manager_.subscribe(shared_from_this(), topic); json response = { {"type", "subscribe_ack"}, {"topic", topic}, {"status", "success"} }; send(response.dump()); } else if (msg.contains("unsubscribe")) { std::string topic = msg["unsubscribe"]; manager_.unsubscribe(shared_from_this(), topic); json response = { {"type", "unsubscribe_ack"}, {"topic", topic}, {"status", "success"} }; send(response.dump()); } else { json response = { {"type", "error"}, {"message", "未知消息类型"} }; send(response.dump()); } do_read(); } catch (const std::exception& e) { std::cerr << "消息处理失败: " << e.what() << std::endl; json error_resp = {{"type", "error"}, {"message", "消息格式错误"}}; send(error_resp.dump()); do_read(); } } void do_write() { ws_.text(ws_.got_text()); ws_.async_write( asio::buffer(write_queue_.front()), beast::bind_front_handler( &WebSocketSession::on_write, shared_from_this() ) ); } void on_write(beast::error_code ec, std::size_t bytes_transferred) { boost::ignore_unused(bytes_transferred); if (ec) { std::cerr << "发送消息失败: " << ec.message() << std::endl; return; } write_queue_.pop_front(); if (!write_queue_.empty()) { do_write(); } } websocket::stream<tcp::socket> ws_; beast::flat_buffer buffer_; SessionManager& manager_; std::deque<std::string> write_queue_; }; class WebSocketServer { public: WebSocketServer(asio::io_context& ioc, tcp::endpoint endpoint) : acceptor_(ioc, endpoint), manager_(std::make_shared<SessionManager>()) { do_accept(); // 模拟定时向news主题推送消息 std::thread push_thread([this]() { int count = 0; while (true) { std::this_thread::sleep_for(std::chrono::seconds(5)); json push_msg = { {"type", "push"}, {"topic", "news"}, {"content", "定时推送消息 " + std::to_string(++count)} }; manager_->publish("news", push_msg.dump()); std::cout << "已向news主题推送消息" << std::endl; } }); push_thread.detach(); } private: void do_accept() { acceptor_.async_accept( beast::bind_front_handler( &WebSocketServer::on_accept, this ) ); } void on_accept(beast::error_code ec, tcp::socket socket) { if (ec) { std::cerr << "接受连接失败: " << ec.message() << std::endl; } else { std::make_shared<WebSocketSession>(std::move(socket), *manager_)->run(); } do_accept(); } tcp::acceptor acceptor_; std::shared_ptr<SessionManager> manager_; }; int main() { std::cout << "WebSocket服务端已启动,监听端口9002..." << std::endl; try { asio::io_context ioc{1}; tcp::endpoint endpoint(tcp::v4(), 9002); WebSocketServer server(ioc, endpoint); ioc.run(); } catch (const std::exception& e) { std::cerr << "服务端启动失败: " << e.what() << std::endl; return EXIT_FAILURE; } return EXIT_SUCCESS; }
服务端核心说明
- 会话管理:通过
SessionManager维护所有连接,线程安全地处理订阅、取消订阅操作。 - 消息响应:解析JSON消息,对合法请求返回对应响应,非法消息返回错误提示,避免客户端无响应等待。
- 订阅推送:模拟定时向指定主题推送消息,遍历订阅该主题的所有会话异步发送。
- 异步发送队列:每个会话维护消息队列,保证异步发送的顺序,避免并发发送冲突。
二、客户端实现(含订阅、消息收发)
修复了原代码中多线程直接操作WebSocket的线程安全问题,使用asio::strand序列化操作,实现完整的订阅、收发逻辑:
#include <boost/asio.hpp> #include <boost/beast.hpp> #include <nlohmann/json.hpp> #include <iostream> #include <memory> #include <string> #include <thread> namespace asio = boost::asio; namespace beast = boost::beast; namespace websocket = beast::websocket; using tcp = asio::ip::tcp; using json = nlohmann::json; class WebSocketClient : public std::enable_shared_from_this<WebSocketClient> { public: WebSocketClient(asio::io_context& ioc) : ioc_(ioc), strand_(asio::make_strand(ioc)) {} void connect(const std::string& host, const std::string& port) { resolver_.async_resolve(host, port, asio::bind_executor(strand_, beast::bind_front_handler( &WebSocketClient::on_resolve, shared_from_this() ) ) ); } // 线程安全的消息发送 void send(const json& msg) { asio::post(strand_, [this, msg_str = msg.dump()]() { bool write_in_progress = !write_queue_.empty(); write_queue_.push_back(msg_str); if (!write_in_progress) { do_write(); } }); } private: void on_resolve(beast::error_code ec, tcp::resolver::results_type results) { if (ec) { std::cerr << "解析地址失败: " << ec.message() << std::endl; return; } asio::async_connect(ws_.next_layer(), results.begin(), results.end(), asio::bind_executor(strand_, beast::bind_front_handler( &WebSocketClient::on_connect, shared_from_this() ) ) ); } void on_connect(beast::error_code ec, tcp::resolver::results_type::endpoint_type) { if (ec) { std::cerr << "连接服务端失败: " << ec.message() << std::endl; return; } ws_.async_handshake("127.0.0.1", "/", asio::bind_executor(strand_, beast::bind_front_handler( &WebSocketClient::on_handshake, shared_from_this() ) ) ); } void on_handshake(beast::error_code ec) { if (ec) { std::cerr << "握手失败: " << ec.message() << std::endl; return; } std::cout << "已连接到服务端" << std::endl; do_read(); } void do_read() { ws_.async_read(buffer_, asio::bind_executor(strand_, beast::bind_front_handler( &WebSocketClient::on_read, shared_from_this() ) ) ); } void on_read(beast::error_code ec, std::size_t bytes_transferred) { boost::ignore_unused(bytes_transferred); if (ec) { if (ec == websocket::error::closed) { std::cerr << "连接已关闭" << std::endl; } else { std::cerr << "读取消息失败: " << ec.message() << std::endl; } return; } try { auto msg_str = beast::buffers_to_string(buffer_.data()); json msg = json::parse(msg_str); buffer_.consume(buffer_.size()); // 处理服务端消息 if (msg["type"] == "echo_response") { std::cout << "\n收到echo响应: " << msg["original"] << std::endl; } else if (msg["type"] == "subscribe_ack") { std::cout << "\n订阅主题" << msg["topic"] << "成功" << std::endl; } else if (msg["type"] == "unsubscribe_ack") { std::cout << "\n取消订阅主题" << msg["topic"] << "成功" << std::endl; } else if (msg["type"] == "push") { std::cout << "\n收到推送消息(" << msg["topic"] << "): " << msg["content"] << std::endl; } else if (msg["type"] == "error") { std::cout << "\n收到错误消息: " << msg["message"] << std::endl; } do_read(); } catch (const std::exception& e) { std::cerr << "解析消息失败: " << e.what() << std::endl; do_read(); } } void do_write() { ws_.text(ws_.got_text()); ws_.async_write(asio::buffer(write_queue_.front()), asio::bind_executor(strand_, beast::bind_front_handler( &WebSocketClient::on_write, shared_from_this() ) ) ); } void on_write(beast::error_code ec, std::size_t bytes_transferred) { boost::ignore_unused(bytes_transferred); if (ec) { std::cerr << "发送消息失败: " << ec.message() << std::endl; return; } write_queue_.pop_front(); if (!write_queue_.empty()) { do_write(); } } asio::io_context& ioc_; asio::strand<asio::io_context::executor_type> strand_; tcp::resolver resolver_{strand_}; websocket::stream<tcp::socket> ws_{strand_}; beast::flat_buffer buffer_; std::deque<std::string> write_queue_; }; int main() { try { asio::io_context ioc; auto client = std::make_shared<WebSocketClient>(ioc); client->connect("127.0.0.1", "9002"); // 启动io线程 std::thread io_thread([&ioc]() { ioc.run(); }); // 控制台输入处理 std::string input; while (true) { std::cout << "\n输入指令(echo <内容>/subscribe <主题>/unsubscribe <主题>/stop): "; std::getline(std::cin, input); if (input == "stop") break; else if (input.substr(0,4) == "echo" && input.size()>5) { client->send({{"echo", input.substr(5)}}); } else if (input.substr(0,9) == "subscribe" && input.size()>10) { client->send({{"subscribe", input.substr(10)}}); } else if (input.substr(0,11) == "unsubscribe" && input.size()>12) { client->send({{"unsubscribe", input.substr(12)}}); } else { std::cout << "指令格式错误,请重新输入" << std::endl; } } // 优雅关闭连接 asio::post(client->ws_.get_executor(), [client]() { client->ws_.close(websocket::close_code::normal); }); io_thread.join(); } catch (const std::exception& e) { std::cerr << "客户端运行错误: " << e.what() << std::endl; return EXIT_FAILURE; } return EXIT_SUCCESS; }
客户端核心说明
- 线程安全:通过
相关产品推荐
相关产品推荐

