switch from aws sdk to qmqtt

This commit is contained in:
Michael Zanetti 2018-08-30 11:59:55 +02:00
parent 637a30700b
commit 096a285212
6 changed files with 145 additions and 209 deletions

2
debian/control vendored
View File

@ -23,7 +23,7 @@ Build-Depends: debhelper (>= 9.0.0),
libavahi-client-dev, libavahi-client-dev,
libavahi-common-dev, libavahi-common-dev,
libssl-dev, libssl-dev,
libaws-iot-device-sdk-cpp, libqmqtt-dev,
dbus-test-runner, dbus-test-runner,
Package: nymea Package: nymea

View File

@ -30,35 +30,21 @@
#include <QSettings> #include <QSettings>
#include <QSslCertificate> #include <QSslCertificate>
#include <QFile> #include <QFile>
#include <QSslKey>
using namespace awsiotsdk;
using namespace awsiotsdk::network;
using namespace awsiotsdk::mqtt;
QHash<quint16, AWSConnector*> AWSConnector::s_requestMap;
// Somehow the linker fails to find this... missing in awsiotsdk?
DisconnectCallbackContextData::~DisconnectCallbackContextData() {}
AWSConnector::AWSConnector(QObject *parent) : QObject(parent) AWSConnector::AWSConnector(QObject *parent) : QObject(parent)
{ {
// Enable some AWS logging (does not regard our logging categories)
std::shared_ptr<awsiotsdk::util::Logging::ConsoleLogSystem> p_log_system =
std::make_shared<awsiotsdk::util::Logging::ConsoleLogSystem>(awsiotsdk::util::Logging::LogLevel::Info);
awsiotsdk::util::Logging::InitializeAWSLogging(p_log_system);
m_disconnectContextData = std::shared_ptr<awsiotsdk::DisconnectCallbackContextData>(new DisconnectContext(this));
m_subscriptionContextData = std::shared_ptr<awsiotsdk::mqtt::SubscriptionHandlerContextData>(new SubscriptionContext(this));
qRegisterMetaType<AWSConnector::PushNotificationsEndpoint>(); qRegisterMetaType<AWSConnector::PushNotificationsEndpoint>();
m_clientName = readSyncedNameCache(); m_clientName = readSyncedNameCache();
} }
AWSConnector::~AWSConnector() AWSConnector::~AWSConnector()
{ {
qCDebug(dcAWS()) << "Stopping AWS connection. This might take a while..."; if (m_client) {
m_client.reset(); m_client->disconnectFromHost();
m_networkConnection.reset(); qCDebug(dcAWS()) << "Disconnected from AWS.";
qCDebug(dcAWS()) << "AWS connection stopped."; emit disconnected();
}
} }
void AWSConnector::connect2AWS(const QString &endpoint, const QString &clientId, const QString &clientName, const QString &caFile, const QString &clientCertFile, const QString &clientPrivKeyFile) void AWSConnector::connect2AWS(const QString &endpoint, const QString &clientId, const QString &clientName, const QString &caFile, const QString &clientCertFile, const QString &clientPrivKeyFile)
@ -71,8 +57,12 @@ void AWSConnector::connect2AWS(const QString &endpoint, const QString &clientId,
m_clientId = clientId; m_clientId = clientId;
m_clientName = clientName; m_clientName = clientName;
m_client.reset(); if (m_client) {
m_networkConnection.reset(); m_shouldReconnect = true;
m_client->disconnectFromHost();
qCDebug(dcAWS()) << "Disconnecting from AWS";
return;
}
doConnect(); doConnect();
} }
@ -81,35 +71,42 @@ void AWSConnector::doConnect()
{ {
m_setupInProgress = true; m_setupInProgress = true;
m_subscriptionCache.clear(); m_subscriptionCache.clear();
m_networkConnection = std::shared_ptr<OpenSSLConnection>(new OpenSSLConnection(
m_currentEndpoint.toStdString(),
8883,
m_caFile.toStdString(),
m_clientCertFile.toStdString(),
m_clientPrivKeyFile.toStdString(),
std::chrono::milliseconds(3000),
std::chrono::milliseconds(3000),
std::chrono::milliseconds(3000),
true
));
m_networkConnection->Initialize();
m_client = MqttClient::Create(m_networkConnection, std::chrono::milliseconds(2800), &onDisconnectedCallback, m_disconnectContextData);
m_client->SetAutoReconnectEnabled(false); QSslConfiguration sslConfig = QSslConfiguration::defaultConfiguration();
m_client->SetMaxReconnectBackoffTimeout(std::chrono::seconds(10)); QFile certFile(m_clientCertFile);
certFile.open(QFile::ReadOnly);
QSslCertificate certificate(certFile.readAll());
qCDebug(dcAWS()) << "Connecting to AWS with ID:" << m_clientId << "endpoint:" << m_currentEndpoint << "Certificate file:" << m_clientCertFile << "Fingerprint:" << getCertificateFingerprint(m_clientCertFile); QFile keyFile(m_clientPrivKeyFile);
m_connectingFuture = QtConcurrent::run([&]() { keyFile.open(QFile::ReadOnly);
ResponseCode rc = m_client->Connect(std::chrono::milliseconds(10000), true, mqtt::Version::MQTT_3_1_1, std::chrono::seconds(30), Utf8String::Create(m_clientId.toStdString()), nullptr, nullptr, nullptr); QSslKey key(keyFile.readAll(), QSsl::Rsa);
if (rc == ResponseCode::MQTT_CONNACK_CONNECTION_ACCEPTED) {
staticMetaObject.invokeMethod(this, "onConnected", Qt::QueuedConnection); sslConfig.setLocalCertificate(certificate);
} else { sslConfig.setPrivateKey(key);
qCWarning(dcAWS) << "Error connecting to AWS. Response code:" << QString::fromStdString(ResponseHelper::ToString(rc)); sslConfig.setPeerVerifyMode(QSslSocket::VerifyNone);
m_client.reset();
m_networkConnection.reset(); QFile caCertFile(m_caFile);
QTimer::singleShot(10000, this, &AWSConnector::doConnect); caCertFile.open(QFile::ReadOnly);
} QSslCertificate caCertificate(caCertFile.readAll());
sslConfig.setCaCertificates({caCertificate});
m_client = new QMQTT::Client(m_currentEndpoint, 8883, sslConfig, true, this);
m_client->setClientId(m_clientId);
m_client->setVersion(QMQTT::V3_1_1);
m_client->setKeepAlive(30*60);
m_client->setCleanSession(true);
m_client->connectToHost();
connect(m_client, &QMQTT::Client::connected, this, &AWSConnector::onConnected);
connect(m_client, &QMQTT::Client::disconnected, this, &AWSConnector::onDisconnected);
connect(m_client, &QMQTT::Client::error, this, [](const QMQTT::ClientError error){
qDebug() << "error" << error;
}); });
connect(m_client, &QMQTT::Client::subscribed, this, &AWSConnector::onSubscribed);
connect(m_client, &QMQTT::Client::received, this, &AWSConnector::onSubscriptionReceived);
connect(m_client, &QMQTT::Client::published, this, &AWSConnector::onPublished);
} }
void AWSConnector::onConnected() void AWSConnector::onConnected()
@ -135,7 +132,7 @@ void AWSConnector::registerDevice()
m_createDeviceId = QUuid::createUuid().toString().remove(QRegExp("[{}]*")); m_createDeviceId = QUuid::createUuid().toString().remove(QRegExp("[{}]*"));
// first subscribe to this tmp id topic // first subscribe to this tmp id topic
m_createDeviceSubscriptionId = subscribe({QString("create/device/%1").arg(m_createDeviceId)}); subscribe({QString("create/device/%1").arg(m_createDeviceId)});
} }
void AWSConnector::onDeviceRegistered(bool needsReconnect) void AWSConnector::onDeviceRegistered(bool needsReconnect)
@ -144,12 +141,7 @@ void AWSConnector::onDeviceRegistered(bool needsReconnect)
if (needsReconnect) { if (needsReconnect) {
qCDebug(dcAWS()) << "Disconnecting from AWS and reconnecting to use new policies"; qCDebug(dcAWS()) << "Disconnecting from AWS and reconnecting to use new policies";
QtConcurrent::run([&]() { m_client->disconnectFromHost();
m_client->Disconnect(std::chrono::milliseconds(500));
m_client.reset();
m_networkConnection.reset();
QTimer::singleShot(1000, this, &AWSConnector::doConnect);
});
return; return;
} }
@ -183,6 +175,9 @@ void AWSConnector::fetchPairings()
void AWSConnector::onPairingsRetrieved(const QVariantMap &pairings) void AWSConnector::onPairingsRetrieved(const QVariantMap &pairings)
{ {
m_setupInProgress = false;
emit connected();
qCDebug(dcAWS) << pairings.value("users").toList().count() << "devices paired in cloud."; qCDebug(dcAWS) << pairings.value("users").toList().count() << "devices paired in cloud.";
if (pairings.value("users").toList().count() > 0) { if (pairings.value("users").toList().count() > 0) {
QStringList topics; QStringList topics;
@ -209,14 +204,12 @@ void AWSConnector::onPairingsRetrieved(const QVariantMap &pairings)
} }
} }
} }
emit pushNotificationEndpointsUpdated(pushNotificationEndpoints);
if (readSyncedNameCache() != m_clientName) { if (readSyncedNameCache() != m_clientName) {
setName(); setName();
} }
m_setupInProgress = false; emit pushNotificationEndpointsUpdated(pushNotificationEndpoints);
emit connected();
requestTURNCredentials(); requestTURNCredentials();
} }
@ -225,17 +218,14 @@ void AWSConnector::disconnectAWS()
{ {
m_shouldReconnect = false; m_shouldReconnect = false;
if (isConnected()) { if (isConnected()) {
m_client->Disconnect(std::chrono::milliseconds(2000)); m_client->disconnectFromHost();
m_client.reset(); qCDebug(dcAWS()) << "Disconnecting from AWS.";
m_networkConnection.reset();
qCDebug(dcAWS()) << "Disconnected from AWS.";
emit disconnected();
} }
} }
bool AWSConnector::isConnected() const bool AWSConnector::isConnected() const
{ {
return m_connectingFuture.isFinished() && m_networkConnection && m_client && m_client->IsConnected() && !m_setupInProgress; return m_client && m_client->isConnectedToHost() && !m_setupInProgress;
} }
void AWSConnector::setDeviceName(const QString &deviceName) void AWSConnector::setDeviceName(const QString &deviceName)
@ -304,21 +294,19 @@ quint16 AWSConnector::publish(const QString &topic, const QVariantMap &message)
qCWarning(dcAWS()) << "Can't publish to AWS: Not connected."; qCWarning(dcAWS()) << "Can't publish to AWS: Not connected.";
return -1; return -1;
} }
QString fullTopic = topic;
QJsonDocument jsonDoc = QJsonDocument::fromVariant(message); QJsonDocument jsonDoc = QJsonDocument::fromVariant(message);
uint16_t packetId = 0; QMQTT::Message msg(0, topic, jsonDoc.toJson(QJsonDocument::Compact), 1);
ResponseCode res = m_client->PublishAsync(Utf8String::Create(fullTopic.toStdString()), false, false, mqtt::QoS::QOS1, jsonDoc.toJson(QJsonDocument::Compact).toStdString(), &publishCallback, packetId); qCDebug(dcAWSTraffic()) << "Publishing:" << topic << jsonDoc.toJson(QJsonDocument::Compact);
qCDebug(dcAWSTraffic()) << "publish call queued with status:" << QString::fromStdString(ResponseHelper::ToString(res)) << packetId << "for topic" << topic << jsonDoc.toJson(); quint16 packetId = m_client->publish(msg);
s_requestMap.insert(packetId, this);
return packetId; return packetId;
} }
void AWSConnector::onDisconnected() void AWSConnector::onDisconnected()
{ {
qCDebug(dcAWS) << "AWS disconnected."; qCDebug(dcAWS) << "AWS disconnected.";
m_client.reset(); m_client->deleteLater();
m_networkConnection.reset(); m_client = nullptr;
emit disconnected(); emit disconnected();
bool needReRegistering = false; bool needReRegistering = false;
@ -373,138 +361,99 @@ void AWSConnector::setName()
publish(QString("%1/device/name").arg(m_clientId), params); publish(QString("%1/device/name").arg(m_clientId), params);
} }
quint16 AWSConnector::subscribe(const QStringList &topics) void AWSConnector::subscribe(const QStringList &topics)
{ {
util::Vector<std::shared_ptr<mqtt::Subscription>> subscription_list;
foreach (const QString &topic, topics) { foreach (const QString &topic, topics) {
if (m_subscriptionCache.contains(topic)) { if (m_subscriptionCache.contains(topic)) {
qCDebug(dcAWS()) << "Already subscribed to topic:" << topic << ". Not resubscribing"; qCDebug(dcAWS()) << "Already subscribed to topic:" << topic << ". Not resubscribing";
continue; continue;
} }
qCDebug(dcAWSTraffic()) << "Topic to subscribe is" << topic; qCDebug(dcAWSTraffic()) << "Topic to subscribe is" << topic;
if (!Subscription::IsValidTopicName(topic.toStdString())) { m_client->subscribe(topic, 1);
qCWarning(dcAWS()) << "Trying to subscribe to invalid topic:" << topic;
continue;
}
auto subscription = mqtt::Subscription::Create(Utf8String::Create(topic.toStdString()), mqtt::QoS::QOS1, &onSubscriptionReceivedCallback, m_subscriptionContextData);
subscription_list.push_back(subscription);
m_subscriptionCache.append(topic); m_subscriptionCache.append(topic);
} }
if (subscription_list.size() == 0) {
return 0;
}
uint16_t packetId;
ResponseCode res = m_client->SubscribeAsync(subscription_list, subscribeCallback, packetId);
qCDebug(dcAWSTraffic()) << "Subscribe call queued with status:" << QString::fromStdString(ResponseHelper::ToString(res)) << "Packet ID:" << packetId;
s_requestMap.insert(packetId, this);
return packetId;
} }
void AWSConnector::publishCallback(uint16_t actionId, ResponseCode rc) void AWSConnector::onPublished(const QMQTT::Message& message, quint16 msgid)
{ {
AWSConnector* obj = s_requestMap.take(actionId); qCDebug(dcAWS()) << "Published message:" << message.topic() << msgid;
if (!obj) {
qCWarning(dcAWS())<< "Received a response callback but don't have an object waiting for it.";
return;
}
switch (rc) {
case ResponseCode::SUCCESS:
qCDebug(dcAWSTraffic()) << "Successfully published" << actionId;
break;
default:
qCDebug(dcAWS())<< "Error publishing data to AWS:" << QString::fromStdString(ResponseHelper::ToString(rc));
}
} }
void AWSConnector::subscribeCallback(uint16_t actionId, ResponseCode rc) void AWSConnector::onSubscribed(const QString& topic, const quint8 qos)
{ {
if (rc != ResponseCode::SUCCESS) { qCDebug(dcAWSTraffic()) << "Subscribed to topic:" << topic << qos;
qCWarning(dcAWS()) << "Error subscribing to" << actionId << QString::fromStdString(ResponseHelper::ToString(rc));
return;
}
AWSConnector *connector = s_requestMap.take(actionId); if (topic.startsWith("create/device/")) {
if (!connector) { qCDebug(dcAWS()) << "Subscribed to create/device/";
qCWarning(dcAWS()) << "Received a subscribe callback but don't have a request id for it.";
return;
}
if (actionId == connector->m_createDeviceSubscriptionId) {
qCDebug(dcAWS()) << "Subscribed to create/device/response";
// We might get this callback even if we didn't explicitly ask for it as the // We might get this callback even if we didn't explicitly ask for it as the
// library automatically resubscribes to all the topics upon reconnect. // library automatically resubscribes to all the topics upon reconnect.
if (!connector->readRegisteredFlag()) { if (!readRegisteredFlag()) {
QVariantMap params; QVariantMap params;
params.insert("id", connector->m_createDeviceId); params.insert("id", m_createDeviceId);
params.insert("UUID", connector->m_clientId); params.insert("UUID", m_clientId);
connector->publish("create/device", params); publish("create/device", params);
} }
return; return;
} }
qCDebug(dcAWSTraffic()) << "Successfully subscribed (actionId:" << actionId << ")";
} }
ResponseCode AWSConnector::onSubscriptionReceivedCallback(util::String topic_name, util::String payload, std::shared_ptr<SubscriptionHandlerContextData> p_app_handler_data) void AWSConnector::onSubscriptionReceived(const QMQTT::Message &message)
{ {
QJsonParseError error; QJsonParseError error;
QJsonDocument jsonDoc = QJsonDocument::fromJson(QByteArray::fromStdString(payload), &error); QJsonDocument jsonDoc = QJsonDocument::fromJson(message.payload(), &error);
if (error.error != QJsonParseError::NoError) { if (error.error != QJsonParseError::NoError) {
qCDebug(dcAWS()) << "Failed to parse JSON from AWS subscription on topic" << QString::fromStdString(topic_name) << ":" << error.errorString() << "\n" << QString::fromStdString(payload); qCDebug(dcAWS()) << "Failed to parse JSON from AWS subscription on topic" << message.topic() << ":" << error.errorString() << "\n" << message.payload();
return ResponseCode::JSON_PARSING_ERROR; return;
} }
qCDebug(dcAWSTraffic()) << "Subscription received: Topic:" << QString::fromStdString(topic_name) << "payload:" << QString::fromStdString(payload); qCDebug(dcAWSTraffic()) << "Subscription received: Topic:" << message.topic() << "payload:" << message.payload();
AWSConnector *connector = dynamic_cast<SubscriptionContext*>(p_app_handler_data.get())->c; QString topic = message.topic();
QString topic = QString::fromStdString(topic_name);
if (topic.startsWith("create/device/")) { if (topic.startsWith("create/device/")) {
int statusCode = jsonDoc.toVariant().toMap().value("result").toMap().value("code").toInt(); int statusCode = jsonDoc.toVariant().toMap().value("result").toMap().value("code").toInt();
switch (statusCode) { switch (statusCode) {
case 201: case 201:
qCDebug(dcAWS()) << "Device successfully registered to the cloud server:" << statusCode << jsonDoc.toVariant().toMap().value("result").toMap().value("message").toString(); qCDebug(dcAWS()) << "Device successfully registered to the cloud server:" << statusCode << jsonDoc.toVariant().toMap().value("result").toMap().value("message").toString();
connector->staticMetaObject.invokeMethod(connector, "onDeviceRegistered", Qt::QueuedConnection, Q_ARG(bool, true)); onDeviceRegistered(true);
return ResponseCode::SUCCESS; return;
case 200: case 200:
qCDebug(dcAWS()) << "Device already known to the cloud server:" << statusCode << jsonDoc.toVariant().toMap().value("result").toMap().value("message").toString(); qCDebug(dcAWS()) << "Device already known to the cloud server:" << statusCode << jsonDoc.toVariant().toMap().value("result").toMap().value("message").toString();
// Ok, we have confirmation that everything went fine and we can proceed, let's remember that to minimize traffic. // Ok, we have confirmation that everything went fine and we can proceed, let's remember that to minimize traffic.
connector->staticMetaObject.invokeMethod(connector, "onDeviceRegistered", Qt::QueuedConnection, Q_ARG(bool, false)); onDeviceRegistered(false);
break; break;
default: default:
qCWarning(dcAWS()) << "Error registering device in the cloud. AWS connetion will not work:" << statusCode << jsonDoc.toVariant().toMap().value("result").toMap().value("message").toString(); qCWarning(dcAWS()) << "Error registering device in the cloud. AWS connetion will not work:" << statusCode << jsonDoc.toVariant().toMap().value("result").toMap().value("message").toString();
return ResponseCode::SUCCESS; return;
} }
} else if (topic == QString("%1/pair/response").arg(connector->m_clientId)) { } else if (topic == QString("%1/pair/response").arg(m_clientId)) {
int statusCode = jsonDoc.toVariant().toMap().value("status").toInt(); int statusCode = jsonDoc.toVariant().toMap().value("status").toInt();
int id = jsonDoc.toVariant().toMap().value("id").toInt(); quint16 id = jsonDoc.toVariant().toMap().value("id").toUInt();
QString message = jsonDoc.toVariant().toMap().value("result").toMap().value("message").toString(); QString message = jsonDoc.toVariant().toMap().value("result").toMap().value("message").toString();
QString userId = connector->m_pairingRequests.take(id); QString userId = m_pairingRequests.take(id);
if (statusCode != 200) { if (statusCode != 200) {
qCWarning(dcAWS()) << "Pairing failed:" << statusCode << message; qCWarning(dcAWS()) << "Pairing failed:" << statusCode << message;
emit connector->devicePaired(userId, statusCode, message); emit devicePaired(userId, statusCode, message);
} else if (!userId.isEmpty()) { } else if (!userId.isEmpty()) {
qCDebug(dcAWS()) << "Pairing response for id:" << userId << statusCode; qCDebug(dcAWS()) << "Pairing response for id:" << userId << statusCode;
emit connector->devicePaired(userId, statusCode, message); emit devicePaired(userId, statusCode, message);
connector->staticMetaObject.invokeMethod(connector, "fetchPairings", Qt::QueuedConnection); fetchPairings();
} else { } else {
qCWarning(dcAWS()) << "Received a pairing response for a transaction we didn't start"; qCWarning(dcAWS()) << "Received a pairing response for a transaction we didn't start";
} }
} else if (topic == QString("%1/device/users/response").arg(connector->m_clientId)) { } else if (topic == QString("%1/device/users/response").arg(m_clientId)) {
connector->staticMetaObject.invokeMethod(connector, "onPairingsRetrieved", Qt::QueuedConnection, Q_ARG(QVariantMap, jsonDoc.toVariant().toMap())); onPairingsRetrieved(jsonDoc.toVariant().toMap());
} else if (topic == QString("%1/device/name/response").arg(connector->m_clientId)) { } else if (topic == QString("%1/device/name/response").arg(m_clientId)) {
qCDebug(dcAWS) << "Set device name in cloud with status:" << jsonDoc.toVariant().toMap().value("status").toInt(); qCDebug(dcAWS) << "Set device name in cloud with status:" << jsonDoc.toVariant().toMap().value("status").toInt();
if (jsonDoc.toVariant().toMap().value("status").toInt() == 200) { if (jsonDoc.toVariant().toMap().value("status").toInt() == 200) {
connector->storeSyncedNameCache(connector->m_clientName); storeSyncedNameCache(m_clientName);
} }
} else if (topic.startsWith(QString("%1/eu-west-1:").arg(connector->m_clientId)) && !topic.contains("reply") && !topic.contains("proxy")) { } else if (topic.startsWith(QString("%1/eu-west-1:").arg(m_clientId)) && !topic.contains("reply") && !topic.contains("proxy")) {
static QHash<QString, QDateTime> dupes; static QHash<QString, QDateTime> dupes;
QString id = jsonDoc.toVariant().toMap().value("id").toString(); QString id = jsonDoc.toVariant().toMap().value("id").toString();
QString type = jsonDoc.toVariant().toMap().value("type").toString(); QString type = jsonDoc.toVariant().toMap().value("type").toString();
if (dupes.contains(id+type)) { if (dupes.contains(id+type)) {
qCDebug(dcAWS()) << "Dropping duplicate packet"; qCDebug(dcAWS()) << "Dropping duplicate packet";
return ResponseCode::SUCCESS; return;
} }
dupes.insert(id+type, QDateTime::currentDateTime()); dupes.insert(id+type, QDateTime::currentDateTime());
foreach (const QString &dupe, dupes.keys()) { foreach (const QString &dupe, dupes.keys()) {
@ -513,17 +462,17 @@ ResponseCode AWSConnector::onSubscriptionReceivedCallback(util::String topic_nam
} }
} }
qCDebug(dcAWS) << "received webrtc handshake message."; qCDebug(dcAWS) << "received webrtc handshake message.";
connector->webRtcHandshakeMessageReceived(topic, jsonDoc.toVariant().toMap()); webRtcHandshakeMessageReceived(topic, jsonDoc.toVariant().toMap());
} else if (topic.startsWith(QString("%1/eu-west-1:").arg(connector->m_clientId)) && topic.contains("reply")) { } else if (topic.startsWith(QString("%1/eu-west-1:").arg(m_clientId)) && topic.contains("reply")) {
// silently drop our own things (should not be subscribed to that in the first place) // silently drop our own things (should not be subscribed to that in the first place)
} else if (topic.startsWith(QString("%1/eu-west-1:").arg(connector->m_clientId)) && topic.contains("proxy")) { } else if (topic.startsWith(QString("%1/eu-west-1:").arg(m_clientId)) && topic.contains("proxy")) {
QString token = jsonDoc.toVariant().toMap().value("token").toString(); QString token = jsonDoc.toVariant().toMap().value("token").toString();
qlonglong timestamp = jsonDoc.toVariant().toMap().value("timestamp").toLongLong(); qlonglong timestamp = jsonDoc.toVariant().toMap().value("timestamp").toLongLong();
static QHash<QString, QDateTime> dupes; static QHash<QString, QDateTime> dupes;
QString packetId = topic + token + QString::number(timestamp); QString packetId = topic + token + QString::number(timestamp);
if (dupes.contains(packetId)) { if (dupes.contains(packetId)) {
qCDebug(dcAWS()) << "Dropping duplicate packet"; qCDebug(dcAWS()) << "Dropping duplicate packet";
return ResponseCode::SUCCESS; return;
} }
dupes.insert(packetId, QDateTime::currentDateTime()); dupes.insert(packetId, QDateTime::currentDateTime());
foreach (const QString &dupe, dupes.keys()) { foreach (const QString &dupe, dupes.keys()) {
@ -532,13 +481,13 @@ ResponseCode AWSConnector::onSubscriptionReceivedCallback(util::String topic_nam
} }
} }
qCDebug(dcAWS) << "Proxy remote connection request received"; qCDebug(dcAWS) << "Proxy remote connection request received";
connector->staticMetaObject.invokeMethod(connector, "proxyConnectionRequestReceived", Qt::QueuedConnection, Q_ARG(QString, token)); proxyConnectionRequestReceived(token);
} else if (topic == QString("%1/notify/response").arg(connector->m_clientId)) { } else if (topic == QString("%1/notify/response").arg(m_clientId)) {
int transactionId = jsonDoc.toVariant().toMap().value("id").toInt(); int transactionId = jsonDoc.toVariant().toMap().value("id").toInt();
int status = jsonDoc.toVariant().toMap().value("status").toInt(); int status = jsonDoc.toVariant().toMap().value("status").toInt();
qCDebug(dcAWS()) << "Push notification reply for transaction" << transactionId << " Status:" << status << jsonDoc.toVariant().toMap().value("message").toString(); qCDebug(dcAWS()) << "Push notification reply for transaction" << transactionId << " Status:" << status << jsonDoc.toVariant().toMap().value("message").toString();
emit connector->pushNotificationSent(transactionId, status); emit pushNotificationSent(transactionId, status);
} else if (topic == QString("%1/notify/info/endpoint").arg(connector->m_clientId)) { } else if (topic == QString("%1/notify/info/endpoint").arg(m_clientId)) {
QVariantMap endpoint = jsonDoc.toVariant().toMap().value("newPushNotificationsEndpoint").toMap(); QVariantMap endpoint = jsonDoc.toVariant().toMap().value("newPushNotificationsEndpoint").toMap();
Q_ASSERT(endpoint.keys().count() == 1); Q_ASSERT(endpoint.keys().count() == 1);
QString cognitoId = endpoint.keys().first(); QString cognitoId = endpoint.keys().first();
@ -546,26 +495,17 @@ ResponseCode AWSConnector::onSubscriptionReceivedCallback(util::String topic_nam
ep.userId = cognitoId; ep.userId = cognitoId;
ep.endpointId = endpoint.value(cognitoId).toMap().value("endpointId").toString(); ep.endpointId = endpoint.value(cognitoId).toMap().value("endpointId").toString();
ep.displayName = endpoint.value(cognitoId).toMap().value("displayName").toString(); ep.displayName = endpoint.value(cognitoId).toMap().value("displayName").toString();
emit connector->pushNotificationEndpointAdded(ep); emit pushNotificationEndpointAdded(ep);
} else if (topic == QString("%1/services/turn/response").arg(connector->m_clientId)) { } else if (topic == QString("%1/services/turn/response").arg(m_clientId)) {
QVariantMap turnCreds = jsonDoc.toVariant().toMap(); QVariantMap turnCreds = jsonDoc.toVariant().toMap();
if (turnCreds.value("result").toMap().value("code").toInt() != 201) { if (turnCreds.value("result").toMap().value("code").toInt() != 201) {
qCWarning(dcAWS()) << "Error retrieving TURN credentials:" << turnCreds.value("result").toMap().value("code").toInt() << turnCreds.value("result").toMap().value("message").toString(); qCWarning(dcAWS()) << "Error retrieving TURN credentials:" << turnCreds.value("result").toMap().value("code").toInt() << turnCreds.value("result").toMap().value("message").toString();
return ResponseCode::SUCCESS; return;
} }
connector->staticMetaObject.invokeMethod(connector, "onTurnCredentialsReceived", Qt::QueuedConnection, Q_ARG(QVariantMap, turnCreds.value("turnCredentials").toMap())); onTurnCredentialsReceived(turnCreds.value("turnCredentials").toMap());
} else { } else {
qCWarning(dcAWS()) << "Unhandled subscription received!" << topic << QString::fromStdString(payload); qCWarning(dcAWS()) << "Unhandled subscription received!" << topic << message.payload();
} }
return ResponseCode::SUCCESS;
}
ResponseCode AWSConnector::onDisconnectedCallback(util::String mqtt_client_id, std::shared_ptr<DisconnectCallbackContextData> p_app_handler_data)
{
Q_UNUSED(mqtt_client_id)
AWSConnector* connector = dynamic_cast<DisconnectContext*>(p_app_handler_data.get())->c;
connector->staticMetaObject.invokeMethod(connector, "onDisconnected", Qt::QueuedConnection);
return ResponseCode::SUCCESS;
} }
void AWSConnector::storeRegisteredFlag(bool registered) void AWSConnector::storeRegisteredFlag(bool registered)
@ -596,7 +536,7 @@ QString AWSConnector::getCertificateFingerprint(const QString &certificateFile)
{ {
QFile certFile(certificateFile); QFile certFile(certificateFile);
if (!certFile.open(QFile::ReadOnly)) { if (!certFile.open(QFile::ReadOnly)) {
qCWarning(dcAWS()) << "Error opening certificate file" << certificateFile; qCWarning(dcAWS()) << "Error openi<ng certificate file" << certificateFile;
return QString(); return QString();
} }
QSslCertificate crt = QSslCertificate(certFile.readAll()); QSslCertificate crt = QSslCertificate(certFile.readAll());

View File

@ -25,18 +25,20 @@
#include <QFuture> #include <QFuture>
#include <QDateTime> #include <QDateTime>
#include "OpenSSL/OpenSSLConnection.hpp" #include <qmqtt/qmqtt.h>
#include <mqtt/Client.hpp>
#include <mqtt/Common.hpp>
#include "util/logging/Logging.hpp"
#include "util/logging/LogMacros.hpp"
#include "util/logging/ConsoleLogSystem.hpp"
class AWSConnector : public QObject, public awsiotsdk::mqtt::SubscriptionHandlerContextData, public awsiotsdk::DisconnectCallbackContextData //#include "OpenSSL/OpenSSLConnection.hpp"
//#include <mqtt/Client.hpp>
//#include <mqtt/Common.hpp>
//#include "util/logging/Logging.hpp"
//#include "util/logging/LogMacros.hpp"
//#include "util/logging/ConsoleLogSystem.hpp"
class AWSConnector : public QObject
{ {
Q_OBJECT Q_OBJECT
public: public:
explicit AWSConnector(QObject *parent = 0); explicit AWSConnector(QObject *parent = nullptr);
~AWSConnector(); ~AWSConnector();
class PushNotificationsEndpoint { class PushNotificationsEndpoint {
@ -74,36 +76,29 @@ signals:
private slots: private slots:
void doConnect(); void doConnect();
void onConnected(); void onConnected();
void onDisconnected();
void onPublished(const QMQTT::Message &message, quint16 msgid);
void onSubscribed(const QString& topic, const quint8 qos);
void onSubscriptionReceived(const QMQTT::Message &message);
void registerDevice(); void registerDevice();
void onDeviceRegistered(bool needsReconnect); void onDeviceRegistered(bool needsReconnect);
void setupSubscriptions(); void setupSubscriptions();
void fetchPairings(); void fetchPairings();
void onPairingsRetrieved(const QVariantMap &pairings); void onPairingsRetrieved(const QVariantMap &pairings);
void setName(); void setName();
void onDisconnected();
void onTurnCredentialsReceived(const QVariantMap &turnCredentials); void onTurnCredentialsReceived(const QVariantMap &turnCredentials);
private: private:
class SubscriptionContext: public awsiotsdk::mqtt::SubscriptionHandlerContextData
{
public:
SubscriptionContext(AWSConnector *connector): c(connector) {}
AWSConnector *c;
};
class DisconnectContext: public awsiotsdk::DisconnectCallbackContextData
{
public:
DisconnectContext(AWSConnector *connector): c(connector) {}
AWSConnector *c;
};
quint16 publish(const QString &topic, const QVariantMap &message); quint16 publish(const QString &topic, const QVariantMap &message);
quint16 subscribe(const QStringList &topics); void subscribe(const QStringList &topics);
static void publishCallback(uint16_t actionId, awsiotsdk::ResponseCode rc); // static void publishCallback(uint16_t actionId, awsiotsdk::ResponseCode rc);
static void subscribeCallback(uint16_t actionId, awsiotsdk::ResponseCode rc); // static void subscribeCallback(uint16_t actionId, awsiotsdk::ResponseCode rc);
static awsiotsdk::ResponseCode onSubscriptionReceivedCallback(awsiotsdk::util::String topic_name, awsiotsdk::util::String payload, // static awsiotsdk::ResponseCode onSubscriptionReceivedCallback(awsiotsdk::util::String topic_name, awsiotsdk::util::String payload,
std::shared_ptr<awsiotsdk::mqtt::SubscriptionHandlerContextData> p_app_handler_data); // std::shared_ptr<awsiotsdk::mqtt::SubscriptionHandlerContextData> p_app_handler_data);
static awsiotsdk::ResponseCode onDisconnectedCallback(awsiotsdk::util::String mqtt_client_id, // static awsiotsdk::ResponseCode onDisconnectedCallback(awsiotsdk::util::String mqtt_client_id,
std::shared_ptr<awsiotsdk::DisconnectCallbackContextData> p_app_handler_data); // std::shared_ptr<awsiotsdk::DisconnectCallbackContextData> p_app_handler_data);
void storeRegisteredFlag(bool registered); void storeRegisteredFlag(bool registered);
bool readRegisteredFlag() const; bool readRegisteredFlag() const;
@ -114,8 +109,9 @@ private:
QString getCertificateFingerprint(const QString &certificateFilePath) const; QString getCertificateFingerprint(const QString &certificateFilePath) const;
private: private:
std::shared_ptr<awsiotsdk::network::OpenSSLConnection> m_networkConnection; // std::shared_ptr<awsiotsdk::network::OpenSSLConnection> m_networkConnection;
std::shared_ptr<awsiotsdk::MqttClient> m_client; // std::shared_ptr<awsiotsdk::MqttClient> m_client;
QMQTT::Client *m_client = nullptr;
QString m_currentEndpoint; QString m_currentEndpoint;
QString m_caFile; QString m_caFile;
QString m_clientCertFile; QString m_clientCertFile;
@ -137,11 +133,11 @@ private:
QStringList m_subscriptionCache; QStringList m_subscriptionCache;
QPair<QVariantMap, QDateTime> m_cachedTURNCredentials; QPair<QVariantMap, QDateTime> m_cachedTURNCredentials;
std::shared_ptr<awsiotsdk::mqtt::SubscriptionHandlerContextData> m_subscriptionContextData; // std::shared_ptr<awsiotsdk::mqtt::SubscriptionHandlerContextData> m_subscriptionContextData;
std::shared_ptr<awsiotsdk::DisconnectCallbackContextData> m_disconnectContextData; // std::shared_ptr<awsiotsdk::DisconnectCallbackContextData> m_disconnectContextData;
static AWSConnector* s_instance; // static AWSConnector* s_instance;
static QHash<quint16, AWSConnector*> s_requestMap; // QHash<quint16, AWSConnector*> m_requestMap;
}; };
Q_DECLARE_METATYPE(AWSConnector::PushNotificationsEndpoint) Q_DECLARE_METATYPE(AWSConnector::PushNotificationsEndpoint)

View File

@ -3,7 +3,7 @@ TARGET = nymea-core
include(../nymea.pri) include(../nymea.pri)
QT += sql QT += sql qmqtt
INCLUDEPATH += $$top_srcdir/libnymea INCLUDEPATH += $$top_srcdir/libnymea
LIBS += -L$$top_builddir/libnymea/ -lnymea -lssl -lcrypto -lavahi-common -lavahi-client LIBS += -L$$top_builddir/libnymea/ -lnymea -lssl -lcrypto -lavahi-common -lavahi-client
@ -97,7 +97,7 @@ HEADERS += nymeacore.h \
tagging/tagsstorage.h \ tagging/tagsstorage.h \
tagging/tag.h \ tagging/tag.h \
jsonrpc/tagshandler.h \ jsonrpc/tagshandler.h \
cloud/cloudtransport.h cloud/cloudtransport.h \
SOURCES += nymeacore.cpp \ SOURCES += nymeacore.cpp \
tcpserver.cpp \ tcpserver.cpp \
@ -154,7 +154,7 @@ SOURCES += nymeacore.cpp \
cloud/awsconnector.cpp \ cloud/awsconnector.cpp \
cloud/cloudmanager.cpp \ cloud/cloudmanager.cpp \
cloud/cloudnotifications.cpp \ cloud/cloudnotifications.cpp \
cloud/OpenSSL/OpenSSLConnection.cpp \ # cloud/OpenSSL/OpenSSLConnection.cpp \
cloud/janusconnector.cpp \ cloud/janusconnector.cpp \
pushbuttondbusservice.cpp \ pushbuttondbusservice.cpp \
hardwaremanagerimplementation.cpp \ hardwaremanagerimplementation.cpp \
@ -179,4 +179,4 @@ SOURCES += nymeacore.cpp \
tagging/tagsstorage.cpp \ tagging/tagsstorage.cpp \
tagging/tag.cpp \ tagging/tag.cpp \
jsonrpc/tagshandler.cpp \ jsonrpc/tagshandler.cpp \
cloud/cloudtransport.cpp cloud/cloudtransport.cpp \

View File

@ -8,11 +8,11 @@ INCLUDEPATH += ../libnymea ../libnymea-core
target.path = /usr/bin target.path = /usr/bin
INSTALLS += target INSTALLS += target
QT *= sql xml websockets bluetooth dbus network QT *= sql xml websockets bluetooth dbus network qmqtt
LIBS += -L$$top_builddir/libnymea/ -lnymea \ LIBS += -L$$top_builddir/libnymea/ -lnymea \
-L$$top_builddir/libnymea-core -lnymea-core \ -L$$top_builddir/libnymea-core -lnymea-core \
-lssl -lcrypto -laws-iot-sdk-cpp -lnymea-remoteproxyclient -lssl -lcrypto -lnymea-remoteproxyclient
# Server files # Server files
include(qtservice/qtservice.pri) include(qtservice/qtservice.pri)

View File

@ -8,7 +8,7 @@ INCLUDEPATH += $$top_srcdir/libnymea \
LIBS += -L$$top_builddir/libnymea/ -lnymea \ LIBS += -L$$top_builddir/libnymea/ -lnymea \
-L$$top_builddir/libnymea-core/ -lnymea-core \ -L$$top_builddir/libnymea-core/ -lnymea-core \
-L$$top_builddir/plugins/mock/ \ -L$$top_builddir/plugins/mock/ \
-lssl -lcrypto -laws-iot-sdk-cpp -lavahi-common -lavahi-client -lnymea-remoteproxyclient -lssl -lcrypto -lavahi-common -lavahi-client -lnymea-remoteproxyclient
SOURCES += ../nymeatestbase.cpp \ SOURCES += ../nymeatestbase.cpp \