VoltMod
C++23 framework for CS2 server plugins
Loading...
Searching...
No Matches
Client.hpp
Go to the documentation of this file.
1#pragma once
2
11#include <atomic>
12#include <chrono>
13#include <concepts>
14#include <condition_variable>
15#include <deque>
16#include <functional>
17#include <future>
18#include <memory>
19#include <mutex>
20#include <optional>
21#include <stdexcept>
22#include <string>
23#include <string_view>
24#include <thread>
25#include <type_traits>
26#include <utility>
27#include <variant>
28
29namespace VoltMod
30{
31
32/**
33 * @brief Async database access over Postgres, MariaDB and SQLite.
34 *
35 * One worker thread owns the connection; jobs run in order and their completions run on the game
36 * thread. A job is a callable over `auto& conn` that returns the same type for every driver.
37 * @ref Run blocks, so it is for load time; gameplay uses @ref RunAsync.
38 */
40{
41public:
42 /** What a job returns; @ref MakeJob makes every driver agree. */
43 template <class Fn>
44 using ResultOf = std::invoke_result_t<Fn&, SqliteConnection&>;
45
46 /** @p scheduler drives per-frame completion delivery and must outlive this object. */
47 explicit Database(Scheduler& scheduler) : _scheduler(scheduler) {}
48 ~Database();
49 Database(const Database&) = delete;
50 Database& operator=(const Database&) = delete;
51
52 /** Start the worker and ping; false on a bad config or an unreachable database. */
53 bool Connect(const DatabaseConfig& config);
54
55 /** Let queued jobs finish within @p stopDeadline, then join the worker; later jobs fail and
56 * undelivered completions are dropped. The destructor calls it. */
57 void Disconnect(std::chrono::milliseconds stopDeadline = std::chrono::seconds(5));
58
59 /** Run @p fn on the worker; @p onDone runs on the game thread later. @p name is a log label. */
60 template <class Fn>
61 void RunAsync(std::string name, Fn fn, std::move_only_function<void(Result<ResultOf<Fn>>)> onDone = {})
62 {
63 if (!onDone)
64 {
65 Enqueue(MakeJob(std::move(name), std::move(fn), {}));
66 return;
67 }
68 Enqueue(MakeJob(
69 std::move(name), std::move(fn), [this, onDone = std::move(onDone)](Result<ResultOf<Fn>> result) mutable {
70 _completions.Push(
71 [onDone = std::move(onDone), result = std::move(result)]() mutable { onDone(std::move(result)); });
72 }));
73 }
74
75 /** @ref RunAsync for a callback taking the value alone; a failure is only logged. */
76 template <class Fn, class OnValue>
77 requires(!std::is_void_v<ResultOf<Fn>> && std::invocable<OnValue&, ResultOf<Fn>> &&
78 !std::invocable<OnValue&, Result<ResultOf<Fn>>>)
79 void RunAsync(std::string name, Fn fn, OnValue onValue)
80 {
81 if (!IsSet(onValue))
82 {
83 RunAsync(std::move(name), std::move(fn));
84 return;
85 }
86 RunAsync(std::move(name), std::move(fn),
87 std::move_only_function<void(Result<ResultOf<Fn>>)>(
88 [onValue = std::move(onValue)](Result<ResultOf<Fn>> result) mutable {
89 if (result)
90 {
91 onValue(std::move(*result));
92 }
93 }));
94 }
95
96 /** @ref RunAsync that blocks; load time only. */
97 template <class Fn>
98 Result<ResultOf<Fn>> Run(std::string name, Fn fn)
99 {
100 // Shared: the worker may still hold the job after this returns.
101 auto promise = std::make_shared<std::promise<Result<ResultOf<Fn>>>>();
102 std::future<Result<ResultOf<Fn>>> answer = promise->get_future();
103 Enqueue(MakeJob(std::move(name), std::move(fn),
104 [promise](Result<ResultOf<Fn>> result) { promise->set_value(std::move(result)); }));
105 return answer.get();
106 }
107
108 /** @ref Run, falling back to @p fallback when the job fails. */
109 template <class Fn>
110 ResultOf<Fn> RunOr(std::string name, Fn fn, ResultOf<Fn> fallback = {})
111 {
112 auto result = Run(std::move(name), std::move(fn));
113 return result ? std::move(*result) : std::move(fallback);
114 }
115
116 /** Run the finished jobs' completions; @ref Connect schedules this every frame. */
117 void DispatchCompletions() { _completions.RunAll(); }
118
119 /** Whether the worker's last job had a live connection; a diagnostic, not a guarantee. */
120 bool IsConnected() const { return _connected.load(std::memory_order_relaxed); }
121
122 /** The configured driver, set by @ref Connect. */
123 Driver GetDriver() const { return _driver; }
124
125private:
126 struct Job
127 {
128 std::string Name; ///< log label only
129 std::move_only_function<void(AnyConnection&)> Run;
130 std::move_only_function<void(Error)> Fail;
131 };
132
133 enum class State
134 {
135 Stopped,
136 Running,
137 Stopping,
138 };
139
140 /** A job handing @p fn's result, or why it failed, to @p deliver on the worker thread; an
141 * empty @p deliver only runs it. */
142 template <class Fn>
143 static Job MakeJob(std::string name, Fn fn, std::move_only_function<void(Result<ResultOf<Fn>>)> deliver)
144 {
145 static_assert(std::same_as<std::invoke_result_t<Fn&, PostgresConnection&>, ResultOf<Fn>> &&
146 std::same_as<std::invoke_result_t<Fn&, MariaDbConnection&>, ResultOf<Fn>>,
147 "a database job must return the same type for every driver");
148 if (!deliver)
149 {
150 return {
151 .Name = std::move(name),
152 .Run = [fn = std::move(fn)](AnyConnection& conn) mutable { Invoke(conn, fn); },
153 .Fail = [](Error) {},
154 };
155 }
156 // Run and Fail both hold it; only one of them delivers.
157 auto shared = std::make_shared<decltype(deliver)>(std::move(deliver));
158 return {
159 .Name = std::move(name),
160 .Run =
161 [fn = std::move(fn), shared](AnyConnection& conn) mutable {
162 if constexpr (std::is_void_v<ResultOf<Fn>>)
163 {
164 Invoke(conn, fn);
165 (*shared)(Result<void>{});
166 }
167 else
168 {
169 (*shared)(Result<ResultOf<Fn>>{Invoke(conn, fn)});
170 }
171 },
172 .Fail = [shared](Error error) { (*shared)(std::unexpected(std::move(error))); },
173 };
174 }
175
176 /** @p fn's result on the open connection; throws without one, as a failing query does. */
177 template <class Fn>
178 static ResultOf<Fn> Invoke(AnyConnection& conn, Fn& fn)
179 {
180 return std::visit(
181 [&fn](auto& open) -> ResultOf<Fn> {
182 if constexpr (std::same_as<std::remove_cvref_t<decltype(open)>, std::monostate>)
183 {
184 throw std::logic_error("no database connection");
185 }
186 else
187 {
188 return fn(open);
189 }
190 },
191 conn);
192 }
193
194 /** False for an empty std::function, which several callers pass. */
195 static bool IsSet(const auto& callback)
196 {
197 if constexpr (requires { static_cast<bool>(callback); })
198 {
199 return static_cast<bool>(callback);
200 }
201 else
202 {
203 return true;
204 }
205 }
206
207 void Enqueue(Job job);
208 void WorkerMain();
209 /** Waits for the next job; nothing once stopping. */
210 std::optional<Job> TakeJob();
211 void RunJob(Job& job);
212 /** Open the connection if needed; its failure log carries no secrets. */
213 bool EnsureOpen();
214 /** The next job reopens it. */
215 void DropConnection();
216
217 Scheduler& _scheduler;
218 DatabaseConfig _config;
219 Driver _driver = Driver::Postgres;
220
221 std::mutex _queueMutex;
222 std::condition_variable _queueCv;
223 std::deque<Job> _queue;
224 State _state = State::Stopped;
225 std::chrono::steady_clock::time_point _stopDeadline{};
226
227 GameThreadQueue _completions;
228 std::thread _worker;
229 Subscription _onFrame;
230
231 /** Written by the worker, read on the game thread. */
232 std::atomic<bool> _connected{false};
233
234 /** Worker thread only. */
235 AnyConnection _connection;
236 std::chrono::steady_clock::time_point _retryAt{};
237};
238
239/**
240 * Apply the `NNNN_*.sql` files in `dir` newer than the database's version, in order, each in its own
241 * transaction under a lock. @ref ResolveDialect fills their placeholders for the live driver. A
242 * missing directory is a no-op; on failure the database stays at the last applied version.
243 */
244MigrationResult RunMigrations(Database& db, std::string_view dir, const MigrationOptions& options = {});
245
246} // namespace VoltMod
Async database access over Postgres, MariaDB and SQLite.
Definition Client.hpp:40
Result< ResultOf< Fn > > Run(std::string name, Fn fn)
Definition Client.hpp:98
Database(const Database &)=delete
ResultOf< Fn > RunOr(std::string name, Fn fn, ResultOf< Fn > fallback={})
Definition Client.hpp:110
bool Connect(const DatabaseConfig &config)
Definition Client.cpp:19
Database(Scheduler &scheduler)
Definition Client.hpp:47
void Disconnect(std::chrono::milliseconds stopDeadline=std::chrono::seconds(5))
Definition Client.cpp:55
bool IsConnected() const
Definition Client.hpp:120
std::invoke_result_t< Fn &, SqliteConnection & > ResultOf
Definition Client.hpp:44
void DispatchCompletions()
Definition Client.hpp:117
void RunAsync(std::string name, Fn fn, OnValue onValue)
Definition Client.hpp:79
void RunAsync(std::string name, Fn fn, std::move_only_function< void(Result< ResultOf< Fn > >)> onDone={})
Definition Client.hpp:61
Driver GetDriver() const
Definition Client.hpp:123
Database & operator=(const Database &)=delete
void Push(std::move_only_function< void()> completion)
One-shot delays and repeating timers, run on the game thread.
Definition Scheduler.hpp:18
void Error(std::format_string< Args... > fmt, Args &&... args)
Definition Log.hpp:97
constexpr std::string_view Name(E value) noexcept
Definition EnumNames.hpp:33
static std::string ReadFile(const std::filesystem::path &path)
Definition Loader.cpp:56
MigrationResult RunMigrations(Database &db, std::string_view dir, const MigrationOptions &options={})
Definition Migrator.cpp:57
std::expected< T, Error > Result
Definition Result.hpp:64
std::variant< std::monostate, PostgresConnection, MariaDbConnection, SqliteConnection > AnyConnection
Connection parameters for every supported backend.
One failure: a code to branch on, text for the log, and an optional translation key.
Definition Result.hpp:42