diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java index 17119232cc..48425d58db 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java @@ -16,9 +16,9 @@ package org.thingsboard.server.actors.device; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.server.common.msg.ruleengine.DeviceAttributesEventNotificationMsg; -import org.thingsboard.server.common.msg.ruleengine.DeviceEdgeUpdateMsg; -import org.thingsboard.server.common.msg.ruleengine.DeviceNameOrTypeUpdateMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceAttributesEventNotificationMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceEdgeUpdateMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.TbActorCtx; import org.thingsboard.server.actors.TbActorException; diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index d87548487c..cd7a7be4df 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -25,10 +25,10 @@ import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections.CollectionUtils; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.LinkedHashMapRemoveEldest; -import org.thingsboard.server.common.msg.ruleengine.DeviceAttributesEventNotificationMsg; -import org.thingsboard.server.common.msg.ruleengine.DeviceCredentialsUpdateNotificationMsg; -import org.thingsboard.server.common.msg.ruleengine.DeviceEdgeUpdateMsg; -import org.thingsboard.server.common.msg.ruleengine.DeviceNameOrTypeUpdateMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceAttributesEventNotificationMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceCredentialsUpdateNotificationMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceEdgeUpdateMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.TbActorCtx; import org.thingsboard.server.actors.shared.AbstractContextAwareMsgProcessor; diff --git a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java b/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java index 829a6dd040..5f88752ff4 100644 --- a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java +++ b/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java @@ -45,7 +45,7 @@ import org.springframework.web.bind.annotation.RestController; import org.springframework.web.context.request.async.DeferredResult; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardThreadFactory; -import org.thingsboard.server.common.msg.ruleengine.DeviceAttributesEventNotificationMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceAttributesEventNotificationMsg; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.StringUtils; 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 fafd2d1c4d..52bf7dbaa2 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 @@ -29,7 +29,7 @@ import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.tuple.ImmutablePair; import org.apache.commons.lang3.tuple.Pair; import org.thingsboard.common.util.JacksonUtil; -import org.thingsboard.server.common.msg.ruleengine.DeviceAttributesEventNotificationMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceAttributesEventNotificationMsg; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/DefaultTbNotificationEntityService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/DefaultTbNotificationEntityService.java index 1b1fa4a4b4..72e8f44020 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/DefaultTbNotificationEntityService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/DefaultTbNotificationEntityService.java @@ -19,7 +19,7 @@ import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import org.thingsboard.common.util.JacksonUtil; -import org.thingsboard.server.common.msg.ruleengine.DeviceCredentialsUpdateNotificationMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceCredentialsUpdateNotificationMsg; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.HasName; diff --git a/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java b/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java index a6bce41ce0..fef96c0fcb 100644 --- a/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java @@ -20,7 +20,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; import org.thingsboard.rule.engine.api.RuleEngineTelemetryService; -import org.thingsboard.server.common.msg.ruleengine.DeviceAttributesEventNotificationMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceAttributesEventNotificationMsg; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; 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 7955b52393..7241c0e571 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 @@ -23,8 +23,8 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Lazy; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; -import org.thingsboard.server.common.msg.ruleengine.DeviceEdgeUpdateMsg; -import org.thingsboard.server.common.msg.ruleengine.DeviceNameOrTypeUpdateMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceEdgeUpdateMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.Device; diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index 8d725b837c..56038214d0 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -279,6 +279,19 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService actorMsg = encodingService.decode(toCoreMsg.getToDeviceActorNotificationMsg().toByteArray()); + if (actorMsg.isPresent()) { + TbActorMsg tbActorMsg = actorMsg.get(); + if (tbActorMsg.getMsgType().equals(MsgType.DEVICE_RPC_REQUEST_TO_DEVICE_ACTOR_MSG)) { + tbCoreDeviceRpcService.forwardRpcRequestToDeviceActor((ToDeviceRpcRequestActorMsg) tbActorMsg); + } else { + log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg.get()); + actorContext.tell(actorMsg.get()); + } + } + callback.onSuccess(); } else if (toCoreMsg.hasNotificationSchedulerServiceMsg()) { TransportProtos.NotificationSchedulerServiceMsg notificationSchedulerServiceMsg = toCoreMsg.getNotificationSchedulerServiceMsg(); log.trace("[{}] Forwarding message to notification scheduler service {}", id, toCoreMsg.getNotificationSchedulerServiceMsg()); @@ -359,12 +372,21 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService actorMsg, TbCallback callback) { + actorMsg.ifPresent(tbActorMsg -> forwardToAppActor(id, tbActorMsg)); + callback.onSuccess(); + } + private void forwardToAppActor(UUID id, TbActorMsg actorMsg) { log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg); actorContext.tell(actorMsg); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/ProtoUtils.java b/application/src/main/java/org/thingsboard/server/service/queue/ProtoUtils.java index 9d6e215c39..539f26ef03 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/ProtoUtils.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/ProtoUtils.java @@ -15,10 +15,10 @@ */ package org.thingsboard.server.service.queue; -import org.thingsboard.server.common.msg.ruleengine.DeviceAttributesEventNotificationMsg; -import org.thingsboard.server.common.msg.ruleengine.DeviceCredentialsUpdateNotificationMsg; -import org.thingsboard.server.common.msg.ruleengine.DeviceEdgeUpdateMsg; -import org.thingsboard.server.common.msg.ruleengine.DeviceNameOrTypeUpdateMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceAttributesEventNotificationMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceCredentialsUpdateNotificationMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceEdgeUpdateMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EdgeId; diff --git a/application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerStats.java b/application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerStats.java index b7f1147ec0..5ff72e97e3 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerStats.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerStats.java @@ -174,10 +174,16 @@ public class TbCoreConsumerStats { toCoreNfComponentLifecycleCounter.increment(); } else if (msg.hasEdgeEventUpdate()) { toCoreNfEdgeEventUpdateCounter.increment(); + } else if (!msg.getEdgeEventUpdateMsg().isEmpty()) { + toCoreNfEdgeEventUpdateCounter.increment(); } else if (msg.hasToEdgeSyncRequest()) { toCoreNfEdgeSyncRequestCounter.increment(); + } else if (!msg.getToEdgeSyncRequestMsg().isEmpty()) { + toCoreNfEdgeSyncRequestCounter.increment(); } else if (msg.hasFromEdgeSyncResponse()) { toCoreNfEdgeSyncResponseCounter.increment(); + } else if (!msg.getFromEdgeSyncResponseMsg().isEmpty()) { + toCoreNfEdgeSyncResponseCounter.increment(); } else if (msg.hasQueueUpdateMsg()) { toCoreNfQueueUpdateCounter.increment(); } else if (msg.hasQueueDeleteMsg()) { diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java index b92aa8cd44..01461f87d3 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java @@ -21,7 +21,7 @@ import org.springframework.stereotype.Service; import org.thingsboard.common.util.DonAsynchron; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardThreadFactory; -import org.thingsboard.server.common.msg.ruleengine.DeviceAttributesEventNotificationMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceAttributesEventNotificationMsg; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntityType; diff --git a/application/src/test/java/org/thingsboard/server/service/queue/ProtoUtilsTest.java b/application/src/test/java/org/thingsboard/server/service/queue/ProtoUtilsTest.java index 0024c6b259..a1ac31a68f 100644 --- a/application/src/test/java/org/thingsboard/server/service/queue/ProtoUtilsTest.java +++ b/application/src/test/java/org/thingsboard/server/service/queue/ProtoUtilsTest.java @@ -17,10 +17,10 @@ package org.thingsboard.server.service.queue; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; -import org.thingsboard.server.common.msg.ruleengine.DeviceAttributesEventNotificationMsg; -import org.thingsboard.server.common.msg.ruleengine.DeviceCredentialsUpdateNotificationMsg; -import org.thingsboard.server.common.msg.ruleengine.DeviceEdgeUpdateMsg; -import org.thingsboard.server.common.msg.ruleengine.DeviceNameOrTypeUpdateMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceAttributesEventNotificationMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceCredentialsUpdateNotificationMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceEdgeUpdateMsg; +import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EdgeId; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceAttributes.java b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceAttributes.java similarity index 98% rename from common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceAttributes.java rename to common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceAttributes.java index f7f042485a..00079558e6 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceAttributes.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceAttributes.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.common.msg.ruleengine; +package org.thingsboard.server.common.msg.rule.engine; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.kv.AttributeKey; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceAttributesEventNotificationMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceAttributesEventNotificationMsg.java similarity index 97% rename from common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceAttributesEventNotificationMsg.java rename to common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceAttributesEventNotificationMsg.java index 93db8efd11..d0124fbb01 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceAttributesEventNotificationMsg.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceAttributesEventNotificationMsg.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.common.msg.ruleengine; +package org.thingsboard.server.common.msg.rule.engine; import lombok.Data; import org.thingsboard.server.common.data.id.DeviceId; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceCredentialsUpdateNotificationMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceCredentialsUpdateNotificationMsg.java similarity index 96% rename from common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceCredentialsUpdateNotificationMsg.java rename to common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceCredentialsUpdateNotificationMsg.java index e4edbc2311..b13026ba38 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceCredentialsUpdateNotificationMsg.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceCredentialsUpdateNotificationMsg.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.common.msg.ruleengine; +package org.thingsboard.server.common.msg.rule.engine; import lombok.Data; import org.thingsboard.server.common.data.id.DeviceId; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceEdgeUpdateMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceEdgeUpdateMsg.java similarity index 95% rename from common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceEdgeUpdateMsg.java rename to common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceEdgeUpdateMsg.java index 99411f7043..99e52e31b0 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceEdgeUpdateMsg.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceEdgeUpdateMsg.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.common.msg.ruleengine; +package org.thingsboard.server.common.msg.rule.engine; import lombok.Data; import org.thingsboard.server.common.data.id.DeviceId; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceMetaData.java b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceMetaData.java similarity index 94% rename from common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceMetaData.java rename to common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceMetaData.java index 21302c0f3c..f7fbb7b1ef 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceMetaData.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceMetaData.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.common.msg.ruleengine; +package org.thingsboard.server.common.msg.rule.engine; import lombok.Data; import org.thingsboard.server.common.data.id.DeviceId; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceNameOrTypeUpdateMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceNameOrTypeUpdateMsg.java similarity index 95% rename from common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceNameOrTypeUpdateMsg.java rename to common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceNameOrTypeUpdateMsg.java index db6ff91b22..13b6a8558f 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/ruleengine/DeviceNameOrTypeUpdateMsg.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceNameOrTypeUpdateMsg.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.common.msg.ruleengine; +package org.thingsboard.server.common.msg.rule.engine; import lombok.Data; import org.thingsboard.server.common.data.id.DeviceId;