Browse Source
Merge pull request #5357 from smatvienko-tb/transport_from_core_skipping_message_without_nodeId
Transport - skipping messages to transport without nodeId
pull/5378/head
Andrew Shvayka
5 years ago
committed by
GitHub
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
3 changed files with
18 additions and
0 deletions
-
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
-
application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java
-
application/src/main/java/org/thingsboard/server/service/transport/DefaultTbCoreToTransportService.java
|
|
|
@ -196,6 +196,13 @@ 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); |
|
|
|
if (callback != null) { |
|
|
|
callback.onSuccess(null); //callback that message already sent, no useful payload expected
|
|
|
|
} |
|
|
|
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); |
|
|
|
|
|
|
|
@ -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(); |
|
|
|
|
|
|
|
@ -54,6 +54,13 @@ public class DefaultTbCoreToTransportService implements TbCoreToTransportService |
|
|
|
|
|
|
|
@Override |
|
|
|
public void process(String nodeId, ToTransportMsg msg, Runnable onSuccess, Consumer<Throwable> 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); |
|
|
|
UUID sessionId = new UUID(msg.getSessionIdMSB(), msg.getSessionIdLSB()); |
|
|
|
log.trace("[{}][{}] Pushing session data to topic: {}", tpi.getFullTopicName(), sessionId, msg); |
|
|
|
|