|
|
|
@ -544,6 +544,9 @@ public class DefaultTransportService implements TransportService { |
|
|
|
|
|
|
|
protected void sendToDeviceActor(TransportProtos.SessionInfoProto sessionInfo, TransportToDeviceActorMsg toDeviceActorMsg, TransportServiceCallback<Void> callback) { |
|
|
|
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, getTenantId(sessionInfo), getDeviceId(sessionInfo)); |
|
|
|
if (log.isTraceEnabled()) { |
|
|
|
log.trace("[{}][{}] Pushing to topic {} message {}", getTenantId(sessionInfo), getDeviceId(sessionInfo), tpi.getFullTopicName(), toDeviceActorMsg); |
|
|
|
} |
|
|
|
tbCoreMsgProducer.send(tpi, |
|
|
|
new TbProtoQueueMsg<>(getRoutingKey(sessionInfo), |
|
|
|
ToCoreMsg.newBuilder().setToDeviceActorMsg(toDeviceActorMsg).build()), callback != null ? |
|
|
|
@ -552,6 +555,9 @@ public class DefaultTransportService implements TransportService { |
|
|
|
|
|
|
|
protected void sendToRuleEngine(TenantId tenantId, TbMsg tbMsg, TbQueueCallback callback) { |
|
|
|
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, tbMsg.getOriginator()); |
|
|
|
if (log.isTraceEnabled()) { |
|
|
|
log.trace("[{}][{}] Pushing to topic {} message {}", tenantId, tbMsg.getOriginator(), tpi.getFullTopicName(), tbMsg); |
|
|
|
} |
|
|
|
ToRuleEngineMsg msg = ToRuleEngineMsg.newBuilder().setTbMsg(TbMsg.toByteString(tbMsg)) |
|
|
|
.setTenantIdMSB(tenantId.getId().getMostSignificantBits()) |
|
|
|
.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()).build(); |
|
|
|
|