co_usb
Loading...
Searching...
No Matches
complete_sequence_awaitable.hpp
Go to the documentation of this file.
1
6#pragma once
7
12#include "co_usb/usb_error.hpp"
13#include <boost/capy/buffers.hpp>
14#include <boost/capy/concept/io_awaitable.hpp>
15#include <boost/capy/continuation.hpp>
16#include <boost/capy/ex/io_env.hpp>
17#include <boost/capy/io_result.hpp>
18#include <cassert>
19#include <coroutine>
20#include <libusb.h>
21#include <memory_resource>
22#include <mutex>
23#include <unordered_map>
24#include <utility>
25
27{
28
29template <detail::TransferSequence TSeq, AnyBufferSequence BuffersTy>
31{
33 using buf_iter_t = decltype(boost::capy::begin(std::declval<BuffersTy const &>()));
34 using buffer_t = boost::capy::buffer_type<BuffersTy>;
35
37 {
38 uint8_t *ptr{nullptr};
39 size_t expected{0};
40 size_t total{0};
41 };
42
44 {
45 explicit await_state_t (std::pmr::memory_resource *memres) : states(memres)
46 {
47 }
48
49 std::mutex mutex;
50 std::pmr::unordered_map<libusb_transfer *, transfer_progress_t> states;
51
52 boost::capy::io_env const *io_env;
53 boost::capy::continuation cont;
54
55 std::error_code ec;
56 size_t err_idx{0};
57
58 size_t in_flight{0};
59 };
60
61 static void transfer_callback (libusb_transfer *tfer)
62 {
64 self_t &self = *static_cast<self_t *>(tfer->user_data);
65 const auto resume_on_zero = [&] ()
66 {
67 if (self.await_state->in_flight == 0)
68 {
69 self.await_state->io_env->executor.post(self.await_state->cont);
70 }
71 };
72 self.await_state->in_flight--;
73 bool ok = false;
74 {
75 std::unique_lock lock{self.await_state->mutex};
76 ok = (bool)self.await_state->ec;
77 }
78 if (ok)
79 {
80 resume_on_zero();
81 return;
82 }
83 size_t subtotal = 0;
84 transfer_progress_t &prg = self.await_state->states.at(tfer);
85 if (tfer->type == LIBUSB_TRANSFER_TYPE_ISOCHRONOUS)
86 {
87 for (int i = 0; i < tfer->num_iso_packets; i++)
88 {
89 subtotal += tfer->iso_packet_desc[i].actual_length;
90 }
91 prg.total += subtotal;
92 for (int i = 0; i < tfer->num_iso_packets; i++)
93 {
94 if (tfer->status != LIBUSB_TRANSFER_COMPLETED) [[unlikely]]
95 {
96 {
97 std::unique_lock lock{self.await_state->mutex};
98 self.await_state->ec = make_transfer_status(
99 static_cast<status>(tfer->iso_packet_desc[i].status));
100 }
101 resume_on_zero();
102 return;
103 }
104 }
105 }
106 else
107 {
108 subtotal = tfer->actual_length;
109 prg.total += subtotal;
110 }
111 if (tfer->status != LIBUSB_TRANSFER_COMPLETED) [[unlikely]]
112 {
113 {
114 std::unique_lock lock{self.await_state->mutex};
115 self.await_state->ec = make_transfer_status(static_cast<status>(tfer->status));
116 }
117 resume_on_zero();
118 return;
119 }
120 if (prg.total >= prg.expected)
121 {
122 if (self.buf_current == self.buf_end) [[unlikely]]
123 {
124 resume_on_zero();
125 return;
126 }
127 buffer_t buf = *self.buf_current;
128 self.buf_current++;
129 prg.ptr = (uint8_t *)buf.data();
130 prg.expected = buf.size();
131 tfer->buffer = prg.ptr;
132 tfer->length = std::min(prg.expected, self.single_transfer_limit);
133 }
134 else
135 {
136 prg.ptr += subtotal;
137 tfer->buffer = prg.ptr;
138 tfer->length = std::min(prg.expected - prg.total, self.single_transfer_limit);
139 }
140 int r = libusb_submit_transfer(tfer);
141 if (r != LIBUSB_SUCCESS) [[unlikely]]
142 {
143 {
144 std::unique_lock lock{self.await_state->mutex};
145 self.await_state->ec = make_usb_error_code(static_cast<usb_error>(r));
146 }
147 resume_on_zero();
148 return;
149 }
150 self.await_state->in_flight++;
151 }
152
154 sequence_view<TSeq> *seq_view,
155 BuffersTy const &buffers, size_t submission_size,
157 : await_state(await_state), view(seq_view), buf_current(boost::capy::begin(buffers)),
158 buf_end(boost::capy::end(buffers)), submission_size(submission_size),
160 {
161 assert(submission_size > 0);
162 assert(single_transfer_limit > 0);
163 assert((submission_size <= seq_view->size()) &&
164 "submission size must not be greater than number of transfers");
165 assert((submission_size <= boost::capy::buffer_length(buffers)) &&
166 "submission size must not be greater than number of buffers");
167 }
168
169 inline bool await_ready ()
170 {
171 return buf_current >= buf_end;
172 }
173
174 inline std::coroutine_handle<> await_suspend (std::coroutine_handle<> h,
175 boost::capy::io_env const *io_env)
176 {
177 if (io_env->stop_token.stop_requested())
178 {
180 return h;
181 }
182 await_state->io_env = io_env;
183 await_state->cont = {h};
185 const tfer_iter_t end = v.end() + submission_size;
186 for (tfer_iter_t iter = v.begin(); iter != end; iter++)
187 {
188 libusb_transfer *tfer = transfer_of(*iter);
190 buffer_t buf = *buf_current;
191 buf_current++;
192 tfer->buffer = (uint8_t *)buf.data();
193 tfer->length = std::min(buf.size(), single_transfer_limit);
194 tfer->user_data = this;
195 tfer->callback = transfer_callback;
196 prg.ptr = (uint8_t *)buf.data();
197 prg.expected = buf.size();
198 }
200 for (tfer_iter_t iter = v.begin(); iter != end; iter++)
201 {
202 libusb_transfer *tfer = transfer_of(*iter);
203 int r = libusb_submit_transfer(tfer);
204 if (r != LIBUSB_SUCCESS) [[unlikely]]
205 {
206 std::unique_lock lock{await_state->mutex};
207 await_state->ec = make_transfer_status(static_cast<status>(r));
208 lock.unlock();
209 break;
210 }
211 }
212 return std::noop_coroutine();
213 }
214
215 inline boost::capy::io_result<size_t> await_resume ()
216 {
217 size_t total{0};
218 for (auto const &[_, state] : await_state->states)
219 {
220 total += state.total;
221 }
222 return {await_state->ec, total};
223 }
224
227
230
233};
234
235static_assert(boost::capy::IoAwaitable<
236 complete_sequence_awaitable<libusb_transfer *, boost::capy::mutable_buffer>>,
237 "Not a proper IoAwaitable");
238static_assert(boost::capy::IoAwaitable<complete_sequence_awaitable<std::vector<libusb_transfer *>,
239 boost::capy::const_buffer>>,
240 "Not a proper IoAwaitable");
241
242} // namespace co_usb::transfer::detail
Unified concept for buffer sequences.
constexpr auto transfer_of(Ty const &tfer_res) -> libusb_transfer *
Uniform accessor for obtaining a libusb_transfer * from a co_usb::transfer::detail::TransferResource.
Definition transfer_sequence.hpp:41
status
Status codes for transfers.
Definition status.hpp:22
Definition complete_io.hpp:21
std::error_code make_transfer_status(status e) noexcept
Definition status.hpp:74
usb_error
USB error enumeration.
Definition usb_error.hpp:20
std::error_code make_usb_error_code(usb_error e) noexcept
Definition usb_error.hpp:85
View of a certain transfer sequence. Provides a uniform range interface.
Status codes for transfers.
boost::capy::continuation cont
Definition complete_sequence_awaitable.hpp:53
size_t in_flight
Definition complete_sequence_awaitable.hpp:58
std::error_code ec
Definition complete_sequence_awaitable.hpp:55
size_t err_idx
Definition complete_sequence_awaitable.hpp:56
await_state_t(std::pmr::memory_resource *memres)
Definition complete_sequence_awaitable.hpp:45
boost::capy::io_env const * io_env
Definition complete_sequence_awaitable.hpp:52
std::pmr::unordered_map< libusb_transfer *, transfer_progress_t > states
Definition complete_sequence_awaitable.hpp:50
std::mutex mutex
Definition complete_sequence_awaitable.hpp:49
uint8_t * ptr
Definition complete_sequence_awaitable.hpp:38
size_t total
Definition complete_sequence_awaitable.hpp:40
size_t expected
Definition complete_sequence_awaitable.hpp:39
Definition complete_sequence_awaitable.hpp:31
size_t submission_size
Definition complete_sequence_awaitable.hpp:231
sequence_view< TSeq > * view
Definition complete_sequence_awaitable.hpp:226
std::coroutine_handle await_suspend(std::coroutine_handle<> h, boost::capy::io_env const *io_env)
Definition complete_sequence_awaitable.hpp:174
complete_sequence_awaitable(await_state_t *await_state, sequence_view< TSeq > *seq_view, BuffersTy const &buffers, size_t submission_size, size_t single_transfer_limit)
Definition complete_sequence_awaitable.hpp:153
await_state_t * await_state
Definition complete_sequence_awaitable.hpp:225
boost::capy::buffer_type< BuffersTy > buffer_t
Definition complete_sequence_awaitable.hpp:34
decltype(boost::capy::begin(std::declval< BuffersTy const & >())) buf_iter_t
Definition complete_sequence_awaitable.hpp:33
buf_iter_t buf_current
Definition complete_sequence_awaitable.hpp:228
static void transfer_callback(libusb_transfer *tfer)
Definition complete_sequence_awaitable.hpp:61
bool await_ready()
Definition complete_sequence_awaitable.hpp:169
size_t single_transfer_limit
Definition complete_sequence_awaitable.hpp:232
sequence_view< TSeq >::iterator_type tfer_iter_t
Definition complete_sequence_awaitable.hpp:32
buf_iter_t buf_end
Definition complete_sequence_awaitable.hpp:229
boost::capy::io_result< size_t > await_resume()
Definition complete_sequence_awaitable.hpp:215
View of a certain transfer sequence. Provides a uniform range interface.
Definition sequence_view.hpp:28
decltype(detail::transfer_begin(std::declval< Seq const & >())) iterator_type
Definition sequence_view.hpp:29
auto begin() const noexcept
Definition sequence_view.hpp:48
auto end() const noexcept
Definition sequence_view.hpp:53
Transfer resource and transfer sequence concepts.
USB error enumeration.