From d8dbf94162c68f3bbaf781ffea01df71cdf47f71 Mon Sep 17 00:00:00 2001 From: AndriiD Date: Fri, 31 Mar 2023 16:27:00 +0300 Subject: [PATCH 1/6] added entity check in rulenode init --- .../org/thingsboard/rule/engine/flow/TbRuleChainInputNode.java | 1 + 1 file changed, 1 insertion(+) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/flow/TbRuleChainInputNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/flow/TbRuleChainInputNode.java index 3c431fd602..ac875688fa 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/flow/TbRuleChainInputNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/flow/TbRuleChainInputNode.java @@ -56,6 +56,7 @@ public class TbRuleChainInputNode implements TbNode { public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { this.config = TbNodeUtils.convert(configuration, TbRuleChainInputNodeConfiguration.class); this.ruleChainId = new RuleChainId(UUID.fromString(config.getRuleChainId())); + ctx.checkTenantEntity(ruleChainId); } @Override From 1b29c97008b4e2b90d5eb63dd8aa6cdf900d9436 Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 5 Apr 2023 13:41:09 +0300 Subject: [PATCH 2/6] Edge Session - improved error logging in case general process is interrupred because of sync started --- .../server/service/edge/rpc/EdgeGrpcSession.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 1717fea812..729de18614 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -423,6 +423,7 @@ public final class EdgeGrpcSession implements Closeable { stopCurrentSendDownlinkMsgsTask(null); } } catch (Exception e) { + log.warn("[{}] Failed to send downlink msgs. Error msg {}", this.sessionId, e.getMessage(), e); stopCurrentSendDownlinkMsgsTask(e); } }; @@ -669,19 +670,18 @@ public final class EdgeGrpcSession implements Closeable { } private void interruptPreviousSendDownlinkMsgsTask() { - String msg = String.format("[%s] Previous send downlink future was not properly completed, stopping it now!", this.sessionId); - stopCurrentSendDownlinkMsgsTask(new RuntimeException(msg)); + log.info("[{}]Previous send downlink future was not properly completed, stopping it now!", this.sessionId); + stopCurrentSendDownlinkMsgsTask(new RuntimeException()); } private void interruptGeneralProcessingOnSync(TenantId tenantId, EdgeId edgeId) { - String msg = String.format("[%s][%s] Sync process started. General processing interrupted!", tenantId, edgeId); - stopCurrentSendDownlinkMsgsTask(new RuntimeException(msg)); + log.info("[{}][{}][{}] Sync process started. General processing interrupted!", this.sessionId, tenantId, edgeId); + stopCurrentSendDownlinkMsgsTask(new RuntimeException()); } public void stopCurrentSendDownlinkMsgsTask(Exception e) { if (sessionState.getSendDownlinkMsgsFuture() != null && !sessionState.getSendDownlinkMsgsFuture().isDone()) { if (e != null) { - log.warn(e.getMessage(), e); sessionState.getSendDownlinkMsgsFuture().setException(e); } else { sessionState.getSendDownlinkMsgsFuture().set(null); From 463ad698c4b426d9668bf4172935c5564d19a0bb Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 5 Apr 2023 15:25:03 +0300 Subject: [PATCH 3/6] Edge telemetry processor - fixed processing of delete attribute request for non device entities --- .../telemetry/BaseTelemetryProcessor.java | 40 +++++++------ .../server/edge/BaseTelemetryEdgeTest.java | 60 +++++++++++++++++++ 2 files changed, 83 insertions(+), 17 deletions(-) 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 e50ffd047a..fece86b92a 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 @@ -18,6 +18,7 @@ package org.thingsboard.server.service.edge.rpc.processor.telemetry; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.FutureCallback; +import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.SettableFuture; import com.google.gson.Gson; @@ -283,26 +284,31 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor { private ListenableFuture processAttributeDeleteMsg(TenantId tenantId, EntityId entityId, AttributeDeleteMsg attributeDeleteMsg, String entityType) { - SettableFuture futureToSet = SettableFuture.create(); + String scope = attributeDeleteMsg.getScope(); List attributeKeys = attributeDeleteMsg.getAttributeNamesList(); - attributesService.removeAll(tenantId, entityId, scope, attributeKeys); - if (EntityType.DEVICE.name().equals(entityType)) { - tbClusterService.pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete( - tenantId, (DeviceId) entityId, scope, attributeKeys), new TbQueueCallback() { - @Override - public void onSuccess(TbQueueMsgMetadata metadata) { - futureToSet.set(null); - } + ListenableFuture> removeAllFuture = attributesService.removeAll(tenantId, entityId, scope, attributeKeys); + return Futures.transformAsync(removeAllFuture, removeAttributes -> { + if (EntityType.DEVICE.name().equals(entityType)) { + SettableFuture futureToSet = SettableFuture.create(); + tbClusterService.pushMsgToCore(DeviceAttributesEventNotificationMsg.onDelete( + tenantId, (DeviceId) entityId, scope, attributeKeys), new TbQueueCallback() { + @Override + public void onSuccess(TbQueueMsgMetadata metadata) { + futureToSet.set(null); + } - @Override - public void onFailure(Throwable t) { - log.error("Can't process attribute delete msg [{}]", attributeDeleteMsg, t); - futureToSet.setException(t); - } - }); - } - return futureToSet; + @Override + public void onFailure(Throwable t) { + log.error("Can't process attribute delete msg [{}]", attributeDeleteMsg, t); + futureToSet.setException(t); + } + }); + return futureToSet; + } else { + return Futures.immediateFuture(null); + } + }, dbCallbackExecutorService); } public EntityDataProto convertTelemetryEventToEntityDataProto(EntityType entityType, diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java index ff6728cecc..8793495b53 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java @@ -15,11 +15,20 @@ */ package org.thingsboard.server.edge; +import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ArrayNode; +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.gson.Gson; +import com.google.gson.reflect.TypeToken; import com.google.protobuf.AbstractMessage; +import org.awaitility.Awaitility; import org.junit.Assert; import org.junit.Test; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; +import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventType; @@ -27,9 +36,11 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.gen.edge.v1.AttributeDeleteMsg; import org.thingsboard.server.gen.edge.v1.DeviceUpdateMsg; import org.thingsboard.server.gen.edge.v1.EntityDataProto; +import org.thingsboard.server.gen.edge.v1.UplinkMsg; import org.thingsboard.server.gen.transport.TransportProtos; import java.util.List; +import java.util.concurrent.TimeUnit; abstract public class BaseTelemetryEdgeTest extends AbstractEdgeTest { @@ -233,4 +244,53 @@ abstract public class BaseTelemetryEdgeTest extends AbstractEdgeTest { Assert.assertEquals("value1", keyValueProto.getStringV()); } + @Test + public void testSendAttributesDeleteRequestToCloud_nonDeviceEntity() throws Exception { + edgeImitator.expectMessageAmount(2); + Asset savedAsset = saveAsset("Delete Attribute Test"); + doPost("/api/edge/" + edge.getUuidId() + "/asset/" + savedAsset.getUuidId(), Asset.class); + Assert.assertTrue(edgeImitator.waitForMessages()); + + final String attributeKey = "key1"; + ObjectNode attributesData = JacksonUtil.OBJECT_MAPPER.createObjectNode(); + attributesData.put(attributeKey, "value1"); + doPost("/api/plugins/telemetry/ASSET/" + savedAsset.getId() + "/attributes/" + DataConstants.SERVER_SCOPE, attributesData); + + // Wait before device attributes saved to database before deleting them + Awaitility.await() + .atMost(10, TimeUnit.SECONDS) + .until(() -> { + String urlTemplate = "/api/plugins/telemetry/ASSET/" + savedAsset.getId() + "/keys/attributes/" + DataConstants.SERVER_SCOPE; + List actualKeys = doGetAsyncTyped(urlTemplate, new TypeReference<>() {}); + return actualKeys != null && !actualKeys.isEmpty() && actualKeys.contains(attributeKey); + }); + + EntityDataProto.Builder builder = EntityDataProto.newBuilder() + .setEntityIdMSB(savedAsset.getUuidId().getMostSignificantBits()) + .setEntityIdLSB(savedAsset.getUuidId().getLeastSignificantBits()) + .setEntityType(savedAsset.getId().getEntityType().name()); + AttributeDeleteMsg.Builder attributeDeleteMsg = AttributeDeleteMsg.newBuilder(); + attributeDeleteMsg.setScope(DataConstants.SERVER_SCOPE); + ArrayNode arrayNode = JacksonUtil.OBJECT_MAPPER.createArrayNode(); + arrayNode.add(attributeKey); + List keys = new Gson().fromJson(arrayNode.toString(), new TypeToken<>(){}.getType()); + attributeDeleteMsg.addAllAttributeNames(keys); + attributeDeleteMsg.build(); + builder.setAttributeDeleteMsg(attributeDeleteMsg); + UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); + uplinkMsgBuilder.addEntityData(builder.build()); + + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + Assert.assertTrue(edgeImitator.waitForResponses()); + + Awaitility.await() + .atMost(10, TimeUnit.SECONDS) + .until(() -> { + String urlTemplate = "/api/plugins/telemetry/ASSET/" + savedAsset.getId() + "/keys/attributes/" + DataConstants.SERVER_SCOPE; + List actualKeys = doGetAsyncTyped(urlTemplate, new TypeReference<>() {}); + return actualKeys != null && actualKeys.isEmpty(); + }); + } + } From 9e189c30e64fcd854378f43dc26e326184470d8b Mon Sep 17 00:00:00 2001 From: Volodymyr Babak Date: Wed, 5 Apr 2023 16:09:41 +0300 Subject: [PATCH 4/6] Remove redundant JSON code --- .../thingsboard/server/edge/BaseTelemetryEdgeTest.java | 8 +------- 1 file changed, 1 insertion(+), 7 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java index 8793495b53..9a5fe26782 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java @@ -17,10 +17,7 @@ package org.thingsboard.server.edge; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.JsonNode; -import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; -import com.google.gson.Gson; -import com.google.gson.reflect.TypeToken; import com.google.protobuf.AbstractMessage; import org.awaitility.Awaitility; import org.junit.Assert; @@ -271,10 +268,7 @@ abstract public class BaseTelemetryEdgeTest extends AbstractEdgeTest { .setEntityType(savedAsset.getId().getEntityType().name()); AttributeDeleteMsg.Builder attributeDeleteMsg = AttributeDeleteMsg.newBuilder(); attributeDeleteMsg.setScope(DataConstants.SERVER_SCOPE); - ArrayNode arrayNode = JacksonUtil.OBJECT_MAPPER.createArrayNode(); - arrayNode.add(attributeKey); - List keys = new Gson().fromJson(arrayNode.toString(), new TypeToken<>(){}.getType()); - attributeDeleteMsg.addAllAttributeNames(keys); + attributeDeleteMsg.addAllAttributeNames(List.of(attributeKey)); attributeDeleteMsg.build(); builder.setAttributeDeleteMsg(attributeDeleteMsg); UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); From 672c05d9b72b5876a28e1c39ad6f5b0cda6209d7 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Thu, 6 Apr 2023 12:59:29 +0300 Subject: [PATCH 5/6] clear duplicated attribute keys before loading them from cache to avoid NPE --- .../org/thingsboard/server/dao/attributes/AttributeUtils.java | 1 + .../server/dao/attributes/CachedAttributesService.java | 4 ++++ 2 files changed, 5 insertions(+) diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributeUtils.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributeUtils.java index 18802130c7..168782c0fa 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributeUtils.java +++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributeUtils.java @@ -21,6 +21,7 @@ import org.thingsboard.server.dao.exception.IncorrectParameterException; import org.thingsboard.server.dao.service.Validator; public class AttributeUtils { + public static void validate(EntityId id, String scope) { Validator.validateId(id.getId(), "Incorrect id " + id); Validator.validateString(scope, "Incorrect scope " + scope); diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java index 017fe4315c..01a6d1201c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java @@ -42,6 +42,7 @@ import java.util.ArrayList; import java.util.Collection; import java.util.HashMap; import java.util.HashSet; +import java.util.LinkedHashSet; import java.util.List; import java.util.Map; import java.util.Objects; @@ -123,6 +124,7 @@ public class CachedAttributesService implements AttributesService { return result; } catch (Throwable e) { cacheTransaction.rollback(); + log.debug("Could not find attribute from cache: [{}] [{}] [{}]", entityId, scope, attributeKey, e); throw e; } }); @@ -132,6 +134,7 @@ public class CachedAttributesService implements AttributesService { @Override public ListenableFuture> find(TenantId tenantId, EntityId entityId, String scope, Collection attributeKeys) { validate(entityId, scope); + attributeKeys = new LinkedHashSet<>(attributeKeys); // deduplicate the attributes attributeKeys.forEach(attributeKey -> Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey)); Map> wrappedCachedAttributes = findCachedAttributes(entityId, scope, attributeKeys); @@ -170,6 +173,7 @@ public class CachedAttributesService implements AttributesService { return mergedAttributes; } catch (Throwable e) { cacheTransaction.rollback(); + log.debug("Could not find attributes from cache: [{}] [{}] [{}]", entityId, scope, notFoundAttributeKeys, e); throw e; } }); From 0bdae38d1db76858278a41aa918851c5ef5cb2f6 Mon Sep 17 00:00:00 2001 From: Vladyslav_Prykhodko Date: Thu, 6 Apr 2023 13:29:04 +0300 Subject: [PATCH 6/6] UI: Fixed updated notification when open notifications popover --- .../websocket/notification-ws.models.ts | 61 +++++++++---------- 1 file changed, 30 insertions(+), 31 deletions(-) diff --git a/ui-ngx/src/app/shared/models/websocket/notification-ws.models.ts b/ui-ngx/src/app/shared/models/websocket/notification-ws.models.ts index da876d0d3f..0acbba1139 100644 --- a/ui-ngx/src/app/shared/models/websocket/notification-ws.models.ts +++ b/ui-ngx/src/app/shared/models/websocket/notification-ws.models.ts @@ -16,7 +16,7 @@ import { BehaviorSubject, ReplaySubject } from 'rxjs'; import { CmdUpdate, CmdUpdateMsg, CmdUpdateType, WebsocketCmd } from '@shared/models/telemetry/telemetry.models'; -import { first, map } from 'rxjs/operators'; +import { map } from 'rxjs/operators'; import { NgZone } from '@angular/core'; import { isDefinedAndNotNull } from '@core/utils'; import { Notification } from '@shared/models/notification.models'; @@ -47,12 +47,19 @@ export class NotificationsUpdate extends CmdUpdate { export class NotificationSubscriber extends WsSubscriber { private notificationCountSubject = new ReplaySubject(1); - private notificationsSubject = new BehaviorSubject(null); + private notificationsSubject = new BehaviorSubject({ + cmdId: 0, + cmdUpdateType: undefined, + errorCode: 0, + errorMsg: '', + notifications: [], + totalUnreadCount: 0 + }); public messageLimit = 10; public notificationCount$ = this.notificationCountSubject.asObservable().pipe(map(msg => msg.totalUnreadCount)); - public notifications$ = this.notificationsSubject.asObservable().pipe(map(msg => msg?.notifications || [])); + public notifications$ = this.notificationsSubject.asObservable().pipe(map(msg => msg.notifications )); public static createNotificationCountSubscription(notificationWsService: NotificationWebsocketService, zone: NgZone): NotificationSubscriber { @@ -109,35 +116,27 @@ export class NotificationSubscriber extends WsSubscriber { } onNotificationsUpdate(message: NotificationsUpdate) { - this.notificationsSubject.asObservable().pipe( - first() - ).subscribe((value) => { - let saveMessage; - if (isDefinedAndNotNull(value) && message.update) { - const findIndex = value.notifications.findIndex(item => item.id.id === message.update.id.id); - if (findIndex !== -1) { - value.notifications.push(message.update); - value.notifications.sort((a, b) => b.createdTime - a.createdTime); - if (value.notifications.length > this.messageLimit) { - value.notifications.pop(); - } - } - saveMessage = value; - } else { - saveMessage = message; - } - if (this.zone) { - this.zone.run( - () => { - this.notificationsSubject.next(saveMessage); - this.notificationCountSubject.next(saveMessage); - } - ); - } else { - this.notificationsSubject.next(saveMessage); - this.notificationCountSubject.next(saveMessage); + const currentNotifications = this.notificationsSubject.value; + let processMessage = message; + if (isDefinedAndNotNull(currentNotifications) && message.update) { + currentNotifications.notifications.unshift(message.update); + if (currentNotifications.notifications.length > this.messageLimit) { + currentNotifications.notifications.pop(); } - }); + processMessage = currentNotifications; + processMessage.totalUnreadCount = message.totalUnreadCount; + } + if (this.zone) { + this.zone.run( + () => { + this.notificationsSubject.next(processMessage); + this.notificationCountSubject.next(processMessage); + } + ); + } else { + this.notificationsSubject.next(processMessage); + this.notificationCountSubject.next(processMessage); + } } }