ublk-cpp v0.0
Loading...
Searching...
No Matches
daemon.hpp
Go to the documentation of this file.
1
6
7#pragma once
8
10#include "ublk/detail/io.hpp"
12#include "ublk/detail/shm.hpp"
13#include "ublk/detail/utils.hpp"
14#include "ublk/handler.hpp"
15#include "ublk/runtime.hpp"
16#include "ublk/ublk_cmd.h"
17#include <condy.hpp>
18#include <stdexcept>
19
20namespace ublk {
21namespace detail {
22
23namespace ex = condy::detail::ex;
24
25inline bool need_recovery(const ublksrv_ctrl_dev_info *info) noexcept {
26 return info->state == UBLK_S_DEV_QUIESCED ||
27 info->state == UBLK_S_DEV_FAIL_IO;
28}
29
30struct daemon_setup_t {
31 template <typename Sched, typename Alloc>
32 ex::task<bool, TaskEnv<Sched, Alloc>> invoke(int control_fd,
33 ublksrv_ctrl_dev_info *info) {
34 if (info->dev_id == -1) {
35 co_await raw::add_dev(control_fd, info);
36 co_return true;
37 }
38
39 auto s = raw::add_dev(control_fd, info) |
40 ex::upon_error([&](std::error_code ec) {
41 if (ec.value() != EEXIST) {
42 throw std::system_error(ec, "add_dev");
43 }
44 return -ec.value();
45 });
46 auto r = co_await std::move(s);
47 if (r != -EEXIST) {
48 co_return true;
49 }
50
51 co_await control_get_dev_info_t{}.invoke<Sched, Alloc>(
52 control_fd, info->dev_id, info);
53
54 if (info->flags & UBLK_F_USER_RECOVERY && need_recovery(info)) {
55 co_await (control_start_user_recovery_t{}.invoke<Sched, Alloc>(
56 control_fd, info->dev_id) |
57 ex::write_env(ex::prop{fetch_dev_info, info}));
58 }
59
60 co_return false;
61 }
62};
63
64struct daemon_configure_t {
65 template <typename Sched, typename Alloc>
66 ex::task<bool, TaskEnv<Sched, Alloc>>
67 invoke(int control_fd, uint32_t dev_id, ublk_params *params) {
68 ublksrv_ctrl_dev_info info;
69 co_await cached_get_dev_info<Sched, Alloc>(control_fd, dev_id, &info);
70
71 if (info.state == UBLK_S_DEV_DEAD) {
72 co_await (control_set_params_t{}.invoke<Sched, Alloc>(
73 control_fd, dev_id, params) |
74 ex::write_env(ex::prop{fetch_dev_info, &info}));
75 co_return true;
76 } else {
77 co_await (control_get_params_t{}.invoke<Sched, Alloc>(
78 control_fd, dev_id, params) |
79 ex::write_env(ex::prop{fetch_dev_info, &info}));
80 co_return false;
81 }
82 }
83};
84
85struct daemon_start_t {
86 template <typename Sched, typename Alloc>
87 ex::task<void, TaskEnv<Sched, Alloc>>
88 invoke(int control_fd, uint32_t dev_id, int32_t daemon_pid) {
89 ublksrv_ctrl_dev_info info;
90 co_await cached_get_dev_info<Sched, Alloc>(control_fd, dev_id, &info);
91
92 if (!need_recovery(&info)) {
93 co_await (control_start_dev_t{}.invoke<Sched, Alloc>(
94 control_fd, dev_id, daemon_pid) |
95 ex::write_env(ex::prop{fetch_dev_info, &info}));
96 } else if (info.flags & UBLK_F_USER_RECOVERY) {
97 co_await (control_end_user_recovery_t{}.invoke<Sched, Alloc>(
98 control_fd, dev_id, daemon_pid) |
99 ex::write_env(ex::prop{fetch_dev_info, &info}));
100 } else {
101 throw std::runtime_error(
102 "device can neither be started nor recovered");
103 }
104 }
105};
106
107struct daemon_run_t {
108 template <typename Sched, typename Alloc, IoHandler Handler>
109 ex::task<void, TaskEnv<Sched, Alloc>>
110 invoke(int control_fd, uint32_t dev_id, Handler *handler,
111 const RuntimeOptions *runtime_options, size_t nr_files,
112 size_t nr_buffers) {
113 auto alloc = co_await ex::read_env(ex::get_allocator);
114
115 ublksrv_ctrl_dev_info info;
116 co_await cached_get_dev_info<Sched, Alloc>(control_fd, dev_id, &info);
117
118 nr_files = std::max<size_t>(nr_files, 1);
119 if (info.flags & (UBLK_F_SUPPORT_ZERO_COPY | UBLK_F_AUTO_BUF_REG)) {
120 nr_buffers = std::max<size_t>(nr_buffers, info.queue_depth);
121 }
122 auto options = runtime_options
123 ? runtime_options->build(info.queue_depth)
124 : condy::RuntimeOptions();
125
126 AllocVector<std::unique_ptr<IoLoop>, decltype(alloc)> loops(alloc);
127 loops.reserve(info.nr_hw_queues);
128 for (uint16_t q_id = 0; q_id < info.nr_hw_queues; q_id++) {
129 cpu_set_t cpuset;
130 co_await (control_get_queue_affinity_t{}.invoke<Sched, Alloc>(
131 control_fd, dev_id, q_id, &cpuset) |
132 ex::write_env(ex::prop{fetch_dev_info, &info}));
133 auto loop =
134 std::make_unique<IoLoop>(options, nr_files, nr_buffers, cpuset);
135 loops.push_back(std::move(loop));
136 options.enable_attach_wq(loops[0]->runtime());
137 }
138
139 std::string path = dev_path(dev_id);
140 int ublkc_fd = co_await condy::async_open(path.c_str(), O_RDWR, 0);
141 auto d = defer([&] noexcept { close(ublkc_fd); });
142
143 ex::simple_counting_scope scope;
144 AllocVector<std::exception_ptr, decltype(alloc)> errs(info.nr_hw_queues,
145 alloc);
146 for (uint16_t q_id = 0; q_id < info.nr_hw_queues; q_id++) {
147 auto sched = loops[q_id]->get_scheduler();
148 auto env = ex::env{ex::prop{ex::get_allocator, alloc},
149 ex::prop{fetch_dev_info, &info}};
150 auto task = io_run_dev_t{}.invoke<decltype(sched), Alloc>(
151 ublkc_fd, q_id, &info, handler) |
152 ex::write_env(std::move(env));
153 auto s = ex::starts_on(sched, std::move(task)) |
154 ex::upon_error(
155 [&, q_id](const std::exception_ptr &ep) noexcept {
156 errs[q_id] = ep;
157 });
158 ex::spawn(std::move(s), scope.get_token());
159 }
160 co_await scope.join();
161
162 for (auto &ep : errs) {
163 if (ep) {
164 std::rethrow_exception(ep);
165 }
166 }
167 }
168};
169
170struct daemon_shm_server_t {
171 template <typename Sched, typename Alloc, ShmHandler Handler>
172 auto invoke(int control_fd, uint32_t dev_id, std::string_view path,
173 Handler *handler, uint32_t flags) {
174 return shm_server_run<Sched, Alloc>(
175 path, Session<Sched, Alloc, Handler>(control_fd, dev_id, flags,
176 *handler));
177 }
178};
179
180} // namespace detail
181} // namespace ublk
Implementation of ublk control command interface.
Implementation of the ublk I/O loop for processing I/O requests on a dedicated scheduler.
Handler concepts for ublk server.
Helpful condy runtime wrapper for running a dedicated I/O loop on a pinned thread.
auto add_dev(Fd fd, ublksrv_ctrl_dev_info *info) noexcept
Add a new ublk device.
Definition raw.hpp:80
The main namespace of the ublk-cpp library.
Definition ublk.hpp:21
constexpr fetch_dev_info_t fetch_dev_info
Query object instance of fetch_dev_info_t.
Definition query.hpp:48
Runtime options based on condy::RuntimeOptions.
Implementation of the shared-memory buffer registration service over a Unix domain socket.
Small utility helpers.