Condy v1.9
C++ Asynchronous System Call Layer for Linux
Loading...
Searching...
No Matches
zcrx.hpp
Go to the documentation of this file.
1
10
11#pragma once
12
15#include "condy/detail/ring.hpp"
17#include "condy/runtime.hpp"
18#include <bit>
19#include <sys/mman.h>
20
21namespace condy {
22
23#if CONDY_URING_VERSION_GE(2, 15) // >= 2.15
24
25class ZeroCopyRxBufferPool;
26
35class ZeroCopyRxBuffer : public detail::ManagedBuffer<ZeroCopyRxBufferPool> {
36public:
37 using Base = detail::ManagedBuffer<ZeroCopyRxBufferPool>;
38 using Base::Base;
39};
40
44struct ZeroCopyRxArea {
47 void *addr = nullptr;
49 size_t size;
50};
51
55struct ZeroCopyRxDMABufArea {
57 int dmabuf_fd;
59 size_t offset;
61 size_t size;
62};
63
77class ZeroCopyRxBufferPool {
78public:
86 ZeroCopyRxBufferPool(uint32_t if_idx, uint32_t if_rxq, uint32_t rq_entries,
87 const ZeroCopyRxArea &area)
88 : ZeroCopyRxBufferPool(*detail::Context::current().runtime(), if_idx,
89 if_rxq, rq_entries, area) {}
90
99 ZeroCopyRxBufferPool(Runtime &runtime, uint32_t if_idx, uint32_t if_rxq,
100 uint32_t rq_entries, const ZeroCopyRxArea &area)
101 : ZeroCopyRxBufferPool(runtime.ring_internal(), if_idx, if_rxq,
102 rq_entries, area, 0) {}
103
104 // Device-less constructor, DO NOT use this in production code if you don't
105 // know what you are doing.
106 ZeroCopyRxBufferPool(Runtime &runtime, uint32_t rq_entries,
107 const ZeroCopyRxArea &area)
108 : ZeroCopyRxBufferPool(runtime.ring_internal(), 0, 0, rq_entries, area,
109 ZCRX_REG_NODEV) {}
110
118 ZeroCopyRxBufferPool(uint32_t if_idx, uint32_t if_rxq, uint32_t rq_entries,
119 const ZeroCopyRxDMABufArea &area)
120 : ZeroCopyRxBufferPool(*detail::Context::current().runtime(), if_idx,
121 if_rxq, rq_entries, area) {}
122
131 ZeroCopyRxBufferPool(Runtime &runtime, uint32_t if_idx, uint32_t if_rxq,
132 uint32_t rq_entries, const ZeroCopyRxDMABufArea &area)
133 : ring_(&runtime.ring_internal()), flags_(0) {
134 bool ok = false;
135 auto d = detail::defer([&]() {
136 if (!ok) {
137 cleanup_();
138 }
139 });
140
141 area_size_ = 0;
142 area_ptr_ = nullptr;
143
144 io_uring_zcrx_area_reg area_reg = {};
145 area_reg.addr = area.offset;
146 area_reg.len = area.size;
147 area_reg.flags = IORING_ZCRX_AREA_DMABUF;
148
149 register_ifq_(if_idx, if_rxq, rq_entries, area_reg,
150 sysconf(_SC_PAGESIZE));
151 ok = true;
152 }
153
154 ~ZeroCopyRxBufferPool() { cleanup_(); }
155
156 CONDY_DELETE_COPY_MOVE(ZeroCopyRxBufferPool);
157
158private:
159 ZeroCopyRxBufferPool(detail::Ring &ring, uint32_t if_idx, uint32_t if_rxq,
160 uint32_t rq_entries, const ZeroCopyRxArea &area,
161 uint32_t flags)
162 : ring_(&ring), flags_(flags) {
163 bool ok = false;
164 auto d = detail::defer([&]() {
165 if (!ok) {
166 cleanup_();
167 }
168 });
169
170 const size_t page_size = sysconf(_SC_PAGESIZE);
171
172 if (area.addr == nullptr) {
173 area_size_ = detail::align_up(area.size, page_size);
174 area_ptr_ = mmap(nullptr, area_size_, PROT_READ | PROT_WRITE,
175 MAP_ANONYMOUS | MAP_PRIVATE, 0, 0);
176 if (area_ptr_ == MAP_FAILED) {
177 throw detail::make_system_error("mmap");
178 }
179
180 io_uring_zcrx_area_reg area_reg = {};
181 area_reg.addr = reinterpret_cast<uint64_t>(area_ptr_);
182 area_reg.len = area_size_;
183 area_reg.flags = 0;
184
185 register_ifq_(if_idx, if_rxq, rq_entries, area_reg, page_size);
186 } else {
187 // Not owned, so we don't track the size for unmapping
188 area_size_ = 0;
189 area_ptr_ = area.addr;
190
191 io_uring_zcrx_area_reg area_reg = {};
192 area_reg.addr = reinterpret_cast<uint64_t>(area_ptr_);
193 area_reg.len = area.size;
194 area_reg.flags = 0;
195
196 register_ifq_(if_idx, if_rxq, rq_entries, area_reg, page_size);
197 }
198
199 ok = true;
200 }
201
202public:
203 uint32_t zcrx_id() const noexcept { return zcrx_id_; }
204
205 ZeroCopyRxBuffer handle_finish(io_uring_cqe *cqe) noexcept {
206 assert(ring_->check_cqe32(cqe) && "Expected big CQE for ZeroCopyRx");
207
208 if (cqe->res < 0) {
209 return ZeroCopyRxBuffer();
210 }
211 io_uring_zcrx_cqe *rcqe =
212 reinterpret_cast<io_uring_zcrx_cqe *>(cqe->big_cqe);
213 void *data = static_cast<char *>(area_ptr_) +
214 (rcqe->off & ~IORING_ZCRX_AREA_MASK);
215 size_t size = static_cast<size_t>(cqe->res);
216 return ZeroCopyRxBuffer(data, size, this);
217 }
218
219 void add_buffer_back(void *ptr, size_t size) noexcept {
220 rq_enqueue_(ptr, size);
221 maybe_flush_rq_();
222 }
223
224private:
225 void cleanup_() noexcept {
226 [[maybe_unused]] int r;
227 if (area_ptr_ != nullptr && area_size_ > 0) {
228 r = munmap(area_ptr_, area_size_);
229 assert(r == 0);
230 }
231 if (rq_ring_.ring_ptr != nullptr) {
232 r = munmap(rq_ring_.ring_ptr, ring_size_);
233 assert(r == 0);
234 }
235 // TODO: Unregister ifq. For now there's no way to unregister ifq, so we
236 // just leak the registration.
237 }
238
239 void register_ifq_(uint32_t if_idx, uint32_t if_rxq, uint32_t rq_entries,
240 io_uring_zcrx_area_reg &area_reg, size_t page_size) {
241 rq_entries = std::bit_ceil(rq_entries);
242 io_uring_region_desc region_reg = {};
243 ring_size_ = get_refill_ring_size_(rq_entries, page_size);
244 region_reg.user_addr = 0;
245 region_reg.size = ring_size_;
246 region_reg.flags = 0;
247
248 io_uring_zcrx_ifq_reg reg = {};
249 reg.if_idx = if_idx;
250 reg.if_rxq = if_rxq;
251 reg.rq_entries = rq_entries;
252 reg.area_ptr = reinterpret_cast<uint64_t>(&area_reg);
253 reg.region_ptr = reinterpret_cast<uint64_t>(&region_reg);
254 reg.flags = flags_;
255
256 int r = io_uring_register_ifq(ring_->ring(), &reg);
257 if (r != 0) {
258 throw detail::make_system_error("io_uring_register_ifq", -r);
259 }
260 // TODO: Unregister ifq if any exception. For now there's no way to
261 // unregister ifq.
262
263 void *ring_ptr = mmap(nullptr, ring_size_, PROT_READ | PROT_WRITE,
264 MAP_SHARED | MAP_POPULATE, ring_->ring()->ring_fd,
265 static_cast<off_t>(region_reg.mmap_offset));
266 if (ring_ptr == MAP_FAILED) {
267 throw detail::make_system_error("mmap");
268 }
269 rq_ring_.khead = (unsigned int *)((char *)ring_ptr + reg.offsets.head);
270 rq_ring_.ktail = (unsigned int *)((char *)ring_ptr + reg.offsets.tail);
271 rq_ring_.rqes =
272 (struct io_uring_zcrx_rqe *)((char *)ring_ptr + reg.offsets.rqes);
273 rq_ring_.rq_tail = 0;
274 rq_ring_.ring_entries = reg.rq_entries;
275 rq_ring_.ring_ptr = ring_ptr;
276
277 zcrx_id_ = reg.zcrx_id;
278 area_token_ = area_reg.rq_area_token;
279 }
280
281 static size_t get_refill_ring_size_(uint32_t rq_entries,
282 size_t page_size) noexcept {
283 size_t ring_size = rq_entries * sizeof(io_uring_zcrx_rqe);
284 ring_size += page_size;
285 ring_size = detail::align_up(ring_size, page_size);
286 return ring_size;
287 }
288
289 size_t rq_nr_queued_() const noexcept {
290 return rq_ring_.rq_tail - io_uring_smp_load_acquire(rq_ring_.khead);
291 }
292
293 void rq_enqueue_(void *ptr, size_t size) noexcept {
294 assert(rq_nr_queued_() < rq_ring_.ring_entries);
295 io_uring_zcrx_rqe *rqe;
296 unsigned rq_mask = rq_ring_.ring_entries - 1;
297 rqe = &rq_ring_.rqes[rq_ring_.rq_tail & rq_mask];
298 rqe->off = (static_cast<char *>(ptr) - static_cast<char *>(area_ptr_)) |
299 area_token_;
300 rqe->len = static_cast<uint32_t>(size);
301 io_uring_smp_store_release(rq_ring_.ktail, ++rq_ring_.rq_tail);
302 }
303
304 void flush_rq_() noexcept {
305 zcrx_ctrl ctrl = {};
306 ctrl.zcrx_id = zcrx_id_;
307 ctrl.op = ZCRX_CTRL_FLUSH_RQ;
308 [[maybe_unused]] int r =
309 io_uring_register_zcrx_ctrl(ring_->ring(), &ctrl);
310 assert(r == 0);
311 }
312
313 void maybe_flush_rq_() noexcept {
314 if (rq_nr_queued_() >= rq_ring_.ring_entries ||
315 (flags_ & ZCRX_REG_NODEV)) {
316 flush_rq_();
317 }
318 }
319
320private:
321 detail::Ring *ring_;
322 size_t area_size_ = 0;
323 void *area_ptr_ = nullptr;
324 size_t ring_size_ = 0;
325 io_uring_zcrx_rq rq_ring_ = {};
326 uint32_t zcrx_id_;
327 uint64_t area_token_;
328 uint32_t flags_;
329};
330
331#endif
332
333} // namespace condy
Basic buffer types and conversion utilities.
The main namespace for the Condy library.
Definition condy.hpp:37
Wrapper classes for liburing interfaces.
Runtime type for running the io_uring event loop.
Internal utility classes and functions used by Condy.