From b31f0748c820aa3e426b50ca6ae1a20de1098ea0 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Wed, 29 Sep 2021 23:50:15 +0300 Subject: [PATCH 1/4] DefaultTbCoreToTransportService - skipping message without nodeId --- .../service/transport/DefaultTbCoreToTransportService.java | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java index 95f7e72ad7..1208518e3c 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java @@ -54,6 +54,10 @@ public class DefaultTbCoreToTransportService implements TbCoreToTransportService @Override public void process(String nodeId, ToTransportMsg msg, Runnable onSuccess, Consumer onFailure) { + if (nodeId == null || nodeId.isEmpty()){ + log.trace("process: skipping message without nodeId [{}], (ToTransportMsg) msg [{}]", nodeId, msg); + return; + } TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_TRANSPORT, nodeId); UUID sessionId = new UUID(msg.getSessionIdMSB(), msg.getSessionIdLSB()); log.trace("[{}][{}] Pushing session data to topic: {}", tpi.getFullTopicName(), sessionId, msg); From 4c0812898cb1a07c87bac73bf6ce3b7fe2f92aec Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Wed, 29 Sep 2021 18:04:29 +0300 Subject: [PATCH 2/4] Skip Notification To Transport without ServiceId --- .../server/service/queue/DefaultTbClusterService.java | 4 ++++ .../server/service/rpc/DefaultTbRuleEngineRpcService.java | 4 ++++ 2 files changed, 8 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java index 6677888854..7ad22a7a8b 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java @@ -196,6 +196,10 @@ public class DefaultTbClusterService implements TbClusterService { @Override public void pushNotificationToTransport(String serviceId, ToTransportMsg response, TbQueueCallback callback) { + if (serviceId == null || serviceId.isEmpty()){ + log.trace("pushNotificationToTransport: skipping message without serviceId [{}], (ToTransportMsg) response [{}]", serviceId, response); + return; + } TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_TRANSPORT, serviceId); log.trace("PUSHING msg: {} to:{}", response, tpi); producerProvider.getTransportNotificationsMsgProducer().send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), response), callback); diff --git a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java index eca8864fda..e0d2cca803 100644 --- a/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java +++ b/application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java @@ -87,6 +87,10 @@ public class DefaultTbRuleEngineRpcService implements TbRuleEngineDeviceRpcServi @Override public void sendRpcReplyToDevice(String serviceId, UUID sessionId, int requestId, String body) { + if (serviceId == null || serviceId.isEmpty()){ + log.trace("sendRpcReplyToDevice: skipping message without serviceId [{}], sessionId[{}], requestId[{}], body[{}]", serviceId, sessionId, requestId, body); + return; + } TransportProtos.ToServerRpcResponseMsg responseMsg = TransportProtos.ToServerRpcResponseMsg.newBuilder() .setRequestId(requestId) .setPayload(body).build(); From d25c605ac617be3ccfeae644fa3f3b22325f4b3d Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 12 Oct 2021 12:36:46 +0300 Subject: [PATCH 3/4] DefaultTbCoreToTransportService: fire "onSuccess send" when skipping message without nodeId --- .../service/transport/DefaultTbCoreToTransportService.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java index 1208518e3c..b4101edc8e 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java @@ -56,6 +56,9 @@ public class DefaultTbCoreToTransportService implements TbCoreToTransportService public void process(String nodeId, ToTransportMsg msg, Runnable onSuccess, Consumer onFailure) { if (nodeId == null || nodeId.isEmpty()){ log.trace("process: skipping message without nodeId [{}], (ToTransportMsg) msg [{}]", nodeId, msg); + if (onSuccess != null) { + onSuccess.run(); + } return; } TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_TRANSPORT, nodeId); From 4ea397b4d59fadaf54585b8b607e66e26a10fe27 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 12 Oct 2021 17:16:14 +0300 Subject: [PATCH 4/4] ClusterService - pushNotificationToTransport - fire callback.onSuccess(null) if serviceId is empty --- .../server/service/queue/DefaultTbClusterService.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java index 7ad22a7a8b..0a82994a42 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java @@ -198,6 +198,9 @@ public class DefaultTbClusterService implements TbClusterService { public void pushNotificationToTransport(String serviceId, ToTransportMsg response, TbQueueCallback callback) { if (serviceId == null || serviceId.isEmpty()){ log.trace("pushNotificationToTransport: skipping message without serviceId [{}], (ToTransportMsg) response [{}]", serviceId, response); + if (callback != null) { + callback.onSuccess(null); //callback that message already sent, no useful payload expected + } return; } TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_TRANSPORT, serviceId);