Boost.Asio是否提供支持同步/异步读写的简单身份流类?
实现Boost.Asio兼容的同步/异步身份流用于测试
好问题!Boost.Asio本身并没有直接提供你想要的这种simple_stream类,但你完全可以基于Asio的核心模型自己实现一个,而且实现起来并不复杂——正好能满足你模拟串口、支持同步/异步读写以及协程的需求。
核心实现思路
我们要创建的流类需要满足Asio的同步流和异步流概念,核心是用一个线程安全的环形缓冲区作为数据中转层:
- 写操作把数据存入缓冲区,同步操作直接返回写入长度,异步操作完成后通知回调
- 读操作从缓冲区取出数据,同步操作会阻塞直到有数据(或者可配置超时),异步操作在有数据时触发回调
- 用Asio的
strand来保证多线程下的操作线程安全,同时兼容协程的异步等待逻辑
完整实现示例
#include <boost/asio.hpp> #include <boost/circular_buffer.hpp> #include <mutex> #include <condition_variable> namespace asio = boost::asio; class simple_identity_stream : public asio::async_write_stream , public asio::async_read_stream { public: explicit simple_identity_stream(asio::io_context& io_ctx) : strand_(asio::make_strand(io_ctx)) , buffer_(1024) // 缓冲区大小可按需调整 {} // 同步写操作 std::size_t write_some(const asio::const_buffer& buf) { std::lock_guard<std::mutex> lock(mutex_); const char* data = asio::buffer_cast<const char*>(buf); std::size_t size = asio::buffer_size(buf); std::size_t written = std::min(size, buffer_.remaining_space()); buffer_.insert(buffer_.end(), data, data + written); // 通知等待读操作的线程/异步任务 cv_.notify_all(); return written; } // 同步读操作 std::size_t read_some(const asio::mutable_buffer& buf) { std::unique_lock<std::mutex> lock(mutex_); // 阻塞直到有数据 cv_.wait(lock, [this] { return !buffer_.empty(); }); char* data = asio::buffer_cast<char*>(buf); std::size_t size = asio::buffer_size(buf); std::size_t read = std::min(size, buffer_.size()); std::copy(buffer_.begin(), buffer_.begin() + read, data); buffer_.erase(buffer_.begin(), buffer_.begin() + read); return read; } // 异步写操作 template <typename WriteHandler> void async_write_some(const asio::const_buffer& buf, WriteHandler&& handler) { asio::post(strand_, [this, buf, handler = std::forward<WriteHandler>(handler)]() mutable { try { std::size_t n = write_some(buf); asio::post(strand_, std::bind(std::move(handler), std::error_code(), n)); } catch (const std::exception& e) { asio::post(strand_, std::bind(std::move(handler), asio::error::operation_aborted, 0)); } }); } // 异步读操作 template <typename ReadHandler> void async_read_some(const asio::mutable_buffer& buf, ReadHandler&& handler) { asio::post(strand_, [this, buf, handler = std::forward<ReadHandler>(handler)]() mutable { try { std::size_t n = read_some(buf); asio::post(strand_, std::bind(std::move(handler), std::error_code(), n)); } catch (const std::exception& e) { asio::post(strand_, std::bind(std::move(handler), asio::error::operation_aborted, 0)); } }); } // 兼容Asio流概念的必要类型定义 using executor_type = asio::strand<asio::io_context::executor_type>; executor_type get_executor() noexcept { return strand_; } private: executor_type strand_; boost::circular_buffer<char> buffer_; std::mutex mutex_; std::condition_variable cv_; };
用法示例
同步读写(符合你的需求)
asio::io_context io_ctx; simple_identity_stream s(io_ctx); // 同步写 std::size_t bytes_written = s.write_some(asio::buffer("test")); // 同步读 char buf[1024]; std::size_t bytes_read = s.read_some(asio::buffer(buf));
异步读写(包括协程支持)
// 使用Asio协程的异步读 asio::co_spawn(io_ctx, [&s]() -> asio::awaitable<void> { char buf[1024]; std::size_t n = co_await asio::async_read_some(s, asio::buffer(buf), asio::use_awaitable); // 处理读取到的数据 }, asio::detached); // 异步写 s.async_write_some(asio::buffer("async test"), [](const std::error_code& ec, std::size_t n) { if (!ec) { // 写操作完成 } }); io_ctx.run();
额外说明
- 你可以根据测试需求扩展这个类,比如添加超时机制、清空缓冲区的方法,或者模拟错误场景(比如返回特定的error_code)
- 环形缓冲区的大小可以根据你的测试数据量调整,避免溢出
- 用
strand保证了异步操作的线程安全,同时兼容Asio的协程模型,因为协程会自动遵循strand的执行顺序
内容的提问来源于stack exchange,提问作者Liarokapis Alexandros
相关产品推荐
相关产品推荐

