无需定时器等待boost::asio::spawn协程结果,求DBus串行化优化方案
问题背景与需求
需要实现一个D-Bus方法,该方法需调用带有yield_context的协程函数,并将该函数的响应作为结果返回。当前存在以下限制:
- 底层系统不允许存在多个待处理消息
- 提供的协程函数库仅通过
yield_context发送消息并等待响应,未实现任何串行化机制,因此协程函数不能在前一次调用返回前再次被调用
现有实现代码
struct RequestResponse { boost::asio::steady_timer& timer; const UnderlyingSystemLibrary::ByteArray& request; std::chrono::milliseconds timeout; boost::system::error_code ec; UnderlyingSystemLibrary::ByteArray response; }; class ExecuteLock { public: ExecuteLock(bool& executing) : executing(executing) { if (!executing) executing = locked = true; else locked = false; } ~ExecuteLock() { if (locked) executing = false; } bool isLocked() { return locked; } private: bool& executing; bool locked; }; ...... /* wrapper function to prohibit multiple in-flight message */ void Wrapper::sendReceiveYield( boost::asio::yield_context yield, std::shared_ptr<RequestResponse> message) { static constexpr size_t limit = 100; ExecuteLock lock(isSending); /* Custom lock to check another coroutine. It'll be released on return */ messageQueue.push(message); /* Push message into message queue */ if (!lock.isLocked()) /* If another coroutine holds lock */ { /* A coroutine that holds lock will do the work */ return; } /* If there is no coroutine which has run already */ while (!messageQueue.empty()) /* Pop and process message until empty */ { auto& message = messageQueue.front(); std::tie(message->ec, message->response) = underlyingSystemLibrary->sendReceiveYield( yield, message->request, message->timeout); /* Call underlying system function with yield context */ message->timer.cancel(); /* Cancel timer after return from underlying system function */ messageQueue.pop(); } } std::pair<boost::system::error_code, UnderlyingSystemLibrary::ByteArray> Wrapper::sendReceiveYield( boost::asio::yield_context yield, const UnderlyingSystemLibrary::ByteArray& request, std::chrono::milliseconds timeout) { static const std::chrono::milliseconds async_timeout(1000); auto executor = boost::asio::get_associated_executor(yield); boost::asio::steady_timer timer(executor); auto message = std::make_shared<RequestResponse>( timer, request, timeout, boost::system::error_code(), UnderlyingSystemLibrary::ByteArray() ); boost::asio::spawn(strand, [this, message](boost::asio::yield_context yield) { sendReceiveYield(yield, message); /* call wrapper function within a strand */ }); auto ec = boost::system::error_code(); timer.expires_after(std::max(timeout, async_timeout)); timer.async_wait(yield[ec]); /* Wait by timer. The timer will be canceled by wrapper function. */ if (ec != boost::asio::error::operation_aborted) { phosphor::logging::log<phosphor::logging::level::ERR>( "transaction timeout"); /* It is error if the timer isn't canceled. */ return std::make_pair(ec, message->response); } /* Return response */ return std::make_pair(message->ec, message->response); }
现有方案说明
当前方案核心思路是将请求入队,由持有锁的协程依次处理队列中的请求,处理完成后取消定时器通知调用方。目前方案可正常运行,但并非最优解,恳请提供此类问题的示例或相关文档参考。
内容的提问来源于stack exchange,提问作者GT Lee
相关产品推荐
相关产品推荐

