fix for potentially missed updates (#1226)
This commit is contained in:
@@ -899,6 +899,7 @@ void KnxHandler::run() {
|
|||||||
if (!m_updatedMessages.empty()) {
|
if (!m_updatedMessages.empty()) {
|
||||||
m_messages->lock();
|
m_messages->lock();
|
||||||
if (m_con->isConnected()) {
|
if (m_con->isConnected()) {
|
||||||
|
time_t maxUpdates = 0;
|
||||||
for (auto it = m_updatedMessages.begin(); it != m_updatedMessages.end(); ) {
|
for (auto it = m_updatedMessages.begin(); it != m_updatedMessages.end(); ) {
|
||||||
const vector<Message*>* messages = m_messages->getByKey(it->first);
|
const vector<Message*>* messages = m_messages->getByKey(it->first);
|
||||||
if (!messages) {
|
if (!messages) {
|
||||||
@@ -906,16 +907,20 @@ void KnxHandler::run() {
|
|||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
for (const auto& message : *messages) {
|
for (const auto& message : *messages) {
|
||||||
if (message->getLastChangeTime() <= 0) {
|
time_t changeTime = message->getLastChangeTime();
|
||||||
|
if (changeTime <= 0) {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
if (changeTime > lastUpdates && changeTime > maxUpdates) {
|
||||||
|
maxUpdates = changeTime;
|
||||||
|
}
|
||||||
const auto mit = m_subscribedMessages.find(message->getKey());
|
const auto mit = m_subscribedMessages.find(message->getKey());
|
||||||
if (mit == m_subscribedMessages.cend()) {
|
if (mit == m_subscribedMessages.cend()) {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
if (!(message->getDataHandlerState()&2)) {
|
if (!(message->getDataHandlerState()&2)) {
|
||||||
message->setDataHandlerState(2, true); // first update still needed
|
message->setDataHandlerState(2, true); // first update still needed
|
||||||
} else if (message->getLastChangeTime() <= lastUpdates) {
|
} else if (changeTime <= lastUpdates) {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
for (auto destFlags : mit->second) {
|
for (auto destFlags : mit->second) {
|
||||||
@@ -936,7 +941,7 @@ void KnxHandler::run() {
|
|||||||
}
|
}
|
||||||
it = m_updatedMessages.erase(it);
|
it = m_updatedMessages.erase(it);
|
||||||
}
|
}
|
||||||
time(&lastUpdates);
|
lastUpdates = maxUpdates == 0 || lastUpdates > maxUpdates ? now : maxUpdates + 1;
|
||||||
} else {
|
} else {
|
||||||
m_updatedMessages.clear();
|
m_updatedMessages.clear();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1068,22 +1068,27 @@ void MqttHandler::run() {
|
|||||||
if (!m_updatedMessages.empty()) {
|
if (!m_updatedMessages.empty()) {
|
||||||
m_messages->lock();
|
m_messages->lock();
|
||||||
if (m_connected) {
|
if (m_connected) {
|
||||||
|
time_t maxUpdates = 0;
|
||||||
for (auto it = m_updatedMessages.begin(); it != m_updatedMessages.end(); ) {
|
for (auto it = m_updatedMessages.begin(); it != m_updatedMessages.end(); ) {
|
||||||
const vector<Message*>* messages = m_messages->getByKey(it->first);
|
const vector<Message*>* messages = m_messages->getByKey(it->first);
|
||||||
if (messages) {
|
if (messages) {
|
||||||
for (const auto& message : *messages) {
|
for (const auto& message : *messages) {
|
||||||
if (message->getLastChangeTime() > 0 && message->isAvailable()
|
time_t changeTime = message->getLastChangeTime();
|
||||||
&& (!g_onlyChanges || message->getLastChangeTime() > lastUpdates)) {
|
if (changeTime > 0 && message->isAvailable()
|
||||||
|
&& (!g_onlyChanges || changeTime > lastUpdates)) {
|
||||||
updates.str("");
|
updates.str("");
|
||||||
updates.clear();
|
updates.clear();
|
||||||
updates << dec;
|
updates << dec;
|
||||||
publishMessage(message, &updates);
|
publishMessage(message, &updates);
|
||||||
}
|
}
|
||||||
|
if (changeTime > lastUpdates && changeTime > maxUpdates) {
|
||||||
|
maxUpdates = changeTime;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
it = m_updatedMessages.erase(it);
|
it = m_updatedMessages.erase(it);
|
||||||
}
|
}
|
||||||
time(&lastUpdates);
|
lastUpdates = maxUpdates == 0 || lastUpdates > maxUpdates ? now : maxUpdates + 1;
|
||||||
} else {
|
} else {
|
||||||
m_updatedMessages.clear();
|
m_updatedMessages.clear();
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user