如何在含同步/异步接口的库中正确使用asio::strand?
我正在通过编写一个基于UDP连接特定设备的库来学习ASIO。设备开机后会进入广播模式,向网络发送身份信息,因此主机需要监听UDP端口,匹配设备格式的数据包,并跟踪已检测到的设备(因为设备连接后仍会持续广播)。我曾用Python异步生成器完成原型开发,在ASIO中,类似的实现类比为asio::ip::tcp::acceptor。由于每个设备数据率较高,我希望自定义acceptor运行在多线程io_context中,因此需要为acceptor及每个设备控制器使用strand来序列化操作。
以下是我当前的示例代码:
template <typename Executor> class Acceptor { std::unordered_set<asio::ip::udp::endpoint> accepted_connections; asio::basic_datagram_socket<asio::ip::udp, Executor> receive_sock; std::shared_ptr<Device<Executor>> accept(asio::error_code& ec) { std::array<std::byte, buffer_size> buffer; typename asio::basic_datagram_socket<asio::ip::udp, Executor>::endpoint_type remote_endpoint; receive_sock.receive_from(asio::buffer(buffer), remote_endpoint, {}, ec); if (!ec) { return {}; } if (accepted_connections.contains(remote_endpoint)) { return {}; } if (auto ret = std::start_lifetime_as<Header>(buffer.data())->validate(); !ret) { ec = ret.error(); return {}; } accepted_connections.emplace(remote_endpoint.remote_endpoint()); return MakeDevice(remote_endpoint, buffer); } };
问题1:如何修改阻塞式accept函数,使其在strand中运行,确保对accepted_connections集合的访问被正确序列化?
我发现strand的defer、dispatch、execute、post等操作都会将函数推入队列延后执行,这会导致无法即时检查error_code(因为函数延后执行时,accept返回时ec变量尚未赋值)。查看asio::ip::tcp::acceptor的accept实现源码也未得到有效信息。
解决方案
要让阻塞式accept的所有逻辑(socket接收、集合访问)都在strand上下文中执行,同时能同步获取结果,核心是用内部strand+同步Promise的方式实现:
- 在Acceptor类内部持有
asio::strand<Executor>,将socket绑定到strand执行器,确保socket操作也受strand序列化。 - 用
std::promise同步获取strand中执行的操作结果,解决即时获取error_code和返回值的问题。
修改后的代码:
template <typename Executor> class Acceptor { asio::strand<Executor> strand_; std::unordered_set<asio::ip::udp::endpoint> accepted_connections; asio::basic_datagram_socket<asio::ip::udp, decltype(strand_)> receive_sock; public: Acceptor(Executor ex) : strand_(std::move(ex)), receive_sock(strand_) {} std::shared_ptr<Device<decltype(strand_)>> accept(asio::error_code& ec) { std::promise<std::pair<asio::error_code, std::shared_ptr<Device<decltype(strand_)>>>> prom; auto fut = prom.get_future(); // 将accept逻辑投递到strand执行 asio::dispatch(strand_, [this, &prom]() { asio::error_code local_ec; std::array<std::byte, buffer_size> buffer; asio::ip::udp::endpoint remote_endpoint; receive_sock.receive_from(asio::buffer(buffer), remote_endpoint, {}, local_ec); std::shared_ptr<Device<decltype(strand_)>> dev; if (!local_ec) { if (!accepted_connections.contains(remote_endpoint)) { if (auto ret = std::start_lifetime_as<Header>(buffer.data())->validate(); ret) { accepted_connections.emplace(remote_endpoint); dev = MakeDevice(remote_endpoint, buffer); } else { local_ec = ret.error(); } } } prom.set_value({local_ec, std::move(dev)}); }); // 同步等待结果返回 auto result = fut.get(); ec = result.first; return result.second; } };
关键说明
- 内部strand确保所有对
accepted_connections的读写和socket操作都被序列化,多线程环境下不会出现并发冲突。 std::promise和std::future解决了strand队列操作无法即时返回结果的问题,调用线程会同步等待strand内的操作完成。- socket绑定到strand执行器,保证即使是阻塞的
receive_from也在strand上下文中执行,避免跨线程操作socket的潜在问题。
问题2:实现async_accept时,如何将strand融入其中以保证accepted_connections集合的操作有序?
我需要使用async_initiate来兼容各种完成令牌,但不确定是在协程中await strand,还是需要修改async_initiate让协程在strand中运行?
解决方案
实现async_accept时,要让整个异步流程(接收、验证、集合操作、回调)都在strand上下文中执行,结合async_initiate和协程的实现方式如下:
- 编写异步协程
async_accept_impl,利用绑定到strand的socket执行异步接收,后续逻辑自然在strand中执行。 - 在
async_initiate内部,通过strand调度协程的启动和完成回调,确保所有操作都被序列化。
修改后的代码:
template <typename Executor> class Acceptor { asio::strand<Executor> strand_; std::unordered_set<asio::ip::udp::endpoint> accepted_connections; asio::basic_datagram_socket<asio::ip::udp, decltype(strand_)> receive_sock; // 异步接收处理协程 asio::awaitable<std::shared_ptr<Device<decltype(strand_)>>> async_accept_impl() { std::array<std::byte, buffer_size> buffer; asio::ip::udp::endpoint remote_endpoint; asio::error_code ec; // 异步接收,socket绑定到strand,操作自动在strand上下文执行 co_await receive_sock.async_receive_from( asio::buffer(buffer), remote_endpoint, asio::redirect_error(asio::use_awaitable, ec) ); if (ec) { co_return nullptr; } // 集合访问在strand上下文内,无需额外同步 if (accepted_connections.contains(remote_endpoint)) { co_return nullptr; } if (auto ret = std::start_lifetime_as<Header>(buffer.data())->validate(); !ret) { co_return nullptr; } accepted_connections.emplace(remote_endpoint); co_return MakeDevice(remote_endpoint, buffer); } public: Acceptor(Executor ex) : strand_(std::move(ex)), receive_sock(strand_) {} // 兼容任意完成令牌的async_accept template <typename CompletionToken> auto async_accept(CompletionToken&& token) { return asio::async_initiate<CompletionToken, void(asio::error_code, std::shared_ptr<Device<decltype(strand_)>>)>( [this](auto&& handler) { // 确保协程在strand上下文启动 asio::dispatch(strand_, [this, handler = std::forward<decltype(handler)>(handler)]() mutable { asio::co_spawn(strand_, async_accept_impl(), [handler = std::move(handler)](std::exception_ptr e, std::shared_ptr<Device<decltype(strand_)>> dev) mutable { asio::error_code ec; if (e) { try { std::rethrow_exception(e); } catch (const asio::system_error& se) { ec = se.code(); } catch (...) { ec = asio::error::operation_aborted; } } // 回调也通过strand调度,保证序列化 asio::dispatch(strand_, [handler = std::move(handler), ec, dev]() mutable { handler(ec, dev); }); } ); }); }, token ); } };
关键说明
- 协程
async_accept_impl直接使用绑定到strand的socket,异步接收操作自动在strand上下文中执行,后续的集合访问和验证逻辑无需额外同步。 - 在
async_initiate内部,通过asio::dispatch(strand_, ...)启动协程,确保整个异步流程都处于strand的序列化上下文。 - 完成回调同样通过strand调度,避免回调在非strand线程执行时引发并发问题。
- 如果不使用协程,也可以采用链式异步操作:先调用
socket.async_receive_from,在其完成回调中通过strand调度集合检查、设备创建逻辑,最后调用完成令牌。
内容的提问来源于stack exchange,提问作者DBS4261

