291 size_t max_pooled_connections = kDefaultMaxPooledHandles)
292 : io_(io), timer_(io), max_pooled_handles_(max_pooled_connections) {
297 const std::lock_guard<std::mutex> lock(global_mutex());
298 curl_global_init(CURL_GLOBAL_ALL);
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);
321 for (
const auto& [fd, state] : sockets_) {
324 state->socket.close(ignored);
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);
333 for (
const auto& t : idle_) curl_easy_cleanup(t.eh);
334 curl_multi_cleanup(multi_);
335 const std::lock_guard<std::mutex> lock(global_mutex());
336 curl_global_cleanup();
350 template <
typename CompletionToken>
352 return net::async_initiate<CompletionToken, void(error_code, Response)>(
353 Initiation{
this}, token, std::move(req));
356 template <
typename CompletionToken>
357 auto async_get(std::string url, CompletionToken&& token) {
359 std::forward<CompletionToken>(token));
362 template <
typename CompletionToken>
363 auto async_post(std::string url, std::string body, CompletionToken&& token) {
366 .body(std::move(body)),
367 std::forward<CompletionToken>(token));
381 : client_(client), req_(std::move(url)) {}
388 req_.
body(std::move(b));
400 req_.
curl(std::move(fn));
404 template <
typename CompletionToken>
405 auto get(CompletionToken&& token) {
407 std::forward<CompletionToken>(token));
409 template <
typename CompletionToken>
410 auto post(CompletionToken&& token) {
412 std::forward<CompletionToken>(token));
414 template <
typename CompletionToken>
415 auto put(CompletionToken&& token) {
417 std::forward<CompletionToken>(token));
419 template <
typename CompletionToken>
420 auto patch(CompletionToken&& token) {
422 std::forward<CompletionToken>(token));
424 template <
typename CompletionToken>
425 auto del(CompletionToken&& token) {
427 std::forward<CompletionToken>(token));
431 template <
typename CompletionToken>
435 std::forward<CompletionToken>(token));
459 using executor_type = net::io_context::executor_type;
460 executor_type get_executor() const noexcept {
461 return self->io_.get_executor();
463 template <
typename H>
464 void operator()(H handler, Request r)
const {
465 self->start(std::move(r), Handler(std::move(handler)));
472 static constexpr size_t kDefaultMaxPooledHandles = 64;
477 static constexpr size_t kHeaderReserve = 1024;
484 explicit Transfer(CURL* e) : eh(e) {}
486 curl_slist* headers =
nullptr;
489 std::string header_buffer;
491 std::list<Transfer>::iterator self;
492 std::uint64_t
id = 0;
496 bool scrub_on_done =
false;
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;
505 curl_socket_t fd = CURL_SOCKET_BAD;
508 bool read_armed =
false;
509 bool write_armed =
false;
512 void start(Request r, Handler h) {
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);
520 transfers_.splice(transfers_.begin(), idle_, idle_.begin());
522 const auto it = transfers_.begin();
524 CURL*
const eh = t.eh;
526 t.handler = std::move(h);
527 t.body = std::move(r.body_);
528 t.id = ++next_transfer_id_;
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()));
537 curl_easy_setopt(eh, CURLOPT_FOLLOWLOCATION,
538 r.max_redirects_ != 0 ? 1L : 0L);
539 curl_easy_setopt(eh, CURLOPT_MAXREDIRS, r.max_redirects_);
541 curl_slist* chunk =
nullptr;
542 for (
const auto& header : r.headers_) {
543 chunk = curl_slist_append(chunk, header.c_str());
546 curl_easy_setopt(eh, CURLOPT_HTTPHEADER, chunk);
551 curl_easy_setopt(eh, CURLOPT_CUSTOMREQUEST,
nullptr);
555 curl_easy_setopt(eh, CURLOPT_POSTFIELDS,
nullptr);
556 curl_easy_setopt(eh, CURLOPT_HTTPGET, 1L);
559 curl_easy_setopt(eh, CURLOPT_CUSTOMREQUEST,
nullptr);
560 curl_easy_setopt(eh, CURLOPT_POST, 1L);
564 curl_easy_setopt(eh, CURLOPT_CUSTOMREQUEST,
"PUT");
565 curl_easy_setopt(eh, CURLOPT_POST, 1L);
569 curl_easy_setopt(eh, CURLOPT_CUSTOMREQUEST,
"PATCH");
570 curl_easy_setopt(eh, CURLOPT_POST, 1L);
574 curl_easy_setopt(eh, CURLOPT_CUSTOMREQUEST,
"DELETE");
575 curl_easy_setopt(eh, CURLOPT_POST, 1L);
582 t.scrub_on_done =
true;
586 auto slot = net::get_associated_cancellation_slot(t.handler);
587 if (slot.is_connected()) {
592 slot.assign([
this, alive = alive_,
id = t.id](net::cancellation_type_t) {
593 if (*alive) cancel(id);
598 const CURLMcode rc = curl_multi_add_handle(multi_, eh);
599 if (rc != CURLM_OK) {
601 Handler handler = std::move(t.handler);
603 fail(std::move(handler),
604 rc == CURLM_OUT_OF_MEMORY ? CURLE_OUT_OF_MEMORY : CURLE_FAILED_INIT);
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 {
615 Response{code, 0, {}, {}});
623 void cancel(std::uint64_t
id) {
624 for (
auto& t : transfers_) {
625 if (t.id !=
id)
continue;
627 curl_multi_remove_handle(multi_, t.eh);
628 if (running_ > 0) --running_;
629 Handler handler = std::move(t.handler);
636 net::post(io_, net::bind_allocator(
637 net::recycling_allocator<void>(),
638 [h = std::move(handler)]()
mutable {
641 Response{CURLE_ABORTED_BY_CALLBACK, 0, {}, {}});
652 void retire(Transfer& t,
bool scrub) {
653 curl_slist_free_all(t.headers);
656 t.scrub_on_done =
false;
658 t.header_buffer.clear();
659 std::string().swap(t.body);
661 curl_easy_reset(t.eh);
662 configure_handle(t.eh);
664 if (idle_.size() < max_pooled_handles_) {
665 idle_.splice(idle_.begin(), transfers_, t.self);
667 curl_easy_cleanup(t.eh);
668 transfers_.erase(t.self);
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
682 curl_easy_setopt(eh, CURLOPT_MAXLIFETIME_CONN, 30L);
685 curl_easy_setopt(eh, CURLOPT_ACCEPT_ENCODING,
"");
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);
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());
702 static size_t header_cb(
char* data,
size_t n,
size_t l, std::string* buf) {
703 if (buf->empty()) buf->reserve(kHeaderReserve);
704 buf->append(data, n * l);
708 static size_t body_write_cb(
char* data,
size_t n,
size_t l, Transfer* t) {
709 if (t->buffer.empty()) {
711 if (curl_easy_getinfo(t->eh, CURLINFO_CONTENT_LENGTH_DOWNLOAD_T, &len) ==
714 t->buffer.reserve(
static_cast<size_t>(len));
717 t->buffer.append(data, n * l);
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;
734 std::make_shared<SocketState>(self->io_,
false);
736 state->socket.open(protocol, ec);
737 if (ec)
return CURL_SOCKET_BAD;
738 const curl_socket_t fd = state->socket.native_handle();
740 self->sockets_[fd] = std::move(state);
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()) {
753 SocketState& state = *it->second;
756 state.socket.close(ignored);
757 if (state.borrowed) close_native(fd);
758 self->sockets_.erase(it);
762 static void close_native(curl_socket_t fd) {
777 SocketState* borrow(curl_socket_t fd) {
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_,
true);
788 state->socket.assign(net::ip::tcp::v4(), dup_fd, ec);
794 SocketState*
const raw = state.get();
795 sockets_[fd] = std::move(state);
800 static int socket_cb(CURL*, curl_socket_t fd,
int what,
void* userp,
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;
808 const auto it = self->sockets_.find(fd);
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);
814 if (what == CURL_POLL_REMOVE) {
816 if (state->borrowed) {
821 state->socket.close(ignored);
822 self->sockets_.erase(fd);
827 self->arm(state->shared_from_this());
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,
838 net::recycling_allocator<void>(),
839 [
this, w = std::weak_ptr<SocketState>(state)](
error_code ec) {
843 if (
auto s = w.lock())
844 on_event(std::move(s), CURL_CSELECT_IN, ec);
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,
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);
860 void on_event(std::shared_ptr<SocketState> state,
int flag,
error_code ec) {
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,
869 if (state->watch != 0) arm(state);
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;
879 if (timeout_ms == 0) {
883 if (!self->kick_pending_) {
884 self->kick_pending_ =
true;
886 net::bind_allocator(net::recycling_allocator<void>(),
887 [self, alive = self->alive_] {
889 self->kick_pending_ = false;
895 const auto deadline = std::chrono::steady_clock::now() +
896 std::chrono::milliseconds(timeout_ms);
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>(),
906 if (ec || !*alive) return;
907 self->timer_armed_ = false;
914 curl_multi_socket_action(multi_, CURL_SOCKET_TIMEOUT, 0, &running_);
918 void check_completions() {
920 while (CURLMsg* msg = curl_multi_info_read(multi_, &msgs_left)) {
921 if (msg->msg != CURLMSG_DONE)
continue;
922 Transfer* t =
nullptr;
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);
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);
941 std::move(handler)(ec, std::move(res));
945 static std::mutex& global_mutex() {
950 net::io_context& io_;
952 net::steady_timer timer_;
954 std::unordered_map<curl_socket_t, std::shared_ptr<SocketState>> sockets_;
955 std::list<Transfer> transfers_;
956 std::list<Transfer> idle_;
957 const size_t max_pooled_handles_;
959 std::uint64_t next_transfer_id_ = 0;
960 bool kick_pending_ =
false;
961 bool timer_armed_ =
false;
964 std::shared_ptr<bool> alive_ = std::make_shared<bool>(
true);