You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于Boost.Asio协程实现单线程WebSocket读写及挂起问题求助

问题描述

我想用C++20协程结合boost::asio实现单线程WebSocket,需求如下:

  • 拆分读写协程:一个负责async_write,一个负责async_read
  • 异常重连:任一协程抛出异常(假设均代表连接断开)时自动尝试重新连接
  • 延迟写入:外部调用ws.Send(data)时,消息先存入writeBuffer队列,直到触发发送逻辑才执行写入

遇到的核心错误:当writeBuffer为空时,尝试用std::suspend_always{}挂起写入协程,触发编译错误:

error C2665: 'boost::asio::detail::awaitable_frame_base<Executor>::await_transform': no overloaded function could convert all the argument types

额外需求:

  • 替换io_context为单线程thread_pool,让协程脱离主线程后台运行
  • 无读写操作时让线程休眠50ms,避免空转浪费资源

解决方案

1. 正确实现协程挂起:使用asio条件变量

不能直接co_await std::suspend_always,因为asio的awaitable框架需要适配的可等待对象。这里用asio::condition_variable_any实现队列通知机制:

  • 给writeBuffer加互斥锁(外部线程调用Write时必须保证线程安全)
  • 缓冲区为空时,通过async_wait挂起协程;有新消息时唤醒协程

2. 修复读取逻辑

原代码中async_read_some的使用不符合WebSocket的读取规范,改用async_read完整读取单条消息,同时用flat_buffer管理读取缓冲区,避免溢出和分片处理错误。

3. 替换io_context为thread_pool并优化线程休眠

用asio::thread_pool(1)创建单线程池替代io_context,后台自动调度协程。通过条件变量和定时器结合,实现无操作时的休眠逻辑。


修改后的完整代码
#include <iostream>
#include <coroutine>
#include <optional>
#include <mutex>
#include <queue>
#include <atomic>

#include <boost/asio.hpp>
#include <boost/beast.hpp>
#include <boost/asio/awaitable.hpp>
#include <boost/asio/experimental/awaitable_operators.hpp>

namespace asio = boost::asio;
namespace beast = boost::beast;
namespace websocket = beast::websocket;
using namespace std::chrono_literals;
using namespace asio::experimental::awaitable_operators;

struct CoroWebsocket {
    CoroWebsocket(std::string host, std::string port)
        : _host(std::move(host))
        , _port(std::move(port))
        , _ioc(1) // 单线程线程池
        , _ws(_ioc)
        , _write_cv(_ioc) {
        asio::co_spawn(_ioc, do_run(), asio::detached);
    }

    ~CoroWebsocket() {
        _stop = true;
        _write_cv.notify_all();
        _ioc.stop();
        _ioc.join();
    }

    void Write(std::string data) {
        std::lock_guard<std::mutex> lock(_write_mutex);
        _writeBuffer.push(std::move(data));
        _write_cv.notify_one(); // 唤醒写入协程
    }

    std::optional<std::string> Read() {
        std::lock_guard<std::mutex> lock(_read_mutex);
        if (_readBuffer.empty())
            return {};
        auto message = std::move(_readBuffer.front());
        _readBuffer.pop();
        return message;
    }

private:
    const std::string _host, _port;
    using tcp = asio::ip::tcp;
    std::queue<std::string>       _writeBuffer;
    std::queue<std::string>       _readBuffer;
    asio::thread_pool             _ioc;
    websocket::stream<tcp::socket> _ws;
    std::mutex _write_mutex, _read_mutex;
    asio::condition_variable_any _write_cv;
    std::atomic<bool> _stop{false};

    asio::awaitable<void> do_run() {
        while (!_stop) {
            try {
                co_await do_connect();
                // 同时运行读写协程,任一结束则触发重连
                co_await (do_write() || do_read());
            } catch (const boost::system::system_error& se) {
                std::cerr << "连接异常,准备重连: " << se.code().message() << std::endl;
                // 重连前休眠1秒,避免频繁重试
                co_await asio::steady_timer(_ioc, 1s).async_wait(asio::use_awaitable);
            }
        }
    }

    asio::awaitable<void> do_connect() {
        tcp::resolver resolver(_ioc);
        auto endpoints = co_await resolver.async_resolve(_host, _port, asio::use_awaitable);

        while (!_stop) {
            try {
                co_await asio::async_connect(_ws.next_layer(), endpoints, asio::use_awaitable);
                _ws.set_option(websocket::stream_base::decorator([](websocket::request_type& req) {
                    req.set(beast::http::field::user_agent, BOOST_BEAST_VERSION_STRING " WsConnect");
                }));
                co_await _ws.async_handshake(_host + ':' + _port, "/", asio::use_awaitable);
                std::cerr << "连接成功" << std::endl;
                co_return; // 连接成功退出重试循环
            } catch (boost::system::system_error const& se) {
                std::cerr << "连接失败: " << se.code().message() << ",1秒后重试" << std::endl;
                co_await asio::steady_timer(_ioc, 1s).async_wait(asio::use_awaitable);
            }
        }
    }

    asio::awaitable<void> do_write() {
        while (!_stop) {
            std::unique_lock<std::mutex> lock(_write_mutex);
            // 等待缓冲区有消息或停止信号
            co_await _write_cv.async_wait(lock, [this] { return !_writeBuffer.empty() || _stop; }, asio::use_awaitable);
            
            if (_stop) co_return;

            // 批量取出所有待发消息,释放锁避免阻塞写入
            std::queue<std::string> messages;
            std::swap(messages, _writeBuffer);
            lock.unlock();

            for (auto& msg : messages) {
                co_await _ws.async_write(asio::buffer(msg), asio::use_awaitable);
            }
        }
    }

    asio::awaitable<void> do_read() {
        beast::flat_buffer buffer;
        while (!_stop) {
            co_await _ws.async_read(buffer, asio::use_awaitable);
            std::lock_guard<std::mutex> lock(_read_mutex);
            _readBuffer.push(beast::buffers_to_string(buffer.data()));
            buffer.consume(buffer.size()); // 清空缓冲区,准备下一次读取
        }
    }
};

关键修改说明

  • 协程挂起修复:用asio::condition_variable_any的async_wait替代std::suspend_always,完美适配asio的awaitable框架
  • 线程安全:读写缓冲区均添加互斥锁,避免外部线程调用Write/Read时的竞态条件
  • thread_pool替换:单线程线程池替代io_context,后台自动运行,无需手动调用Run
  • 读取逻辑修复:用async_read完整读取WebSocket消息,避免分片处理错误,flat_buffer自动管理内存
  • 优雅停止:添加_stop原子变量,析构时安全停止所有协程和线程池
  • 重连优化:连接失败/异常后添加1秒休眠,避免频繁重试消耗资源

内容的提问来源于stack exchange,提问作者Some Dinosaur

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 09:43:16