From b7265cb6824f14db7047353d61d2a5cd34447a20 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 25 Jan 2023 12:23:59 +0200 Subject: [PATCH 01/16] Added check if entity exists before processing telemetry update msg --- .../edge/rpc/processor/BaseEdgeProcessor.java | 21 ++++++ .../relation/BaseRelationProcessor.java | 28 -------- .../telemetry/BaseTelemetryProcessor.java | 66 ++++++++++--------- 3 files changed, 56 insertions(+), 59 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java index 4b94cee592..36077d3fd5 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java @@ -485,4 +485,25 @@ public abstract class BaseEdgeProcessor { } return customerId; } + + protected boolean isEntityExists(TenantId tenantId, EntityId entityId) { + switch (entityId.getEntityType()) { + case DEVICE: + return deviceService.findDeviceById(tenantId, new DeviceId(entityId.getId())) != null; + case ASSET: + return assetService.findAssetById(tenantId, new AssetId(entityId.getId())) != null; + case ENTITY_VIEW: + return entityViewService.findEntityViewById(tenantId, new EntityViewId(entityId.getId())) != null; + case CUSTOMER: + return customerService.findCustomerById(tenantId, new CustomerId(entityId.getId())) != null; + case USER: + return userService.findUserById(tenantId, new UserId(entityId.getId())) != null; + case DASHBOARD: + return dashboardService.findDashboardById(tenantId, new DashboardId(entityId.getId())) != null; + case EDGE: + return edgeService.findEdgeById(tenantId, new EdgeId(entityId.getId())) != null; + default: + return false; + } + } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/relation/BaseRelationProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/relation/BaseRelationProcessor.java index e1c937eaad..36a5decc3f 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/relation/BaseRelationProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/relation/BaseRelationProcessor.java @@ -20,16 +20,9 @@ import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.common.data.id.AssetId; -import org.thingsboard.server.common.data.id.CustomerId; -import org.thingsboard.server.common.data.id.DashboardId; -import org.thingsboard.server.common.data.id.DeviceId; -import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; -import org.thingsboard.server.common.data.id.EntityViewId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.gen.edge.v1.RelationUpdateMsg; @@ -80,25 +73,4 @@ public abstract class BaseRelationProcessor extends BaseEdgeProcessor { return Futures.immediateFailedFuture(e); } } - - private boolean isEntityExists(TenantId tenantId, EntityId entityId) { - switch (entityId.getEntityType()) { - case DEVICE: - return deviceService.findDeviceById(tenantId, new DeviceId(entityId.getId())) != null; - case ASSET: - return assetService.findAssetById(tenantId, new AssetId(entityId.getId())) != null; - case ENTITY_VIEW: - return entityViewService.findEntityViewById(tenantId, new EntityViewId(entityId.getId())) != null; - case CUSTOMER: - return customerService.findCustomerById(tenantId, new CustomerId(entityId.getId())) != null; - case USER: - return userService.findUserById(tenantId, new UserId(entityId.getId())) != null; - case DASHBOARD: - return dashboardService.findDashboardById(tenantId, new DashboardId(entityId.getId())) != null; - case EDGE: - return edgeService.findEdgeById(tenantId, new EdgeId(entityId.getId())) != null; - default: - return false; - } - } } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/telemetry/BaseTelemetryProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/telemetry/BaseTelemetryProcessor.java index 30979187d1..a3d3e3c280 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/telemetry/BaseTelemetryProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/telemetry/BaseTelemetryProcessor.java @@ -90,42 +90,46 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor { log.trace("[{}] processTelemetryMsg [{}]", tenantId, entityData); List> result = new ArrayList<>(); EntityId entityId = constructEntityId(entityData.getEntityType(), entityData.getEntityIdMSB(), entityData.getEntityIdLSB()); - if ((entityData.hasPostAttributesMsg() || entityData.hasPostTelemetryMsg() || entityData.hasAttributesUpdatedMsg()) && entityId != null) { - Pair pair = getBaseMsgMetadataAndCustomerId(tenantId, entityId); - TbMsgMetaData metaData = pair.getKey(); - CustomerId customerId = pair.getValue(); - metaData.putValue(DataConstants.MSG_SOURCE_KEY, getMsgSourceKey()); - if (entityData.hasPostAttributesMsg()) { - result.add(processPostAttributes(tenantId, customerId, entityId, entityData.getPostAttributesMsg(), metaData)); - } - if (entityData.hasAttributesUpdatedMsg()) { - metaData.putValue("scope", entityData.getPostAttributeScope()); - result.add(processAttributesUpdate(tenantId, customerId, entityId, entityData.getAttributesUpdatedMsg(), metaData)); - } - if (entityData.hasPostTelemetryMsg()) { - result.add(processPostTelemetry(tenantId, customerId, entityId, entityData.getPostTelemetryMsg(), metaData)); - } - if (EntityType.DEVICE.equals(entityId.getEntityType())) { - DeviceId deviceId = new DeviceId(entityId.getId()); + if (entityId != null && isEntityExists(tenantId, entityId)) { + if ((entityData.hasPostAttributesMsg() || entityData.hasPostTelemetryMsg() || entityData.hasAttributesUpdatedMsg())) { + Pair pair = getBaseMsgMetadataAndCustomerId(tenantId, entityId); + TbMsgMetaData metaData = pair.getKey(); + CustomerId customerId = pair.getValue(); + metaData.putValue(DataConstants.MSG_SOURCE_KEY, getMsgSourceKey()); + if (entityData.hasPostAttributesMsg()) { + result.add(processPostAttributes(tenantId, customerId, entityId, entityData.getPostAttributesMsg(), metaData)); + } + if (entityData.hasAttributesUpdatedMsg()) { + metaData.putValue("scope", entityData.getPostAttributeScope()); + result.add(processAttributesUpdate(tenantId, customerId, entityId, entityData.getAttributesUpdatedMsg(), metaData)); + } + if (entityData.hasPostTelemetryMsg()) { + result.add(processPostTelemetry(tenantId, customerId, entityId, entityData.getPostTelemetryMsg(), metaData)); + } + if (EntityType.DEVICE.equals(entityId.getEntityType())) { + DeviceId deviceId = new DeviceId(entityId.getId()); - long currentTs = System.currentTimeMillis(); + long currentTs = System.currentTimeMillis(); - TransportProtos.DeviceActivityProto deviceActivityMsg = TransportProtos.DeviceActivityProto.newBuilder() - .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) - .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) - .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) - .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) - .setLastActivityTime(currentTs).build(); + TransportProtos.DeviceActivityProto deviceActivityMsg = TransportProtos.DeviceActivityProto.newBuilder() + .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) + .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) + .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) + .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) + .setLastActivityTime(currentTs).build(); - log.trace("[{}][{}] device activity time is going to be updated, ts {}", tenantId, deviceId, currentTs); + log.trace("[{}][{}] device activity time is going to be updated, ts {}", tenantId, deviceId, currentTs); - TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, deviceId); - tbCoreMsgProducer.send(tpi, new TbProtoQueueMsg<>(deviceId.getId(), - TransportProtos.ToCoreMsg.newBuilder().setDeviceActivityMsg(deviceActivityMsg).build()), null); + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, deviceId); + tbCoreMsgProducer.send(tpi, new TbProtoQueueMsg<>(deviceId.getId(), + TransportProtos.ToCoreMsg.newBuilder().setDeviceActivityMsg(deviceActivityMsg).build()), null); + } } - } - if (entityData.hasAttributeDeleteMsg()) { - result.add(processAttributeDeleteMsg(tenantId, entityId, entityData.getAttributeDeleteMsg(), entityData.getEntityType())); + if (entityData.hasAttributeDeleteMsg()) { + result.add(processAttributeDeleteMsg(tenantId, entityId, entityData.getAttributeDeleteMsg(), entityData.getEntityType())); + } + } else { + log.warn("Skipping telemetry update msg because entity doesn't exists on edge, {}", entityData); } return result; } From 096458d268d93c3beddae3b4f8b79ccfd06b5f07 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 25 Jan 2023 15:38:53 +0200 Subject: [PATCH 02/16] Updated edge root rule chain; --- .../rule_chains/edge_root_rule_chain.json | 74 +++++++++++-------- 1 file changed, 43 insertions(+), 31 deletions(-) diff --git a/application/src/main/data/json/tenant/edge_management/rule_chains/edge_root_rule_chain.json b/application/src/main/data/json/tenant/edge_management/rule_chains/edge_root_rule_chain.json index f908b1661e..69eae724a2 100644 --- a/application/src/main/data/json/tenant/edge_management/rule_chains/edge_root_rule_chain.json +++ b/application/src/main/data/json/tenant/edge_management/rule_chains/edge_root_rule_chain.json @@ -6,7 +6,8 @@ "firstRuleNodeId": null, "root": true, "debugMode": false, - "configuration": null + "configuration": null, + "externalId": null }, "metadata": { "firstNodeIndex": 0, @@ -23,7 +24,8 @@ "configuration": { "persistAlarmRulesState": false, "fetchAlarmRulesStateOnStart": false - } + }, + "externalId": null }, { "additionalInfo": { @@ -35,7 +37,8 @@ "debugMode": false, "configuration": { "defaultTTL": 0 - } + }, + "externalId": null }, { "additionalInfo": { @@ -47,7 +50,8 @@ "debugMode": false, "configuration": { "scope": "CLIENT_SCOPE" - } + }, + "externalId": null }, { "additionalInfo": { @@ -59,7 +63,8 @@ "debugMode": false, "configuration": { "version": 0 - } + }, + "externalId": null }, { "additionalInfo": { @@ -73,7 +78,8 @@ "scriptLang": "TBEL", "jsScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);", "tbelScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);" - } + }, + "externalId": null }, { "additionalInfo": { @@ -87,7 +93,8 @@ "scriptLang": "TBEL", "jsScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);", "tbelScript": "return '\\nIncoming message:\\n' + JSON.stringify(msg) + '\\nIncoming metadata:\\n' + JSON.stringify(metadata);" - } + }, + "externalId": null }, { "additionalInfo": { @@ -99,19 +106,34 @@ "debugMode": false, "configuration": { "timeoutInSeconds": 60 - } + }, + "externalId": null }, { "additionalInfo": { - "layoutX": 1129, - "layoutY": 52 + "layoutX": 1126, + "layoutY": 104 }, "type": "org.thingsboard.rule.engine.edge.TbMsgPushToCloudNode", "name": "Push to cloud", "debugMode": false, "configuration": { "scope": "SERVER_SCOPE" - } + }, + "externalId": null + }, + { + "additionalInfo": { + "layoutX": 826, + "layoutY": 601 + }, + "type": "org.thingsboard.rule.engine.edge.TbMsgPushToCloudNode", + "name": "Push to cloud", + "debugMode": false, + "configuration": { + "scope": "SERVER_SCOPE" + }, + "externalId": null } ], "connections": [ @@ -132,24 +154,14 @@ }, { "fromIndex": 3, - "toIndex": 6, - "type": "RPC Request to Device" - }, - { - "fromIndex": 3, - "toIndex": 5, - "type": "Other" + "toIndex": 1, + "type": "Post telemetry" }, { "fromIndex": 3, "toIndex": 2, "type": "Post attributes" }, - { - "fromIndex": 3, - "toIndex": 1, - "type": "Post telemetry" - }, { "fromIndex": 3, "toIndex": 4, @@ -157,23 +169,23 @@ }, { "fromIndex": 3, - "toIndex": 7, - "type": "Attributes Updated" + "toIndex": 5, + "type": "Other" }, { "fromIndex": 3, - "toIndex": 7, - "type": "Attributes Deleted" + "toIndex": 6, + "type": "RPC Request to Device" }, { "fromIndex": 3, - "toIndex": 7, - "type": "Timeseries Deleted" + "toIndex": 8, + "type": "Attributes Deleted" }, { "fromIndex": 3, - "toIndex": 7, - "type": "Timeseries Updated" + "toIndex": 8, + "type": "Attributes Updated" } ], "ruleChainConnections": null From 39ce423ef98259561ba86e73a7d25ad006536087 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 25 Jan 2023 15:40:26 +0200 Subject: [PATCH 03/16] Moved edge root rule chain from tenant dir into edge dir --- .../rule_chains/edge_root_rule_chain.json | 0 .../thingsboard/server/service/install/InstallScripts.java | 5 ++--- 2 files changed, 2 insertions(+), 3 deletions(-) rename application/src/main/data/json/{tenant/edge_management => edge}/rule_chains/edge_root_rule_chain.json (100%) diff --git a/application/src/main/data/json/tenant/edge_management/rule_chains/edge_root_rule_chain.json b/application/src/main/data/json/edge/rule_chains/edge_root_rule_chain.json similarity index 100% rename from application/src/main/data/json/tenant/edge_management/rule_chains/edge_root_rule_chain.json rename to application/src/main/data/json/edge/rule_chains/edge_root_rule_chain.json diff --git a/application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java b/application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java index 38ff5aacec..893ef5ec20 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java +++ b/application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java @@ -59,6 +59,7 @@ public class InstallScripts { public static final String JSON_DIR = "json"; public static final String SYSTEM_DIR = "system"; public static final String TENANT_DIR = "tenant"; + public static final String EDGE_DIR = "edge"; public static final String DEVICE_PROFILE_DIR = "device_profile"; public static final String DEMO_DIR = "demo"; public static final String RULE_CHAINS_DIR = "rule_chains"; @@ -68,8 +69,6 @@ public class InstallScripts { public static final String MODELS_DIR = "models"; public static final String CREDENTIALS_DIR = "credentials"; - public static final String EDGE_MANAGEMENT = "edge_management"; - public static final String JSON_EXT = ".json"; public static final String XML_EXT = ".xml"; @@ -103,7 +102,7 @@ public class InstallScripts { } private Path getEdgeRuleChainsDir() { - return Paths.get(getDataDir(), JSON_DIR, TENANT_DIR, EDGE_MANAGEMENT, RULE_CHAINS_DIR); + return Paths.get(getDataDir(), JSON_DIR, EDGE_DIR, RULE_CHAINS_DIR); } public String getDataDir() { From 7203a0b5c86357721f2229a0f46c747806dc3273 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 25 Jan 2023 15:56:36 +0200 Subject: [PATCH 04/16] Updated edge root rule chain Save Client Attributes rule node config --- .../main/data/json/edge/rule_chains/edge_root_rule_chain.json | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/application/src/main/data/json/edge/rule_chains/edge_root_rule_chain.json b/application/src/main/data/json/edge/rule_chains/edge_root_rule_chain.json index 69eae724a2..ec1341cc71 100644 --- a/application/src/main/data/json/edge/rule_chains/edge_root_rule_chain.json +++ b/application/src/main/data/json/edge/rule_chains/edge_root_rule_chain.json @@ -49,7 +49,8 @@ "name": "Save Client Attributes", "debugMode": false, "configuration": { - "scope": "CLIENT_SCOPE" + "scope": "CLIENT_SCOPE", + "notifyDevice": "false" }, "externalId": null }, From 4228baa33164f11eeaea6b7a33dbcb88c0004411 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Fri, 24 Mar 2023 10:51:03 +0200 Subject: [PATCH 05/16] Alarm processor - remove usage of deprecate createOrUpdateAlarm method --- .../service/edge/rpc/fetch/GeneralEdgeEventFetcher.java | 9 +-------- .../edge/rpc/processor/alarm/BaseAlarmProcessor.java | 5 ++--- .../pages/rulechain/rulechains-table-config.resolver.ts | 4 ++-- 3 files changed, 5 insertions(+), 13 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java index 327184e6a9..05b973b63c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java @@ -21,7 +21,6 @@ import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; -import org.thingsboard.server.common.data.page.SortOrder; import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.dao.edge.EdgeEventService; @@ -33,13 +32,7 @@ public class GeneralEdgeEventFetcher implements EdgeEventFetcher { @Override public PageLink getPageLink(int pageSize) { - return new TimePageLink( - pageSize, - 0, - null, - new SortOrder("createdTime", SortOrder.Direction.ASC), - queueStartTs, - null); + return new TimePageLink(pageSize, 0, null, null, queueStartTs, null); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java index d85750f7d2..c695ca84c6 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java @@ -19,10 +19,10 @@ import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.springframework.stereotype.Component; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.alarm.Alarm; +import org.thingsboard.server.common.data.alarm.AlarmCreateOrUpdateActiveRequest; import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmStatus; import org.thingsboard.server.common.data.edge.EdgeEventActionType; @@ -31,7 +31,6 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; -import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; import java.util.UUID; @@ -68,7 +67,7 @@ public abstract class BaseAlarmProcessor extends BaseEdgeProcessor { existentAlarm.setAckTs(alarmUpdateMsg.getAckTs()); existentAlarm.setEndTs(alarmUpdateMsg.getEndTs()); existentAlarm.setDetails(JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails())); - alarmService.createOrUpdateAlarm(existentAlarm); + alarmService.createAlarm(AlarmCreateOrUpdateActiveRequest.fromAlarm(existentAlarm)); break; case ALARM_ACK_RPC_MESSAGE: if (existentAlarm != null) { diff --git a/ui-ngx/src/app/modules/home/pages/rulechain/rulechains-table-config.resolver.ts b/ui-ngx/src/app/modules/home/pages/rulechain/rulechains-table-config.resolver.ts index d44fc67764..5360a97054 100644 --- a/ui-ngx/src/app/modules/home/pages/rulechain/rulechains-table-config.resolver.ts +++ b/ui-ngx/src/app/modules/home/pages/rulechain/rulechains-table-config.resolver.ts @@ -142,11 +142,11 @@ export class RuleChainsTableConfigResolver implements Resolve('root', 'rulechain.edge-template-root', '60px', + new EntityTableColumn('root', 'rulechain.edge-template-root', '70px', entity => { return checkBoxCell(entity.root); }), - new EntityTableColumn('assignToEdge', 'rulechain.assign-to-edge', '60px', + new EntityTableColumn('assignToEdge', 'rulechain.assign-to-edge', '70px', entity => { return checkBoxCell(this.isAutoAssignToEdgeRuleChain(entity)); }) From fc4be34d55971616535081e9444022d41fb1537f Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Fri, 24 Mar 2023 14:22:34 +0200 Subject: [PATCH 06/16] Added edgeAlarmId into AlarmCreateOrUpdateActiveRequest. BaseAlarmProcessor - use alarmId from cloud --- .../processor/alarm/BaseAlarmProcessor.java | 62 ++++++++++--------- .../server/edge/BaseAlarmEdgeTest.java | 6 ++ .../AlarmCreateOrUpdateActiveRequest.java | 8 +++ .../server/dao/sql/alarm/JpaAlarmDao.java | 2 +- 4 files changed, 49 insertions(+), 29 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java index c695ca84c6..84a768a886 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java @@ -42,51 +42,57 @@ public abstract class BaseAlarmProcessor extends BaseEdgeProcessor { log.trace("[{}] processAlarmMsg [{}]", tenantId, alarmUpdateMsg); EntityId originatorId = getAlarmOriginator(tenantId, alarmUpdateMsg.getOriginatorName(), EntityType.valueOf(alarmUpdateMsg.getOriginatorType())); + AlarmId alarmId = new AlarmId(new UUID(alarmUpdateMsg.getIdMSB(), alarmUpdateMsg.getIdLSB())); if (originatorId == null) { log.warn("Originator not found for the alarm msg {}", alarmUpdateMsg); return Futures.immediateFuture(null); } try { - Alarm existentAlarm = alarmService.findLatestActiveByOriginatorAndType(tenantId, originatorId, alarmUpdateMsg.getType()); switch (alarmUpdateMsg.getMsgType()) { case ENTITY_CREATED_RPC_MESSAGE: case ENTITY_UPDATED_RPC_MESSAGE: - if (existentAlarm == null || existentAlarm.getStatus().isCleared()) { - existentAlarm = new Alarm(); - existentAlarm.setTenantId(tenantId); - existentAlarm.setType(alarmUpdateMsg.getName()); - existentAlarm.setOriginator(originatorId); - existentAlarm.setSeverity(AlarmSeverity.valueOf(alarmUpdateMsg.getSeverity())); - existentAlarm.setStartTs(alarmUpdateMsg.getStartTs()); - existentAlarm.setClearTs(alarmUpdateMsg.getClearTs()); - existentAlarm.setPropagate(alarmUpdateMsg.getPropagate()); + Alarm alarm = alarmService.findAlarmById(tenantId, alarmId); + if (alarm == null) { + alarm = new Alarm(); + alarm.setTenantId(tenantId); + alarm.setType(alarmUpdateMsg.getName()); + alarm.setOriginator(originatorId); + alarm.setSeverity(AlarmSeverity.valueOf(alarmUpdateMsg.getSeverity())); + alarm.setStartTs(alarmUpdateMsg.getStartTs()); } var alarmStatus = AlarmStatus.valueOf(alarmUpdateMsg.getStatus()); - existentAlarm.setCleared(alarmStatus.isCleared()); - existentAlarm.setAcknowledged(alarmStatus.isAck()); - existentAlarm.setAckTs(alarmUpdateMsg.getAckTs()); - existentAlarm.setEndTs(alarmUpdateMsg.getEndTs()); - existentAlarm.setDetails(JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails())); - alarmService.createAlarm(AlarmCreateOrUpdateActiveRequest.fromAlarm(existentAlarm)); - break; + alarm.setClearTs(alarmUpdateMsg.getClearTs()); + alarm.setPropagate(alarmUpdateMsg.getPropagate()); + alarm.setCleared(alarmStatus.isCleared()); + alarm.setAcknowledged(alarmStatus.isAck()); + alarm.setAckTs(alarmUpdateMsg.getAckTs()); + alarm.setEndTs(alarmUpdateMsg.getEndTs()); + alarm.setDetails(JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails())); + alarmService.createAlarm(AlarmCreateOrUpdateActiveRequest.fromAlarm(alarm, null, alarmId)); + return Futures.immediateFuture(null); case ALARM_ACK_RPC_MESSAGE: - if (existentAlarm != null) { - alarmService.acknowledgeAlarm(tenantId, existentAlarm.getId(), alarmUpdateMsg.getAckTs()); + Alarm alarmToAck = alarmService.findAlarmById(tenantId, alarmId); + if (alarmToAck != null) { + alarmService.acknowledgeAlarm(tenantId, alarmId, alarmUpdateMsg.getAckTs()); } - break; + return Futures.immediateFuture(null); case ALARM_CLEAR_RPC_MESSAGE: - if (existentAlarm != null) { - alarmService.clearAlarm(tenantId, existentAlarm.getId(), - alarmUpdateMsg.getAckTs(), JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails())); + Alarm alarmToClear = alarmService.findAlarmById(tenantId, alarmId); + if (alarmToClear != null) { + alarmService.clearAlarm(tenantId, alarmId, alarmUpdateMsg.getClearTs(), + JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails())); } - break; + return Futures.immediateFuture(null); case ENTITY_DELETED_RPC_MESSAGE: - if (existentAlarm != null) { - alarmService.delAlarm(tenantId, existentAlarm.getId()); + Alarm alarmToDelete = alarmService.findAlarmById(tenantId, alarmId); + if (alarmToDelete != null) { + alarmService.delAlarm(tenantId, alarmId); } - break; + return Futures.immediateFuture(null); + case UNRECOGNIZED: + default: + return handleUnsupportedMsgType(alarmUpdateMsg.getMsgType()); } - return Futures.immediateFuture(null); } catch (Exception e) { log.error("[{}] Failed to process alarm update msg [{}]", tenantId, alarmUpdateMsg, e); return Futures.immediateFailedFuture(e); diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java index 1cf2cd44bb..4c73128f3b 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java @@ -25,6 +25,7 @@ import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmStatus; +import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; @@ -33,6 +34,7 @@ import org.thingsboard.server.gen.edge.v1.UplinkMsg; import java.util.List; import java.util.Optional; +import java.util.UUID; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @@ -42,8 +44,11 @@ abstract public class BaseAlarmEdgeTest extends AbstractEdgeTest { public void testSendAlarmToCloud() throws Exception { Device device = saveDeviceOnCloudAndVerifyDeliveryToEdge(); + UUID alarmUUID = UUID.randomUUID(); UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); AlarmUpdateMsg.Builder alarmUpdateMgBuilder = AlarmUpdateMsg.newBuilder(); + alarmUpdateMgBuilder.setIdMSB(alarmUUID.getMostSignificantBits()); + alarmUpdateMgBuilder.setIdLSB(alarmUUID.getLeastSignificantBits()); alarmUpdateMgBuilder.setName("alarm from edge"); alarmUpdateMgBuilder.setStatus(AlarmStatus.ACTIVE_UNACK.name()); alarmUpdateMgBuilder.setSeverity(AlarmSeverity.CRITICAL.name()); @@ -65,6 +70,7 @@ abstract public class BaseAlarmEdgeTest extends AbstractEdgeTest { Optional foundAlarm = alarms.stream().filter(alarm -> alarm.getType().equals("alarm from edge")).findAny(); Assert.assertTrue(foundAlarm.isPresent()); AlarmInfo alarmInfo = foundAlarm.get(); + Assert.assertEquals(new AlarmId(alarmUUID), alarmInfo.getId()); Assert.assertEquals(device.getId(), alarmInfo.getOriginator()); Assert.assertEquals(AlarmStatus.ACTIVE_UNACK, alarmInfo.getStatus()); Assert.assertEquals(AlarmSeverity.CRITICAL, alarmInfo.getSeverity()); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/alarm/AlarmCreateOrUpdateActiveRequest.java b/common/data/src/main/java/org/thingsboard/server/common/data/alarm/AlarmCreateOrUpdateActiveRequest.java index ecf882e2c9..0f3dd5b636 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/alarm/AlarmCreateOrUpdateActiveRequest.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/alarm/AlarmCreateOrUpdateActiveRequest.java @@ -19,6 +19,7 @@ import com.fasterxml.jackson.databind.JsonNode; import io.swagger.annotations.ApiModelProperty; import lombok.Builder; import lombok.Data; +import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -61,11 +62,17 @@ public class AlarmCreateOrUpdateActiveRequest implements AlarmModificationReques private UserId userId; + private AlarmId edgeAlarmId; + public static AlarmCreateOrUpdateActiveRequest fromAlarm(Alarm a) { return fromAlarm(a, null); } public static AlarmCreateOrUpdateActiveRequest fromAlarm(Alarm a, UserId userId) { + return fromAlarm(a, userId, null); + } + + public static AlarmCreateOrUpdateActiveRequest fromAlarm(Alarm a, UserId userId, AlarmId edgeAlarmId) { return AlarmCreateOrUpdateActiveRequest.builder() .tenantId(a.getTenantId()) .customerId(a.getCustomerId()) @@ -81,6 +88,7 @@ public class AlarmCreateOrUpdateActiveRequest implements AlarmModificationReques .propagateToTenant(a.isPropagateToTenant()) .propagateRelationTypes(a.getPropagateRelationTypes()).build()) .userId(userId) + .edgeAlarmId(edgeAlarmId) .build(); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java index f1abd07d95..321208eb8f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/alarm/JpaAlarmDao.java @@ -245,7 +245,7 @@ public class JpaAlarmDao extends JpaAbstractDao implements A return toAlarmApiResult(alarmRepository.createOrUpdateActiveAlarm( request.getTenantId().getId(), request.getCustomerId() != null ? request.getCustomerId().getId() : CustomerId.NULL_UUID, - UUID.randomUUID(), + request.getEdgeAlarmId() != null ? request.getEdgeAlarmId().getId() : UUID.randomUUID(), System.currentTimeMillis(), request.getOriginator().getId(), request.getOriginator().getEntityType().ordinal(), From 19dcc85262bbb8e76b5f57f3961d9494e992040c Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 28 Mar 2023 12:30:06 +0300 Subject: [PATCH 07/16] Fixed circular dependency issue - WebSocketService and TbEntityDataSubscriptionService --- .../subscription/DefaultTbEntityDataSubscriptionService.java | 1 + 1 file changed, 1 insertion(+) diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java index 9fa8429f78..f001cc4822 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java @@ -95,6 +95,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc private final Map> subscriptionsBySessionId = new ConcurrentHashMap<>(); @Autowired + @Lazy private WebSocketService wsService; @Autowired From 98ee96256075719c5189005f8da49110e0684136 Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Tue, 28 Mar 2023 13:42:19 +0300 Subject: [PATCH 08/16] tbel: - ver 1.0.6 - public interface ExecutionObject -> int getId(); to int getExecutionObjectId(); --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index b577bf4142..fb1aafd6b6 100755 --- a/pom.xml +++ b/pom.xml @@ -78,7 +78,7 @@ 3.5.5 3.21.9 1.42.1 - 1.0.5 + 1.0.6 1.18.18 1.2.4 1.2.5 From f7a1a7eabee0176be94f8f2aa87e0c3103e79435 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Tue, 28 Mar 2023 15:05:57 +0300 Subject: [PATCH 09/16] Code review changes --- .../service/edge/rpc/fetch/GeneralEdgeEventFetcher.java | 9 ++++++++- .../DefaultTbEntityDataSubscriptionService.java | 1 - 2 files changed, 8 insertions(+), 2 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java index 05b973b63c..327184e6a9 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java @@ -21,6 +21,7 @@ import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.common.data.page.SortOrder; import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.dao.edge.EdgeEventService; @@ -32,7 +33,13 @@ public class GeneralEdgeEventFetcher implements EdgeEventFetcher { @Override public PageLink getPageLink(int pageSize) { - return new TimePageLink(pageSize, 0, null, null, queueStartTs, null); + return new TimePageLink( + pageSize, + 0, + null, + new SortOrder("createdTime", SortOrder.Direction.ASC), + queueStartTs, + null); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java index f001cc4822..9fa8429f78 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java @@ -95,7 +95,6 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc private final Map> subscriptionsBySessionId = new ConcurrentHashMap<>(); @Autowired - @Lazy private WebSocketService wsService; @Autowired From bacda7e3f6752904de9b6af7a7cb30ccfd28cec4 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Fri, 31 Mar 2023 12:32:19 +0300 Subject: [PATCH 10/16] Base Edge Processor - Added check for tenant entity --- .../server/service/edge/rpc/processor/BaseEdgeProcessor.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java index 4c8cf585a5..ac7c791338 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/BaseEdgeProcessor.java @@ -488,6 +488,8 @@ public abstract class BaseEdgeProcessor { protected boolean isEntityExists(TenantId tenantId, EntityId entityId) { switch (entityId.getEntityType()) { + case TENANT: + return tenantService.findTenantById(tenantId) != null; case DEVICE: return deviceService.findDeviceById(tenantId, new DeviceId(entityId.getId())) != null; case ASSET: From 8d4f1e31dc6693c3d560ecc94453e94893dae046 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Mon, 3 Apr 2023 13:41:02 +0300 Subject: [PATCH 11/16] Code review updates - remove redundant alarmService.findById - use correct alarm service method based on msg type --- .../processor/alarm/BaseAlarmProcessor.java | 23 +++++++++++-------- .../server/edge/BaseAlarmEdgeTest.java | 17 ++++++++++++++ .../engine/edge/AbstractTbMsgPushNode.java | 17 +++++++++++++- 3 files changed, 46 insertions(+), 11 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java index 84a768a886..be64f475a2 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/alarm/BaseAlarmProcessor.java @@ -25,6 +25,7 @@ import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.AlarmCreateOrUpdateActiveRequest; import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmStatus; +import org.thingsboard.server.common.data.alarm.AlarmUpdateRequest; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.EntityId; @@ -51,15 +52,13 @@ public abstract class BaseAlarmProcessor extends BaseEdgeProcessor { switch (alarmUpdateMsg.getMsgType()) { case ENTITY_CREATED_RPC_MESSAGE: case ENTITY_UPDATED_RPC_MESSAGE: - Alarm alarm = alarmService.findAlarmById(tenantId, alarmId); - if (alarm == null) { - alarm = new Alarm(); - alarm.setTenantId(tenantId); - alarm.setType(alarmUpdateMsg.getName()); - alarm.setOriginator(originatorId); - alarm.setSeverity(AlarmSeverity.valueOf(alarmUpdateMsg.getSeverity())); - alarm.setStartTs(alarmUpdateMsg.getStartTs()); - } + Alarm alarm = new Alarm(); + alarm.setId(alarmId); + alarm.setTenantId(tenantId); + alarm.setType(alarmUpdateMsg.getName()); + alarm.setOriginator(originatorId); + alarm.setSeverity(AlarmSeverity.valueOf(alarmUpdateMsg.getSeverity())); + alarm.setStartTs(alarmUpdateMsg.getStartTs()); var alarmStatus = AlarmStatus.valueOf(alarmUpdateMsg.getStatus()); alarm.setClearTs(alarmUpdateMsg.getClearTs()); alarm.setPropagate(alarmUpdateMsg.getPropagate()); @@ -68,7 +67,11 @@ public abstract class BaseAlarmProcessor extends BaseEdgeProcessor { alarm.setAckTs(alarmUpdateMsg.getAckTs()); alarm.setEndTs(alarmUpdateMsg.getEndTs()); alarm.setDetails(JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails())); - alarmService.createAlarm(AlarmCreateOrUpdateActiveRequest.fromAlarm(alarm, null, alarmId)); + if (UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE.equals(alarmUpdateMsg.getMsgType())) { + alarmService.createAlarm(AlarmCreateOrUpdateActiveRequest.fromAlarm(alarm, null, alarmId)); + } else { + alarmService.updateAlarm(AlarmUpdateRequest.fromAlarm(alarm)); + } return Futures.immediateFuture(null); case ALARM_ACK_RPC_MESSAGE: Alarm alarmToAck = alarmService.findAlarmById(tenantId, alarmId); diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java index 4c73128f3b..c06093535a 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseAlarmEdgeTest.java @@ -19,6 +19,7 @@ import com.fasterxml.jackson.core.type.TypeReference; import com.google.protobuf.AbstractMessage; import org.junit.Assert; import org.junit.Test; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.alarm.Alarm; @@ -91,12 +92,28 @@ abstract public class BaseAlarmEdgeTest extends AbstractEdgeTest { Assert.assertTrue(latestMessage instanceof AlarmUpdateMsg); AlarmUpdateMsg alarmUpdateMsg = (AlarmUpdateMsg) latestMessage; Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, alarmUpdateMsg.getMsgType()); + Assert.assertEquals(savedAlarm.getUuidId().getMostSignificantBits(), alarmUpdateMsg.getIdMSB()); + Assert.assertEquals(savedAlarm.getUuidId().getLeastSignificantBits(), alarmUpdateMsg.getIdLSB()); Assert.assertEquals(savedAlarm.getType(), alarmUpdateMsg.getType()); Assert.assertEquals(savedAlarm.getName(), alarmUpdateMsg.getName()); Assert.assertEquals(device.getName(), alarmUpdateMsg.getOriginatorName()); Assert.assertEquals(savedAlarm.getStatus().name(), alarmUpdateMsg.getStatus()); Assert.assertEquals(savedAlarm.getSeverity().name(), alarmUpdateMsg.getSeverity()); + // update alarm + String updatedDetails = "{\"testKey\":\"testValue\"}"; + savedAlarm.setDetails(JacksonUtil.OBJECT_MAPPER.readTree(updatedDetails)); + edgeImitator.expectMessageAmount(1); + savedAlarm = doPost("/api/alarm", savedAlarm, Alarm.class); + Assert.assertTrue(edgeImitator.waitForMessages()); + latestMessage = edgeImitator.getLatestMessage(); + Assert.assertTrue(latestMessage instanceof AlarmUpdateMsg); + alarmUpdateMsg = (AlarmUpdateMsg) latestMessage; + Assert.assertEquals(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE, alarmUpdateMsg.getMsgType()); + Assert.assertEquals(savedAlarm.getUuidId().getMostSignificantBits(), alarmUpdateMsg.getIdMSB()); + Assert.assertEquals(savedAlarm.getUuidId().getLeastSignificantBits(), alarmUpdateMsg.getIdLSB()); + Assert.assertEquals(updatedDetails, alarmUpdateMsg.getDetails()); + // ack alarm edgeImitator.expectMessageAmount(1); doPost("/api/alarm/" + savedAlarm.getUuidId() + "/ack"); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java index 6d8063b8b5..7be7dae2de 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java @@ -74,7 +74,8 @@ public abstract class AbstractTbMsgPushNode entityBody = new HashMap<>(); @@ -107,6 +108,20 @@ public abstract class AbstractTbMsgPushNode Date: Mon, 3 Apr 2023 19:23:13 +0200 Subject: [PATCH 12/16] fixed upgrade script --- .../DefaultSystemDataLoaderService.java | 23 ++++++------------- .../service/install/InstallScripts.java | 18 +++++++-------- 2 files changed, 15 insertions(+), 26 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java b/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java index a182a4d83f..bd68aad316 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java @@ -22,7 +22,6 @@ import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.Getter; import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; @@ -91,6 +90,7 @@ import org.thingsboard.server.dao.device.DeviceProfileService; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.notification.NotificationSettingsService; +import org.thingsboard.server.dao.notification.NotificationTargetService; import org.thingsboard.server.dao.queue.QueueService; import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.settings.AdminSettingsService; @@ -177,6 +177,9 @@ public class DefaultSystemDataLoaderService implements SystemDataLoaderService { @Autowired private NotificationSettingsService notificationSettingsService; + @Autowired + private NotificationTargetService notificationTargetService; + @Bean protected BCryptPasswordEncoder passwordEncoder() { return new BCryptPasswordEncoder(); @@ -679,27 +682,15 @@ public class DefaultSystemDataLoaderService implements SystemDataLoaderService { @Override public void createDefaultNotificationConfigs() { - try { - log.info("Creating default notification configs for system admin"); + log.info("Creating default notification configs for system admin"); + if (notificationTargetService.findNotificationTargetsByTenantId(TenantId.SYS_TENANT_ID, new PageLink(1)).getTotalElements() == 0) { notificationSettingsService.createDefaultNotificationConfigs(TenantId.SYS_TENANT_ID); - } catch (Exception e) { - if (StringUtils.contains(e.getMessage(), "already exists")) { - log.info("Default notification configs are already present for system admin, skipping"); - } else { - throw e; - } } PageDataIterable tenants = new PageDataIterable<>(tenantService::findTenantsIds, 500); log.info("Creating default notification configs for all tenants"); for (TenantId tenantId : tenants) { - try { + if (notificationTargetService.findNotificationTargetsByTenantId(tenantId, new PageLink(1)).getTotalElements() == 0) { notificationSettingsService.createDefaultNotificationConfigs(tenantId); - } catch (Exception e) { - if (StringUtils.contains(e.getMessage(), "already exists")) { - log.info("Default notification configs are already present for tenant {}, skipping", tenantId); - } else { - throw e; - } } } } diff --git a/application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java b/application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java index 0ccb201cb9..de467f572a 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java +++ b/application/src/main/java/org/thingsboard/server/service/install/InstallScripts.java @@ -293,17 +293,15 @@ public class InstallScripts { } private void doSaveLwm2mResource(TbResource resource) throws ThingsboardException { - try { - log.trace("Executing saveResource [{}]", resource); - if (StringUtils.isEmpty(resource.getData())) { - throw new DataValidationException("Resource data should be specified!"); - } - toLwm2mResource(resource); + log.trace("Executing saveResource [{}]", resource); + if (StringUtils.isEmpty(resource.getData())) { + throw new DataValidationException("Resource data should be specified!"); + } + toLwm2mResource(resource); + TbResource foundResource = + resourceService.getResource(TenantId.SYS_TENANT_ID, ResourceType.LWM2M_MODEL, resource.getResourceKey()); + if (foundResource == null) { resourceService.saveResource(resource); - } catch (DataValidationException e) { - log.debug("[{}] {}", resource.getFileName(), e.getMessage()); - } catch (Exception ex) { - throw ex; } } } From 6b93de8bef4a1ddc427be28a57570912cb8dac7a Mon Sep 17 00:00:00 2001 From: Andrew Shvayka Date: Tue, 4 Apr 2023 18:10:19 +0300 Subject: [PATCH 13/16] Update pom.xml --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index fb1aafd6b6..91bc3832b1 100755 --- a/pom.xml +++ b/pom.xml @@ -78,7 +78,7 @@ 3.5.5 3.21.9 1.42.1 - 1.0.6 + 1.0.7 1.18.18 1.2.4 1.2.5 From 09b3299015703ba18af4d99c282bad66f1d07010 Mon Sep 17 00:00:00 2001 From: Andrew Shvayka Date: Tue, 4 Apr 2023 18:15:21 +0300 Subject: [PATCH 14/16] Revert "[Fix bug][3.5] tbel json decode id" --- pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pom.xml b/pom.xml index 91bc3832b1..b577bf4142 100755 --- a/pom.xml +++ b/pom.xml @@ -78,7 +78,7 @@ 3.5.5 3.21.9 1.42.1 - 1.0.7 + 1.0.5 1.18.18 1.2.4 1.2.5 From 19e58193b52da1afe5b6eda55c95729faa5559f6 Mon Sep 17 00:00:00 2001 From: nickAS21 Date: Wed, 5 Apr 2023 13:32:57 +0300 Subject: [PATCH 15/16] sparkplug: change locale --- ui-ngx/src/assets/locale/locale.constant-en_US.json | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index 7d9eed195b..c4e9b7f46b 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -1396,10 +1396,10 @@ "create-new-device-profile": "Create a new one!", "mqtt-device-topic-filters": "MQTT device topic filters", "mqtt-device-topic-filters-unique": "MQTT device topic filters need to be unique.", - "mqtt-device-topic-filters-spark-plug": "MQTT device topic filters SparkPlug.", - "mqtt-device-topic-filters-spark-plug-hint": "Default - telemetry. Example: namespace/group_id/message_type/edge_node_id/[device_id].", - "mqtt-device-topic-filters-spark-plug-attribute-metric-names": "SparkPlug attributes metric names", - "mqtt-device-topic-filters-spark-plug-attribute-metric-names-hint": "Names of SparkPlug metrics that will be stored as device attributes. All other metrics will be stored as device telemetry", + "mqtt-device-topic-filters-spark-plug": "MQTT Sparkplug B Edge of Network (EoN) node.", + "mqtt-device-topic-filters-spark-plug-hint": "Allow connections from EoN nodes with Sparkplug B payload and topic format.", + "mqtt-device-topic-filters-spark-plug-attribute-metric-names": "SparkPlug metrics to store as attributes.", + "mqtt-device-topic-filters-spark-plug-attribute-metric-names-hint": "Names of SparkPlug metrics that will be stored as device attributes. All other metrics will be stored as device telemetry.", "mqtt-device-payload-type": "MQTT device payload", "mqtt-device-payload-type-json": "JSON", "mqtt-device-payload-type-proto": "Protobuf", From b1884f675a62c61b81c3e1331c560b964edc4da1 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Wed, 5 Apr 2023 15:43:09 +0300 Subject: [PATCH 16/16] Update service and TBEL version improvement --- .../NewPlatformVersionTriggerProcessor.java | 2 +- .../service/update/DefaultUpdateService.java | 16 +++++++++------- .../server/common/data/UpdateMessage.java | 12 ++++++++---- pom.xml | 2 +- 4 files changed, 19 insertions(+), 13 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/NewPlatformVersionTriggerProcessor.java b/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/NewPlatformVersionTriggerProcessor.java index ba346caab8..030409d014 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/NewPlatformVersionTriggerProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/rule/trigger/NewPlatformVersionTriggerProcessor.java @@ -44,7 +44,7 @@ public class NewPlatformVersionTriggerProcessor implements NotificationRuleTrigg @Override public NotificationInfo constructNotificationInfo(NewPlatformVersionTrigger trigger, NewPlatformVersionNotificationRuleTriggerConfig triggerConfig) { return NewPlatformVersionNotificationInfo.builder() - .message(trigger.getMessage().getMessage()) + .message("New version available - " + trigger.getMessage().getLatestVersion()) .build(); } diff --git a/application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java b/application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java index d9f49ae4d3..46f3d89a50 100644 --- a/application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java +++ b/application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java @@ -21,6 +21,7 @@ import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.info.BuildProperties; import org.springframework.stereotype.Service; import org.springframework.web.client.RestTemplate; import org.thingsboard.common.util.ThingsBoardThreadFactory; @@ -57,6 +58,9 @@ public class DefaultUpdateService implements UpdateService { @Value("${updates.enabled}") private boolean updatesEnabled; + @Autowired(required = false) + private BuildProperties buildProperties; + @Autowired private NotificationRuleProcessingService notificationRuleProcessingService; @@ -73,14 +77,11 @@ public class DefaultUpdateService implements UpdateService { @PostConstruct private void init() { - updateMessage = new UpdateMessage("", false, ""); + version = buildProperties != null ? buildProperties.getVersion() : "unknown"; + updateMessage = new UpdateMessage(false, version, "", ""); if (updatesEnabled) { try { platform = System.getProperty("platform", "unknown"); - version = getClass().getPackage().getImplementationVersion(); - if (version == null) { - version = "unknown"; - } instanceId = parseInstanceId(); checkUpdatesFuture = scheduler.scheduleAtFixedRate(checkUpdatesRunnable, 0, 1, TimeUnit.HOURS); } catch (Exception e) { @@ -131,9 +132,10 @@ public class DefaultUpdateService implements UpdateService { JsonNode response = restClient.postForObject(UPDATE_SERVER_BASE_URL + "/api/thingsboard/updates", request, JsonNode.class); UpdateMessage prevUpdateMessage = updateMessage; updateMessage = new UpdateMessage( - response.get("message").asText(), response.get("updateAvailable").asBoolean(), - version + version, + "3.5.0", + "https://thingsboard.io/docs/user-guide/install/pe/upgrade-instructions" ); if (updateMessage.isUpdateAvailable() && !updateMessage.equals(prevUpdateMessage)) { notificationRuleProcessingService.process(TenantId.SYS_TENANT_ID, NewPlatformVersionTrigger.builder() diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/UpdateMessage.java b/common/data/src/main/java/org/thingsboard/server/common/data/UpdateMessage.java index b5ecdcbf79..aecfcc16b4 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/UpdateMessage.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/UpdateMessage.java @@ -23,10 +23,14 @@ import lombok.Data; @Data public class UpdateMessage { - @ApiModelProperty(position = 1, value = "The message about new platform update available.") - private final String message; - @ApiModelProperty(position = 2, value = "'True' if new platform update is available.") + @ApiModelProperty(position = 1, value = "'True' if new platform update is available.") private final boolean isUpdateAvailable; - @ApiModelProperty(position = 3, value = "Current ThingsBoard version.") + @ApiModelProperty(position = 2, value = "Current ThingsBoard version.") private final String currentVersion; + @ApiModelProperty(position = 3, value = "Latest ThingsBoard version.") + private final String latestVersion; + @ApiModelProperty(position = 4, value = "Upgrade instructions URL.") + private final String upgradeInstructionsUrl; + + } diff --git a/pom.xml b/pom.xml index b577bf4142..fb1aafd6b6 100755 --- a/pom.xml +++ b/pom.xml @@ -78,7 +78,7 @@ 3.5.5 3.21.9 1.42.1 - 1.0.5 + 1.0.6 1.18.18 1.2.4 1.2.5