GKD.RoboCtrl
RoboMaster Linux 电控:异步 IO、设备驱动与机器人控制
载入中...
搜索中...
未找到
write_queue.hpp
1#pragma once
2
3#include <asio.hpp>
4#include <asio/use_awaitable.hpp>
5
6#include <algorithm>
7#include <cstddef>
8#include <cstdint>
9#include <deque>
10#include <exception>
11#include <functional>
12#include <iterator>
13#include <memory>
14#include <optional>
15#include <span>
16#include <stdexcept>
17#include <string>
18#include <utility>
19#include <vector>
20
21#include "core/async.hpp"
22#include "core/logger.h"
23
24namespace roboctrl::io {
25
29class write_queue_full final : public std::runtime_error {
30public:
31 explicit write_queue_full(std::size_t queued_bytes)
32 : std::runtime_error{
33 "asynchronous write queue is full (" +
34 std::to_string(queued_bytes) + " bytes queued)"} {}
35};
36
44public:
45 using writer_type = std::function<awaitable<void>(std::span<const std::byte>)>;
46
47 explicit write_queue(writer_type writer) : writer_{std::move(writer)} {}
48
49 write_queue(const write_queue&) = delete;
50 write_queue& operator=(const write_queue&) = delete;
51
52 awaitable<void> send(std::span<const std::byte> data,
53 std::shared_ptr<void> keepalive = {}) {
54 enqueue(data, std::move(keepalive), std::nullopt);
55 return completed();
56 }
57
65 awaitable<void> send_latest(std::uint64_t key,
66 std::span<const std::byte> data,
67 std::shared_ptr<void> keepalive = {}) {
68 enqueue(data, std::move(keepalive), key);
69 return completed();
70 }
71
76 [[nodiscard]] bool failed() const noexcept {
77 return write_failure_ != nullptr;
78 }
79
80private:
81 static awaitable<void> completed() {
82 co_return;
83 }
84
85 void enqueue(std::span<const std::byte> data,
86 std::shared_ptr<void> keepalive,
87 std::optional<std::uint64_t> replace_key) {
88 if (write_failure_) {
89 std::rethrow_exception(write_failure_);
90 }
91 if (data.empty()) {
92 return;
93 }
94
95 if (replace_key) {
96 const auto first = std::find_if(queue_.begin(), queue_.end(), [&](const item& queued) {
97 return queued.replace_key == replace_key;
98 });
99 if (first != queue_.end()) {
100 std::size_t matching_bytes = 0;
101 for (const auto& queued : queue_) {
102 if (queued.replace_key == replace_key) {
103 matching_bytes += queued.data->size();
104 }
105 }
106 const auto retained_bytes = queued_bytes_ - matching_bytes;
107 if (data.size() > max_queued_bytes_ - retained_bytes) {
108 throw write_queue_full{queued_bytes_};
109 }
110
111 auto replacement =
112 std::make_shared<std::vector<std::byte>>(data.begin(), data.end());
113 queued_bytes_ -= first->data->size();
114 first->data = std::move(replacement);
115 first->keepalive = std::move(keepalive);
116 queued_bytes_ += first->data->size();
117
118 for (auto it = std::next(first); it != queue_.end();) {
119 if (it->replace_key == replace_key) {
120 queued_bytes_ -= it->data->size();
121 it = queue_.erase(it);
122 } else {
123 ++it;
124 }
125 }
126 return;
127 }
128 }
129
130 if (data.size() > max_queued_bytes_ - queued_bytes_) {
131 throw write_queue_full{queued_bytes_};
132 }
133
134 auto replacement = std::make_shared<std::vector<std::byte>>(data.begin(), data.end());
135 auto owner_keepalive = keepalive;
136 queue_.push_back(item{
137 std::move(replacement), std::move(keepalive), replace_key});
138 queued_bytes_ += data.size();
139 if (!writing_) {
140 writing_ = true;
141 // Keep the owner alive for the complete drain coroutine. A
142 // per-frame keepalive alone can be released immediately after
143 // the final write, before drain() performs its final member
144 // accesses.
145 drain_keepalive_ = std::move(owner_keepalive);
146 roboctrl::spawn(drain());
147 }
148 }
149
150 awaitable<void> drain() {
151 auto drain_keepalive = drain_keepalive_;
152 try {
153 while (!queue_.empty()) {
154 auto frame = std::move(queue_.front());
155 queue_.pop_front();
156 queued_bytes_ -= frame.data->size();
157 co_await writer_(std::span<const std::byte>{frame.data->data(), frame.data->size()});
158 }
159 } catch (const asio::system_error& error) {
160 if (error.code() != asio::error::operation_aborted) {
161 write_failure_ = std::current_exception();
162 logger::instance().log_error(
163 "asynchronous write failed: {}", error.what());
164 }
165 queue_.clear();
166 queued_bytes_ = 0;
167 } catch (const std::exception& error) {
168 queue_.clear();
169 queued_bytes_ = 0;
170 write_failure_ = std::current_exception();
171 roboctrl::logger::instance().log_error(
172 "asynchronous write failed: {}", error.what());
173 } catch (...) {
174 queue_.clear();
175 queued_bytes_ = 0;
176 write_failure_ = std::current_exception();
177 roboctrl::logger::instance().log_error(
178 "asynchronous write failed with an unknown exception");
179 }
180 writing_ = false;
181 drain_keepalive_.reset();
182 }
183
184 writer_type writer_;
185 struct item {
186 std::shared_ptr<std::vector<std::byte>> data;
187 std::shared_ptr<void> keepalive;
188 std::optional<std::uint64_t> replace_key;
189 };
190
191 std::deque<item> queue_;
192 std::size_t queued_bytes_{0};
193 bool writing_{false};
194 std::shared_ptr<void> drain_keepalive_;
195 std::exception_ptr write_failure_;
196
197 static constexpr std::size_t max_queued_bytes_ = 64 * 1024;
198};
199
200} // namespace roboctrl::io
异步任务上下文组件。
写队列无法接纳新的独立报文。
将同一异步流上的写操作串行化。
bool failed() const noexcept
writer 是否已经失败。
awaitable< void > send_latest(std::uint64_t key, std::span< const std::byte > data, std::shared_ptr< void > keepalive={})
按 key 入队,并用新数据替换同 key 的尚未发送项。
用于日志输出的组件。
auto spawn(task_context::task_type &&task)
添加一个协程任务到全局任务上下文中执行。
Definition async.hpp:191
asio::awaitable< T > awaitable
协程任务类型。
Definition async.hpp:46