mirror of
https://github.com/xmrig/xmrig.git
synced 2024-11-05 16:07:42 +00:00
Added send with callback.
This commit is contained in:
parent
3752551e53
commit
83a5923568
9 changed files with 127 additions and 72 deletions
|
@ -29,6 +29,9 @@
|
||||||
#include "rapidjson/fwd.h"
|
#include "rapidjson/fwd.h"
|
||||||
|
|
||||||
|
|
||||||
|
#include <functional>
|
||||||
|
|
||||||
|
|
||||||
namespace xmrig {
|
namespace xmrig {
|
||||||
|
|
||||||
|
|
||||||
|
@ -51,6 +54,8 @@ public:
|
||||||
EXT_MAX
|
EXT_MAX
|
||||||
};
|
};
|
||||||
|
|
||||||
|
using Callback = std::function<void(const rapidjson::Value &result, bool success, uint64_t elapsed)>;
|
||||||
|
|
||||||
virtual ~IClient() = default;
|
virtual ~IClient() = default;
|
||||||
|
|
||||||
virtual bool disconnect() = 0;
|
virtual bool disconnect() = 0;
|
||||||
|
@ -64,6 +69,7 @@ public:
|
||||||
virtual const Pool &pool() const = 0;
|
virtual const Pool &pool() const = 0;
|
||||||
virtual const String &ip() const = 0;
|
virtual const String &ip() const = 0;
|
||||||
virtual int id() const = 0;
|
virtual int id() const = 0;
|
||||||
|
virtual int64_t send(const rapidjson::Value &obj, Callback callback) = 0;
|
||||||
virtual int64_t send(const rapidjson::Value &obj) = 0;
|
virtual int64_t send(const rapidjson::Value &obj) = 0;
|
||||||
virtual int64_t sequence() const = 0;
|
virtual int64_t sequence() const = 0;
|
||||||
virtual int64_t submit(const JobResult &result) = 0;
|
virtual int64_t submit(const JobResult &result) = 0;
|
||||||
|
|
|
@ -26,6 +26,7 @@
|
||||||
#include "base/kernel/interfaces/IClientListener.h"
|
#include "base/kernel/interfaces/IClientListener.h"
|
||||||
#include "base/net/stratum/BaseClient.h"
|
#include "base/net/stratum/BaseClient.h"
|
||||||
#include "base/net/stratum/SubmitResult.h"
|
#include "base/net/stratum/SubmitResult.h"
|
||||||
|
#include "rapidjson/document.h"
|
||||||
|
|
||||||
|
|
||||||
namespace xmrig {
|
namespace xmrig {
|
||||||
|
@ -42,6 +43,32 @@ xmrig::BaseClient::BaseClient(int id, IClientListener *listener) :
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
bool xmrig::BaseClient::handleResponse(int64_t id, const rapidjson::Value &result, const rapidjson::Value &error)
|
||||||
|
{
|
||||||
|
if (id == 1) {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
auto it = m_callbacks.find(id);
|
||||||
|
if (it != m_callbacks.end()) {
|
||||||
|
const uint64_t elapsed = Chrono::steadyMSecs() - it->second.ts;
|
||||||
|
|
||||||
|
if (error.IsObject()) {
|
||||||
|
it->second.callback(error, false, elapsed);
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
it->second.callback(result, true, elapsed);
|
||||||
|
}
|
||||||
|
|
||||||
|
m_callbacks.erase(it);
|
||||||
|
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
bool xmrig::BaseClient::handleSubmitResponse(int64_t id, const char *error)
|
bool xmrig::BaseClient::handleSubmitResponse(int64_t id, const char *error)
|
||||||
{
|
{
|
||||||
auto it = m_results.find(id);
|
auto it = m_results.find(id);
|
||||||
|
|
|
@ -32,6 +32,7 @@
|
||||||
#include "base/kernel/interfaces/IClient.h"
|
#include "base/kernel/interfaces/IClient.h"
|
||||||
#include "base/net/stratum/Job.h"
|
#include "base/net/stratum/Job.h"
|
||||||
#include "base/net/stratum/Pool.h"
|
#include "base/net/stratum/Pool.h"
|
||||||
|
#include "base/tools/Chrono.h"
|
||||||
|
|
||||||
|
|
||||||
namespace xmrig {
|
namespace xmrig {
|
||||||
|
@ -70,8 +71,17 @@ protected:
|
||||||
ReconnectingState
|
ReconnectingState
|
||||||
};
|
};
|
||||||
|
|
||||||
|
struct SendResult
|
||||||
|
{
|
||||||
|
inline SendResult(Callback &&callback) : callback(callback), ts(Chrono::steadyMSecs()) {}
|
||||||
|
|
||||||
|
Callback callback;
|
||||||
|
const uint64_t ts;
|
||||||
|
};
|
||||||
|
|
||||||
inline bool isQuiet() const { return m_quiet || m_failures >= m_retries; }
|
inline bool isQuiet() const { return m_quiet || m_failures >= m_retries; }
|
||||||
|
|
||||||
|
bool handleResponse(int64_t id, const rapidjson::Value &result, const rapidjson::Value &error);
|
||||||
bool handleSubmitResponse(int64_t id, const char *error = nullptr);
|
bool handleSubmitResponse(int64_t id, const char *error = nullptr);
|
||||||
|
|
||||||
bool m_quiet = false;
|
bool m_quiet = false;
|
||||||
|
@ -82,6 +92,7 @@ protected:
|
||||||
Job m_job;
|
Job m_job;
|
||||||
Pool m_pool;
|
Pool m_pool;
|
||||||
SocketState m_state = UnconnectedState;
|
SocketState m_state = UnconnectedState;
|
||||||
|
std::map<int64_t, SendResult> m_callbacks;
|
||||||
std::map<int64_t, SubmitResult> m_results;
|
std::map<int64_t, SubmitResult> m_results;
|
||||||
String m_ip;
|
String m_ip;
|
||||||
uint64_t m_retryPause = 5000;
|
uint64_t m_retryPause = 5000;
|
||||||
|
|
|
@ -137,6 +137,16 @@ const char *xmrig::Client::tlsVersion() const
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
int64_t xmrig::Client::send(const rapidjson::Value &obj, Callback callback)
|
||||||
|
{
|
||||||
|
assert(obj["id"] == sequence());
|
||||||
|
|
||||||
|
m_callbacks.insert({ sequence(), std::move(callback) });
|
||||||
|
|
||||||
|
return send(obj);
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
int64_t xmrig::Client::send(const rapidjson::Value &obj)
|
int64_t xmrig::Client::send(const rapidjson::Value &obj)
|
||||||
{
|
{
|
||||||
using namespace rapidjson;
|
using namespace rapidjson;
|
||||||
|
@ -736,6 +746,10 @@ void xmrig::Client::parseNotification(const char *method, const rapidjson::Value
|
||||||
|
|
||||||
void xmrig::Client::parseResponse(int64_t id, const rapidjson::Value &result, const rapidjson::Value &error)
|
void xmrig::Client::parseResponse(int64_t id, const rapidjson::Value &result, const rapidjson::Value &error)
|
||||||
{
|
{
|
||||||
|
if (handleResponse(id, result, error)) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
if (error.IsObject()) {
|
if (error.IsObject()) {
|
||||||
const char *message = error["message"].GetString();
|
const char *message = error["message"].GetString();
|
||||||
|
|
||||||
|
|
|
@ -77,6 +77,7 @@ protected:
|
||||||
bool isTLS() const override;
|
bool isTLS() const override;
|
||||||
const char *tlsFingerprint() const override;
|
const char *tlsFingerprint() const override;
|
||||||
const char *tlsVersion() const override;
|
const char *tlsVersion() const override;
|
||||||
|
int64_t send(const rapidjson::Value &obj, Callback callback) override;
|
||||||
int64_t send(const rapidjson::Value &obj) override;
|
int64_t send(const rapidjson::Value &obj) override;
|
||||||
int64_t submit(const JobResult &result) override;
|
int64_t submit(const JobResult &result) override;
|
||||||
void connect() override;
|
void connect() override;
|
||||||
|
|
|
@ -58,6 +58,7 @@ protected:
|
||||||
inline const char *mode() const override { return "daemon"; }
|
inline const char *mode() const override { return "daemon"; }
|
||||||
inline const char *tlsFingerprint() const override { return m_tlsFingerprint; }
|
inline const char *tlsFingerprint() const override { return m_tlsFingerprint; }
|
||||||
inline const char *tlsVersion() const override { return m_tlsVersion; }
|
inline const char *tlsVersion() const override { return m_tlsVersion; }
|
||||||
|
inline int64_t send(const rapidjson::Value &, Callback) override { return -1; }
|
||||||
inline int64_t send(const rapidjson::Value &) override { return -1; }
|
inline int64_t send(const rapidjson::Value &) override { return -1; }
|
||||||
inline void deleteLater() override { delete this; }
|
inline void deleteLater() override { delete this; }
|
||||||
inline void tick(uint64_t) override {}
|
inline void tick(uint64_t) override {}
|
||||||
|
|
|
@ -122,7 +122,7 @@ bool xmrig::SelfSelectClient::parseResponse(int64_t id, rapidjson::Value &result
|
||||||
|
|
||||||
submitBlockTemplate(result);
|
submitBlockTemplate(result);
|
||||||
|
|
||||||
m_listener->onJobReceived(this, m_job, rapidjson::Value{});
|
|
||||||
|
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
@ -207,7 +207,9 @@ void xmrig::SelfSelectClient::submitBlockTemplate(rapidjson::Value &result)
|
||||||
|
|
||||||
JsonRequest::create(doc, sequence(), "block_template", params);
|
JsonRequest::create(doc, sequence(), "block_template", params);
|
||||||
|
|
||||||
send(doc);
|
send(doc, [this](const rapidjson::Value &result, bool success, uint64_t elapsed) {
|
||||||
|
m_listener->onJobReceived(this, m_job, rapidjson::Value{});
|
||||||
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|
|
@ -58,6 +58,7 @@ protected:
|
||||||
inline const Pool &pool() const override { return m_client->pool(); }
|
inline const Pool &pool() const override { return m_client->pool(); }
|
||||||
inline const String &ip() const override { return m_client->ip(); }
|
inline const String &ip() const override { return m_client->ip(); }
|
||||||
inline int id() const override { return m_client->id(); }
|
inline int id() const override { return m_client->id(); }
|
||||||
|
inline int64_t send(const rapidjson::Value &obj, Callback callback) override { return m_client->send(obj, callback); }
|
||||||
inline int64_t send(const rapidjson::Value &obj) override { return m_client->send(obj); }
|
inline int64_t send(const rapidjson::Value &obj) override { return m_client->send(obj); }
|
||||||
inline int64_t sequence() const override { return m_client->sequence(); }
|
inline int64_t sequence() const override { return m_client->sequence(); }
|
||||||
inline int64_t submit(const JobResult &result) override { return m_client->submit(result); }
|
inline int64_t submit(const JobResult &result) override { return m_client->submit(result); }
|
||||||
|
|
|
@ -35,34 +35,26 @@ namespace xmrig {
|
||||||
class SubmitResult
|
class SubmitResult
|
||||||
{
|
{
|
||||||
public:
|
public:
|
||||||
inline SubmitResult() :
|
SubmitResult() = default;
|
||||||
reqId(0),
|
|
||||||
seq(0),
|
|
||||||
actualDiff(0),
|
|
||||||
diff(0),
|
|
||||||
elapsed(0),
|
|
||||||
m_start(0)
|
|
||||||
{}
|
|
||||||
|
|
||||||
inline SubmitResult(int64_t seq, uint64_t diff, uint64_t actualDiff, int64_t reqId = 0) :
|
inline SubmitResult(int64_t seq, uint64_t diff, uint64_t actualDiff, int64_t reqId = 0) :
|
||||||
reqId(reqId),
|
reqId(reqId),
|
||||||
seq(seq),
|
seq(seq),
|
||||||
actualDiff(actualDiff),
|
actualDiff(actualDiff),
|
||||||
diff(diff),
|
diff(diff),
|
||||||
elapsed(0),
|
|
||||||
m_start(Chrono::steadyMSecs())
|
m_start(Chrono::steadyMSecs())
|
||||||
{}
|
{}
|
||||||
|
|
||||||
inline void done() { elapsed = Chrono::steadyMSecs() - m_start; }
|
inline void done() { elapsed = Chrono::steadyMSecs() - m_start; }
|
||||||
|
|
||||||
int64_t reqId;
|
int64_t reqId = 0;
|
||||||
int64_t seq;
|
int64_t seq = 0;
|
||||||
uint64_t actualDiff;
|
uint64_t actualDiff = 0;
|
||||||
uint64_t diff;
|
uint64_t diff = 0;
|
||||||
uint64_t elapsed;
|
uint64_t elapsed = 0;
|
||||||
|
|
||||||
private:
|
private:
|
||||||
uint64_t m_start;
|
uint64_t m_start = 0;
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|
||||||
|
|
Loading…
Reference in a new issue