cofetch
Chainable, high-performance async HTTP client for C++ event loops
Loading...
Searching...
No Matches
cofetch.h
Go to the documentation of this file.
1#pragma once
2// cofetch: async HTTP client on top of libcurl's multi interface and ASIO.
3//
4// One implementation, any ASIO completion token. C++17 and up; the
5// co_await interface additionally needs C++20.
6//
7// // fluent chain, finished by the HTTP verb (zero-overhead hot path):
8// http.request(url)
9// .headers({"content-type: application/json"})
10// .body(payload)
11// .post([](std::error_code ec, const cofetch::Response& res) {});
12//
13// // std::future:
14// auto fut = http.async_get(url, asio::use_future);
15//
16// // .then()-style chaining (see test/cofetch_tests.cpp):
17// http.async_get(url, asio::deferred)(asio::deferred(next))(handler);
18//
19// // C++20 coroutine:
20// auto res = co_await http.async_get(url, asio::use_awaitable);
21//
22// Drive it with io_context::run(), or io_context::poll() in a busy loop.
23// Requests are cancellable through asio's cancellation slots
24// (asio::cancel_after, asio::bind_cancellation_slot, co_spawn).
25// Not thread-safe: run the client and its io_context on one thread.
26// Define COFETCH_USE_BOOST_ASIO to build on Boost.Asio instead of
27// standalone asio. The API is identical; cofetch::error_code follows the
28// flavor (std::error_code, or boost::system::error_code) — curl category.
29#if defined(COFETCH_USE_BOOST_ASIO)
30#include <boost/asio.hpp>
31#else
32#include <asio.hpp>
33#endif
34//
35#include <curl/curl.h>
36
37#include <chrono>
38#include <cstdint>
39#include <functional>
40#include <list>
41#include <memory>
42#include <mutex>
43#include <optional>
44#include <string>
45#include <string_view>
46#include <system_error>
47#include <unordered_map>
48#include <utility>
49#include <vector>
50
51#if !defined(_WIN32)
52#include <fcntl.h> // F_DUPFD_CLOEXEC: watching descriptors curl owns
53#include <unistd.h> // ::close
54#endif
55
56namespace cofetch {
57
58#if defined(COFETCH_USE_BOOST_ASIO)
59namespace net = boost::asio;
60// The error type follows the asio flavor so completion tokens recognise
61// it (Boost.Asio only unwraps boost::system::error_code).
62using error_code = boost::system::error_code;
63using error_category = boost::system::error_category;
64#else
65namespace net = asio;
66using error_code = std::error_code;
67using error_category = std::error_category;
68#endif
69
74 class category final : public error_category {
75 public:
76 const char* name() const noexcept override { return "curl"; }
77 std::string message(int ev) const override {
78 return curl_easy_strerror(static_cast<CURLcode>(ev));
79 }
80 };
81 static category instance;
82 return instance;
83}
84
85inline error_code make_error_code(CURLcode code) {
86 return {static_cast<int>(code), curl_category()};
87}
88
89namespace detail {
90
91// HTTP field names are case-insensitive (RFC 9110 §5.1), so the header map
92// hashes and compares them without regard to case. ASCII-only folding — field
93// names are ASCII tokens — which also sidesteps std::tolower's locale and
94// signed-char pitfalls. Both are transparent, so on C++20 the map accepts
95// std::string_view lookups without building a temporary std::string.
96inline char ascii_lower(char c) {
97 return (c >= 'A' && c <= 'Z') ? static_cast<char>(c - 'A' + 'a') : c;
98}
99struct CiHash {
100 using is_transparent = void;
101 size_t operator()(std::string_view s) const noexcept {
102 size_t h = 0;
103 for (char c : s) h = h * 31 + static_cast<unsigned char>(ascii_lower(c));
104 return h;
105 }
106};
107struct CiEqual {
108 using is_transparent = void;
109 bool operator()(std::string_view a, std::string_view b) const noexcept {
110 if (a.size() != b.size()) return false;
111 for (size_t i = 0; i < a.size(); ++i) {
112 if (ascii_lower(a[i]) != ascii_lower(b[i])) return false;
113 }
114 return true;
115 }
116};
117
118} // namespace detail
119
120class Response {
121 public:
122 // Case-insensitive field-name -> value map (see headers()).
123 using Headers = std::unordered_map<std::string, std::string, detail::CiHash,
125
126 Response() = default;
127 Response(CURLcode curl_code, long http_code, std::string data,
128 std::string header_data)
129 : curl_code_(curl_code),
130 http_code_(http_code),
131 data_(std::move(data)),
132 header_data_(std::move(header_data)) {}
133
137 bool is_ok() const {
138 return curl_code_ == CURLE_OK && http_code_ >= 200 && http_code_ < 300;
139 }
140
145 const char* error() const { return curl_easy_strerror(curl_code_); }
146
156 Headers headers() const {
157 Headers out;
158 for_each_field([&](std::string_view name, std::string_view value) {
159 const auto [it, inserted] = out.try_emplace(std::string(name), value);
160 if (!inserted) {
161 it->second += ", ";
162 it->second += value;
163 }
164 });
165 return out;
166 }
167
173 std::optional<std::string> header(std::string_view name) const {
174 std::optional<std::string> found;
175 for_each_field([&](std::string_view field, std::string_view value) {
176 if (!detail::CiEqual{}(field, name)) return;
177 if (found) {
178 *found += ", ";
179 *found += value;
180 } else {
181 found = std::string(value);
182 }
183 });
184 return found;
185 }
186
187 CURLcode curl_code_ = CURLE_OK;
188 long http_code_ = 0;
189 std::string data_;
190 std::string header_data_;
191
192 private:
193 // Iterate the "name: value" fields of the last response block in
194 // header_data_, trimmed, skipping the status line and blank separators.
195 // The views point into header_data_, so they are valid only for the
196 // lifetime of this Response.
197 template <typename F>
198 void for_each_field(F&& f) const {
199 std::string_view sv(header_data_);
200 // Each hop's block opens with a status line. Field names are tokens and
201 // cannot contain '/', so a line starting with "HTTP/" is only ever a
202 // status line: the last one marks where the final response begins.
203 const size_t last_status = sv.rfind("\nHTTP/");
204 size_t pos = (last_status == std::string_view::npos) ? 0 : last_status + 1;
205 while (pos < sv.size()) {
206 const size_t nl = sv.find('\n', pos);
207 std::string_view line =
208 sv.substr(pos, (nl == std::string_view::npos ? sv.size() : nl) - pos);
209 pos = (nl == std::string_view::npos) ? sv.size() : nl + 1;
210 if (!line.empty() && line.back() == '\r') line.remove_suffix(1);
211 const size_t colon = line.find(':');
212 if (colon == std::string_view::npos) continue; // status line or blank
213 f(trim(line.substr(0, colon)), trim(line.substr(colon + 1)));
214 }
215 }
216
217 // Strip leading/trailing HTTP optional whitespace (space and htab).
218 static std::string_view trim(std::string_view s) {
219 const auto b = s.find_first_not_of(" \t");
220 if (b == std::string_view::npos) return {};
221 return s.substr(b, s.find_last_not_of(" \t") - b + 1);
222 }
223};
224
228class Request {
229 public:
230 enum class Method { GET, POST, PUT, PATCH, DEL };
231
232 explicit Request(std::string url) : url_(std::move(url)) {}
233
235 method_ = m;
236 return *this;
237 }
238 Request& headers(std::vector<std::string> h) {
239 headers_ = std::move(h);
240 return *this;
241 }
242 Request& body(std::string b) {
243 body_ = std::move(b);
244 return *this;
245 }
251 Request& timeout(std::chrono::milliseconds t) {
252 timeout_ = t;
253 return *this;
254 }
259 Request& follow_redirects(long max = 30) {
260 max_redirects_ = max;
261 return *this;
262 }
269 Request& curl(std::function<void(CURL*)> fn) {
270 curl_setup_ = std::move(fn);
271 return *this;
272 }
273
274 std::string url_;
276 std::vector<std::string> headers_;
277 std::string body_;
278 std::chrono::milliseconds timeout_{std::chrono::seconds{5}};
279 long max_redirects_ = 0; // 0: do not follow redirects
280 std::function<void(CURL*)> curl_setup_;
281};
282
283class Client {
284 public:
290 explicit Client(net::io_context& io,
291 size_t max_pooled_connections = kDefaultMaxPooledHandles)
292 : io_(io), timer_(io), max_pooled_handles_(max_pooled_connections) {
293 {
294 // Reference counted inside libcurl, but only thread-safe from 7.84
295 // on: serialise clients constructed on different threads at once
296 // (one event loop per core).
297 const std::lock_guard<std::mutex> lock(global_mutex());
298 curl_global_init(CURL_GLOBAL_ALL);
299 }
300 multi_ = curl_multi_init();
301 curl_multi_setopt(multi_, CURLMOPT_SOCKETFUNCTION, socket_cb);
302 curl_multi_setopt(multi_, CURLMOPT_SOCKETDATA, this);
303 curl_multi_setopt(multi_, CURLMOPT_TIMERFUNCTION, timer_cb);
304 curl_multi_setopt(multi_, CURLMOPT_TIMERDATA, this);
305 }
306
307 Client(const Client&) = delete;
308 Client& operator=(const Client&) = delete;
309
316 *alive_ = false;
317 timer_.cancel();
318 // Sever every socket first: pending waits complete with
319 // operation_aborted (and find no state), and curl's shutdown attempts
320 // below fail at once instead of waiting on peers.
321 for (const auto& [fd, state] : sockets_) {
322 state->watch = 0;
323 error_code ignored;
324 state->socket.close(ignored);
325 }
326 // Transfers still in flight: detach them so their handles and header
327 // lists can be freed. The handlers go down with transfers_, uninvoked.
328 for (const auto& t : transfers_) {
329 curl_multi_remove_handle(multi_, t.eh);
330 curl_slist_free_all(t.headers);
331 curl_easy_cleanup(t.eh);
332 }
333 for (const auto& t : idle_) curl_easy_cleanup(t.eh);
334 curl_multi_cleanup(multi_); // closes cached connections (close_socket_cb)
335 const std::lock_guard<std::mutex> lock(global_mutex());
336 curl_global_cleanup();
337 }
338
350 template <typename CompletionToken>
351 auto async_perform(Request req, CompletionToken&& token) {
352 return net::async_initiate<CompletionToken, void(error_code, Response)>(
353 Initiation{this}, token, std::move(req));
354 }
355
356 template <typename CompletionToken>
357 auto async_get(std::string url, CompletionToken&& token) {
358 return async_perform(Request(std::move(url)),
359 std::forward<CompletionToken>(token));
360 }
361
362 template <typename CompletionToken>
363 auto async_post(std::string url, std::string body, CompletionToken&& token) {
364 return async_perform(Request(std::move(url))
365 .method(Request::Method::POST)
366 .body(std::move(body)),
367 std::forward<CompletionToken>(token));
368 }
369
379 public:
380 RequestBuilder(Client& client, std::string url)
381 : client_(client), req_(std::move(url)) {}
382
383 RequestBuilder& headers(std::vector<std::string> h) {
384 req_.headers(std::move(h));
385 return *this;
386 }
387 RequestBuilder& body(std::string b) {
388 req_.body(std::move(b));
389 return *this;
390 }
391 RequestBuilder& timeout(std::chrono::milliseconds t) {
392 req_.timeout(t);
393 return *this;
394 }
396 req_.follow_redirects(max);
397 return *this;
398 }
399 RequestBuilder& curl(std::function<void(CURL*)> fn) {
400 req_.curl(std::move(fn));
401 return *this;
402 }
403
404 template <typename CompletionToken>
405 auto get(CompletionToken&& token) {
406 return perform(Request::Method::GET,
407 std::forward<CompletionToken>(token));
408 }
409 template <typename CompletionToken>
410 auto post(CompletionToken&& token) {
411 return perform(Request::Method::POST,
412 std::forward<CompletionToken>(token));
413 }
414 template <typename CompletionToken>
415 auto put(CompletionToken&& token) {
416 return perform(Request::Method::PUT,
417 std::forward<CompletionToken>(token));
418 }
419 template <typename CompletionToken>
420 auto patch(CompletionToken&& token) {
421 return perform(Request::Method::PATCH,
422 std::forward<CompletionToken>(token));
423 }
424 template <typename CompletionToken>
425 auto del(CompletionToken&& token) {
426 return perform(Request::Method::DEL,
427 std::forward<CompletionToken>(token));
428 }
429
430 private:
431 template <typename CompletionToken>
432 auto perform(Request::Method m, CompletionToken&& token) {
433 req_.method(m);
434 return client_.async_perform(std::move(req_),
435 std::forward<CompletionToken>(token));
436 }
437
438 Client& client_;
439 Request req_;
440 };
441
445 RequestBuilder request(std::string url) {
446 return RequestBuilder(*this, std::move(url));
447 }
448
449 int pending_requests() const { return running_; }
450
451 private:
452 using Handler = net::any_completion_handler<void(error_code, Response)>;
453
454 // Initiation for async_perform. Exposing the io_context executor lets
455 // executor-aware tokens work — asio::cancel_after, for one, builds its
456 // timeout timer on Initiation::executor_type.
457 struct Initiation {
458 Client* self;
459 using executor_type = net::io_context::executor_type;
460 executor_type get_executor() const noexcept {
461 return self->io_.get_executor();
462 }
463 template <typename H>
464 void operator()(H handler, Request r) const {
465 self->start(std::move(r), Handler(std::move(handler)));
466 }
467 };
468
469 // Default cap on idle easy handles kept for reuse; beyond this they are
470 // freed so a burst of concurrent requests does not pin memory forever.
471 // Overridable per client through the constructor.
472 static constexpr size_t kDefaultMaxPooledHandles = 64;
473
474 // Pre-size for the response header block: curl delivers it one line per
475 // callback, and growing a std::string from empty takes several
476 // reallocations to reach the few hundred bytes a typical block needs.
477 static constexpr size_t kHeaderReserve = 1024;
478
479 // One transfer's state. Nodes live in transfers_ while in flight and
480 // move (spliced, never reallocated) to idle_ afterwards, together with
481 // their configured easy handle, so a request costs neither a node nor a
482 // handle allocation once the client is warm.
483 struct Transfer {
484 explicit Transfer(CURL* e) : eh(e) {}
485 CURL* eh;
486 curl_slist* headers = nullptr;
487 std::string body;
488 std::string buffer;
489 std::string header_buffer;
490 Handler handler;
491 std::list<Transfer>::iterator self;
492 std::uint64_t id = 0;
493 // Set when Request::curl ran on this handle: unknown options must be
494 // wiped (curl_easy_reset + configure_handle) before the handle is
495 // pooled, or they would leak into whatever request draws it next.
496 bool scrub_on_done = false;
497 };
498
499 struct SocketState : std::enable_shared_from_this<SocketState> {
500 SocketState(net::io_context& io, bool borrowed_fd)
501 : socket(io), borrowed(borrowed_fd) {}
502 net::ip::tcp::socket socket;
503 // The descriptor as curl knows it. For our own sockets that is
504 // socket.native_handle(); for borrowed ones socket holds a dup.
505 curl_socket_t fd = CURL_SOCKET_BAD;
506 const bool borrowed;
507 int watch = 0; // current CURL_POLL_* interest
508 bool read_armed = false;
509 bool write_armed = false;
510 };
511
512 void start(Request r, Handler h) {
513 if (idle_.empty()) {
514 CURL* const eh = curl_easy_init();
515 if (eh == nullptr) return fail(std::move(h), CURLE_FAILED_INIT);
516 configure_handle(eh);
517 transfers_.emplace_front(eh);
518 } else {
519 // Idle nodes keep their configured handle: no allocation, no reset.
520 transfers_.splice(transfers_.begin(), idle_, idle_.begin());
521 }
522 const auto it = transfers_.begin();
523 Transfer& t = *it;
524 CURL* const eh = t.eh;
525 t.self = it;
526 t.handler = std::move(h);
527 t.body = std::move(r.body_);
528 t.id = ++next_transfer_id_;
529
530 curl_easy_setopt(eh, CURLOPT_URL, r.url_.c_str());
531 curl_easy_setopt(eh, CURLOPT_WRITEDATA, &t);
532 curl_easy_setopt(eh, CURLOPT_HEADERDATA, &t.header_buffer);
533 curl_easy_setopt(eh, CURLOPT_PRIVATE, &t);
534 curl_easy_setopt(eh, CURLOPT_TIMEOUT_MS,
535 static_cast<long>(r.timeout_.count()));
536 // Always set: clears the previous transfer's values on pooled handles.
537 curl_easy_setopt(eh, CURLOPT_FOLLOWLOCATION,
538 r.max_redirects_ != 0 ? 1L : 0L);
539 curl_easy_setopt(eh, CURLOPT_MAXREDIRS, r.max_redirects_);
540
541 curl_slist* chunk = nullptr;
542 for (const auto& header : r.headers_) {
543 chunk = curl_slist_append(chunk, header.c_str());
544 }
545 // Always set: clears the previous transfer's list on pooled handles.
546 curl_easy_setopt(eh, CURLOPT_HTTPHEADER, chunk);
547 t.headers = chunk;
548
549 switch (r.method_) {
551 curl_easy_setopt(eh, CURLOPT_CUSTOMREQUEST, nullptr);
552 // Drop the previous transfer's body pointer (freed with it) so a
553 // Request::curl hook turning this into an upload cannot read it.
554 // POSTFIELDS selects POST, so it goes before HTTPGET.
555 curl_easy_setopt(eh, CURLOPT_POSTFIELDS, nullptr);
556 curl_easy_setopt(eh, CURLOPT_HTTPGET, 1L);
557 break;
559 curl_easy_setopt(eh, CURLOPT_CUSTOMREQUEST, nullptr);
560 curl_easy_setopt(eh, CURLOPT_POST, 1L);
561 set_body(t);
562 break;
564 curl_easy_setopt(eh, CURLOPT_CUSTOMREQUEST, "PUT");
565 curl_easy_setopt(eh, CURLOPT_POST, 1L);
566 set_body(t);
567 break;
569 curl_easy_setopt(eh, CURLOPT_CUSTOMREQUEST, "PATCH");
570 curl_easy_setopt(eh, CURLOPT_POST, 1L);
571 set_body(t);
572 break;
574 curl_easy_setopt(eh, CURLOPT_CUSTOMREQUEST, "DELETE");
575 curl_easy_setopt(eh, CURLOPT_POST, 1L);
576 set_body(t);
577 break;
578 }
579
580 // The escape hatch runs last so it can override anything above.
581 if (r.curl_setup_) {
582 t.scrub_on_done = true;
583 r.curl_setup_(eh);
584 }
585
586 auto slot = net::get_associated_cancellation_slot(t.handler);
587 if (slot.is_connected()) {
588 // Look the transfer up by id at emit time: the slot outlives the
589 // transfer (asio only guarantees clearing on handler destruction),
590 // so a late emit must find nothing rather than follow a dangling
591 // pointer. alive_ covers emits after ~Client.
592 slot.assign([this, alive = alive_, id = t.id](net::cancellation_type_t) {
593 if (*alive) cancel(id);
594 });
595 }
596
597 // curl schedules the kickstart itself through the timer callback.
598 const CURLMcode rc = curl_multi_add_handle(multi_, eh);
599 if (rc != CURLM_OK) {
600 // Out of memory, realistically. Never leave a handler uncalled.
601 Handler handler = std::move(t.handler);
602 retire(t, /*scrub=*/true);
603 fail(std::move(handler),
604 rc == CURLM_OUT_OF_MEMORY ? CURLE_OUT_OF_MEMORY : CURLE_FAILED_INIT);
605 }
606 }
607
608 // Complete h with a transport error. Posted, not invoked: asio forbids
609 // running a completion handler from inside the initiating function.
610 void fail(Handler h, CURLcode code) {
611 net::post(io_, net::bind_allocator(net::recycling_allocator<void>(),
612 [h = std::move(h), code]() mutable {
613 std::move(h)(
614 make_error_code(code),
615 Response{code, 0, {}, {}});
616 }));
617 }
618
619 // Cooperative cancellation, reached through the completion handler's
620 // associated cancellation slot. Any cancellation type aborts: the
621 // transfer is torn down and the handler completes with
622 // operation_aborted. No-op when the transfer already completed.
623 void cancel(std::uint64_t id) {
624 for (auto& t : transfers_) {
625 if (t.id != id) continue;
626 // Also discards any DONE message this handle queued in the multi.
627 curl_multi_remove_handle(multi_, t.eh);
628 if (running_ > 0) --running_;
629 Handler handler = std::move(t.handler);
630 // Severed mid-flight: scrub before the handle is reused.
631 retire(t, /*scrub=*/true);
632 // Unlike normal completions this one is posted, not invoked: we are
633 // inside the cancellation emit, and completing here can destroy the
634 // very signal being emitted (asio::cancel_after owns its signal in
635 // the operation state the completion frees).
636 net::post(io_, net::bind_allocator(
637 net::recycling_allocator<void>(),
638 [h = std::move(handler)]() mutable {
639 std::move(h)(
640 error_code(net::error::operation_aborted),
641 Response{CURLE_ABORTED_BY_CALLBACK, 0, {}, {}});
642 }));
643 return;
644 }
645 }
646
647 // Done with t (handler already moved out, handle no longer in the multi):
648 // return the node and its easy handle to the idle pool for the next
649 // request, or free them when the pool is full. scrub wipes the handle's
650 // options first — needed after Request::curl ran on it, or when the
651 // transfer was cut short.
652 void retire(Transfer& t, bool scrub) {
653 curl_slist_free_all(t.headers);
654 t.headers = nullptr;
655 t.id = 0;
656 t.scrub_on_done = false;
657 t.buffer.clear(); // moved-from: valid but unspecified until cleared
658 t.header_buffer.clear();
659 std::string().swap(t.body); // a large upload must not sit in the pool
660 if (scrub) {
661 curl_easy_reset(t.eh);
662 configure_handle(t.eh);
663 }
664 if (idle_.size() < max_pooled_handles_) {
665 idle_.splice(idle_.begin(), transfers_, t.self);
666 } else {
667 curl_easy_cleanup(t.eh);
668 transfers_.erase(t.self);
669 }
670 }
671
672 // Request-independent options, set once per easy handle. Everything a
673 // transfer can vary must be (re)set in start() — pooled handles are
674 // reused without curl_easy_reset.
675 void configure_handle(CURL* eh) {
676 curl_easy_setopt(eh, CURLOPT_WRITEFUNCTION, body_write_cb);
677 curl_easy_setopt(eh, CURLOPT_HEADERFUNCTION, header_cb);
678 curl_easy_setopt(eh, CURLOPT_NOSIGNAL, 1L);
679#if LIBCURL_VERSION_NUM >= 0x075000 // 7.80.0
680 // Cap connection reuse age; on older libcurl the option is absent and
681 // pooled connections simply live longer.
682 curl_easy_setopt(eh, CURLOPT_MAXLIFETIME_CONN, 30L);
683#endif
684 // "" advertises every decoder curl was built with (gzip, br, ...).
685 curl_easy_setopt(eh, CURLOPT_ACCEPT_ENCODING, "");
686 // Wait for an in-progress connection to the same host and multiplex over
687 // it (HTTP/2) instead of racing to open one connection per request.
688 curl_easy_setopt(eh, CURLOPT_PIPEWAIT, 1L);
689 curl_easy_setopt(eh, CURLOPT_BUFFERSIZE, 512L * 1024L);
690 curl_easy_setopt(eh, CURLOPT_OPENSOCKETFUNCTION, open_socket_cb);
691 curl_easy_setopt(eh, CURLOPT_OPENSOCKETDATA, this);
692 curl_easy_setopt(eh, CURLOPT_CLOSESOCKETFUNCTION, close_socket_cb);
693 curl_easy_setopt(eh, CURLOPT_CLOSESOCKETDATA, this);
694 }
695
696 void set_body(const Transfer& t) {
697 curl_easy_setopt(t.eh, CURLOPT_POSTFIELDSIZE,
698 static_cast<long>(t.body.size()));
699 curl_easy_setopt(t.eh, CURLOPT_POSTFIELDS, t.body.c_str());
700 }
701
702 static size_t header_cb(char* data, size_t n, size_t l, std::string* buf) {
703 if (buf->empty()) buf->reserve(kHeaderReserve); // no-op once sized
704 buf->append(data, n * l);
705 return n * l;
706 }
707
708 static size_t body_write_cb(char* data, size_t n, size_t l, Transfer* t) {
709 if (t->buffer.empty()) {
710 curl_off_t len = 0;
711 if (curl_easy_getinfo(t->eh, CURLINFO_CONTENT_LENGTH_DOWNLOAD_T, &len) ==
712 CURLE_OK &&
713 len > 0) {
714 t->buffer.reserve(static_cast<size_t>(len));
715 }
716 }
717 t->buffer.append(data, n * l);
718 return n * l;
719 }
720
721 // curl asks us (not the OS directly) for sockets, so every fd it uses is
722 // backed by an ASIO object we can async_wait on. Cross-platform, no epoll.
723 static curl_socket_t open_socket_cb(void* clientp, curlsocktype purpose,
724 curl_sockaddr* address) {
725 auto* const self = static_cast<Client*>(clientp);
726 if (purpose != CURLSOCKTYPE_IPCXN) return CURL_SOCKET_BAD;
727 net::ip::tcp protocol = net::ip::tcp::v4();
728 if (address->family == AF_INET6) {
729 protocol = net::ip::tcp::v6();
730 } else if (address->family != AF_INET) {
731 return CURL_SOCKET_BAD;
732 }
733 auto state =
734 std::make_shared<SocketState>(self->io_, /*borrowed_fd=*/false);
735 error_code ec;
736 state->socket.open(protocol, ec);
737 if (ec) return CURL_SOCKET_BAD;
738 const curl_socket_t fd = state->socket.native_handle();
739 state->fd = fd;
740 self->sockets_[fd] = std::move(state);
741 return fd;
742 }
743
744 static int close_socket_cb(void* clientp, curl_socket_t fd) {
745 auto* const self = static_cast<Client*>(clientp);
746 const auto it = self->sockets_.find(fd);
747 if (it == self->sockets_.end()) {
748 // Opened past open_socket_cb (a Request::curl hook replaced it):
749 // still ours to close, as the callback contract says.
750 close_native(fd);
751 return 0;
752 }
753 SocketState& state = *it->second;
754 state.watch = 0;
755 error_code ignored;
756 state.socket.close(ignored);
757 if (state.borrowed) close_native(fd); // we only closed our dup
758 self->sockets_.erase(it);
759 return 0;
760 }
761
762 static void close_native(curl_socket_t fd) {
763#if defined(_WIN32)
764 ::closesocket(fd);
765#else
766 ::close(fd);
767#endif
768 }
769
770 // curl also hands us descriptors it created itself: the threaded
771 // resolver's wake-up pipe, c-ares' UDP sockets. Watch a dup of those,
772 // so the lifetime of curl's own descriptor stays curl's business (it may
773 // close it any time after CURL_POLL_REMOVE, and the number may be
774 // reused). Ignoring them keeps things working — curl's timer polls the
775 // resolver at 1 ms doubling up to 250 ms — but every hostname lookup is
776 // then noticed late, up to twice its real duration.
777 SocketState* borrow(curl_socket_t fd) {
778#if defined(_WIN32)
779 (void)fd;
780 return nullptr; // no dup for SOCKETs; curl's polling covers it
781#else
782 const int dup_fd = ::fcntl(fd, F_DUPFD_CLOEXEC, 0);
783 if (dup_fd < 0) return nullptr;
784 auto state = std::make_shared<SocketState>(io_, /*borrowed_fd=*/true);
785 error_code ec;
786 // The protocol is a label here (the descriptor may be a socketpair end
787 // or an eventfd); only the reactor registration matters to async_wait.
788 state->socket.assign(net::ip::tcp::v4(), dup_fd, ec);
789 if (ec) { // LCOV_EXCL_START: reactor registration failing (ENOMEM)
790 ::close(dup_fd);
791 return nullptr;
792 } // LCOV_EXCL_STOP
793 state->fd = fd;
794 SocketState* const raw = state.get();
795 sockets_[fd] = std::move(state);
796 return raw;
797#endif
798 }
799
800 static int socket_cb(CURL*, curl_socket_t fd, int what, void* userp,
801 void* socketp) {
802 auto* const self = static_cast<Client*>(userp);
803 auto* state = static_cast<SocketState*>(socketp);
804 if (state == nullptr) {
805 if (what == CURL_POLL_REMOVE) return 0; // never watched: nothing to do
806 // First notification for this descriptor: attach the state so curl
807 // hands it back on later calls and we skip the lookup.
808 const auto it = self->sockets_.find(fd);
809 state =
810 (it != self->sockets_.end()) ? it->second.get() : self->borrow(fd);
811 if (state == nullptr) return 0;
812 curl_multi_assign(self->multi_, fd, state);
813 }
814 if (what == CURL_POLL_REMOVE) {
815 state->watch = 0;
816 if (state->borrowed) {
817 // curl closes its descriptor right after this; drop our dup now (a
818 // later re-add borrows afresh). Our own sockets stay: the
819 // connection is merely idle in curl's cache.
820 error_code ignored;
821 state->socket.close(ignored);
822 self->sockets_.erase(fd);
823 }
824 return 0;
825 }
826 state->watch = what;
827 self->arm(state->shared_from_this());
828 return 0;
829 }
830
831 void arm(const std::shared_ptr<SocketState>& state) {
832 if (!state->socket.is_open()) return;
833 if ((state->watch & CURL_POLL_IN) && !state->read_armed) {
834 state->read_armed = true;
835 state->socket.async_wait(
836 net::ip::tcp::socket::wait_read,
837 net::bind_allocator(
838 net::recycling_allocator<void>(),
839 [this, w = std::weak_ptr<SocketState>(state)](error_code ec) {
840 // Lock before touching the client: a state outlives
841 // ~Client only inside a queued completion like this one,
842 // and then the lock fails.
843 if (auto s = w.lock())
844 on_event(std::move(s), CURL_CSELECT_IN, ec);
845 }));
846 }
847 if ((state->watch & CURL_POLL_OUT) && !state->write_armed) {
848 state->write_armed = true;
849 state->socket.async_wait(
850 net::ip::tcp::socket::wait_write,
851 net::bind_allocator(
852 net::recycling_allocator<void>(),
853 [this, w = std::weak_ptr<SocketState>(state)](error_code ec) {
854 if (auto s = w.lock())
855 on_event(std::move(s), CURL_CSELECT_OUT, ec);
856 }));
857 }
858 }
859
860 void on_event(std::shared_ptr<SocketState> state, int flag, error_code ec) {
861 // The shared_ptr keeps the state alive across the socket_action call
862 // below, which may close this very socket via close_socket_cb (or drop
863 // a borrowed one on CURL_POLL_REMOVE).
864 (flag == CURL_CSELECT_IN ? state->read_armed : state->write_armed) = false;
865 if (ec == net::error::operation_aborted) return;
866 curl_multi_socket_action(multi_, state->fd, ec ? CURL_CSELECT_ERR : flag,
867 &running_);
868 check_completions();
869 if (state->watch != 0) arm(state);
870 }
871
872 static int timer_cb(CURLM*, long timeout_ms, void* userp) {
873 auto* const self = static_cast<Client*>(userp);
874 if (timeout_ms < 0) {
875 self->timer_.cancel();
876 self->timer_armed_ = false;
877 return 0;
878 }
879 if (timeout_ms == 0) {
880 // "Act as soon as possible" — the common per-transfer kick. A plain
881 // post (deduplicated) is much cheaper than rescheduling the timer,
882 // and we may not call curl back from inside its own callback.
883 if (!self->kick_pending_) {
884 self->kick_pending_ = true;
885 net::post(self->io_,
886 net::bind_allocator(net::recycling_allocator<void>(),
887 [self, alive = self->alive_] {
888 if (!*alive) return;
889 self->kick_pending_ = false;
890 self->kick();
891 }));
892 }
893 return 0;
894 }
895 const auto deadline = std::chrono::steady_clock::now() +
896 std::chrono::milliseconds(timeout_ms);
897 // A pending wait that fires no later than the new deadline is good
898 // enough: a kick() finding nothing due is a cheap no-op, while
899 // rescheduling reprograms the timer every time.
900 if (self->timer_armed_ && self->timer_.expiry() <= deadline) return 0;
901 self->timer_.expires_at(deadline);
902 self->timer_armed_ = true;
903 self->timer_.async_wait(
904 net::bind_allocator(net::recycling_allocator<void>(),
905 [self, alive = self->alive_](error_code ec) {
906 if (ec || !*alive) return;
907 self->timer_armed_ = false;
908 self->kick();
909 }));
910 return 0;
911 }
912
913 void kick() {
914 curl_multi_socket_action(multi_, CURL_SOCKET_TIMEOUT, 0, &running_);
915 check_completions();
916 }
917
918 void check_completions() {
919 int msgs_left = 0;
920 while (CURLMsg* msg = curl_multi_info_read(multi_, &msgs_left)) {
921 if (msg->msg != CURLMSG_DONE) continue;
922 Transfer* t = nullptr;
923 long http_code = 0;
924 // msg must not be dereferenced after curl_multi_remove_handle().
925 const CURLcode curl_code = msg->data.result;
926 CURL* const eh = msg->easy_handle;
927 curl_easy_getinfo(eh, CURLINFO_PRIVATE, &t);
928 curl_easy_getinfo(eh, CURLINFO_RESPONSE_CODE, &http_code);
929 curl_multi_remove_handle(multi_, eh);
930
931 Response res{curl_code, http_code, std::move(t->buffer),
932 std::move(t->header_buffer)};
933 Handler handler = std::move(t->handler);
934 retire(*t, t->scrub_on_done);
935
936 const error_code ec =
937 curl_code == CURLE_OK ? error_code{} : make_error_code(curl_code);
938 // Single-threaded by contract: the handler's executor is this
939 // io_context, where we already are — invoke without the
940 // type-erased dispatch hop.
941 std::move(handler)(ec, std::move(res));
942 }
943 }
944
945 static std::mutex& global_mutex() {
946 static std::mutex m;
947 return m;
948 }
949
950 net::io_context& io_;
951 CURLM* multi_;
952 net::steady_timer timer_;
953 // Keyed by the descriptor curl reports (SocketState::fd).
954 std::unordered_map<curl_socket_t, std::shared_ptr<SocketState>> sockets_;
955 std::list<Transfer> transfers_; // in flight
956 std::list<Transfer> idle_; // warm nodes + handles for reuse (LIFO)
957 const size_t max_pooled_handles_; // cap on idle_; set by the constructor
958 int running_ = 0;
959 std::uint64_t next_transfer_id_ = 0;
960 bool kick_pending_ = false;
961 bool timer_armed_ = false;
962 // Outlives the client inside posted/timed kicks: they bail out when the
963 // client is gone instead of touching a destroyed multi handle.
964 std::shared_ptr<bool> alive_ = std::make_shared<bool>(true);
965};
966
967} // namespace cofetch
Fluent builder bound to this client.
Definition cofetch.h:378
RequestBuilder & headers(std::vector< std::string > h)
Definition cofetch.h:383
RequestBuilder(Client &client, std::string url)
Definition cofetch.h:380
auto patch(CompletionToken &&token)
Definition cofetch.h:420
RequestBuilder & curl(std::function< void(CURL *)> fn)
Definition cofetch.h:399
auto post(CompletionToken &&token)
Definition cofetch.h:410
auto del(CompletionToken &&token)
Definition cofetch.h:425
RequestBuilder & body(std::string b)
Definition cofetch.h:387
RequestBuilder & follow_redirects(long max=30)
Definition cofetch.h:395
auto get(CompletionToken &&token)
Definition cofetch.h:405
RequestBuilder & timeout(std::chrono::milliseconds t)
Definition cofetch.h:391
auto put(CompletionToken &&token)
Definition cofetch.h:415
Definition cofetch.h:283
RequestBuilder request(std::string url)
Start a fluent request chain: request(url).body(...).post(token).
Definition cofetch.h:445
Client & operator=(const Client &)=delete
~Client()
Destroy the client.
Definition cofetch.h:315
auto async_perform(Request req, CompletionToken &&token)
Start a transfer described by req.
Definition cofetch.h:351
Client(net::io_context &io, size_t max_pooled_connections=kDefaultMaxPooledHandles)
Construct a client driven by io.
Definition cofetch.h:290
int pending_requests() const
Definition cofetch.h:449
auto async_post(std::string url, std::string body, CompletionToken &&token)
Definition cofetch.h:363
auto async_get(std::string url, CompletionToken &&token)
Definition cofetch.h:357
Client(const Client &)=delete
Value-type description of a request; pass to Client::async_perform.
Definition cofetch.h:228
Request(std::string url)
Definition cofetch.h:232
Request & headers(std::vector< std::string > h)
Definition cofetch.h:238
long max_redirects_
Definition cofetch.h:279
std::string body_
Definition cofetch.h:277
Method
Definition cofetch.h:230
Request & body(std::string b)
Definition cofetch.h:242
Request & method(Method m)
Definition cofetch.h:234
Request & follow_redirects(long max=30)
Follow HTTP 3xx redirects, at most max hops (the transfer fails with CURLE_TOO_MANY_REDIRECTS beyond ...
Definition cofetch.h:259
std::chrono::milliseconds timeout_
Definition cofetch.h:278
Method method_
Definition cofetch.h:275
std::function< void(CURL *)> curl_setup_
Definition cofetch.h:280
std::vector< std::string > headers_
Definition cofetch.h:276
Request & curl(std::function< void(CURL *)> fn)
Escape hatch: fn runs on the underlying easy handle after cofetch's own options, so it can set (or ov...
Definition cofetch.h:269
std::string url_
Definition cofetch.h:274
Request & timeout(std::chrono::milliseconds t)
Whole-transfer timeout (default 5s), millisecond resolution.
Definition cofetch.h:251
Definition cofetch.h:120
Response(CURLcode curl_code, long http_code, std::string data, std::string header_data)
Definition cofetch.h:127
long http_code_
Definition cofetch.h:188
const char * error() const
Human readable description of the transport error ("No error" when the transfer itself succeeded).
Definition cofetch.h:145
CURLcode curl_code_
Definition cofetch.h:187
Response()=default
Headers headers() const
Parse the response headers into a case-insensitive name->value map.
Definition cofetch.h:156
bool is_ok() const
True when the transfer succeeded and the HTTP status is 2xx.
Definition cofetch.h:137
std::string header_data_
Definition cofetch.h:190
std::optional< std::string > header(std::string_view name) const
Case-insensitive lookup of one field without building the full map; repeats are comma-combined.
Definition cofetch.h:173
std::unordered_map< std::string, std::string, detail::CiHash, detail::CiEqual > Headers
Definition cofetch.h:124
std::string data_
Definition cofetch.h:189
char ascii_lower(char c)
Definition cofetch.h:96
Definition cofetch.h:56
error_code make_error_code(CURLcode code)
Definition cofetch.h:85
const error_category & curl_category()
std::error_category for libcurl transport errors (CURLcode values).
Definition cofetch.h:73
std::error_category error_category
Definition cofetch.h:67
std::error_code error_code
Definition cofetch.h:66
Definition cofetch.h:107
void is_transparent
Definition cofetch.h:108
bool operator()(std::string_view a, std::string_view b) const noexcept
Definition cofetch.h:109
Definition cofetch.h:99
size_t operator()(std::string_view s) const noexcept
Definition cofetch.h:101
void is_transparent
Definition cofetch.h:100