track changes only per data sink (#1226)
This commit is contained in:
@@ -75,8 +75,11 @@ bool datahandler_register(UserInfo* userInfo, BusHandler* busHandler, MessageMap
|
|||||||
return success;
|
return success;
|
||||||
}
|
}
|
||||||
|
|
||||||
void DataSink::notifyUpdate(Message* message) {
|
void DataSink::notifyUpdate(Message* message, bool changed) {
|
||||||
if (message && message->hasLevel(m_levels)) {
|
if (message && message->hasLevel(m_levels)) {
|
||||||
|
if (m_changedOnly && !changed) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
m_updatedMessages[message->getKey()]++;
|
m_updatedMessages[message->getKey()]++;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -143,8 +143,9 @@ class DataSink : virtual public DataHandler {
|
|||||||
* Constructor.
|
* Constructor.
|
||||||
* @param userInfo the @a UserInfo instance.
|
* @param userInfo the @a UserInfo instance.
|
||||||
* @param user the user name for determining the allowed access levels (fall back to default levels).
|
* @param user the user name for determining the allowed access levels (fall back to default levels).
|
||||||
|
* @param changedOnly whether to handle changed messages only in the updates.
|
||||||
*/
|
*/
|
||||||
DataSink(const UserInfo* userInfo, const string& user) {
|
DataSink(const UserInfo* userInfo, const string& user, bool changedOnly) : m_changedOnly(changedOnly) {
|
||||||
m_levels = userInfo->getLevels(userInfo->hasUser(user) ? user : "");
|
m_levels = userInfo->getLevels(userInfo->hasUser(user) ? user : "");
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -159,8 +160,9 @@ class DataSink : virtual public DataHandler {
|
|||||||
/**
|
/**
|
||||||
* Notify the sink of an updated @a Message (not necessarily changed though).
|
* Notify the sink of an updated @a Message (not necessarily changed though).
|
||||||
* @param message the updated @a Message.
|
* @param message the updated @a Message.
|
||||||
|
* @param changed whether the message data changed since the last notification.
|
||||||
*/
|
*/
|
||||||
virtual void notifyUpdate(Message* message);
|
virtual void notifyUpdate(Message* message, bool changed);
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Notify the sink of the latest update check result.
|
* Notify the sink of the latest update check result.
|
||||||
@@ -178,6 +180,9 @@ class DataSink : virtual public DataHandler {
|
|||||||
/** the allowed access levels. */
|
/** the allowed access levels. */
|
||||||
string m_levels;
|
string m_levels;
|
||||||
|
|
||||||
|
/** whether to handle changed messages only in the updates. */
|
||||||
|
bool m_changedOnly;
|
||||||
|
|
||||||
/** a map of updated @p Message keys. */
|
/** a map of updated @p Message keys. */
|
||||||
map<uint64_t, int> m_updatedMessages;
|
map<uint64_t, int> m_updatedMessages;
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -163,7 +163,7 @@ bool knxhandler_register(UserInfo* userInfo, BusHandler* busHandler, MessageMap*
|
|||||||
}
|
}
|
||||||
|
|
||||||
KnxHandler::KnxHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap* messages)
|
KnxHandler::KnxHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap* messages)
|
||||||
: DataSink(userInfo, "knx"), DataSource(busHandler), WaitThread(), m_messages(messages),
|
: DataSink(userInfo, "knx", true), DataSource(busHandler), WaitThread(), m_messages(messages),
|
||||||
m_start(0), m_lastUpdateCheckResult("."),
|
m_start(0), m_lastUpdateCheckResult("."),
|
||||||
m_lastScanStatus(SCAN_STATUS_NONE), m_scanFinishReceived(false), m_lastErrorLogTime(0) {
|
m_lastScanStatus(SCAN_STATUS_NONE), m_scanFinishReceived(false), m_lastErrorLogTime(0) {
|
||||||
m_con = KnxConnection::create(g_url);
|
m_con = KnxConnection::create(g_url);
|
||||||
@@ -709,7 +709,7 @@ void KnxHandler::handleGroupTelegram(knx_addr_t src, knx_addr_t dest, int len, c
|
|||||||
#define UPTIME_INTERVAL 3600
|
#define UPTIME_INTERVAL 3600
|
||||||
|
|
||||||
void KnxHandler::run() {
|
void KnxHandler::run() {
|
||||||
time_t lastTaskRun, now, lastSignal = 0, lastUptime = 0, lastUpdates = 0;
|
time_t lastTaskRun, now, lastSignal = 0, lastUptime = 0;
|
||||||
bool signal = false;
|
bool signal = false;
|
||||||
result_t result = RESULT_OK;
|
result_t result = RESULT_OK;
|
||||||
time(&now);
|
time(&now);
|
||||||
@@ -899,7 +899,6 @@ 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) {
|
||||||
@@ -911,18 +910,10 @@ void KnxHandler::run() {
|
|||||||
if (changeTime <= 0) {
|
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)) {
|
|
||||||
message->setDataHandlerState(2, true); // first update still needed
|
|
||||||
} else if (changeTime <= lastUpdates) {
|
|
||||||
continue;
|
|
||||||
}
|
|
||||||
for (auto destFlags : mit->second) {
|
for (auto destFlags : mit->second) {
|
||||||
auto sit = m_subscribedGroups.find(destFlags);
|
auto sit = m_subscribedGroups.find(destFlags);
|
||||||
if (sit == m_subscribedGroups.end()) {
|
if (sit == m_subscribedGroups.end()) {
|
||||||
@@ -941,7 +932,6 @@ void KnxHandler::run() {
|
|||||||
}
|
}
|
||||||
it = m_updatedMessages.erase(it);
|
it = m_updatedMessages.erase(it);
|
||||||
}
|
}
|
||||||
lastUpdates = maxUpdates == 0 || lastUpdates > maxUpdates ? now : maxUpdates + 1;
|
|
||||||
} else {
|
} else {
|
||||||
m_updatedMessages.clear();
|
m_updatedMessages.clear();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -352,8 +352,9 @@ void MainLoop::run() {
|
|||||||
m_messages->lock();
|
m_messages->lock();
|
||||||
m_messages->findAll("", "", "*", false, true, true, true, true, true, sinkSince, now, false, &messages);
|
m_messages->findAll("", "", "*", false, true, true, true, true, true, sinkSince, now, false, &messages);
|
||||||
for (const auto message : messages) {
|
for (const auto message : messages) {
|
||||||
|
bool changed = message->getLastChangeTime() >= sinkSince;
|
||||||
for (const auto dataSink : dataSinks) {
|
for (const auto dataSink : dataSinks) {
|
||||||
dataSink->notifyUpdate(message);
|
dataSink->notifyUpdate(message, changed);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
m_messages->unlock();
|
m_messages->unlock();
|
||||||
|
|||||||
@@ -368,7 +368,8 @@ string removeTrailingNonTopicPart(const string& str) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
MqttHandler::MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap* messages)
|
MqttHandler::MqttHandler(UserInfo* userInfo, BusHandler* busHandler, MessageMap* messages)
|
||||||
: DataSink(userInfo, "mqtt"), DataSource(busHandler), WaitThread(), m_messages(messages), m_connected(false),
|
: DataSink(userInfo, "mqtt", g_onlyChanges), DataSource(busHandler), WaitThread(),
|
||||||
|
m_messages(messages), m_connected(false),
|
||||||
m_lastUpdateCheckResult("."), m_lastScanStatus(SCAN_STATUS_NONE) {
|
m_lastUpdateCheckResult("."), m_lastScanStatus(SCAN_STATUS_NONE) {
|
||||||
m_definitionsSince = 0;
|
m_definitionsSince = 0;
|
||||||
m_client = nullptr;
|
m_client = nullptr;
|
||||||
@@ -696,7 +697,7 @@ void splitFields(const string& str, vector<string>* row) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
void MqttHandler::run() {
|
void MqttHandler::run() {
|
||||||
time_t lastTaskRun, now, start, lastSignal = 0, lastUpdates = 0;
|
time_t lastTaskRun, now, start, lastSignal = 0;
|
||||||
bool signal = false;
|
bool signal = false;
|
||||||
bool globalHasName = m_globalTopic.has("name");
|
bool globalHasName = m_globalTopic.has("name");
|
||||||
string signalTopic = m_globalTopic.get("", "signal");
|
string signalTopic = m_globalTopic.get("", "signal");
|
||||||
@@ -1068,27 +1069,21 @@ 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) {
|
||||||
time_t changeTime = message->getLastChangeTime();
|
time_t changeTime = message->getLastChangeTime();
|
||||||
if (changeTime > 0 && message->isAvailable()
|
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);
|
||||||
}
|
}
|
||||||
lastUpdates = maxUpdates == 0 || lastUpdates > maxUpdates ? now : maxUpdates + 1;
|
|
||||||
} else {
|
} else {
|
||||||
m_updatedMessages.clear();
|
m_updatedMessages.clear();
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user