45 using writer_type = std::function<awaitable<void>(std::span<const std::byte>)>;
47 explicit write_queue(writer_type writer) : writer_{std::move(writer)} {}
53 std::shared_ptr<void> keepalive = {}) {
54 enqueue(data, std::move(keepalive), std::nullopt);
66 std::span<const std::byte> data,
67 std::shared_ptr<void> keepalive = {}) {
68 enqueue(data, std::move(keepalive), key);
76 [[nodiscard]]
bool failed() const noexcept {
77 return write_failure_ !=
nullptr;
85 void enqueue(std::span<const std::byte> data,
86 std::shared_ptr<void> keepalive,
87 std::optional<std::uint64_t> replace_key) {
89 std::rethrow_exception(write_failure_);
96 const auto first = std::find_if(queue_.begin(), queue_.end(), [&](
const item& queued) {
97 return queued.replace_key == replace_key;
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();
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_};
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();
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);
130 if (data.size() > max_queued_bytes_ - queued_bytes_) {
131 throw write_queue_full{queued_bytes_};
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();
145 drain_keepalive_ = std::move(owner_keepalive);
150 awaitable<void> drain() {
151 auto drain_keepalive = drain_keepalive_;
153 while (!queue_.empty()) {
154 auto frame = std::move(queue_.front());
156 queued_bytes_ -= frame.data->size();
157 co_await writer_(std::span<const std::byte>{frame.data->data(), frame.data->size()});
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());
167 }
catch (
const std::exception& error) {
170 write_failure_ = std::current_exception();
171 roboctrl::logger::instance().log_error(
172 "asynchronous write failed: {}", error.what());
176 write_failure_ = std::current_exception();
177 roboctrl::logger::instance().log_error(
178 "asynchronous write failed with an unknown exception");
181 drain_keepalive_.reset();
186 std::shared_ptr<std::vector<std::byte>> data;
187 std::shared_ptr<void> keepalive;
188 std::optional<std::uint64_t> replace_key;
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_;
197 static constexpr std::size_t max_queued_bytes_ = 64 * 1024;