自定义boost::asio异步操作实现及相关技术问题咨询
我正尝试实现自己的boost::asio异步操作,使其在io_context.run()期间异步执行。同时需要实现一个可取消的操作,该操作会等待某个条件/谓词满足(或boost::signals2::signal被触发),之后在io_context线程中调用完成处理器。想知道Boost中是否已提供此类功能?
我参考Boost文档写了以下代码:
#include <boost/asio.hpp> #include <iostream> /// From boost documentation: /// — If the initiating function is not a member function, /// the associated executor is that returned by the get_executor member function /// of the first argument to the initiating function. struct executor_owner { using executor_type = boost::asio::io_context; executor_type* m_executor = nullptr; explicit executor_owner(executor_type& ex) : m_executor(&ex) { } executor_type& get_executor() BOOST_ASIO_NOEXCEPT { std::cout << "get_executor() called." << std::endl; return *m_executor; } }; /// My initiating function template<class CompletionToken> auto async_xyz(executor_owner executorOwner, CompletionToken&& token) { using completion_handler_t = typename boost::asio::async_completion<CompletionToken, void(executor_owner&)>::completion_handler_type; return boost::asio::async_initiate<CompletionToken, void(executor_owner&)>( [](completion_handler_t completion_handler, executor_owner& owner) { // initiate the operation and cause completion_handler to be invoked // with the result std::cout << "I'm performing the async func operation here!" << std::endl; completion_handler(); }, token, executorOwner); } int main() { boost::asio::io_context io_context; executor_owner executorOwner(io_context); async_xyz(executorOwner, []() { std::cout << "completion handler here!" << std::endl; }); // I expect that completion handler will be called asynchronously, during io_context::run(). // Unfortunately, completion handler has already been called. // std::cout << "io_context begin." << std::endl; // io_context.run(); // It doesn't matter. completion handler // std::cout << "io_context done." << std::endl; return 0; }
输出结果:
I'm performing the async func operation here! completion handler here!
我的问题:
- 如何实现符合
boost::asio设计规范的自定义异步操作? - 为什么我的executor没有按照Boost文档所述通过
get_executor()方法获取? boost::asio::async_initiate内部是否会将操作提交到io_context,还是需要我在自定义操作实现中手动调用io_context::post(...)来触发完成处理器?- 如何实现一个简单的等待操作,用于等待谓词满足或信号触发?
- 如何为此类操作添加超时/取消功能?
使用环境:Boost 1.82、C14(优先考虑C11解决方案)
1. 符合Boost.Asio规范的自定义异步操作实现方式
需遵循组合异步操作的设计模式,核心要点:
- 利用
boost::asio::async_completion和async_initiate处理CompletionToken(支持回调、future等多种形式) - 确保操作异步性:必须将完成处理器提交到关联的executor(如io_context)执行,禁止直接同步调用
- 遵循Executor模型:通过
get_associated_executor获取与完成处理器关联的executor,或从操作第一个参数(如你的executor_owner)获取 - 保证异常安全:操作过程中若发生异常,需正确传递给完成处理器,避免资源泄漏
2. get_executor()未被调用的原因
你的代码中async_initiate的lambda直接同步调用了completion_handler,且未触发Executor关联逻辑:
async_initiate本身不会自动获取executor,需要手动在lambda中通过boost::asio::get_associated_executor获取,或显式调用executor_owner的get_executor()- 你没有利用
executor_owner将任务提交到io_context,因此get_executor()没有触发场景
修正示例:
auto ex = boost::asio::get_associated_executor(completion_handler, owner.get_executor()); ex.post([completion_handler]() mutable { completion_handler(); });
3. async_initiate是否自动提交到io_context?
不会。async_initiate仅负责处理CompletionToken的类型擦除,将用户提供的token转换为统一的完成处理器类型,不会自动将操作或完成处理器提交到executor。必须手动调用executor的post()/dispatch()/defer()方法,将完成处理器的执行提交到io_context线程,才能实现异步效果。
你的代码直接同步调用completion_handler(),因此会在调用async_xyz的线程立即执行,而非等待io_context.run()。
4. 实现等待谓词/信号触发的异步操作
方式1:等待谓词满足
使用定时器定期检查谓词,条件满足后触发完成处理器:
template <typename Predicate, typename CompletionToken> auto async_wait_predicate(boost::asio::io_context& io, Predicate pred, CompletionToken&& token) { using completion_handler_t = typename boost::asio::async_completion<CompletionToken, void()>::completion_handler_type; return boost::asio::async_initiate<CompletionToken, void()>( [&io, pred](completion_handler_t handler) mutable { auto check_predicate = [&io, pred, handler](const boost::system::error_code& ec) mutable { if (!ec && pred()) { boost::asio::post(io, std::move(handler)); return; } // 10ms后再次检查 auto timer = std::make_shared<boost::asio::steady_timer>(io); timer->expires_after(std::chrono::milliseconds(10)); timer->async_wait(std::move(check_predicate)); }; check_predicate({}); }, token ); }
方式2:等待boost::signals2::signal触发
将完成处理器绑定到信号,触发时提交到io_context执行:
template <typename CompletionToken> auto async_wait_signal(boost::asio::io_context& io, boost::signals2::signal<void()>& sig, CompletionToken&& token) { using completion_handler_t = typename boost::asio::async_completion<CompletionToken, void()>::completion_handler_type; return boost::asio::async_initiate<CompletionToken, void()>( [&io, &sig](completion_handler_t handler) { sig.connect([&io, handler]() mutable { boost::asio::post(io, std::move(handler)); }); }, token ); }
5. 添加超时/取消功能
超时功能
在等待操作中添加定时器,超时后触发取消逻辑:
template <typename Predicate, typename CompletionToken> auto async_wait_predicate_with_timeout(boost::asio::io_context& io, Predicate pred, std::chrono::milliseconds timeout, CompletionToken&& token) { using completion_handler_t = typename boost::asio::async_completion<CompletionToken, void(boost::system::error_code)>::completion_handler_type; return boost::asio::async_initiate<CompletionToken, void(boost::system::error_code)>( [&io, pred, timeout](completion_handler_t handler) mutable { auto timer = std::make_shared<boost::asio::steady_timer>(io); auto timeout_timer = std::make_shared<boost::asio::steady_timer>(io, timeout); auto check_predicate = [&io, pred, handler, timer, timeout_timer](const boost::system::error_code& ec) mutable { if (ec == boost::asio::error::operation_aborted) { handler(boost::asio::error::operation_aborted); return; } if (pred()) { timeout_timer->cancel(); boost::asio::post(io, [handler]() { handler({}); }); return; } timer->expires_after(std::chrono::milliseconds(10)); timer->async_wait(std::move(check_predicate)); }; timeout_timer->async_wait([handler, timer](const boost::system::error_code& ec) { if (!ec) { timer->cancel(); handler(boost::asio::error::timed_out); } }); check_predicate({}); }, token ); }
取消功能
使用boost::asio::cancellation_signal,触发信号时终止等待操作:
template <typename Predicate, typename CompletionToken> auto async_wait_predicate_cancellable(boost::asio::io_context& io, Predicate pred, boost::asio::cancellation_signal& cancel_sig, CompletionToken&& token) { using completion_handler_t = typename boost::asio::async_completion<CompletionToken, void(boost::system::error_code)>::completion_handler_type; return boost::asio::async_initiate<CompletionToken, void(boost::system::error_code)>( [&io, pred, &cancel_sig](completion_handler_t handler) mutable { auto timer = std::make_shared<boost::asio::steady_timer>(io); auto slot = cancel_sig.connect([timer]() { timer->cancel(); }); auto check_predicate = [&io, pred, handler, timer, slot](const boost::system::error_code& ec) mutable { if (ec == boost::asio::error::operation_aborted) { handler(boost::asio::error::operation_aborted); return; } if (pred()) { boost::asio::post(io, [handler]() { handler({}); }); return; } timer->expires_after(std::chrono::milliseconds(10)); timer->async_wait(std::move(check_predicate)); }; check_predicate({}); }, token ); }
内容的提问来源于stack exchange,提问作者Michał Jaroń

