20#include <linux/ublk_cmd.h>
22#include <sys/signalfd.h>
27size_t queue_depth = 32;
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;
37std::vector<std::unique_ptr<condy::Runtime>> queue_runtimes;
45bool multi_queues() {
return num_queues > 1; }
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);
52auto ublk_ctrl_cmd(
int fd,
int cmd_op,
const ublksrv_ctrl_cmd &ctrl) {
55 [&](io_uring_sqe *sqe) { std::memcpy(sqe->cmd, &ctrl,
sizeof(ctrl)); },
59auto ublk_io_cmd(
int cmd_op,
const ublksrv_io_cmd &cmd) {
62 [&](io_uring_sqe *sqe) { std::memcpy(sqe->cmd, &cmd,
sizeof(cmd)); },
68 ublksrv_ctrl_cmd ctrl;
72 std::cerr << std::format(
"Failed to open {}: {}\n", CTRL_FILE, ctrl_fd);
76 ublksrv_ctrl_dev_info info = {};
78 info.nr_hw_queues = num_queues;
79 info.queue_depth = queue_depth;
80 info.max_io_buf_bytes = BUF_SIZE;
82 ctrl.dev_id = info.dev_id;
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);
88 std::cerr << std::format(
"ADD_DEV failed: {}\n", r);
93 std::string ublkc_path = std::format(
"/dev/ublkc{}", dev_id);
96 std::cerr << std::format(
"Failed to open {}: {}\n", ublkc_path,
100 std::cout << std::format(
"ublk-loop: /dev/ublkc{} created\n", dev_id);
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;
113 ctrl.dev_id = dev_id;
115 ctrl.addr =
reinterpret_cast<uint64_t
>(¶ms);
116 ctrl.len =
sizeof(params);
117 r =
co_await ublk_ctrl_cmd(ctrl_fd, UBLK_U_CMD_SET_PARAMS, ctrl);
119 std::cerr << std::format(
"SET_PARAMS failed: {}\n", r);
123 std::cout << std::format(
"ublk-loop: /dev/ublkc{} configured\n", dev_id);
127 ublksrv_io_desc *io_desc_base) {
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;
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) {
143 std::cerr << std::format(
"FETCH_REQ failed: {}\n", r);
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;
152 case UBLK_IO_OP_READ:
156 std::cerr << std::format(
"async_read backing file failed: {}\n",
160 case UBLK_IO_OP_WRITE:
164 std::cerr << std::format(
165 "async_write backing file failed: {}\n", r);
168 case UBLK_IO_OP_FLUSH:
171 std::cerr << std::format(
"async_fsync failed: {}\n", r);
175 std::cerr << std::format(
"Unknown op: {}\n", op);
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) {
189 std::cerr << std::format(
"COMMIT_AND_FETCH_REQ failed: {}\n", r);
196 size_t io_desc_base_size = queue_depth *
sizeof(ublksrv_io_desc);
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");
204 auto *io_desc_base =
reinterpret_cast<ublksrv_io_desc *
>(addr);
207 int r = fd_table.init(&ublkc_fd, 1);
209 std::cerr << std::format(
"fd_table.init failed: {}\n", r);
213 std::vector<condy::Task<void>> tasks;
214 tasks.reserve(queue_depth);
215 for (
size_t tag = 0; tag < queue_depth; tag++) {
218 for (
auto &t : tasks) {
223 munmap(io_desc_base, io_desc_base_size);
228 ublksrv_ctrl_cmd ctrl;
230 co_await setup_device();
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");
240 std::vector<condy::Task<void>> tasks;
242 if (multi_queues()) {
243 tasks.reserve(num_queues);
244 for (
size_t q_id = 0; q_id < num_queues; q_id++) {
253 ctrl.dev_id = dev_id;
255 ctrl.data[0] = getpid();
256 r =
co_await ublk_ctrl_cmd(ctrl_fd, UBLK_U_CMD_START_DEV, ctrl);
258 std::cerr << std::format(
"START_DEV failed: {}\n", r);
261 std::cout << std::format(
"ublk-loop: /dev/ublkb{} started\n", dev_id);
265 std::cout << std::format(
266 "ublk-loop: received signal {}, shutting down...\n", si.ssi_signo);
269 ctrl.dev_id = dev_id;
271 r =
co_await ublk_ctrl_cmd(ctrl_fd, UBLK_U_CMD_STOP_DEV, ctrl);
273 std::cerr << std::format(
"STOP_DEV failed: {}\n", r);
276 std::cout <<
"ublk-loop: device stopped\n";
278 for (
auto &t : tasks) {
282 munmap(bufs_base, num_queues * queue_depth * BUF_SIZE);
286 ctrl.dev_id = dev_id;
288 r =
co_await ublk_ctrl_cmd(ctrl_fd, UBLK_U_CMD_DEL_DEV, ctrl);
290 std::cerr << std::format(
"DEL_DEV failed: {}\n", r);
292 std::cout << std::format(
"ublk-loop: device /dev/ublkb{} deleted\n",
298void usage(
const char *prog) {
299 std::cerr << std::format(
"Usage: {} [OPTIONS] <backing-file>\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",
308int main(
int argc,
char *argv[])
noexcept(
false) {
310 while ((opt = getopt(argc, argv,
"d:n:i:h")) != -1) {
313 queue_depth = std::stoull(optarg);
316 num_queues = std::stoull(optarg);
319 dev_id = std::stoul(optarg);
324 return opt ==
'h' ? 0 : 1;
328 if (optind >= argc) {
332 std::string backing_path = argv[optind];
334 backing_fd = open(backing_path.c_str(), O_RDWR);
335 if (backing_fd < 0) {
336 std::perror(
"open backing file");
342 sigaddset(&mask, SIGINT);
343 sigaddset(&mask, SIGTERM);
344 sigprocmask(SIG_BLOCK, &mask,
nullptr);
346 signal_fd = signalfd(-1, &mask, SFD_NONBLOCK);
348 std::perror(
"signalfd");
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);
359 options.
sq_size(multi_queues() ? 32 : queue_depth);
362 std::vector<std::thread> queue_threads;
363 if (multi_queues()) {
366 queue_options.
sq_size(queue_depth);
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(); });
379 for (
size_t i = 0; i < queue_runtimes.size(); i++) {
380 queue_runtimes[i]->allow_exit();
381 queue_threads[i].join();
Coroutine type used to define a coroutine function.
The event loop runtime for executing asynchronous.
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.
Task< T, Allocator > co_spawn(Runtime &runtime, Coro< T, Allocator > coro) noexcept
Spawn a coroutine as a task in the given runtime.
auto fixed(int fd)
Mark a file descriptor as fixed for io_uring operations.
MutableBuffer buffer(void *data, size_t size) noexcept
Create a buffer object from various data sources.
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.
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.