Multi brokers MQTT semble fonctionner

This commit is contained in:
2025-05-03 15:01:06 -04:00
parent 9e43e9c748
commit e298577c2f
17 changed files with 780 additions and 558 deletions
+1 -11
View File
@@ -4,8 +4,7 @@
CCANDataLogger::CCANDataLogger():
mTopicDeviceString("")
{
mMQTTCLient = 0;
mCANMsgList.clear();
mCANMsgList.clear();
}
@@ -51,11 +50,6 @@ int CCANDataLogger::SetMQTTTopicDevice(QString DeviceString)
return RET_OK;
}
int CCANDataLogger::SetMQTTClient(CMQTTClientWrapper *MQTTClient)
{
mMQTTCLient = MQTTClient;
return RET_OK;
}
QList<CMQTTMessage> *CCANDataLogger::GetMQTTMessagesList()
{
@@ -115,10 +109,6 @@ QList<CMQTTMessage> *CCANDataLogger::GetMQTTMessagesList()
mCANMsgList.clear();
// if(mMQTTCLient != 0)
// {
// mMQTTCLient->NewMQTTMessages(mMQTTMsgList);
// }
return &mMQTTMsgList;
}
@@ -18,12 +18,10 @@ public:
//MQTT logging
QString mTopicDeviceString;
QList<CMQTTMessage> mMQTTMsgList;
CMQTTClientWrapper *mMQTTCLient;
QList<CCANMessage> mCANMsgList;
CMQTTMessage GetMQTTMessage(CCANMessage* Message, bool Format = false);
int SetMQTTTopicDevice(QString DeviceString);
int SetMQTTClient(CMQTTClientWrapper *MQTTClient);
QList<CMQTTMessage> *GetMQTTMessagesList();
+1 -3
View File
@@ -13,7 +13,7 @@ CCANDevice::CCANDevice(QObject *parent)
mProgramPtr = 0;
}
CCANDevice::CCANDevice(CCANDeviceConfig &SysConfig, CMQTTClientWrapper *MQTTClient, QString DeviceTopicPrefix)
CCANDevice::CCANDevice(CCANDeviceConfig &SysConfig, QString DeviceTopicPrefix)
{
mMessageList.clear();
mMessagesListLoaded = false;
@@ -21,7 +21,6 @@ CCANDevice::CCANDevice(CCANDeviceConfig &SysConfig, CMQTTClientWrapper *MQTTClie
mProgramPtr = 0;
mDeviceConfigInfo = SysConfig;
mCANMQTTClient = MQTTClient;
mDeviceTopicPrefix = DeviceTopicPrefix;
mCANDriverIF = 0;
@@ -101,7 +100,6 @@ int CCANDevice::Init()
#else
mCANDataLogger.SetMQTTTopicDevice(QString("CANBus/%1/").arg(mDeviceConfigInfo.mDeviceName));
#endif
mCANDataLogger.SetMQTTClient(mCANMQTTClient);
mProgramPtr->SetCANConnectionStatusRequest(true);
mProgramPtr->UpdateCANModuleStatusRequest(mDeviceConfigInfo.mDeviceName,"Connecté","NOUPDATE");
+1 -2
View File
@@ -24,7 +24,7 @@ class CCANDevice : public QObject
Q_OBJECT
public:
explicit CCANDevice(QObject *parent = 0);
CCANDevice(CCANDeviceConfig &SysConfig, CMQTTClientWrapper* MQTTClient = 0, QString DeviceTopicPrefix="");
CCANDevice(CCANDeviceConfig &SysConfig, QString DeviceTopicPrefix="");
~CCANDevice();
int Init(QString DatabaseFileName, TPCANHandle CANDeviceID, TPCANBaudrate CANDeviceBaudRate, QString DevDescription, QString DeviceName, unsigned int DevicePollPeriod);
@@ -40,7 +40,6 @@ public:
CCANAnalyzer mCANAnalyzer; //The module that handles the USB puck and decodes the data
CCANDatabase mCANDatabase; //The device's database loaded from dbc file
CCANDataLogger mCANDataLogger;
CMQTTClientWrapper *mCANMQTTClient;
QString mDeviceTopicPrefix;
CCANWatchdog mCANWatchdog;
@@ -175,11 +175,12 @@ void CMQTTClientWrapper::StateChanged()
{
//Connection attempt failed, just restart the timer...
mMQTTReconnectTimer->start(MQTT_CLIENT_RECONNECT_TIMEOUT);
CGeneralMessagesLogDispatcher::instance()->AddLogMessage("Client MQTT déconnecté pendant une reconnexion. ","CMQTTClientWrapper",true,1);
CGeneralMessagesLogDispatcher::instance()->AddLogMessage(QString("Client MQTT %1 déconnecté pendant une reconnexion. ").arg(mMQTTParams.mMQTTBrokerHostName),"CMQTTClientWrapper",true,1);
}
else
{
CGeneralMessagesLogDispatcher::instance()->AddLogMessage("Client MQTT déconnecté.","CMQTTClientWrapper",true,1);
CGeneralMessagesLogDispatcher::instance()->AddLogMessage(QString("Client MQTT %1 déconnecté.").arg(mMQTTParams.mMQTTBrokerHostName),"CMQTTClientWrapper",true,1);
mProgramPtr->SetMQTTConnectionSatusRequest(false);
#ifndef ENABLE_DEVELOPMENT_DEBUG_TOOLS
mMQTTReconnectTimer->start(MQTT_CLIENT_RECONNECT_TIMEOUT);
@@ -195,7 +196,7 @@ void CMQTTClientWrapper::StateChanged()
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","CMQTTClientWrapper",true,1);
CGeneralMessagesLogDispatcher::instance()->AddLogMessage(QString("Passage en mode buffering des messages MQTT pour %1").arg(mMQTTParams.mMQTTBrokerHostName),"CMQTTClientWrapper",true,1);
mBufferingModeText = "Buffering";
UpdateGUIBufferingStatus();
}
@@ -207,11 +208,11 @@ void CMQTTClientWrapper::StateChanged()
mProgramPtr->SetMQTTConnectionSatusRequest(true);
mMQTTReconnectTimer->stop();
mIsClientConnecting = false;
CGeneralMessagesLogDispatcher::instance()->AddLogMessage("Client MQTT connecté.","CMQTTClientWrapper",true,1);
CGeneralMessagesLogDispatcher::instance()->AddLogMessage(QString("Client MQTT %1 connecté.").arg(mMQTTParams.mMQTTBrokerHostName),"CMQTTClientWrapper",true,1);
if(mMQTTMessagesQueue.isEmpty() == false)
{
mMessagesQueueMode = MQTT_QUEUE_MSG_MODE; //Stay in (or enter) queue mode until we empty the buffer
CGeneralMessagesLogDispatcher::instance()->AddLogMessage("FIFO non vide, passage au mode de vidage de la FIFO","CMQTTClientWrapper",true,2);
CGeneralMessagesLogDispatcher::instance()->AddLogMessage(QString("FIFO non vide, passage au mode de vidage de la FIFO pour %1").arg(mMQTTParams.mMQTTBrokerHostName),"CMQTTClientWrapper",true,2);
mBufferingModeText = "Buffering";
mCircularBufferStatusText = QString("%2/%1 messages (0%3\%)").arg(MQTT_CLIENT_MSG_QUEUE_SIZE).arg(mMQTTMessagesQueue.size()).arg((mMQTTMessagesQueue.size()/MQTT_CLIENT_MSG_QUEUE_SIZE)*100);
@@ -234,7 +235,8 @@ void CMQTTClientWrapper::StateChanged()
case QMqttClient::Connecting:
{
mIsClientConnecting = true;
CGeneralMessagesLogDispatcher::instance()->AddLogMessage("Client MQTT en cours de connexion... ","CMQTTClientWrapper",true,1);
CGeneralMessagesLogDispatcher::instance()->AddLogMessage(QString("Client MQTT %1 en cours de connexion... ").arg(mMQTTParams.mMQTTBrokerHostName),"CMQTTClientWrapper",true,1);
mProgramPtr->SetMQTTConnectionSatusRequest(false);
break;
}
}
@@ -414,6 +416,32 @@ void CMQTTClientWrapper::MQTTClientError(QMqttClient::ClientError error)
CGeneralMessagesLogDispatcher::instance()->AddLogMessage(QString("Erreur du client MQTT: %1").arg(error),"CMQTTClientWrapper",true,2);
}
bool CMQTTClientWrapper::IsMQTTClientConnected()
{
if(mMQTTClient.state() != QMqttClient::Connected)
return false;
else
return true;
}
QString CMQTTClientWrapper::GetMQTTClientConnectionState()
{
QString StateText;
if(mMQTTClient.state() == QMqttClient::Connected)
{
StateText = QString("%1 : Connecté").arg(mMQTTParams.mMQTTBrokerHostName);
}
else if(mMQTTClient.state() == QMqttClient::Connecting)
{
StateText = QString("%1 : En connexion...").arg(mMQTTParams.mMQTTBrokerHostName);
}
else
{
StateText = QString("%1 : Déconnecté").arg(mMQTTParams.mMQTTBrokerHostName);
}
return StateText;
}
//This function is used only when flushing the msg queue
void CMQTTClientWrapper::MQTTMessageSent(qint32 MsgID)
{
@@ -41,6 +41,8 @@ public:
int StartMQTTClient();
int StopMQTTClient();
quint64 GetMQTTServerPresenceCANMask();
bool IsMQTTClientConnected();
QString GetMQTTClientConnectionState();
#ifdef ENABLE_DEVELOPMENT_DEBUG_TOOLS
int ForceMQTTClientDisconnection(bool Disconnect);
@@ -13,7 +13,7 @@ CGeneralStatusPage::CGeneralStatusPage(QWidget *parent) :
connect(ui->mClearGenMsgTxtBtn,&QPushButton::clicked,this,&CGeneralStatusPage::ClearGenMsgAreaBtnPressed);
connect(ui->mQuitAppBtn,&QPushButton::clicked,this,&CGeneralStatusPage::QuitAppBtnPressed);
SetMQTTConnectionStatus(false);
//SetMQTTConnectionStatus(false);
SetCANConnectionStatus(false);
ui->mCANModuleStatusTableWdgt->setColumnCount(3);
@@ -155,6 +155,11 @@ int CGeneralStatusPage::SetMQTTConnectionStatus(bool Connected)
return RET_OK;
}
int CGeneralStatusPage::SetMQTTConnectionStatus(QString Status)
{
ui->mClientMQTTConnStatLbl->setText(Status);
}
int CGeneralStatusPage::SetCANConnectionStatus(bool Connected)
{
@@ -46,6 +46,7 @@ public:
int AddGeneralMsgBoxLineEntry(QString LineTxt);
int SetMQTTConnectionStatus(bool Connected);
int SetMQTTConnectionStatus(QString Status);
int SetCANConnectionStatus(bool Connected);
int UpdateCANModuleStatus(QString ModuleName, QString ModuleStatus, QString Buffer);
int ClearCANModuleStatusTable();
+12 -27
View File
@@ -47,31 +47,13 @@
<string>Nettoyer</string>
</property>
</widget>
<widget class="QLabel" name="mClientMQTTLbl">
<widget class="QLabel" name="mClientMQTTConnStatLbl">
<property name="geometry">
<rect>
<x>1200</x>
<y>210</y>
<width>101</width>
<height>16</height>
</rect>
</property>
<property name="font">
<font>
<pointsize>12</pointsize>
</font>
</property>
<property name="text">
<string>Client MQTT:</string>
</property>
</widget>
<widget class="QLabel" name="mClientMQTTConnStatLbl">
<property name="geometry">
<rect>
<x>1300</x>
<y>210</y>
<width>121</width>
<height>16</height>
<width>251</width>
<height>131</height>
</rect>
</property>
<property name="font">
@@ -84,6 +66,9 @@
<property name="text">
<string>Déconnecté</string>
</property>
<property name="alignment">
<set>Qt::AlignLeading|Qt::AlignLeft|Qt::AlignTop</set>
</property>
</widget>
<widget class="QTableWidget" name="mCANModuleStatusTableWdgt">
<property name="geometry">
@@ -131,8 +116,8 @@
<widget class="QLabel" name="mInternetConnectedLbl">
<property name="geometry">
<rect>
<x>1230</x>
<y>240</y>
<x>1200</x>
<y>180</y>
<width>61</width>
<height>16</height>
</rect>
@@ -149,8 +134,8 @@
<widget class="QLabel" name="mInternetPresentStatLbl">
<property name="geometry">
<rect>
<x>1300</x>
<y>240</y>
<x>1270</x>
<y>180</y>
<width>121</width>
<height>16</height>
</rect>
@@ -254,8 +239,8 @@
<widget class="QGroupBox" name="mSystemStateGroupBox">
<property name="geometry">
<rect>
<x>1210</x>
<y>290</y>
<x>1200</x>
<y>370</y>
<width>201</width>
<height>121</height>
</rect>
+69 -13
View File
@@ -23,7 +23,7 @@ COtarcikCan::COtarcikCan(QObject *parent) : QObject(parent)
{
mGPTimer = new QTimer;
connect(mGPTimer,SIGNAL(timeout()),this,SLOT(GPTimerExpired()));
mCANBusMQTTClient.mProgramPtr = this;
// mCANBusMQTTClient.mProgramPtr = this;
mWatchdogTimer = new QTimer;
connect(mWatchdogTimer,&QTimer::timeout,this,&COtarcikCan::WatchdogUpdateTimerExpired);
mWatchdogTimer->setSingleShot(false);
@@ -41,7 +41,14 @@ COtarcikCan::~COtarcikCan()
}
mCANDevicesList.clear();
mCANBusMQTTClient.DisconnectFromBroker();
while(!mCANBusMQTTClientList.isEmpty())
{
mCANBusMQTTClientList.first()->DisconnectFromBroker();
delete mCANBusMQTTClientList.takeFirst();
}
// mCANBusMQTTClient.DisconnectFromBroker();
delete mGPTimer;
@@ -84,12 +91,21 @@ int COtarcikCan::Start()
mCloudLoggingParamsList = mSystemConfig.GetCloudParams();
mMainWindow.mDataLoggingSettingsPage->SetCloudParams(mCloudLoggingParamsList);
mCANBusMQTTClient.SetMQTTParams(mCloudLoggingParamsList->at(0)); //TODO: Fix that
mCANBusMQTTClient.SetCANDevicesList(&mCANDevicesList);
for(int i = 0; i < mCloudLoggingParamsList->size(); i++)
{
CMQTTClientWrapper *NewMQTTWrapper = new CMQTTClientWrapper;
NewMQTTWrapper->mProgramPtr = this;
NewMQTTWrapper->SetMQTTParams(mCloudLoggingParamsList->at(i));
NewMQTTWrapper->SetCANDevicesList(&mCANDevicesList);
#ifdef ENABLE_CHIPSET_DRIVER
mCANBusMQTTClient.SetCPUInterface(&mCPUInterface);
NewMQTTWrapper->SetCPUInterface(&mCPUInterface);
#endif
mCANBusMQTTClient.SetMQTTServerPresenceCANBit(mSystemConfig.GetDeviceDetectionConfig()->mMQTTDetectionCANStatusBit);
NewMQTTWrapper->SetMQTTServerPresenceCANBit(mSystemConfig.GetDeviceDetectionConfig()->mMQTTDetectionCANStatusBit);
mCANBusMQTTClientList.append(NewMQTTWrapper);
}
// mCANBusMQTTClient.SetMQTTParams(mCloudLoggingParamsList->at(0)); //TODO: Fix that
// mCANBusMQTTClient.SetCANDevicesList(&mCANDevicesList);
mGeneralSystemParams = *mSystemConfig.GetGeneralSystemSettings();
mMainWindow.mDataLoggingSettingsPage->SetGeneralSettingsParams(&mGeneralSystemParams);
@@ -107,7 +123,8 @@ int COtarcikCan::Start()
mInternetMonitor.Start(mSystemConfig.GetDeviceDetectionConfig()->mInternetDetectionCANStatusBit);
mLANDevicesPresenceMonitor.Start(mSystemConfig.GetDeviceDetectionConfig()->GetLANDevicesConfigList());
mSystemConfig.mDeviceDetectionParams.SetCANPresenceMonitors(&mCANBusMQTTClient,&mLANDevicesPresenceMonitor,&mInternetMonitor);
//mSystemConfig.mDeviceDetectionParams.SetCANPresenceMonitors(&mCANBusMQTTClient,&mLANDevicesPresenceMonitor,&mInternetMonitor); TODO: Add a bit for each MQTT client?
mSystemConfig.mDeviceDetectionParams.SetCANPresenceMonitors(mCANBusMQTTClientList.at(0),&mLANDevicesPresenceMonitor,&mInternetMonitor);
for(int i = 0; i < mCANDevicesList.size(); i++)
{
mCANDevicesList.at(i)->StartWatchdog(mSystemConfig.GetDeviceDetectionConfig());
@@ -140,7 +157,11 @@ int COtarcikCan::Start()
CGeneralMessagesLogDispatcher::instance()->AddLogMessage(QString("Démarrage du logiciel OtarcikCAN"),"CPCANInterface");
// mCANBusMQTTClient.ConnectToBroker();
mCANBusMQTTClient.StartMQTTClient();
//mCANBusMQTTClient.StartMQTTClient();
for(int i = 0; i < mCANBusMQTTClientList.size(); i++)
{
mCANBusMQTTClientList.at(i)->StartMQTTClient();
}
mMainWindow.mCANbusSettingsPage->SetDevicesList(&mCANDevicesList);
connect(&mInternetMonitor,&CInternetMonitor::InternetStateChanged,mMainWindow.mGeneralStatusPage,&CGeneralStatusPage::InternetStatusChanged);
@@ -206,7 +227,7 @@ int COtarcikCan::PopulateCANDevicesList(QList<CCANDeviceConfig *> *CANDeviceConf
for(int i = 0; i < CANDeviceConfigList->size(); i++)
{
CCANDevice *NewDevice = new CCANDevice(*CANDeviceConfigList->at(i),&mCANBusMQTTClient,mSystemConfig.mCloudLoggingParamsList[0]->mMQTTTopicPrefix); //TODO fix the cloud
CCANDevice *NewDevice = new CCANDevice(*CANDeviceConfigList->at(i),mSystemConfig.mCloudLoggingParamsList[0]->mMQTTTopicPrefix); //TODO fix the cloud
NewDevice->mProgramPtr = this;
NewDevice->Init();
mCANDevicesList.append(NewDevice);
@@ -227,7 +248,28 @@ int COtarcikCan::SaveCloudLoggingConfigRequest(QList<CCloudParams *> *CloudParam
mSystemConfig.SetCloudParamsList(CloudParams);
mCloudLoggingParamsList = mSystemConfig.GetCloudParams();
mCANBusMQTTClient.SetMQTTParams(mCloudLoggingParamsList->at(0)); //TODO: Fix
while(!mCANBusMQTTClientList.isEmpty())
{
mCANBusMQTTClientList.first()->DisconnectFromBroker();
delete mCANBusMQTTClientList.takeFirst();
}
for(int i = 0; i < CloudParams->size(); i++)
{
CMQTTClientWrapper *NewMQTTWrapper = new CMQTTClientWrapper;
NewMQTTWrapper->mProgramPtr = this;
NewMQTTWrapper->SetMQTTParams(CloudParams->at(i));
NewMQTTWrapper->SetCANDevicesList(&mCANDevicesList);
#ifdef ENABLE_CHIPSET_DRIVER
NewMQTTWrapper->SetCPUInterface(&mCPUInterface);
#endif
NewMQTTWrapper->SetMQTTServerPresenceCANBit(mSystemConfig.GetDeviceDetectionConfig()->mMQTTDetectionCANStatusBit);
mCANBusMQTTClientList.append(NewMQTTWrapper);
NewMQTTWrapper->StartMQTTClient();
}
mSystemConfig.mDeviceDetectionParams.SetCANPresenceMonitors(mCANBusMQTTClientList.at(0),&mLANDevicesPresenceMonitor,&mInternetMonitor);
//mCANBusMQTTClient.SetMQTTParams(mCloudLoggingParamsList->at(0)); //TODO: Fix
if(mSystemConfig.SaveConfig() == RET_OK)
{
@@ -284,7 +326,11 @@ int COtarcikCan::SaveDeviceDetectionSettingsRequest(CDeviceDetectionConfig *Devi
mLANDevicesPresenceMonitor.Stop();
mLANDevicesPresenceMonitor.Start(DeviceDetectconfig->GetLANDevicesConfigList());
mInternetMonitor.UpdateCANReportingBit(DeviceDetectconfig->mInternetDetectionCANStatusBit);
mCANBusMQTTClient.SetMQTTServerPresenceCANBit(DeviceDetectconfig->mMQTTDetectionCANStatusBit);
for(int i = 0; i < mCANBusMQTTClientList.size(); i++)
{
mCANBusMQTTClientList.at(i)->SetMQTTServerPresenceCANBit(DeviceDetectconfig->mMQTTDetectionCANStatusBit);
}
// mCANBusMQTTClient.SetMQTTServerPresenceCANBit(DeviceDetectconfig->mMQTTDetectionCANStatusBit);
return RET_OK;
}
@@ -300,8 +346,18 @@ int COtarcikCan::SetCANConnectionStatusRequest(bool Connected)
int COtarcikCan::SetMQTTConnectionSatusRequest(bool Connected)
{
QString MQTTStates("Clients MQTT:\n");
for(int i = 0; i < mCANBusMQTTClientList.size(); i++)
{
CMQTTClientWrapper *Wrapper = mCANBusMQTTClientList.at(i);
MQTTStates.append(Wrapper->GetMQTTClientConnectionState());
MQTTStates.append("\n");
}
mMainWindow.mDataLoggingSettingsPage->SetMQTTPresenceStatus(Connected);
return mMainWindow.mGeneralStatusPage->SetMQTTConnectionStatus(Connected);
mMainWindow.mGeneralStatusPage->SetMQTTConnectionStatus(MQTTStates);
return RET_OK;
}
int COtarcikCan::UpdateCANModuleStatusRequest(QString ModuleName, QString ModuleStatus, QString Buffer)
@@ -350,7 +406,7 @@ void COtarcikCan::QuitApplicationRequest()
int COtarcikCan::ForceMQTTDisconnect(bool Disconnect)
{
mCANBusMQTTClient.ForceMQTTClientDisconnection(Disconnect);
// mCANBusMQTTClient.ForceMQTTClientDisconnection(Disconnect);
}
#endif
+2 -1
View File
@@ -25,7 +25,8 @@ public:
~COtarcikCan();
CMainWindow mMainWindow;
CSystemConfig mSystemConfig;
CMQTTClientWrapper mCANBusMQTTClient;
/// CMQTTClientWrapper mCANBusMQTTClient;
QList<CMQTTClientWrapper*> mCANBusMQTTClientList;
QTimer *mGPTimer;
#ifdef ENABLE_CHIPSET_DRIVER
CComputerBoardInterface mCPUInterface;