Buffering MQTT fonctionnel

This commit is contained in:
2023-06-29 09:16:50 -04:00
parent 3a44111079
commit b631eab073
11 changed files with 778 additions and 28 deletions
@@ -15,12 +15,14 @@ CMQTTClientWrapper::CMQTTClientWrapper()
mMQTTQueueFlushTimer = new QTimer;
mMQTTQueueFlushTimer->setSingleShot(false);
mMQTTQueueFlushTimer->setInterval(MQTT_CLIENT_MSG_QUEUE_FLUSH_TIMEOUT);
connect(mMQTTQueueFlushTimer,&QTimer::timeout,this,&CMQTTClientWrapper::MQTTQueueFlushTimerExipred);
mMQTTQueueFlushTimer->stop();
mProgramPtr = 0;
mMessagesQueueMode = MQTT_DROP_MSG_MODE;
mDisconnectionIsVoluntary = false;
mIsClientConnecting = false;
}
CMQTTClientWrapper::~CMQTTClientWrapper()
@@ -36,6 +38,14 @@ int CMQTTClientWrapper::SetMQTTParams(CCloudParams *Params)
return RET_OK;
}
int CMQTTClientWrapper::StartMQTTClient()
{
mMQTTRefreshTimer->start(mMQTTParams.mMQTTTransmitTimeout);
mMessagesQueueMode = MQTT_QUEUE_MSG_MODE;
ConnectToBroker();
return RET_OK;
}
int CMQTTClientWrapper::ConnectToBroker()
{
//Setup the client before connecting.
@@ -97,37 +107,46 @@ void CMQTTClientWrapper::StateChanged()
{
case QMqttClient::Disconnected:
{
qDebug("MQTT client Disconnected");
mProgramPtr->SetMQTTConnectionSatusRequest(false);
mMQTTRefreshTimer->stop();
mMQTTReconnectTimer->start(MQTT_CLIENT_RECONNECT_TIMEOUT);
mMQTTQueueFlushTimer->stop();
if(mDisconnectionIsVoluntary == false)
if(mIsClientConnecting)
{
mMessagesQueueMode = MQTT_QUEUE_MSG_MODE; //We're disconnected, queue all the messages.
//Connection attempt failed, just restart the timer...
mMQTTReconnectTimer->start(MQTT_CLIENT_RECONNECT_TIMEOUT);
}
else
{
CGeneralMessagesLogDispatcher::instance()->AddLogMessage("Client MQTT déconnecté.",true,1);
mProgramPtr->SetMQTTConnectionSatusRequest(false);
mMQTTReconnectTimer->start(MQTT_CLIENT_RECONNECT_TIMEOUT);
mMQTTQueueFlushTimer->stop();
if(mDisconnectionIsVoluntary == false)
{
mMessagesQueueMode = MQTT_QUEUE_MSG_MODE; //We're disconnected, queue all the messages.
CGeneralMessagesLogDispatcher::instance()->AddLogMessage("Passage en mode buffering des messages MQTT",true,2);
}
}
break;
}
case QMqttClient::Connected:
{
mProgramPtr->SetMQTTConnectionSatusRequest(true);
mMQTTRefreshTimer->start(mMQTTParams.mMQTTTransmitTimeout);
mMQTTReconnectTimer->stop();
mIsClientConnecting = false;
CGeneralMessagesLogDispatcher::instance()->AddLogMessage("Client MQTT connecté.",true,1);
if(mMQTTMessagesQueue.isEmpty() == false)
{
mMessagesQueueMode = MQTT_QUEUE_MSG_MODE; //Stay in (or enter) queue mode until we empty the buffer
mMQTTQueueFlushTimer->start();
CGeneralMessagesLogDispatcher::instance()->AddLogMessage("FIFO non vide, passage au mode de vidage de la FIFO",true,2);
}
else
{
mMessagesQueueMode = MQTT_TRANSMIT_MSG_MODE;
}
qDebug("MQTT client Connected");
break;
}
case QMqttClient::Connecting:
{
mIsClientConnecting = true;
qDebug("MQTT client Connecting...");
break;
}
@@ -172,7 +191,6 @@ void CMQTTClientWrapper::MQTTSendTimerExpired()
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);
@@ -185,14 +203,19 @@ void CMQTTClientWrapper::MQTTSendTimerExpired()
{
CCANDevice *Device = mCANDevicesList->at(j);
QList<CMQTTMessage> *MessagesList = Device->GetMQTTMessagesList();
for(int i = 0; i < MessagesList->size(); i++)
if(MessagesList != 0)
{
CMQTTMessage *NewMsg = new CMQTTMessage(MessagesList->at(i).mMessageTopic,MessagesList->at(i).mMessagePayload);
if(mMQTTMessagesQueue.size() >= MQTT_CLIENT_MSG_QUEUE_SIZE)
for(int i = 0; i < MessagesList->size(); i++)
{
delete mMQTTMessagesQueue.takeFirst();
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);
QString LogMsg = QString("Ajout d'un message MQTT à la FIFO. Topic: %1 Payload: %2 FIFO size: %3").arg(MessagesList->at(i).mMessageTopic).arg(MessagesList->at(i).mMessagePayload).arg(mMQTTMessagesQueue.size());
CGeneralMessagesLogDispatcher::instance()->AddLogMessage(LogMsg,true,3);
}
mMQTTMessagesQueue.append(NewMsg);
}
break;
}
@@ -219,16 +242,16 @@ void CMQTTClientWrapper::MQTTQueueFlushTimerExipred()
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);
QString LogMsg = QString("Envoi d'un message MQTT provenant du buffer. Topic: %1 Payload: %2 Buffer Size: %3").arg(Msg->mMessageTopic).arg(Msg->mMessagePayload).arg(mMQTTMessagesQueue.size());
CGeneralMessagesLogDispatcher::instance()->AddLogMessage(LogMsg,true,3);
delete Msg;
delete Msg; //free memory
if(mMQTTMessagesQueue.isEmpty())
{
mMQTTQueueFlushTimer->stop();
mMessagesQueueMode = MQTT_TRANSMIT_MSG_MODE;
CGeneralMessagesLogDispatcher::instance()->AddLogMessage("Tous les messages MQTT de la FIFO ont été envoyés au serveur",true,2);
}
}