Condy v1.9
C++ Asynchronous System Call Layer for Linux
Loading...
Searching...
No Matches
ublk-loop.cpp
Go to the documentation of this file.
1
10
11#include <condy.hpp>
12#include <csignal>
13#include <cstdint>
14#include <cstdio>
15#include <cstdlib>
16#include <cstring>
17#include <fcntl.h>
18#include <format>
19#include <iostream>
20#include <linux/ublk_cmd.h>
21#include <sys/mman.h>
22#include <sys/signalfd.h>
23#include <unistd.h>
24#include <vector>
25
26size_t num_queues = 1;
27size_t queue_depth = 32;
28uint32_t dev_id = -1;
29
30constexpr auto *CTRL_FILE = "/dev/ublk-control";
31constexpr size_t SECTOR_SIZE = 512;
32constexpr size_t MAX_SECTORS = 128;
33constexpr size_t BUF_SIZE = MAX_SECTORS * SECTOR_SIZE;
34constexpr size_t BS_SHIFT = 9;
35constexpr int FIXED_FD = 0;
36
37std::vector<std::unique_ptr<condy::Runtime>> queue_runtimes;
38size_t dev_sectors;
39int backing_fd;
40int signal_fd;
41int ctrl_fd;
42int ublkc_fd;
43void *bufs_base;
44
45bool multi_queues() { return num_queues > 1; }
46
47off_t io_desc_offset(size_t q_id) {
48 return UBLKSRV_CMD_BUF_OFFSET +
49 q_id * UBLK_MAX_QUEUE_DEPTH * sizeof(ublksrv_io_desc);
50}
51
52auto ublk_ctrl_cmd(int fd, int cmd_op, const ublksrv_ctrl_cmd &ctrl) {
54 cmd_op, fd,
55 [&](io_uring_sqe *sqe) { std::memcpy(sqe->cmd, &ctrl, sizeof(ctrl)); },
57}
58
59auto ublk_io_cmd(int cmd_op, const ublksrv_io_cmd &cmd) {
61 cmd_op, condy::fixed(FIXED_FD),
62 [&](io_uring_sqe *sqe) { std::memcpy(sqe->cmd, &cmd, sizeof(cmd)); },
64}
65
66condy::Coro<void> setup_device() {
67 int r;
68 ublksrv_ctrl_cmd ctrl;
69
70 ctrl_fd = co_await condy::async_open(CTRL_FILE, O_RDWR, 0);
71 if (ctrl_fd < 0) {
72 std::cerr << std::format("Failed to open {}: {}\n", CTRL_FILE, ctrl_fd);
73 std::exit(1);
74 }
75
76 ublksrv_ctrl_dev_info info = {};
77 info.dev_id = dev_id;
78 info.nr_hw_queues = num_queues;
79 info.queue_depth = queue_depth;
80 info.max_io_buf_bytes = BUF_SIZE;
81 ctrl = {};
82 ctrl.dev_id = info.dev_id;
83 ctrl.queue_id = -1;
84 ctrl.addr = reinterpret_cast<uint64_t>(&info);
85 ctrl.len = sizeof(info);
86 r = co_await ublk_ctrl_cmd(ctrl_fd, UBLK_U_CMD_ADD_DEV, ctrl);
87 if (r < 0) {
88 std::cerr << std::format("ADD_DEV failed: {}\n", r);
89 std::exit(1);
90 }
91 dev_id = info.dev_id; // Get actual dev_id
92
93 std::string ublkc_path = std::format("/dev/ublkc{}", dev_id);
94 ublkc_fd = co_await condy::async_open(ublkc_path.c_str(), O_RDWR, 0);
95 if (ublkc_fd < 0) {
96 std::cerr << std::format("Failed to open {}: {}\n", ublkc_path,
97 ublkc_fd);
98 std::exit(1);
99 }
100 std::cout << std::format("ublk-loop: /dev/ublkc{} created\n", dev_id);
101
102 ublk_params params = {};
103 params.len = sizeof(params);
104 params.types = UBLK_PARAM_TYPE_BASIC;
105 params.basic.attrs = UBLK_ATTR_VOLATILE_CACHE;
106 params.basic.logical_bs_shift = BS_SHIFT;
107 params.basic.physical_bs_shift = BS_SHIFT;
108 params.basic.io_opt_shift = BS_SHIFT;
109 params.basic.io_min_shift = BS_SHIFT;
110 params.basic.max_sectors = MAX_SECTORS;
111 params.basic.dev_sectors = dev_sectors;
112 ctrl = {};
113 ctrl.dev_id = dev_id;
114 ctrl.queue_id = -1;
115 ctrl.addr = reinterpret_cast<uint64_t>(&params);
116 ctrl.len = sizeof(params);
117 r = co_await ublk_ctrl_cmd(ctrl_fd, UBLK_U_CMD_SET_PARAMS, ctrl);
118 if (r < 0) {
119 std::cerr << std::format("SET_PARAMS failed: {}\n", r);
120 std::exit(1);
121 }
122
123 std::cout << std::format("ublk-loop: /dev/ublkc{} configured\n", dev_id);
124}
125
126condy::Coro<void> io_loop(size_t q_id, size_t tag,
127 ublksrv_io_desc *io_desc_base) {
128 int r;
129 ublksrv_io_cmd cmd;
130
131 void *buf = static_cast<char *>(bufs_base) +
132 q_id * (queue_depth * BUF_SIZE) + tag * BUF_SIZE;
133 ublksrv_io_desc *io_desc = io_desc_base + tag;
134
135 cmd = {};
136 cmd.q_id = q_id;
137 cmd.tag = tag;
138 cmd.addr = reinterpret_cast<uint64_t>(buf);
139 r = co_await ublk_io_cmd(UBLK_U_IO_FETCH_REQ, cmd);
140 if (r == UBLK_IO_RES_ABORT) {
141 co_return;
142 } else if (r < 0) {
143 std::cerr << std::format("FETCH_REQ failed: {}\n", r);
144 exit(1);
145 }
146
147 while (true) {
148 uint8_t op = ublksrv_get_op(io_desc);
149 size_t start = io_desc->start_sector * SECTOR_SIZE;
150 size_t size = io_desc->nr_sectors * SECTOR_SIZE;
151 switch (op) {
152 case UBLK_IO_OP_READ:
153 r = co_await condy::async_read(backing_fd, condy::buffer(buf, size),
154 start);
155 if (r < 0) {
156 std::cerr << std::format("async_read backing file failed: {}\n",
157 r);
158 }
159 break;
160 case UBLK_IO_OP_WRITE:
161 r = co_await condy::async_write(backing_fd,
162 condy::buffer(buf, size), start);
163 if (r < 0) {
164 std::cerr << std::format(
165 "async_write backing file failed: {}\n", r);
166 }
167 break;
168 case UBLK_IO_OP_FLUSH:
169 r = co_await condy::async_fsync(backing_fd, IORING_FSYNC_DATASYNC);
170 if (r < 0) {
171 std::cerr << std::format("async_fsync failed: {}\n", r);
172 }
173 break;
174 default:
175 std::cerr << std::format("Unknown op: {}\n", op);
176 r = -EINVAL;
177 break;
178 }
179
180 cmd = {};
181 cmd.q_id = q_id;
182 cmd.tag = tag;
183 cmd.result = static_cast<int32_t>(r);
184 cmd.addr = reinterpret_cast<uint64_t>(buf);
185 r = co_await ublk_io_cmd(UBLK_U_IO_COMMIT_AND_FETCH_REQ, cmd);
186 if (r == UBLK_IO_RES_ABORT) {
187 co_return;
188 } else if (r < 0) {
189 std::cerr << std::format("COMMIT_AND_FETCH_REQ failed: {}\n", r);
190 exit(1);
191 }
192 }
193}
194
195condy::Coro<void> io_queue(size_t q_id) {
196 size_t io_desc_base_size = queue_depth * sizeof(ublksrv_io_desc);
197 void *addr =
198 mmap(nullptr, io_desc_base_size, PROT_READ, MAP_SHARED | MAP_POPULATE,
199 ublkc_fd, io_desc_offset(q_id));
200 if (addr == MAP_FAILED) {
201 std::perror("mmap io_desc_base");
202 std::exit(1);
203 }
204 auto *io_desc_base = reinterpret_cast<ublksrv_io_desc *>(addr);
205
206 auto &fd_table = condy::current_runtime().fd_table();
207 int r = fd_table.init(&ublkc_fd, 1);
208 if (r < 0) {
209 std::cerr << std::format("fd_table.init failed: {}\n", r);
210 std::exit(1);
211 }
212
213 std::vector<condy::Task<void>> tasks;
214 tasks.reserve(queue_depth);
215 for (size_t tag = 0; tag < queue_depth; tag++) {
216 tasks.push_back(condy::co_spawn(io_loop(q_id, tag, io_desc_base)));
217 }
218 for (auto &t : tasks) {
219 co_await t;
220 }
221
222 fd_table.destroy();
223 munmap(io_desc_base, io_desc_base_size);
224}
225
226condy::Coro<void> co_main() {
227 int r;
228 ublksrv_ctrl_cmd ctrl;
229
230 co_await setup_device();
231
232 bufs_base =
233 mmap(nullptr, num_queues * queue_depth * BUF_SIZE,
234 PROT_READ | PROT_WRITE, MAP_PRIVATE | MAP_ANONYMOUS, -1, 0);
235 if (bufs_base == MAP_FAILED) {
236 std::perror("mmap bufs_base");
237 std::exit(1);
238 }
239
240 std::vector<condy::Task<void>> tasks;
241
242 if (multi_queues()) {
243 tasks.reserve(num_queues);
244 for (size_t q_id = 0; q_id < num_queues; q_id++) {
245 tasks.push_back(
246 condy::co_spawn(*queue_runtimes[q_id], io_queue(q_id)));
247 }
248 } else {
249 tasks.push_back(condy::co_spawn(io_queue(0)));
250 }
251
252 ctrl = {};
253 ctrl.dev_id = dev_id;
254 ctrl.queue_id = -1;
255 ctrl.data[0] = getpid();
256 r = co_await ublk_ctrl_cmd(ctrl_fd, UBLK_U_CMD_START_DEV, ctrl);
257 if (r < 0) {
258 std::cerr << std::format("START_DEV failed: {}\n", r);
259 std::exit(1);
260 }
261 std::cout << std::format("ublk-loop: /dev/ublkb{} started\n", dev_id);
262
263 signalfd_siginfo si;
264 co_await condy::async_read(signal_fd, condy::buffer(&si, sizeof(si)), 0);
265 std::cout << std::format(
266 "ublk-loop: received signal {}, shutting down...\n", si.ssi_signo);
267
268 ctrl = {};
269 ctrl.dev_id = dev_id;
270 ctrl.queue_id = -1;
271 r = co_await ublk_ctrl_cmd(ctrl_fd, UBLK_U_CMD_STOP_DEV, ctrl);
272 if (r < 0) {
273 std::cerr << std::format("STOP_DEV failed: {}\n", r);
274 exit(1);
275 }
276 std::cout << "ublk-loop: device stopped\n";
277
278 for (auto &t : tasks) {
279 co_await t;
280 }
281
282 munmap(bufs_base, num_queues * queue_depth * BUF_SIZE);
283 co_await condy::async_close(ublkc_fd);
284
285 ctrl = {};
286 ctrl.dev_id = dev_id;
287 ctrl.queue_id = -1;
288 r = co_await ublk_ctrl_cmd(ctrl_fd, UBLK_U_CMD_DEL_DEV, ctrl);
289 if (r < 0) {
290 std::cerr << std::format("DEL_DEV failed: {}\n", r);
291 }
292 std::cout << std::format("ublk-loop: device /dev/ublkb{} deleted\n",
293 dev_id);
294
295 co_await condy::async_close(ctrl_fd);
296}
297
298void usage(const char *prog) {
299 std::cerr << std::format("Usage: {} [OPTIONS] <backing-file>\n"
300 "Options:\n"
301 " -d NUM Queue depth (default: 32)\n"
302 " -n NUM Number of queues (default: 1)\n"
303 " -i NUM Device ID (default: auto)\n"
304 " -h Show this help\n",
305 prog);
306}
307
308int main(int argc, char *argv[]) noexcept(false) {
309 int opt;
310 while ((opt = getopt(argc, argv, "d:n:i:h")) != -1) {
311 switch (opt) {
312 case 'd':
313 queue_depth = std::stoull(optarg);
314 break;
315 case 'n':
316 num_queues = std::stoull(optarg);
317 break;
318 case 'i':
319 dev_id = std::stoul(optarg);
320 break;
321 case 'h':
322 default:
323 usage(argv[0]);
324 return opt == 'h' ? 0 : 1;
325 }
326 }
327
328 if (optind >= argc) {
329 usage(argv[0]);
330 return 1;
331 }
332 std::string backing_path = argv[optind];
333
334 backing_fd = open(backing_path.c_str(), O_RDWR);
335 if (backing_fd < 0) {
336 std::perror("open backing file");
337 return 1;
338 }
339
340 sigset_t mask;
341 sigemptyset(&mask);
342 sigaddset(&mask, SIGINT);
343 sigaddset(&mask, SIGTERM);
344 sigprocmask(SIG_BLOCK, &mask, nullptr);
345
346 signal_fd = signalfd(-1, &mask, SFD_NONBLOCK);
347 if (signal_fd < 0) {
348 std::perror("signalfd");
349 exit(1);
350 }
351
352 off_t file_size = lseek(backing_fd, 0, SEEK_END);
353 dev_sectors = static_cast<size_t>(file_size) / SECTOR_SIZE;
354 std::cout << std::format("ublk-loop: backing file {} MiB ({} sectors)\n",
355 file_size / 1024 / 1024, dev_sectors);
356
357 condy::RuntimeOptions options;
358 options.enable_sqe128();
359 options.sq_size(multi_queues() ? 32 : queue_depth);
360 condy::Runtime runtime(options);
361
362 std::vector<std::thread> queue_threads;
363 if (multi_queues()) {
364 condy::RuntimeOptions queue_options;
365 queue_options.enable_attach_wq(runtime);
366 queue_options.sq_size(queue_depth);
367
368 queue_runtimes.reserve(num_queues);
369 queue_threads.reserve(num_queues);
370 for (size_t i = 0; i < num_queues; i++) {
371 queue_runtimes.push_back(
372 std::make_unique<condy::Runtime>(queue_options));
373 queue_threads.emplace_back([&, i] { queue_runtimes[i]->run(); });
374 }
375 }
376
377 condy::sync_wait(runtime, co_main());
378
379 for (size_t i = 0; i < queue_runtimes.size(); i++) {
380 queue_runtimes[i]->allow_exit();
381 queue_threads[i].join();
382 }
383
384 close(signal_fd);
385 close(backing_fd);
386 return 0;
387}
Coroutine type used to define a coroutine function.
Definition coro.hpp:25
The event loop runtime for executing asynchronous.
Definition runtime.hpp:36
Main include file for the Condy library.
auto async_read(Fd fd, const Buffer &buf, __u64 offset, int flags=0)
See io_uring_prep_read.
auto async_fsync(Fd fd, unsigned fsync_flags)
See io_uring_prep_fsync.
T sync_wait(Runtime &runtime, Coro< T, Allocator > coro)
Synchronously wait for a coroutine to complete in the given runtime.
Definition sync_wait.hpp:25
Task< T, Allocator > co_spawn(Runtime &runtime, Coro< T, Allocator > coro) noexcept
Spawn a coroutine as a task in the given runtime.
Definition task.hpp:101
auto fixed(int fd)
Mark a file descriptor as fixed for io_uring operations.
Definition helpers.hpp:70
MutableBuffer buffer(void *data, size_t size) noexcept
Create a buffer object from various data sources.
Definition buffers.hpp:86
auto async_write(Fd fd, const Buffer &buf, __u64 offset, int flags=0)
See io_uring_prep_write.
auto async_close(int fd)
See io_uring_prep_close.
auto async_open(const char *path, int flags, mode_t mode)
See io_uring_prep_openat.
auto & current_runtime() noexcept
Get the current runtime.
Definition runtime.hpp:452
auto async_uring_cmd(int cmd_op, Fd fd, CmdFunc &&cmd_func, Args &&...handler_args)
See io_uring_prep_uring_cmd.
Self & enable_sqe128()
See IORING_SETUP_SQE128.
Self & enable_attach_wq(Runtime &other)
See IORING_SETUP_ATTACH_WQ.
Self & sq_size(size_t v)
Set SQ size.
A simple CQE handler that extracts the result from the CQE without any additional processing.