diff --git a/src/ebusd/mainloop.cpp b/src/ebusd/mainloop.cpp index 6d41e67e..844bea53 100644 --- a/src/ebusd/mainloop.cpp +++ b/src/ebusd/mainloop.cpp @@ -359,12 +359,14 @@ void MainLoop::run() { time(&now); if (!dataSinks.empty()) { messages.clear(); + m_messages->lock(); m_messages->findAll("", "", "*", false, true, true, true, true, true, sinkSince, now, false, &messages); for (const auto message : messages) { for (const auto dataSink : dataSinks) { dataSink->notifyUpdate(message); } } + m_messages->unlock(); sinkSince = now; } if (netMessage == nullptr) { diff --git a/src/ebusd/mqtthandler.cpp b/src/ebusd/mqtthandler.cpp index 1ded7a24..78df67af 100644 --- a/src/ebusd/mqtthandler.cpp +++ b/src/ebusd/mqtthandler.cpp @@ -747,8 +747,8 @@ void MqttHandler::run() { } } if (!m_updatedMessages.empty()) { + m_messages->lock(); if (m_connected) { - m_messages->lock(); for (auto it = m_updatedMessages.begin(); it != m_updatedMessages.end(); ) { const vector* messages = m_messages->getByKey(it->first); if (messages) { @@ -764,11 +764,11 @@ void MqttHandler::run() { } it = m_updatedMessages.erase(it); } - m_messages->unlock(); time(&lastUpdates); } else { m_updatedMessages.clear(); } + m_messages->unlock(); } if ((!m_connected && !Wait(5)) || (needsWait && !Wait(1))) { break;