VoltMod
C++23 framework for CS2 server plugins
Loading...
Searching...
No Matches
HttpClient.cpp
Go to the documentation of this file.
4#include <atomic>
5#include <chrono>
6#include <condition_variable>
7#include <cpr/cpr.h>
8#include <deque>
9#include <mutex>
10#include <optional>
11#include <thread>
12#include <utility>
13#include <vector>
14
15namespace VoltMod
16{
17
18static constexpr size_t MaxWorkers = 4;
19/** Requests beyond this many waiting fail at once. */
20static constexpr size_t MaxWaiting = 64;
21
22static HttpResult ToResult(cpr::Response&& response)
23{
24 if (response.error)
25 {
26 return {.Error = std::move(response.error.message)};
27 }
28 return {.Ok = true, .StatusCode = static_cast<long>(response.status_code), .Body = std::move(response.text)};
29}
30
31/** Worker thread: touches nothing but the request. */
32static HttpResult Perform(const HttpRequest& request, const std::atomic_bool& stopped)
33{
34 const cpr::Url url{request.Url};
35 const cpr::Header headers(request.Headers.begin(), request.Headers.end());
36 const cpr::Timeout timeout{std::chrono::milliseconds{request.TimeoutMs}};
37 // False aborts the transfer, so Stop does not wait out a stalled endpoint.
38 const cpr::ProgressCallback progress{[&stopped](auto&&...) { return !stopped; }};
39
40 switch (request.Method)
41 {
42 case HttpMethod::Get:
43 return ToResult(cpr::Get(url, headers, timeout, progress));
45 return ToResult(cpr::Post(url, cpr::Body{request.Body}, headers, timeout, progress));
46 case HttpMethod::Put:
47 return ToResult(cpr::Put(url, cpr::Body{request.Body}, headers, timeout, progress));
49 return ToResult(cpr::Patch(url, cpr::Body{request.Body}, headers, timeout, progress));
51 return ToResult(cpr::Delete(url, cpr::Body{request.Body}, headers, timeout, progress));
52 }
53 return {.Error = "unsupported HTTP method"};
54}
55
61
62/** Shared with the workers; Mutex guards the queue and the worker list. */
64{
65 std::mutex Mutex;
66 std::condition_variable Wake;
67 std::deque<WaitingRequest> Waiting;
68 std::vector<std::thread> Workers;
69 size_t IdleWorkers = 0;
71 std::atomic_bool Stopped = false; // never cleared: its plugin is unloading
72
73 void Work()
74 {
75 while (std::optional<WaitingRequest> next = Take())
76 {
77 Complete(std::move(next->OnComplete), Perform(next->Request, Stopped));
78 }
79 }
80
81 /** Runs @p onComplete with @p result on the game thread's next frame. */
83 {
84 Completions.Push([onComplete = std::move(onComplete), result = std::move(result)] {
85 if (onComplete)
86 {
88 }
89 });
90 }
91
92 /** Waits for the next request; nothing once stopped. */
93 std::optional<WaitingRequest> Take()
94 {
95 std::unique_lock lock(Mutex);
97 Wake.wait(lock, [this] { return Stopped || !Waiting.empty(); });
99 if (Stopped)
100 {
101 return std::nullopt;
102 }
103 WaitingRequest next = std::move(Waiting.front());
104 Waiting.pop_front();
105 return next;
106 }
107};
108
110 : _requests(std::make_unique<Requests>()), _onFrame(scheduler.EveryFrame([this] { RunCompletions(); }))
111{}
112
114{
115 Stop();
116}
117
119{
120 Requests& requests = *_requests;
121 if (requests.Stopped)
122 {
123 Log::Warn("http: dropped a request to '{}' because the client is stopped.", request.Url);
124 return;
125 }
126
127 {
128 std::lock_guard lock(requests.Mutex);
129 if (requests.Waiting.size() >= MaxWaiting)
130 {
131 Log::Warn("http: refused a request to '{}': {} requests are already waiting.", request.Url, MaxWaiting);
132 requests.Complete(std::move(onComplete), {.Error = "too many requests waiting"});
133 return;
134 }
135
136 requests.Waiting.push_back({std::move(request), std::move(onComplete)});
137 // Stop joins every worker before Requests goes away.
138 if (requests.IdleWorkers == 0 && requests.Workers.size() < MaxWorkers)
139 {
140 requests.Workers.emplace_back([&requests] { requests.Work(); });
141 }
142 }
143 requests.Wake.notify_one();
144}
145
147{
148 Requests& requests = *_requests;
149 {
150 std::lock_guard lock(requests.Mutex);
151 requests.Stopped = true;
152 }
153 requests.Wake.notify_all();
154 for (std::thread& worker : requests.Workers)
155 {
156 worker.join();
157 }
158
159 requests.Workers.clear();
160 requests.Waiting.clear();
161 requests.Completions.Clear();
162}
163
164void HttpClient::RunCompletions()
165{
166 _requests->Completions.RunAll();
167}
168
169} // namespace VoltMod
Callbacks worker threads hand to the game thread: Push from any thread, RunAll on the game thread.
void Push(std::move_only_function< void()> completion)
void Send(HttpRequest request, HttpCompletion onComplete)
HttpClient(Scheduler &scheduler)
One-shot delays and repeating timers, run on the game thread.
Definition Scheduler.hpp:18
void Warn(std::format_string< Args... > fmt, Args &&... args)
Definition Log.hpp:88
static HttpResult Perform(const HttpRequest &request, const std::atomic_bool &stopped)
static std::string ReadFile(const std::filesystem::path &path)
Definition Loader.cpp:56
static HttpResult ToResult(cpr::Response &&response)
std::function< void(const HttpResult &)> HttpCompletion
static constexpr size_t MaxWaiting
static constexpr size_t MaxWorkers
std::vector< std::thread > Workers
std::deque< WaitingRequest > Waiting
std::condition_variable Wake
void Complete(HttpCompletion onComplete, HttpResult result)
std::optional< WaitingRequest > Take()
HttpCompletion OnComplete