Kurlyk
Loading...
Searching...
No Matches
HttpRequestManager.hpp
Go to the documentation of this file.
1#pragma once
2#ifndef KURLYK_HEADER_KURLYK_HTTP_HTTP_REQUEST_MANAGER_HPP_INCLUDED
3#define KURLYK_HEADER_KURLYK_HTTP_HTTP_REQUEST_MANAGER_HPP_INCLUDED
4
7
8#include <type_traits>
9
16
17
18namespace kurlyk {
19
24 public:
25
29 static HttpRequestManager* instance = new HttpRequestManager();
30 return *instance;
31 }
32
38 std::unique_ptr<HttpRequest> request_ptr,
39 HttpResponseCallback callback) {
40 return submit_request(std::move(request_ptr), std::move(callback)).accepted;
41 }
42
48 std::unique_ptr<HttpRequest> request_ptr,
49 HttpResponseCallback callback) {
50 std::lock_guard<std::mutex> lock(m_mutex);
51 if (m_shutdown) {
53 }
54
55 const std::size_t queue_limit = m_max_pending_requests.load();
56 if (queue_limit && m_pending_requests.size() >= queue_limit) {
58 }
59
60# if __cplusplus >= 201402L
61 auto context = std::make_unique<HttpRequestContext>(std::move(request_ptr), std::move(callback));
62# else
63 auto context = std::unique_ptr<HttpRequestContext>(
64 new HttpRequestContext(std::move(request_ptr), std::move(callback)));
65# endif
66 context->in_flight_token = m_next_in_flight_token.fetch_add(1, std::memory_order_relaxed);
67 m_pending_requests.push_back(std::move(context));
68 return SubmitResult{true, std::error_code()};
69 }
70
77 HttpRateLimitHandlePtr create_rate_limit(long requests_per_period, long period_ms, bool sequential = false) {
78 return m_rate_limiter.create_limit_handle(requests_per_period, period_ms, sequential);
79 }
80
84 return m_rate_limiter.get_limit(limit_id);
85 }
86
90 bool remove_limit(long limit_id) {
91 return m_rate_limiter.remove_limit(limit_id);
92 }
93
97 return m_rate_limiter.remove_limit(limit);
98 }
99
108 const HttpRateLimitHandlePtr& general_limit,
109 const HttpRateLimitHandlePtr& specific_limit,
110 uint64_t in_flight_token,
111 const std::string& general_key,
112 const std::string& specific_key) {
113 return m_rate_limiter.allow_request(
114 general_limit, specific_limit, in_flight_token, general_key, specific_key);
115 }
116
124 const HttpRateLimitHandlePtr& general_limit,
125 const HttpRateLimitHandlePtr& specific_limit,
126 uint64_t in_flight_token,
127 const std::string& general_key,
128 const std::string& specific_key) {
129 m_rate_limiter.release_request(
130 general_limit, specific_limit, in_flight_token, general_key, specific_key);
131 }
132
140 template<typename Duration = std::chrono::milliseconds>
142 const HttpRateLimitHandlePtr& general_limit,
143 const HttpRateLimitHandlePtr& specific_limit,
144 const std::string& general_key,
145 const std::string& specific_key
146 ) {
147 return m_rate_limiter.time_until_next_allowed<Duration>(
148 general_limit, specific_limit, general_key, specific_key);
149 }
150
154 return m_request_id_counter++;
155 }
156
159 uint64_t generate_group_id() {
160 return m_group_id_counter++;
161 }
162
168
171 std::size_t max_pending_requests() const {
172 return m_max_pending_requests.load();
173 }
174
178 bool has_requests_by_group_id(uint64_t group_id) const {
179 return group_request_count(group_id) != 0;
180 }
181
185 std::size_t group_request_count(uint64_t group_id) const {
186 std::lock_guard<std::mutex> lock(m_mutex);
187 return group_request_count_unlocked(group_id);
188 }
189
193 void wait_requests_by_group_id(uint64_t group_id, std::function<void()> callback) {
194 if (m_shutdown || group_id == 0) {
195 if (callback) callback();
196 return;
197 }
198
199 bool invoke_now = false;
200 {
201 std::lock_guard<std::mutex> lock(m_mutex);
202 if (group_request_count_unlocked(group_id) == 0) {
203 invoke_now = true;
204 } else {
205 m_group_waiters[group_id].push_back(std::move(callback));
206 }
207 }
208
209 if (invoke_now && callback) {
210 callback();
211 }
212 }
213
217 void cancel_request_by_id(uint64_t request_id, std::function<void()> callback) {
218 if (m_shutdown || request_id == 0) {
219 if (callback) callback();
220 return;
221 }
222 std::lock_guard<std::mutex> lock(m_mutex);
223 m_requests_to_cancel_by_id[request_id].push_back(std::move(callback));
224 }
225
229 void cancel_requests_by_group_id(uint64_t group_id, std::function<void()> callback) {
230 if (m_shutdown || group_id == 0) {
231 if (callback) callback();
232 return;
233 }
234 bool invoke_now = false;
235 {
236 std::lock_guard<std::mutex> lock(m_mutex);
237 if (group_request_count_unlocked(group_id) == 0) {
238 invoke_now = true;
239 } else {
240 m_groups_to_cancel[group_id].push_back(std::move(callback));
241 }
242 }
243 if (invoke_now && callback) {
244 callback();
245 }
246 }
247
262
265 void shutdown() override {
266 if (m_shutdown.exchange(true)) {
268 return;
269 }
274 }
275
278 bool is_loaded() const override {
279 std::lock_guard<std::mutex> lock(m_mutex);
280 return
281 !m_pending_requests.empty() ||
282 !m_failed_requests.empty() ||
283 !m_active_request_batches.empty() ||
285 !m_groups_to_cancel.empty();
286 }
287
288 private:
289 mutable std::mutex m_mutex;
290 std::list<std::unique_ptr<HttpRequestContext>> m_pending_requests;
291 std::list<std::unique_ptr<HttpRequestContext>> m_failed_requests;
292 std::list<std::unique_ptr<HttpBatchRequestHandler>> m_active_request_batches;
293 using callback_list_t = std::list<std::function<void()>>;
294 using cancel_map_t = std::unordered_map<uint64_t, callback_list_t>;
299 std::atomic<uint64_t> m_next_in_flight_token{1};
300 std::atomic<uint64_t> m_request_id_counter = ATOMIC_VAR_INIT(1);
301 std::atomic<uint64_t> m_group_id_counter = ATOMIC_VAR_INIT(1);
302 std::atomic<bool> m_shutdown = ATOMIC_VAR_INIT(false);
303 std::atomic<std::size_t> m_max_pending_requests = ATOMIC_VAR_INIT(0);
304
305 std::size_t group_request_count_unlocked(uint64_t group_id) const {
306 if (group_id == 0) return 0;
307
308 std::size_t count = 0;
309 for (const auto& context : m_pending_requests) {
310 if (context && context->request && context->request->group_id == group_id) {
311 ++count;
312 }
313 }
314 for (const auto& context : m_failed_requests) {
315 if (context && context->request && context->request->group_id == group_id) {
316 ++count;
317 }
318 }
319 for (const auto& batch : m_active_request_batches) {
320 if (batch) {
321 count += batch->group_request_count(group_id);
322 }
323 }
324 return count;
325 }
326
328 cancel_map_t ready_waiters;
329 {
330 std::lock_guard<std::mutex> lock(m_mutex);
331 auto it = m_group_waiters.begin();
332 while (it != m_group_waiters.end()) {
333 if (group_request_count_unlocked(it->first) != 0) {
334 ++it;
335 continue;
336 }
337 ready_waiters.emplace(it->first, std::move(it->second));
338 it = m_group_waiters.erase(it);
339 }
340 }
341
342 for (const auto& item : ready_waiters) {
343 for (const auto& callback : item.second) {
344 if (callback) callback();
345 }
346 }
347 }
348
350 cancel_map_t waiters;
351 {
352 std::lock_guard<std::mutex> lock(m_mutex);
353 waiters = std::move(m_group_waiters);
354 m_group_waiters.clear();
355 }
356
357 for (const auto& item : waiters) {
358 for (const auto& callback : item.second) {
359 if (callback) callback();
360 }
361 }
362 }
363
366 std::unique_lock<std::mutex> lock(m_mutex);
367 if (m_pending_requests.empty()) return;
368
369 std::vector<std::unique_ptr<HttpRequestContext>> pending_request;
370 std::vector<std::unique_ptr<HttpRequestContext>> failed_requests;
371
372 auto it = m_pending_requests.begin();
373 while (it != m_pending_requests.end()) {
374 auto& context = *it;
375 auto& request = context->request;
376 // Check if the request is valid.
377 if (!request) {
378 failed_requests.push_back(std::move(context));
379 it = m_pending_requests.erase(it);
380 continue;
381 }
382
383 // Set up completion callback before allow_request so it is armed
384 // even if an exception occurs after the token is committed.
385 auto general_limit = request->general_rate_limit;
386 auto specific_limit = request->specific_rate_limit;
387 uint64_t token = context->in_flight_token;
388 auto general_key = request->general_rate_limit_key;
389 auto specific_key = request->specific_rate_limit_key;
390
391 // Preserve any previous on_complete so a retry that already owns
392 // a sequential token does not lose its cleanup callback.
393 auto old_on_complete = std::move(context->on_complete);
394 context->on_complete = [this, general_limit, specific_limit, token, general_key, specific_key]() {
395 m_rate_limiter.release_request(general_limit, specific_limit, token, general_key, specific_key);
396 };
397
398 // Check if the request is allowed by the rate limiter.
399 const bool allowed = m_rate_limiter.allow_request(
400 general_limit,
401 specific_limit,
402 token,
403 request->general_rate_limit_key,
404 request->specific_rate_limit_key);
405 if (!allowed) {
406 context->on_complete = std::move(old_on_complete);
407 ++it;
408 continue;
409 }
410
411 pending_request.push_back(std::move(context));
412 it = m_pending_requests.erase(it);
413 }
414 lock.unlock();
415
416 // Handle failed requests by calling their callback with a 400 status.
417 if (!failed_requests.empty()) {
418 for (const auto &context : failed_requests) {
419# if __cplusplus >= 201402L
420 auto response = std::make_unique<HttpResponse>();
421# else
422 auto response = std::unique_ptr<HttpResponse>(new HttpResponse());
423# endif
424 const long BAD_REQUEST = 400;
425 response->error_code = utils::make_error_code(CURLE_OK);
426 response->status_code = BAD_REQUEST;
427 response->ready = true;
428 context->callback(std::move(response));
429 context->complete();
430 }
431 failed_requests.clear();
432 }
433
434 // If there are ready requests, create a new HttpBatchRequestHandler to manage them.
435 if (pending_request.empty()) return;
436 try {
437# if __cplusplus >= 201402L
438 m_active_request_batches.push_back(std::make_unique<HttpBatchRequestHandler>(pending_request));
439# else
440 m_active_request_batches.push_back(std::unique_ptr<HttpBatchRequestHandler>(new HttpBatchRequestHandler(pending_request)));
441# endif
442 } catch (...) {
443 for (auto& context : pending_request) {
444 if (!context || !context->callback) continue;
445# if __cplusplus >= 201402L
446 auto response = std::make_unique<HttpResponse>();
447# else
448 auto response = std::unique_ptr<HttpResponse>(new HttpResponse());
449# endif
451 response->status_code = 499; // Client closed request
452 response->ready = true;
453 context->callback(std::move(response));
454 context->complete();
455 }
456 return;
457 }
458 }
459
462 auto it = m_active_request_batches.begin();
463 while (it != m_active_request_batches.end()) {
464 auto& request = *it;
465 if (!request->process()) {
466 ++it;
467 continue;
468 }
469 auto failed_requests = request->extract_failed_requests();
470 it = m_active_request_batches.erase(it);
471 for (auto& request : failed_requests) {
472 m_failed_requests.push_back(std::move(request));
473 }
474 }
475 }
476
479 auto it = m_failed_requests.begin();
480 while (it != m_failed_requests.end()) {
481 auto& request_context = *it;
482 if (!request_context || !request_context->request) {
483 it = m_failed_requests.erase(it);
484 continue;
485 }
486
487 const auto now = std::chrono::steady_clock::now();
488 const auto duration = std::chrono::duration_cast<std::chrono::milliseconds>(now - request_context->start_time);
489 const auto& retry_delay_ms = request_context->request->retry_delay_ms;
490 if (duration.count() >= retry_delay_ms) {
491 std::unique_lock<std::mutex> lock(m_mutex);
492 m_pending_requests.push_back(std::move(request_context));
493 lock.unlock();
494 it = m_failed_requests.erase(it);
495 continue;
496 }
497 ++it;
498 }
499 }
500
503 std::unique_lock<std::mutex> lock(m_mutex);
504 if (m_requests_to_cancel_by_id.empty() && m_groups_to_cancel.empty()) return;
505
506 auto requests_to_cancel = std::move(m_requests_to_cancel_by_id);
507 auto groups_to_cancel = std::move(m_groups_to_cancel);
509 m_groups_to_cancel.clear();
510
511 std::list<std::unique_ptr<HttpRequestContext>> canceled_pending_requests;
512 auto pending_it = m_pending_requests.begin();
513 while (pending_it != m_pending_requests.end()) {
514 const auto& ctx = *pending_it;
515 if (!matches_cancel(ctx, requests_to_cancel, groups_to_cancel)) {
516 ++pending_it;
517 continue;
518 }
519 canceled_pending_requests.push_back(std::move(*pending_it));
520 pending_it = m_pending_requests.erase(pending_it);
521 }
522 lock.unlock();
523
524 for (const auto& request_context : canceled_pending_requests) {
525 request_context->callback(make_cancelled_response());
526 request_context->complete();
527 }
528
529 auto failed_it = m_failed_requests.begin();
530 while (failed_it != m_failed_requests.end()) {
531 const auto& request_context = *failed_it;
532 if (!matches_cancel(request_context, requests_to_cancel, groups_to_cancel)) {
533 ++failed_it;
534 continue;
535 }
536 request_context->callback(make_cancelled_response());
537 request_context->complete();
538 failed_it = m_failed_requests.erase(failed_it);
539 }
540
541 for (const auto &handler : m_active_request_batches) {
542 handler->cancel_requests(requests_to_cancel, groups_to_cancel);
543 }
544
545 invoke_cancel_callbacks(requests_to_cancel);
546 invoke_cancel_callbacks(groups_to_cancel);
547 }
548
549 static bool matches_cancel(
550 const std::unique_ptr<HttpRequestContext>& ctx,
551 const cancel_map_t& requests_to_cancel,
552 const cancel_map_t& groups_to_cancel) {
553 if (!ctx || !ctx->request) return false;
554 const auto request_id = ctx->request->request_id;
555 const auto group_id = ctx->request->group_id;
556 return
557 (request_id != 0 && requests_to_cancel.count(request_id) > 0) ||
558 (group_id != 0 && groups_to_cancel.count(group_id) > 0);
559 }
560
562# if __cplusplus >= 201402L
563 auto response = std::make_unique<HttpResponse>();
564# else
565 auto response = std::unique_ptr<HttpResponse>(new HttpResponse());
566# endif
568 response->status_code = 499;
569 response->ready = true;
570 return response;
571 }
572
573 static void invoke_cancel_callbacks(const cancel_map_t& requests_to_cancel) {
574 for (const auto& request : requests_to_cancel) {
575 for (const auto& callback : request.second) {
576 if (callback) callback();
577 }
578 }
579 }
580
583 std::unique_lock<std::mutex> lock(m_mutex);
584 auto pending_requests = std::move(m_pending_requests);
585 auto failed_requests = std::move(m_failed_requests);
586 m_pending_requests.clear();
587 m_failed_requests.clear();
588 lock.unlock();
589
590 for (const auto &request_context : pending_requests) {
591# if __cplusplus >= 201402L
592 auto response = std::make_unique<HttpResponse>();
593# else
594 auto response = std::unique_ptr<HttpResponse>(new HttpResponse());
595# endif
596 const long CANCELED_REQUEST_CODE = 499;
597 response->error_code = utils::make_error_code(CURLE_OK);
598 response->status_code = CANCELED_REQUEST_CODE;
599 response->ready = true;
600 request_context->callback(std::move(response));
601 request_context->complete();
602 }
603 for (const auto &request_context : failed_requests) {
604# if __cplusplus >= 201402L
605 auto response = std::make_unique<HttpResponse>();
606# else
607 auto response = std::unique_ptr<HttpResponse>(new HttpResponse());
608# endif
609 const long CANCELED_REQUEST_CODE = 499;
610 response->error_code = utils::make_error_code(CURLE_OK);
611 response->status_code = CANCELED_REQUEST_CODE;
612 response->ready = true;
613 request_context->callback(std::move(response));
614 request_context->complete();
615 }
616 }
617
620 curl_global_init(CURL_GLOBAL_ALL);
621 }
622
625 curl_global_cleanup();
626 }
627
630
633
634 }; // HttpRequestManager
635
636}; // namespace kurlyk
637
638#endif // KURLYK_HEADER_KURLYK_HTTP_HTTP_REQUEST_MANAGER_HPP_INCLUDED
Manages multiple asynchronous HTTP requests using libcurl's multi interface.
Defines the RateLimitDelay result type for time-until-allowed queries.
Defines RAII handle for keeping HTTP rate limits alive.
Defines the HttpRateLimiter class for managing rate limits on HTTP requests.
Defines the HttpRequestContext class for managing HTTP request context, including retries and timing.
Handles multiple asynchronous HTTP requests using libcurl's multi interface.
Manages rate limits for HTTP requests.
Represents the context of an HTTP request, including the request object, callback function,...
HttpRequestManager(const HttpRequestManager &)=delete
Deleted copy constructor to enforce the singleton pattern.
cancel_map_t m_group_waiters
Map of group IDs to callbacks waiting until a group becomes idle.
std::size_t max_pending_requests() const
Returns the current maximum pending request count.
std::atomic< bool > m_shutdown
Flag indicating if shutdown has been requested.
static bool matches_cancel(const std::unique_ptr< HttpRequestContext > &ctx, const cancel_map_t &requests_to_cancel, const cancel_map_t &groups_to_cancel)
uint64_t generate_group_id()
Generates a new group ID.
bool has_requests_by_group_id(uint64_t group_id) const
Checks whether pending, failed, or active requests exist for a group.
void set_max_pending_requests(std::size_t max_pending_requests)
Sets the maximum number of pending requests accepted into the global queue.
cancel_map_t m_groups_to_cancel
Map of group IDs to their associated cancellation callbacks.
std::list< std::unique_ptr< HttpRequestContext > > m_failed_requests
List of failed HTTP requests for retrying.
void process_active_requests()
Processes active requests, moving failed ones to the failed requests list for retrying.
std::atomic< uint64_t > m_group_id_counter
Atomic counter for group IDs.
bool remove_limit(const HttpRateLimitHandlePtr &limit)
Releases manager-owned handle for the specified rate-limit handle.
std::size_t group_request_count(uint64_t group_id) const
Counts pending, failed, and active requests for a group.
void cancel_request_by_id(uint64_t request_id, std::function< void()> callback)
Cancels one request by request ID.
std::atomic< uint64_t > m_next_in_flight_token
Atomic counter for sequential rate-limit tokens.
static HttpRequestManager & get_instance()
Get the singleton instance of HttpRequestManager.
std::list< std::unique_ptr< HttpBatchRequestHandler > > m_active_request_batches
List of currently active HTTP request batches.
bool is_loaded() const override
Checks if there are active, pending, or failed requests.
std::atomic< uint64_t > m_request_id_counter
Atomic counter for unique request IDs.
void process_retry_failed_requests()
Attempts to retry failed requests if their retry delay has passed.
HttpRateLimitHandlePtr get_rate_limit(long limit_id)
Returns a handle for a registered rate limit ID.
std::atomic< std::size_t > m_max_pending_requests
Maximum number of requests accepted into the pending queue, or zero if unbounded.
void wait_requests_by_group_id(uint64_t group_id, std::function< void()> callback)
Registers a callback invoked after all requests from a group finish.
virtual ~HttpRequestManager()
Private destructor to clean up global resources.
RateLimitDelay< Duration > time_until_next_allowed(const HttpRateLimitHandlePtr &general_limit, const HttpRateLimitHandlePtr &specific_limit, const std::string &general_key, const std::string &specific_key)
Calculates delay until request is allowed by two handles.
SubmitResult submit_request(std::unique_ptr< HttpRequest > request_ptr, HttpResponseCallback callback)
Attempts to enqueue a new HTTP request and reports the admission result.
HttpRateLimiter m_rate_limiter
Rate limiter for controlling request frequency.
std::list< std::function< void()> > callback_list_t
std::mutex m_mutex
Mutex to protect access to manager state shared with public entry points.
std::size_t group_request_count_unlocked(uint64_t group_id) const
void process_cancel_requests()
Processes and cancels HTTP requests based on their request or group IDs.
bool remove_limit(long limit_id)
Removes an existing rate limit with the specified identifier.
void cancel_requests_by_group_id(uint64_t group_id, std::function< void()> callback)
Cancels all requests with the specified group ID.
bool allow_request(const HttpRateLimitHandlePtr &general_limit, const HttpRateLimitHandlePtr &specific_limit, uint64_t in_flight_token, const std::string &general_key, const std::string &specific_key)
Checks if a request is allowed by two optional rate-limit handles.
cancel_map_t m_requests_to_cancel_by_id
Map of request IDs to their associated cancellation callbacks.
static HttpResponsePtr make_cancelled_response()
void release_request(const HttpRateLimitHandlePtr &general_limit, const HttpRateLimitHandlePtr &specific_limit, uint64_t in_flight_token, const std::string &general_key, const std::string &specific_key)
Releases in-flight tokens for sequential rate limits.
std::unordered_map< uint64_t, callback_list_t > cancel_map_t
uint64_t generate_request_id()
Generates a new unique request ID.
bool add_request(std::unique_ptr< HttpRequest > request_ptr, HttpResponseCallback callback)
Adds a new HTTP request to the manager.
void shutdown() override
Shuts down the request manager, clearing all active and pending requests.
HttpRequestManager & operator=(const HttpRequestManager &)=delete
Deleted copy assignment operator to enforce the singleton pattern.
static void invoke_cancel_callbacks(const cancel_map_t &requests_to_cancel)
HttpRequestManager()
Private constructor to initialize global resources (e.g., cURL).
void cleanup_pending_requests()
Cleans up pending requests, marking each as failed and invoking its callback.
HttpRateLimitHandlePtr create_rate_limit(long requests_per_period, long period_ms, bool sequential=false)
Creates a rate-limit handle with specified parameters.
void process() override
Processes all requests in the manager.
void process_pending_requests()
Processes all pending requests, moving valid requests to active batches or marking them as failed.
std::list< std::unique_ptr< HttpRequestContext > > m_pending_requests
List of pending HTTP requests awaiting processing.
Represents an HTTP response.
Interface for modules managed by NetworkWorker (e.g., HTTP, WebSocket).
@ ShuttingDown
Operation was rejected because the owning subsystem is shutting down.
@ AbortedDuringDestruction
Request handler was destroyed before completion, causing the request to abort.
@ QueueLimitExceeded
Operation was rejected because the bounded queue is already full.
@ CancelledByUser
Request was cancelled explicitly by the user via cancel().
std::error_code make_error_code(ClientError e)
Creates a std::error_code from a ClientError value.
Primary namespace for the Kurlyk library, encompassing initialization, request management,...
std::function< void(HttpResponsePtr response)> HttpResponseCallback
Callback invoked with an HTTP response.
std::unique_ptr< HttpResponse > HttpResponsePtr
Owning pointer to an HTTP response.
std::shared_ptr< HttpRateLimitHandle > HttpRateLimitHandlePtr
Shared RAII handle for HTTP rate limits.
Result type for time-until-allowed queries.
Represents the synchronous result of trying to enqueue or submit work.
bool accepted
Indicates whether the work item was accepted for processing.