#include "EXClient.h" #include "EXMgr.h" #include "Util.h" #include namespace { const qint64 MAX_BUFFER = 20000000LL; } class BadServerReply : public Exception { public: using Exception::Exception; /// bring in c'tor ~BadServerReply(); }; BadServerReply::~BadServerReply() {} // for vtable EXClient::EXClient(EXMgr *mgr, qint64 id, const QString &host, quint16 tport, quint16 sport) : QObject(nullptr), id(id), host(host), tport(tport), sport(sport), mgr(mgr) { Debug() << __FUNCTION__ << " host:" << host << " t:" << tport << " s:" << sport; _thread.setObjectName(QString("%1 %2").arg("EXClient").arg(host)); } EXClient::~EXClient() { Debug() << __FUNCTION__ << " host:" << host; stop(); } /// this should only be called from our thread, because it accesses socket which should only be touched from thread QString EXClient::prettyName() const { bool dontTouchSocket = false; if (_thread.isRunning() && QThread::currentThread() != &_thread) { Warning() << __PRETTY_FUNCTION__ << " called from another thread! FIXME!"; dontTouchSocket = true; } QString type = socket && !dontTouchSocket ? (dynamic_cast(socket) ? "SSL" : "TCP") : "(NoSocket)"; QString port = socket && !dontTouchSocket && socket->peerPort() ? QString(":%1").arg(socket->peerPort()) : ""; QString ip = socket && !dontTouchSocket && !socket->peerAddress().isNull() ? socket->peerAddress().toString() : ""; return QString("%1 %2 %3%4").arg(type).arg(host).arg(ip).arg(port); } bool EXClient::isGood() const { return _thread.isRunning() && status == Connected && info.isValid(); } bool EXClient::isStale() const { return isGood() && Util::getTime() - lastGood > stale_threshold; } void EXClient::start() { if (_thread.isRunning()) return; Debug() << host << " starting thread"; moveToThread(&_thread); connect(&_thread, &QThread::started, this, &EXClient::on_started); connect(&_thread, &QThread::finished, this, &EXClient::on_finished); connect(this, &EXClient::sendRequest, this, &EXClient::_sendRequest); _thread.start(); } void EXClient::stop() { /// Disconnect this externally originating signal before stopping to /// ensure no new signals get sent to us after we switch back to the /// main thread. disconnect(this, &EXClient::sendRequest, this, &EXClient::_sendRequest); if (_thread.isRunning()) { Debug() << host << " thread is running, joining thread"; _thread.quit(); _thread.wait(); } disconnect(&_thread, &QThread::started, this, &EXClient::on_started); disconnect(&_thread, &QThread::finished, this, &EXClient::on_finished); } // runs in thread void EXClient::on_started() { Debug() << "started"; reconnect(); } // runs in thread void EXClient::on_finished() { killSocket(); moveToThread(qApp->thread()); Debug() << "finished."; } void EXClient::killSocket() { if (socket && socket->state() != QAbstractSocket::UnconnectedState) { Debug() << host << " aborting connection"; boilerplate_disconnect(); } if (socket) { delete socket; socket = nullptr; } status = NotConnected; } void EXClient::reconnect() { killSocket(); lastConnectionAttempt = Util::getTime(); if (sport && QSslSocket::supportsSsl()) { QSslSocket *ssl; socket = ssl = new QSslSocket(this); auto conf = ssl->sslConfiguration(); conf.setPeerVerifyMode(QSslSocket::VerifyNone); ssl->setSslConfiguration(conf); connect(ssl, &QSslSocket::encrypted, this, [this]{ Debug() << prettyName() << " connected encrypted"; on_connected(); }); connect(socket, SIGNAL(error(QAbstractSocket::SocketError)), this, SLOT(on_error(QAbstractSocket::SocketError))); connect(socket, SIGNAL(stateChanged(QAbstractSocket::SocketState)), this, SLOT(on_socketState(QAbstractSocket::SocketState))); socket->setSocketOption(QAbstractSocket::KeepAliveOption, true); // from Qt docs: required on Windows ssl->connectToHostEncrypted(host, static_cast(sport)); } else if (tport) { socket = new QTcpSocket(this); connect(socket, &QAbstractSocket::connected, this, [this]{ Debug() << prettyName() << " connected"; on_connected(); }); connect(socket, SIGNAL(error(QAbstractSocket::SocketError)), this, SLOT(on_error(QAbstractSocket::SocketError))); connect(socket, SIGNAL(stateChanged(QAbstractSocket::SocketState)), this, SLOT(on_socketState(QAbstractSocket::SocketState))); socket->setSocketOption(QAbstractSocket::KeepAliveOption, true); // from Qt docs: required on Windows socket->connectToHost(host, static_cast(tport)); } else { Error() << "Cannot connect to " << host << "; no TCP port defined and SSL is disabled on this install"; } } void EXClient::on_socketState(QAbstractSocket::SocketState s) { Debug() << prettyName() << " socket state: " << s; switch (s) { case QAbstractSocket::ConnectedState: status = Connected; break; case QAbstractSocket::HostLookupState: case QAbstractSocket::ConnectingState: status = Connecting; break; case QAbstractSocket::UnconnectedState: case QAbstractSocket::ClosingState: default: status = NotConnected; break; } } bool EXClient::_sendRequest(qint64 id, const QString &method, const QVariantList ¶ms) { if (status != Connected || !socket) { Error() << __FUNCTION__ << " method: " << method << "; Not connected!"; return false; } while (idMethodMap.size() > 20000) { // prevent memory leaks in case of misbehaving server idMethodMap.erase(idMethodMap.begin()); } idMethodMap[id] = method; return do_write(makeRequestData(id, method, params)); } void EXClient::boilerplate_disconnect() { status = status == Bad ? Bad : NotConnected; // try and keep Bad status around so EXMgr can decide when to reconnect based on it if (socket) socket->abort(); // this will set status too because state change, but we set it first above to be paranoid } bool EXClient::do_write(const QByteArray & data) { auto data2write = writeBackLog + data; qint64 written = socket->write(data2write); if (written < 0) { Error() << __FUNCTION__ << " error on write " << socket->error() << " (" << socket->errorString() << ") "; boilerplate_disconnect(); return false; } else if (written < data2write.length()) { writeBackLog = data2write.mid(int(written)); } nSent += written; if (writeBackLog.length() > MAX_BUFFER) { Error() << __FUNCTION__ << " MAX_BUFFER reached on write (" << MAX_BUFFER << ")"; boilerplate_disconnect(); return false; } return true; } void EXClient::kill_pingTimer() { if (pingTimer) { delete pingTimer; pingTimer = nullptr; } } void EXClient::start_pingTimer() { kill_pingTimer(); pingTimer = new QTimer(this); pingTimer->setSingleShot(false); connect(pingTimer, SIGNAL(timeout()), this, SLOT(on_pingTimer())); pingTimer->start(pingtime_ms/* 1 minute */ / 2); } void EXClient::on_pingTimer() { if (Util::getTime() - lastGood > pingtime_ms) // only ping if we've been idle for longer than 1 minute emit sendRequest(mgr->newId(), "server.ping"); } void EXClient::on_connected() { // runs in thread Debug() << __FUNCTION__; connect(socket, SIGNAL(readyRead()), this, SLOT(on_readyRead())); if (dynamic_cast(socket)) { // for some reason Qt can't find this old-style signal for QSslSocket so we do the below. // Additionally, bytesWritten is never emitted for QSslSocket, violating OOP! Thanks Qt. :P connect(socket, SIGNAL(encryptedBytesWritten(qint64)), this, SLOT(on_bytesWritten())); } else { connect(socket, SIGNAL(bytesWritten()), this, SLOT(on_bytesWritten())); } connect(socket, &QAbstractSocket::disconnected, this, [this]{ Debug() << prettyName() << " socket disconnected"; kill_pingTimer(); emit lostConnection(this); idMethodMap.clear(); // todo: put stuff to queue up a reconnect sometime later? }); emit newConnection(this); start_pingTimer(); } /* static */ EXResponse EXResponse::fromJson(const QString &json) { const auto m = Util::Json::parseString(json).toMap(); const auto jsonrpc = m.value("jsonrpc", "").toString(); if (jsonrpc != "2.0") throw BadServerReply(QString("Unexpected or missing jsonrpc version: \"%1\"").arg(jsonrpc)); const qint64 id = m.value("id", -1).toLongLong(); QString method = m.value("method", "").toString(); if (id < 0 && method.isEmpty()) throw BadServerReply("Bad server reply, missing required id field in JSON"); QVariantMap err = m.value("error", QVariantMap()).toMap(); if (!err.isEmpty()) { // error reply const int code = err.value("code", 123456789).toInt(); const QString message = err.value("message").toString(); if (code == 123456789 || message.isEmpty()) throw BadServerReply("Bad server reply, error field in JSON is not of the expected format"); return EXResponse{ jsonrpc, id, method, QVariant(), code, message }; } QVariant result = m.value("result", QVariant()); if (result.isNull()) { result = m.value("params", QVariant()); } return EXResponse{ jsonrpc, id, method, result }; } QString EXResponse::toString() const { return QString("jsonrpc: %1 ; id: %2 ; method: %3 ; result: %4 ; error code: %5 ; error message: %6") .arg(jsonRpcVersion).arg(id).arg(method).arg(result.isNull() ? "(null)" : Util::Json::toString(result, true)) .arg(errorCode).arg(errorMessage); } void EXResponse::validate() { if (!errorMessage.isEmpty()) return; if (method == "server.version" && !result.isNull()) { QVariantList l = result.toList(); if (l.count() < 2 || l[0].toString().isNull() || l[1].toString().isNull()) throw BadServerReply(QString("%1 expected string list of size 2").arg(method)); return; // ok } else if (method == "blockchain.headers.subscribe" && !result.isNull()) { QVariantMap m = result.toMap(); QVariantList l = result.toList(); if (m.isEmpty() && !l.isEmpty()) { // spontaneous "subscribe" callbacks pass a list containing a dict rather than a straight up dict, // so mogrify ourselves to always contain the dict m = l.last().toMap(); result = m; // save back result as a map rather than a list } if (m.isEmpty() || m.count() < 2 || m.value("height", -1).toInt() < 0 || m.value("hex", "").toString().isEmpty()) { throw BadServerReply(QString("%1 expected map with 'height' and 'hex'").arg(method)); } return; // ok } else if (method == "server.ping") { // always accept return; } throw BadServerReply(QString("Unexpected method \"%1\", and/or incomplete/missing results").arg(method)); } void EXClient::on_readyRead() { Debug() << __FUNCTION__; try { while (socket->canReadLine()) { auto data = socket->readLine(); nReceived += data.length(); auto line = data.trimmed(); Debug() << "Got: " << line; auto resp = EXResponse::fromJson(line); auto meth = resp.id > 0 ? idMethodMap.take(resp.id) : resp.method; if (meth.isEmpty()) { throw BadServerReply(QString("Unexpected/unknown message id (%1) in server reply").arg(resp.id)); } resp.method = meth; resp.validate(); // may throw, may modify resp Debug() << "Parsed response: " << resp.toString(); lastGood = Util::getTime(); emit gotResponse(this, resp); } if (socket->bytesAvailable() > MAX_BUFFER) { // bad server.. sending us garbage data not containing newlines. Kill connection. throw BadServerReply(QString("Server has sent us more than %1 bytes without a newline! Bad server?").arg(MAX_BUFFER)); } } catch (const Exception &e) { Error() << "Error reading/parsing response: " << e.what(); boilerplate_disconnect(); status = Bad; } } void EXClient::on_bytesWritten() { Debug() << __FUNCTION__; if (!writeBackLog.isEmpty() && status == Connected && socket) { Debug() << "writeBackLog size: " << writeBackLog.length(); do_write(); } } void EXClient::on_error(QAbstractSocket::SocketError err) { Warning() << prettyName() << ": error " << err << " (" << (socket ? socket->errorString() : "(null)") << ")"; boilerplate_disconnect(); // todo: put stuff to queue up a reconnect sometime later? } /* static */ QByteArray EXClient::makeRequestData(qint64 id, const QString &method, const QVariantList ¶ms) { QVariantMap m; m["id"] = id; m["method"] = method; m["params"] = params; try { static const QChar nl(012); return QString("%1%2").arg(Util::Json::toString(m, true)).arg(nl).toUtf8(); } catch (const Util::Json::Error &e) { Error() << __FUNCTION__ << ": " << e.what(); } return QByteArray(); }