json_rpc_request fixes

This commit is contained in:
SChernykh
2026-04-26 19:35:13 +02:00
parent cd1ea288df
commit 4ebae25093
3 changed files with 87 additions and 48 deletions
+85 -46
View File
@@ -41,7 +41,7 @@ struct CurlContext : public nocopy_nomove
return reinterpret_cast<CurlContext*>(ctx)->on_timer(multi, timeout_ms); return reinterpret_cast<CurlContext*>(ctx)->on_timer(multi, timeout_ms);
} }
static size_t write_func(const void* buffer, size_t size, size_t count, void* ctx) static size_t write_func(char* buffer, size_t size, size_t count, void* ctx)
{ {
return reinterpret_cast<CurlContext*>(ctx)->on_write(buffer, size, count); return reinterpret_cast<CurlContext*>(ctx)->on_write(buffer, size, count);
} }
@@ -49,15 +49,13 @@ struct CurlContext : public nocopy_nomove
int on_socket(CURL* easy, curl_socket_t s, int action); int on_socket(CURL* easy, curl_socket_t s, int action);
int on_timer(CURLM* multi, long timeout_ms); int on_timer(CURLM* multi, long timeout_ms);
static void on_timeout(uv_handle_t* req); static void on_timeout(uv_timer_t* req);
size_t on_write(const void* buffer, size_t size, size_t count); size_t on_write(const char* buffer, size_t size, size_t count);
static void curl_perform(uv_poll_t* req, int status, int events); static void curl_perform(uv_poll_t* req, int status, int events);
void check_multi_info(); void check_multi_info();
static void on_close(uv_handle_t* h);
void shutdown(); void shutdown();
std::vector<std::pair<curl_socket_t, uv_poll_t*>> m_pollHandles; std::vector<std::pair<curl_socket_t, uv_poll_t*>> m_pollHandles;
@@ -66,7 +64,7 @@ struct CurlContext : public nocopy_nomove
CallbackBase* m_closeCallback; CallbackBase* m_closeCallback;
uv_loop_t* m_loop; uv_loop_t* m_loop;
uv_timer_t m_timer; uv_timer_t* m_timer;
CURLM* m_multiHandle; CURLM* m_multiHandle;
CURL* m_handle; CURL* m_handle;
@@ -88,7 +86,7 @@ CurlContext::CurlContext(const std::string& address, int port, const std::string
: m_callback(cb) : m_callback(cb)
, m_closeCallback(close_cb) , m_closeCallback(close_cb)
, m_loop(loop) , m_loop(loop)
, m_timer{} , m_timer(nullptr)
, m_multiHandle(nullptr) , m_multiHandle(nullptr)
, m_handle(nullptr) , m_handle(nullptr)
, m_req(req) , m_req(req)
@@ -97,8 +95,6 @@ CurlContext::CurlContext(const std::string& address, int port, const std::string
, m_connectedTime(0) , m_connectedTime(0)
, m_proxy(proxy) , m_proxy(proxy)
{ {
BACKGROUND_JOB_START(CurlContext);
m_pollHandles.reserve(2); m_pollHandles.reserve(2);
{ {
@@ -132,18 +128,10 @@ CurlContext::CurlContext(const std::string& address, int port, const std::string
m_url = buf; m_url = buf;
} }
int err = uv_timer_init(m_loop, &m_timer);
if (err) {
LOGERR(1, "uv_timer_init failed, error " << uv_err_name(err));
throw std::runtime_error("uv_timer_init failed");
}
m_timer.data = this;
m_multiHandle = curl_multi_init(); m_multiHandle = curl_multi_init();
if (!m_multiHandle) { if (!m_multiHandle) {
static constexpr char msg[] = "curl_multi_init() failed"; static constexpr char msg[] = "curl_multi_init() failed";
LOGERR(1, msg); LOGERR(1, msg);
uv_close(reinterpret_cast<uv_handle_t*>(&m_timer), nullptr);
throw std::runtime_error(msg); throw std::runtime_error(msg);
} }
@@ -153,6 +141,7 @@ CurlContext::CurlContext(const std::string& address, int port, const std::string
if (r != CURLM_OK) { \ if (r != CURLM_OK) { \
static constexpr char msg[] = "curl_multi_setopt(" #__VA_ARGS__ ") failed"; \ static constexpr char msg[] = "curl_multi_setopt(" #__VA_ARGS__ ") failed"; \
LOGERR(1, msg << ": " << curl_multi_strerror(r)); \ LOGERR(1, msg << ": " << curl_multi_strerror(r)); \
curl_multi_cleanup(m_multiHandle); \
throw std::runtime_error(msg); \ throw std::runtime_error(msg); \
} \ } \
} while (0) } while (0)
@@ -168,7 +157,6 @@ CurlContext::CurlContext(const std::string& address, int port, const std::string
static constexpr char msg[] = "curl_easy_init() failed"; static constexpr char msg[] = "curl_easy_init() failed";
LOGERR(1, msg); LOGERR(1, msg);
curl_multi_cleanup(m_multiHandle); curl_multi_cleanup(m_multiHandle);
uv_close(reinterpret_cast<uv_handle_t*>(&m_timer), nullptr);
throw std::runtime_error(msg); throw std::runtime_error(msg);
} }
@@ -178,6 +166,9 @@ CurlContext::CurlContext(const std::string& address, int port, const std::string
if (r != CURLE_OK) { \ if (r != CURLE_OK) { \
static constexpr char msg[] = "curl_easy_setopt(" #__VA_ARGS__ ") failed"; \ static constexpr char msg[] = "curl_easy_setopt(" #__VA_ARGS__ ") failed"; \
LOGERR(1, msg << ": " << curl_easy_strerror(r)); \ LOGERR(1, msg << ": " << curl_easy_strerror(r)); \
curl_easy_cleanup(m_handle); \
curl_multi_cleanup(m_multiHandle); \
curl_slist_free_all(m_headers); \
throw std::runtime_error(msg); \ throw std::runtime_error(msg); \
} \ } \
} while (0) } while (0)
@@ -194,7 +185,11 @@ CurlContext::CurlContext(const std::string& address, int port, const std::string
#endif #endif
curl_easy_setopt_checked(m_handle, CURLOPT_URL, m_url.c_str()); curl_easy_setopt_checked(m_handle, CURLOPT_URL, m_url.c_str());
curl_easy_setopt_checked(m_handle, CURLOPT_POSTFIELDS, m_req.c_str());
if (!m_req.empty()) {
curl_easy_setopt_checked(m_handle, CURLOPT_POSTFIELDS, m_req.c_str());
}
curl_easy_setopt_checked(m_handle, CURLOPT_CONNECTTIMEOUT, timeout); curl_easy_setopt_checked(m_handle, CURLOPT_CONNECTTIMEOUT, timeout);
curl_easy_setopt_checked(m_handle, CURLOPT_TIMEOUT, timeout * 10); curl_easy_setopt_checked(m_handle, CURLOPT_TIMEOUT, timeout * 10);
@@ -223,7 +218,7 @@ CurlContext::CurlContext(const std::string& address, int port, const std::string
curl_easy_setopt_checked(m_handle, CURLOPT_SSL_VERIFYHOST, 0L); curl_easy_setopt_checked(m_handle, CURLOPT_SSL_VERIFYHOST, 0L);
if (!ssl_fingerprint.empty()) { if (!ssl_fingerprint.empty()) {
char buf[64] = {}; char buf[128] = {};
log::Stream s(buf); log::Stream s(buf);
s << "sha256//" << ssl_fingerprint; s << "sha256//" << ssl_fingerprint;
@@ -232,15 +227,32 @@ CurlContext::CurlContext(const std::string& address, int port, const std::string
} }
#endif #endif
m_timer = new uv_timer_t{};
const int err = uv_timer_init(m_loop, m_timer);
if (err) {
LOGERR(1, "uv_timer_init failed, error " << uv_err_name(err));
delete m_timer;
curl_easy_cleanup(m_handle);
curl_multi_cleanup(m_multiHandle);
curl_slist_free_all(m_headers);
throw std::runtime_error("uv_timer_init failed");
}
m_timer->data = this;
CURLMcode curl_err = curl_multi_add_handle(m_multiHandle, m_handle); CURLMcode curl_err = curl_multi_add_handle(m_multiHandle, m_handle);
if (curl_err != CURLM_OK) { if (curl_err != CURLM_OK) {
LOGERR(1, "curl_multi_add_handle failed: " << curl_multi_strerror(curl_err)); LOGERR(1, "curl_multi_add_handle failed: " << curl_multi_strerror(curl_err));
curl_easy_cleanup(m_handle); curl_easy_cleanup(m_handle);
curl_multi_cleanup(m_multiHandle); curl_multi_cleanup(m_multiHandle);
uv_close(reinterpret_cast<uv_handle_t*>(&m_timer), nullptr); curl_slist_free_all(m_headers);
uv_close(reinterpret_cast<uv_handle_t*>(m_timer), [](uv_handle_t* h) { delete reinterpret_cast<uv_timer_t*>(h); });
throw std::runtime_error("curl_multi_add_handle failed"); throw std::runtime_error("curl_multi_add_handle failed");
} }
BACKGROUND_JOB_START(CurlContext);
m_startTime = microseconds_since_epoch(); m_startTime = microseconds_since_epoch();
} }
@@ -251,7 +263,12 @@ CurlContext::~CurlContext()
if (m_error.empty() && !m_response.empty()) { if (m_error.empty() && !m_response.empty()) {
const uint64_t t = (m_proxy.empty() && m_connectedTime) ? m_connectedTime : microseconds_since_epoch(); const uint64_t t = (m_proxy.empty() && m_connectedTime) ? m_connectedTime : microseconds_since_epoch();
tcp_ping = static_cast<double>(t - m_startTime) / 1000.0; tcp_ping = static_cast<double>(t - m_startTime) / 1000.0;
(*m_callback)(m_response.data(), m_response.size(), tcp_ping);
try {
(*m_callback)(m_response.data(), m_response.size(), tcp_ping);
}
catch (...) {
}
} }
delete m_callback; delete m_callback;
@@ -264,10 +281,13 @@ CurlContext::~CurlContext()
} }
} }
(*m_closeCallback)(m_error.c_str(), m_error.length(), tcp_ping); try {
delete m_closeCallback; (*m_closeCallback)(m_error.c_str(), m_error.length(), tcp_ping);
}
catch (...) {
}
curl_slist_free_all(m_headers); delete m_closeCallback;
BACKGROUND_JOB_STOP(CurlContext); BACKGROUND_JOB_STOP(CurlContext);
} }
@@ -348,12 +368,17 @@ int CurlContext::on_socket(CURL* /*easy*/, curl_socket_t s, int action)
int CurlContext::on_timer(CURLM* /*multi*/, long timeout_ms) int CurlContext::on_timer(CURLM* /*multi*/, long timeout_ms)
{ {
if (!m_timer) {
LOGERR(1, "on_timer: timer is not initialized/already destroyed");
return -1;
}
if (timeout_ms < 0) { if (timeout_ms < 0) {
uv_timer_stop(&m_timer); uv_timer_stop(m_timer);
return 0; return 0;
} }
const int result = uv_timer_start(&m_timer, reinterpret_cast<uv_timer_cb>(on_timeout), timeout_ms, 0); const int result = uv_timer_start(m_timer, on_timeout, timeout_ms, 0);
if (result < 0) { if (result < 0) {
LOGERR(1, "uv_timer_start failed with error " << uv_err_name(result)); LOGERR(1, "uv_timer_start failed with error " << uv_err_name(result));
return -1; return -1;
@@ -362,7 +387,7 @@ int CurlContext::on_timer(CURLM* /*multi*/, long timeout_ms)
return 0; return 0;
} }
void CurlContext::on_timeout(uv_handle_t* req) void CurlContext::on_timeout(uv_timer_t* req)
{ {
CurlContext* ctx = reinterpret_cast<CurlContext*>(req->data); CurlContext* ctx = reinterpret_cast<CurlContext*>(req->data);
@@ -379,15 +404,14 @@ void CurlContext::on_timeout(uv_handle_t* req)
} }
} }
size_t CurlContext::on_write(const void* buffer, size_t size, size_t count) size_t CurlContext::on_write(const char* buffer, size_t size, size_t count)
{ {
if (!m_connectedTime) { if (!m_connectedTime) {
m_connectedTime = microseconds_since_epoch(); m_connectedTime = microseconds_since_epoch();
} }
const size_t realsize = size * count; const size_t realsize = size * count;
const char* p = reinterpret_cast<const char*>(buffer); m_response.insert(m_response.end(), buffer, buffer + realsize);
m_response.insert(m_response.end(), p, p + realsize);
return realsize; return realsize;
} }
@@ -452,23 +476,15 @@ void CurlContext::check_multi_info()
curl_multi_remove_handle(m_multiHandle, m_handle); curl_multi_remove_handle(m_multiHandle, m_handle);
curl_easy_cleanup(m_handle); curl_easy_cleanup(m_handle);
curl_multi_cleanup(m_multiHandle); curl_multi_cleanup(m_multiHandle);
m_handle = nullptr;
m_multiHandle = nullptr;
return; return;
} }
} }
} }
void CurlContext::on_close(uv_handle_t* h)
{
CurlContext* ctx = reinterpret_cast<CurlContext*>(h->data);
h->data = nullptr;
if (ctx->m_timer.data) {
return;
}
delete ctx;
}
void CurlContext::shutdown() void CurlContext::shutdown()
{ {
for (const auto& p : m_pollHandles) { for (const auto& p : m_pollHandles) {
@@ -477,8 +493,27 @@ void CurlContext::shutdown()
} }
m_pollHandles.clear(); m_pollHandles.clear();
if (m_timer.data && !uv_is_closing(reinterpret_cast<uv_handle_t*>(&m_timer))) { if (m_timer->data && !uv_is_closing(reinterpret_cast<uv_handle_t*>(m_timer))) {
uv_close(reinterpret_cast<uv_handle_t*>(&m_timer), on_close); if (m_handle) {
if (m_multiHandle) {
curl_multi_remove_handle(m_multiHandle, m_handle);
}
curl_easy_cleanup(m_handle);
m_handle = nullptr;
}
if (m_multiHandle) {
curl_multi_cleanup(m_multiHandle);
m_multiHandle = nullptr;
}
curl_slist_free_all(m_headers);
uv_close(reinterpret_cast<uv_handle_t*>(m_timer), [](uv_handle_t* h) {
CurlContext* ctx = reinterpret_cast<CurlContext*>(h->data);
delete reinterpret_cast<uv_timer_t*>(h);
delete ctx;
});
} }
} }
@@ -496,7 +531,11 @@ void Call(const std::string& address, int port, const std::string& req, const st
} }
catch (const std::exception& e) { catch (const std::exception& e) {
const char* msg = e.what(); const char* msg = e.what();
(*close_cb)(msg, strlen(msg), 0.0); try {
(*close_cb)(msg, strlen(msg), 0.0);
}
catch (...) {
}
delete cb; delete cb;
delete close_cb; delete close_cb;
} }
+1 -1
View File
@@ -32,7 +32,7 @@ FORCEINLINE void call(const std::string& address, int port, const std::string& r
typedef Callback<void, const char*, size_t, double>::Derived<T> CallbackT; typedef Callback<void, const char*, size_t, double>::Derived<T> CallbackT;
typedef Callback<void, const char*, size_t, double>::Derived<U> CallbackU; typedef Callback<void, const char*, size_t, double>::Derived<U> CallbackU;
Call(address, port, req, auth, proxy, ssl, ssl_fingerprint, new CallbackT(std::move(cb)), new CallbackU(std::move(close_cb)), loop); Call(address, port, req, auth, proxy, ssl, ssl_fingerprint, new CallbackT(std::forward<T>(cb)), new CallbackU(std::forward<U>(close_cb)), loop);
} }
} // namespace JSONRPCRequest } // namespace JSONRPCRequest
+1 -1
View File
@@ -361,7 +361,7 @@ struct Callback
template<typename T> template<typename T>
struct Derived : public Base struct Derived : public Base
{ {
explicit FORCEINLINE Derived(T&& cb) : m_cb(std::move(cb)) {} explicit FORCEINLINE Derived(T&& cb) : m_cb(std::forward<T>(cb)) {}
R operator()(Args... args) const override { return m_cb(args...); } R operator()(Args... args) const override { return m_cb(args...); }
private: private: