使用可等待运算符无法取消azmq::async_receive操作的问题
注:不确定这是Asio Bug、Azmq Bug还是自身使用问题。
我正在使用Azmq库从socket读取数据,该库的自定义socket类型azmq::socket是boost::asio::socket的封装。执行流程如下:
参考定时器函数
auto timer(std::chrono::steady_clock::duration dur) -> asio::awaitable<void> { asio::steady_timer timer(co_await asio::this_coro::executor); timer.expires_after(dur); co_await timer.async_wait(asio::use_awaitable); }
合并等待的核心代码
auto val = co_await(azmq::async_receive(socket_, asio::buffer(buffer), asio::use_awaitable) || timer(std::chrono::seconds{1}));
Azmq async_receive 的相关实现
顶层async_receive模板
template<class CompletionToken, class MutableBufferSequence> auto async_receive(azmq::socket &socket, MutableBufferSequence const &buffers, CompletionToken &&token) -> BOOST_ASIO_INITFN_RESULT_TYPE(CompletionToken, void(boost::system::error_code, size_t)) { return boost::asio::async_initiate<CompletionToken, void(boost::system::error_code, size_t)>(async_receive_initiation<MutableBufferSequence>{socket, buffers}, token); }
异步操作初始化器
template<typename MutableBufferSequence> struct async_receive_initiation { azmq::socket &socket; MutableBufferSequence const &buffers; template<typename CompletionHandler> void operator()(CompletionHandler &&completion_handler) { auto executor = boost::asio::get_associated_executor( completion_handler, socket.get_executor()); socket.async_receive(buffers, boost::asio::bind_executor(executor, std::bind(std::forward<CompletionHandler>(completion_handler), std::placeholders::_1, std::placeholders::_2))); } };
底层socket的async_receive实现
template<typename MessageReadHandler> void async_receive(MessageReadHandler && handler, flags_type flags = 0) { using type = detail::receive_op<MessageReadHandler>; get_service().enqueue<type>(get_implementation(), detail::socket_service::op_type::read_op, std::forward<MessageReadHandler>(handler), flags); }
问题现象
定时器触发后,async_receive操作并未被取消。查阅Asio文档未找到解决办法,尝试修改完成令牌、给异步处理结构体添加取消函数均无效。
最小复现代码
#pragma once #include <azmq/socket.hpp> #include <boost/asio.hpp> #include <boost/asio/experimental/awaitable_operators.hpp> #include <iostream> namespace asio = boost::asio; using namespace asio::experimental::awaitable_operators; auto timer(std::chrono::steady_clock::duration dur) -> asio::awaitable<void> { asio::steady_timer timer(co_await asio::this_coro::executor); timer.expires_after(dur); co_await timer.async_wait(asio::use_awaitable); } auto foo(asio::io_context &ctx) -> asio::awaitable<void> { azmq::sub_socket socket_{ctx}; socket_ = azmq::sub_socket(socket_.get_io_context(), true); boost::system::error_code error_code; std::string const socket_path{"ipc:///tmp/not_real"}; if (socket_.connect(socket_path, error_code)) { std::exit(-1); } if (socket_.set_option(azmq::socket::subscribe(""), error_code)) { std::exit(-1); } std::array<std::byte, 1024> buffer{}; std::cout << "waiting for async_receive" << std::endl; co_await(azmq::async_receive(socket_, asio::buffer(buffer), asio::use_awaitable) || timer(std::chrono::seconds{1})); std::cout << "async_receive finished" << std::endl; } auto main() -> int { boost::asio::io_context ctx{}; boost::asio::co_spawn(ctx, foo(ctx), asio::detached); ctx.run(); return 0; }
问题分析与解决思路
Asio的awaitable_operators::||操作符在其中一个awaitable完成时,会尝试取消另一个可取消的异步操作,但这个机制生效的前提是异步操作实现了Asio的取消模型:
- 异步操作需要关联取消令牌(Cancellation Token),并监听取消请求;
- 底层操作需要支持主动取消,或者在执行前检查取消状态。
从Azmq的代码实现来看,问题出在:
azmq::async_receive的初始化器没有获取或关联完成处理程序的取消槽(cancellation slot),无法接收Asio的取消信号;- 底层的socket服务只是将接收操作入队,没有提供取消已入队操作的接口,也不会在操作执行前检查是否已被取消。
可行解决方案
修改Azmq库以支持取消:
在async_receive_initiation的operator()中,获取完成处理程序的关联取消槽,将其与socket的操作绑定。如果Azmq的socket服务支持取消已入队操作,就在取消触发时调用对应接口;否则,在操作执行前检查取消状态,若已取消则直接返回boost::asio::error::operation_aborted。示例修改方向:
template<typename CompletionHandler> void operator()(CompletionHandler &&completion_handler) { auto executor = boost::asio::get_associated_executor( completion_handler, socket.get_executor()); auto cancellation_slot = boost::asio::get_associated_cancellation_slot( completion_handler, socket.get_executor().context()); // 封装处理程序,添加取消检查或绑定取消回调 auto wrapped_handler = [completion_handler = std::forward<CompletionHandler>(completion_handler), cancellation_slot](boost::system::error_code ec, size_t size) mutable { if (cancellation_slot.is_cancelled()) { ec = boost::asio::error::operation_aborted; } std::move(completion_handler)(ec, size); }; socket.async_receive(buffers, boost::asio::bind_executor(executor, std::move(wrapped_handler))); }注意:这只是简化示例,实际需要Azmq底层支持取消已入队的read_op才能完全生效。
手动封装支持取消的接收操作:
若无法修改Azmq库,可以手动用asio::cancellation_signal实现取消逻辑:在定时器完成时发送取消信号,同时在接收操作的处理程序中检查该信号。不过这种方式需要手动管理信号,且无法真正取消Azmq底层的等待操作,只能提前终止处理流程。使用Azmq原生超时选项(仅同步场景适用):
Azmq的socket支持ZMQ_RCVTIMEO选项设置同步接收的超时,但这不适用于异步接收场景,无法直接解决当前问题。
内容的提问来源于stack exchange,提问作者magni_mar

