Début implémentation buffering des messages MQTT lors d'un blackout internet
This commit is contained in:
@@ -12,7 +12,15 @@ CMQTTClientWrapper::CMQTTClientWrapper()
|
||||
mMQTTReconnectTimer = new QTimer;
|
||||
mMQTTReconnectTimer->setSingleShot(true);
|
||||
connect(mMQTTReconnectTimer,&QTimer::timeout,this,&CMQTTClientWrapper::MQTTReconnectTimerExpired);
|
||||
|
||||
mMQTTQueueFlushTimer = new QTimer;
|
||||
mMQTTQueueFlushTimer->setSingleShot(false);
|
||||
connect(mMQTTQueueFlushTimer,&QTimer::timeout,this,&CMQTTClientWrapper::MQTTQueueFlushTimerExipred);
|
||||
mMQTTQueueFlushTimer->stop();
|
||||
|
||||
mProgramPtr = 0;
|
||||
mMessagesQueueMode = MQTT_DROP_MSG_MODE;
|
||||
mDisconnectionIsVoluntary = false;
|
||||
}
|
||||
|
||||
CMQTTClientWrapper::~CMQTTClientWrapper()
|
||||
@@ -37,6 +45,7 @@ int CMQTTClientWrapper::ConnectToBroker()
|
||||
mMQTTClient.setPort(mMQTTParams.mMQTTBrokerPort);
|
||||
mMQTTClient.setPassword(mMQTTParams.mMQTTBrokerPassword);
|
||||
mMQTTClient.setUsername(mMQTTParams.mMQTTBrokerUserName);
|
||||
mDisconnectionIsVoluntary = false;
|
||||
|
||||
mMQTTClient.connectToHost();
|
||||
|
||||
@@ -46,6 +55,9 @@ int CMQTTClientWrapper::ConnectToBroker()
|
||||
int CMQTTClientWrapper::DisconnectFromBroker()
|
||||
{
|
||||
mMQTTClient.disconnectFromHost();
|
||||
mDisconnectionIsVoluntary = true;
|
||||
mMessagesQueueMode = MQTT_DROP_MSG_MODE; //It's a voluntary disconnection... don't queue the CAN messages.
|
||||
|
||||
return RET_OK;
|
||||
}
|
||||
|
||||
@@ -89,6 +101,11 @@ void CMQTTClientWrapper::StateChanged()
|
||||
mProgramPtr->SetMQTTConnectionSatusRequest(false);
|
||||
mMQTTRefreshTimer->stop();
|
||||
mMQTTReconnectTimer->start(MQTT_CLIENT_RECONNECT_TIMEOUT);
|
||||
mMQTTQueueFlushTimer->stop();
|
||||
if(mDisconnectionIsVoluntary == false)
|
||||
{
|
||||
mMessagesQueueMode = MQTT_QUEUE_MSG_MODE; //We're disconnected, queue all the messages.
|
||||
}
|
||||
break;
|
||||
}
|
||||
case QMqttClient::Connected:
|
||||
@@ -96,6 +113,16 @@ void CMQTTClientWrapper::StateChanged()
|
||||
mProgramPtr->SetMQTTConnectionSatusRequest(true);
|
||||
mMQTTRefreshTimer->start(mMQTTParams.mMQTTTransmitTimeout);
|
||||
mMQTTReconnectTimer->stop();
|
||||
if(mMQTTMessagesQueue.isEmpty() == false)
|
||||
{
|
||||
mMessagesQueueMode = MQTT_QUEUE_MSG_MODE; //Stay in (or enter) queue mode until we empty the buffer
|
||||
mMQTTQueueFlushTimer->start();
|
||||
}
|
||||
else
|
||||
{
|
||||
mMessagesQueueMode = MQTT_TRANSMIT_MSG_MODE;
|
||||
}
|
||||
|
||||
qDebug("MQTT client Connected");
|
||||
break;
|
||||
}
|
||||
@@ -122,7 +149,7 @@ int CMQTTClientWrapper::SetCANDevicesList(QList<CCANDevice *> *List)
|
||||
|
||||
void CMQTTClientWrapper::MQTTSendTimerExpired()
|
||||
{
|
||||
if(mMQTTClient.state() != QMqttClient::Connected)
|
||||
if(mMessagesQueueMode == MQTT_DROP_MSG_MODE)
|
||||
{
|
||||
return;
|
||||
}
|
||||
@@ -134,19 +161,41 @@ void CMQTTClientWrapper::MQTTSendTimerExpired()
|
||||
//Send the CANbus devices messsages
|
||||
for(int j = 0; j < mCANDevicesList->size(); j++)
|
||||
{
|
||||
CCANDevice *Device = mCANDevicesList->at(j);
|
||||
QList<CMQTTMessage> *MessagesList = Device->GetMQTTMessagesList();
|
||||
if(MessagesList != 0)
|
||||
switch(mMessagesQueueMode)
|
||||
{
|
||||
for(int i = 0; i < MessagesList->size(); i++)
|
||||
case MQTT_TRANSMIT_MSG_MODE:
|
||||
{
|
||||
qint32 res = mMQTTClient.publish(MessagesList->at(i).mMessageTopic,MessagesList->at(i).mMessagePayload.toLocal8Bit(),0,true);
|
||||
qDebug("%s : %s",qPrintable(MessagesList->at(i).mMessageTopic), qPrintable(MessagesList->at(i).mMessagePayload));
|
||||
QString LogMsg = QString("Envoi d'un message MQTT. Topic: %1 Payload: %2").arg(MessagesList->at(i).mMessageTopic).arg(MessagesList->at(i).mMessagePayload);
|
||||
CGeneralMessagesLogDispatcher::instance()->AddLogMessage(LogMsg,true,3);
|
||||
CCANDevice *Device = mCANDevicesList->at(j);
|
||||
QList<CMQTTMessage> *MessagesList = Device->GetMQTTMessagesList();
|
||||
if(MessagesList != 0)
|
||||
{
|
||||
for(int i = 0; i < MessagesList->size(); i++)
|
||||
{
|
||||
qint32 res = mMQTTClient.publish(MessagesList->at(i).mMessageTopic,MessagesList->at(i).mMessagePayload.toLocal8Bit(),0,true);
|
||||
qDebug("%s : %s",qPrintable(MessagesList->at(i).mMessageTopic), qPrintable(MessagesList->at(i).mMessagePayload));
|
||||
QString LogMsg = QString("Envoi d'un message MQTT. Topic: %1 Payload: %2").arg(MessagesList->at(i).mMessageTopic).arg(MessagesList->at(i).mMessagePayload);
|
||||
CGeneralMessagesLogDispatcher::instance()->AddLogMessage(LogMsg,true,3);
|
||||
|
||||
}
|
||||
qDebug("Sent %d MQTT messages",MessagesList->size());
|
||||
}
|
||||
break;
|
||||
}
|
||||
case MQTT_QUEUE_MSG_MODE:
|
||||
{
|
||||
CCANDevice *Device = mCANDevicesList->at(j);
|
||||
QList<CMQTTMessage> *MessagesList = Device->GetMQTTMessagesList();
|
||||
for(int i = 0; i < MessagesList->size(); i++)
|
||||
{
|
||||
CMQTTMessage *NewMsg = new CMQTTMessage(MessagesList->at(i).mMessageTopic,MessagesList->at(i).mMessagePayload);
|
||||
if(mMQTTMessagesQueue.size() >= MQTT_CLIENT_MSG_QUEUE_SIZE)
|
||||
{
|
||||
delete mMQTTMessagesQueue.takeFirst();
|
||||
}
|
||||
mMQTTMessagesQueue.append(NewMsg);
|
||||
}
|
||||
break;
|
||||
}
|
||||
qDebug("Sent %d MQTT messages",MessagesList->size());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -158,3 +207,28 @@ void CMQTTClientWrapper::MQTTReconnectTimerExpired()
|
||||
{
|
||||
ConnectToBroker();
|
||||
}
|
||||
|
||||
void CMQTTClientWrapper::MQTTQueueFlushTimerExipred()
|
||||
{
|
||||
if(mMQTTMessagesQueue.isEmpty()) //Shouldn't happen... but just to be safe
|
||||
{
|
||||
mMQTTQueueFlushTimer->stop();
|
||||
mMessagesQueueMode = MQTT_TRANSMIT_MSG_MODE;
|
||||
return;
|
||||
}
|
||||
|
||||
CMQTTMessage *Msg = mMQTTMessagesQueue.takeFirst();
|
||||
qint32 res = mMQTTClient.publish(Msg->mMessageTopic,Msg->mMessagePayload.toLocal8Bit(),0,true);
|
||||
qDebug("Flushing MQTT Msg queue... %s : %s",qPrintable(Msg->mMessageTopic), qPrintable(Msg->mMessagePayload));
|
||||
QString LogMsg = QString("Envoi d'un message MQTT de la queue. Topic: %1 Payload: %2").arg(Msg->mMessageTopic).arg(Msg->mMessagePayload);
|
||||
CGeneralMessagesLogDispatcher::instance()->AddLogMessage(LogMsg,true,3);
|
||||
|
||||
delete Msg;
|
||||
|
||||
if(mMQTTMessagesQueue.isEmpty())
|
||||
{
|
||||
mMQTTQueueFlushTimer->stop();
|
||||
mMessagesQueueMode = MQTT_TRANSMIT_MSG_MODE;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -17,6 +17,13 @@ class CMQTTClientWrapper : public QObject
|
||||
{
|
||||
Q_OBJECT
|
||||
public:
|
||||
typedef enum
|
||||
{
|
||||
MQTT_DROP_MSG_MODE,
|
||||
MQTT_TRANSMIT_MSG_MODE,
|
||||
MQTT_QUEUE_MSG_MODE
|
||||
}eMQTTMsgQueuingMode;
|
||||
|
||||
CMQTTClientWrapper();
|
||||
~CMQTTClientWrapper();
|
||||
int SetMQTTParams(CCloudParams *Params);
|
||||
@@ -30,17 +37,22 @@ public:
|
||||
// QString mMQTTClientID;
|
||||
QTimer *mMQTTRefreshTimer;
|
||||
QTimer *mMQTTReconnectTimer;
|
||||
QTimer *mMQTTQueueFlushTimer;
|
||||
eMQTTMsgQueuingMode mMessagesQueueMode;
|
||||
bool mDisconnectionIsVoluntary;
|
||||
|
||||
|
||||
private:
|
||||
QMqttClient mMQTTClient;
|
||||
CCloudParams mMQTTParams;
|
||||
QList<CCANDevice*> *mCANDevicesList;
|
||||
QList<CMQTTMessage*> mMQTTMessagesQueue;
|
||||
|
||||
public slots:
|
||||
void StateChanged();
|
||||
void MQTTSendTimerExpired();
|
||||
void MQTTReconnectTimerExpired();
|
||||
void MQTTQueueFlushTimerExipred();
|
||||
|
||||
};
|
||||
|
||||
|
||||
Reference in New Issue
Block a user