diff --git a/application/src/main/data/json/edge/rule_chains/edge_root_rule_chain.json b/application/src/main/data/json/edge/rule_chains/edge_root_rule_chain.json index 6701b59e0e..ef0bf2698f 100644 --- a/application/src/main/data/json/edge/rule_chains/edge_root_rule_chain.json +++ b/application/src/main/data/json/edge/rule_chains/edge_root_rule_chain.json @@ -119,7 +119,7 @@ "type": "org.thingsboard.rule.engine.edge.TbMsgPushToCloudNode", "name": "Push to cloud", "configuration": { - "scope": "SERVER_SCOPE" + "scope": "CLIENT_SCOPE" }, "externalId": null }, 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 9f0e184853..a970097548 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 @@ -74,6 +74,8 @@ import org.thingsboard.server.gen.edge.v1.RequestMsgType; import org.thingsboard.server.gen.edge.v1.ResourceUpdateMsg; import org.thingsboard.server.gen.edge.v1.ResponseMsg; import org.thingsboard.server.gen.edge.v1.RuleChainMetadataRequestMsg; +import org.thingsboard.server.gen.edge.v1.RuleChainMetadataUpdateMsg; +import org.thingsboard.server.gen.edge.v1.RuleChainUpdateMsg; import org.thingsboard.server.gen.edge.v1.SyncCompletedMsg; import org.thingsboard.server.gen.edge.v1.UplinkMsg; import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg; @@ -820,6 +822,16 @@ public abstract class EdgeGrpcSession implements Closeable { result.add(ctx.getAssetProcessor().processAssetMsgFromEdge(edge.getTenantId(), edge, assetUpdateMsg)); } } + if (uplinkMsg.getRuleChainUpdateMsgCount() > 0) { + for (RuleChainUpdateMsg ruleChainUpdateMsg : uplinkMsg.getRuleChainUpdateMsgList()) { + result.add(ctx.getRuleChainProcessor().processRuleChainMsgFromEdge(edge.getTenantId(), edge, ruleChainUpdateMsg)); + } + } + if (uplinkMsg.getRuleChainMetadataUpdateMsgCount() > 0) { + for (RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg : uplinkMsg.getRuleChainMetadataUpdateMsgList()) { + result.add(ctx.getRuleChainProcessor().processRuleChainMetadataMsgFromEdge(edge.getTenantId(), edge, ruleChainMetadataUpdateMsg)); + } + } if (uplinkMsg.getEntityViewUpdateMsgCount() > 0) { for (EntityViewUpdateMsg entityViewUpdateMsg : uplinkMsg.getEntityViewUpdateMsgList()) { result.add(ctx.getEntityViewProcessor().processEntityViewMsgFromEdge(edge.getTenantId(), edge, entityViewUpdateMsg)); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/edge/EdgeEntityProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/edge/EdgeEntityProcessor.java index ddbe4810df..77fa31c028 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/edge/EdgeEntityProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/edge/EdgeEntityProcessor.java @@ -49,8 +49,12 @@ public class EdgeEntityProcessor extends BaseEdgeProcessor { @Override public ListenableFuture processEntityNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { try { + EdgeId originatorEdgeId = safeGetEdgeId(edgeNotificationMsg.getOriginatorEdgeIdMSB(), edgeNotificationMsg.getOriginatorEdgeIdLSB()); EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()); EdgeId edgeId = new EdgeId(new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB())); + if (edgeId.equals(originatorEdgeId)) { + return Futures.immediateFuture(null); + } switch (actionType) { case ASSIGNED_TO_CUSTOMER: { CustomerId customerId = JacksonUtil.fromString(edgeNotificationMsg.getBody(), CustomerId.class); diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/rule/BaseRuleChainProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/rule/BaseRuleChainProcessor.java new file mode 100644 index 0000000000..00b9c732aa --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/rule/BaseRuleChainProcessor.java @@ -0,0 +1,82 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.service.edge.rpc.processor.rule; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.data.util.Pair; +import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.id.RuleChainId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.rule.RuleChain; +import org.thingsboard.server.common.data.rule.RuleChainMetaData; +import org.thingsboard.server.common.data.rule.RuleChainType; +import org.thingsboard.server.common.data.rule.RuleNode; +import org.thingsboard.server.dao.service.DataValidator; +import org.thingsboard.server.gen.edge.v1.RuleChainMetadataUpdateMsg; +import org.thingsboard.server.gen.edge.v1.RuleChainUpdateMsg; +import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; + +import java.util.function.Function; + +@Slf4j +public class BaseRuleChainProcessor extends BaseEdgeProcessor { + + @Autowired + private DataValidator ruleChainValidator; + + protected Pair saveOrUpdateRuleChain(TenantId tenantId, RuleChainId ruleChainId, RuleChainUpdateMsg ruleChainUpdateMsg, RuleChainType ruleChainType) { + boolean created = false; + RuleChain ruleChainFromDb = edgeCtx.getRuleChainService().findRuleChainById(tenantId, ruleChainId); + if (ruleChainFromDb == null) { + created = true; + } + + RuleChain ruleChain = JacksonUtil.fromString(ruleChainUpdateMsg.getEntity(), RuleChain.class, true); + if (ruleChain == null) { + throw new RuntimeException("[{" + tenantId + "}] ruleChainUpdateMsg {" + ruleChainUpdateMsg + "} cannot be converted to rule chain"); + } + boolean isRoot = ruleChain.isRoot(); + if (RuleChainType.CORE.equals(ruleChainType)) { + ruleChain.setRoot(false); + } else { + ruleChain.setRoot(ruleChainFromDb == null ? false : ruleChainFromDb.isRoot()); + } + ruleChain.setType(ruleChainType); + + ruleChainValidator.validate(ruleChain, RuleChain::getTenantId); + if (created) { + ruleChain.setId(ruleChainId); + } + edgeCtx.getRuleChainService().saveRuleChain(ruleChain, true, false); + return Pair.of(created, isRoot); + } + + protected void saveOrUpdateRuleChainMetadata(TenantId tenantId, RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg) { + RuleChainMetaData ruleChainMetadata = JacksonUtil.fromString(ruleChainMetadataUpdateMsg.getEntity(), RuleChainMetaData.class, true); + if (ruleChainMetadata == null) { + throw new RuntimeException("[{" + tenantId + "}] ruleChainMetadataUpdateMsg {" + ruleChainMetadataUpdateMsg + "} cannot be converted to rule chain metadata"); + } + if (!ruleChainMetadata.getNodes().isEmpty()) { + ruleChainMetadata.setVersion(null); + for (RuleNode ruleNode : ruleChainMetadata.getNodes()) { + ruleNode.setRuleChainId(null); + ruleNode.setId(null); + } + edgeCtx.getRuleChainService().saveRuleChainMetaData(tenantId, ruleChainMetadata, Function.identity(), true); + } + } +} diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/rule/RuleChainEdgeProcessor.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/rule/RuleChainEdgeProcessor.java index 772300bded..06fb4c37a2 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/rule/RuleChainEdgeProcessor.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/rule/RuleChainEdgeProcessor.java @@ -15,29 +15,123 @@ */ package org.thingsboard.server.service.edge.rpc.processor.rule; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; +import org.springframework.data.util.Pair; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.id.RuleChainId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainMetaData; +import org.thingsboard.server.common.data.rule.RuleChainType; +import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.gen.edge.v1.DownlinkMsg; import org.thingsboard.server.gen.edge.v1.RuleChainMetadataUpdateMsg; import org.thingsboard.server.gen.edge.v1.RuleChainUpdateMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.edge.EdgeMsgConstructorUtils; -import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; + +import java.util.UUID; import static org.thingsboard.server.dao.edge.EdgeServiceImpl.EDGE_IS_ROOT_BODY_KEY; @Slf4j @Component @TbCoreComponent -public class RuleChainEdgeProcessor extends BaseEdgeProcessor { +public class RuleChainEdgeProcessor extends BaseRuleChainProcessor { + + public ListenableFuture processRuleChainMsgFromEdge(TenantId tenantId, Edge edge, RuleChainUpdateMsg ruleChainUpdateMsg) { + log.trace("[{}] executing processRuleChainMsgFromEdge [{}] from edge [{}]", tenantId, ruleChainUpdateMsg, edge.getName()); + RuleChainId ruleChainId = new RuleChainId(new UUID(ruleChainUpdateMsg.getIdMSB(), ruleChainUpdateMsg.getIdLSB())); + try { + edgeSynchronizationManager.getEdgeId().set(edge.getId()); + + switch (ruleChainUpdateMsg.getMsgType()) { + case ENTITY_CREATED_RPC_MESSAGE: + case ENTITY_UPDATED_RPC_MESSAGE: + return saveOrUpdateRuleChain(tenantId, ruleChainId, ruleChainUpdateMsg, edge); + case ENTITY_DELETED_RPC_MESSAGE: + RuleChain ruleChainToDelete = edgeCtx.getRuleChainService().findRuleChainById(tenantId, ruleChainId); + if (ruleChainToDelete != null) { + edgeCtx.getRuleChainService().unassignRuleChainFromEdge(tenantId, ruleChainId, edge.getId(), false); + } + return Futures.immediateFuture(null); + case UNRECOGNIZED: + default: + return handleUnsupportedMsgType(ruleChainUpdateMsg.getMsgType()); + } + } catch (DataValidationException e) { + if (e.getMessage().contains("limit reached")) { + log.warn("[{}] Number of allowed rule chains violated {}", tenantId, ruleChainUpdateMsg, e); + return Futures.immediateFuture(null); + } else { + return Futures.immediateFailedFuture(e); + } + } finally { + edgeSynchronizationManager.getEdgeId().remove(); + } + } + + private ListenableFuture saveOrUpdateRuleChain(TenantId tenantId, RuleChainId ruleChainId, RuleChainUpdateMsg ruleChainUpdateMsg, Edge edge) { + try { + Pair resultPair = super.saveOrUpdateRuleChain(tenantId, ruleChainId, ruleChainUpdateMsg, RuleChainType.EDGE); + Boolean created = resultPair.getFirst(); + if (created) { + createRelationFromEdge(tenantId, edge.getId(), ruleChainId); + pushRuleChainCreatedEventToRuleEngine(tenantId, edge, ruleChainId, ruleChainUpdateMsg.getEntity()); + edgeCtx.getRuleChainService().assignRuleChainToEdge(tenantId, ruleChainId, edge.getId()); + } + Boolean isRoot = resultPair.getSecond(); + if (isRoot) { + edge = edgeCtx.getEdgeService().findEdgeById(tenantId, edge.getId()); + edgeCtx.getEdgeService().setEdgeRootRuleChain(tenantId, edge, ruleChainId); + } + } catch (Exception e) { + log.error("Failed to save or update rule chain", e); + return Futures.immediateFailedFuture(e); + } + return Futures.immediateFuture(null); + } + + private void pushRuleChainCreatedEventToRuleEngine(TenantId tenantId, Edge edge, RuleChainId ruleChainId, String ruleChainAsString) { + try { + TbMsgMetaData msgMetaData = getEdgeActionTbMsgMetaData(edge, null); + pushEntityEventToRuleEngine(tenantId, ruleChainId, null, TbMsgType.ENTITY_CREATED, ruleChainAsString, msgMetaData); + } catch (Exception e) { + log.warn("[{}][{}] Failed to push rule chain action to rule engine: {}", tenantId, ruleChainId, TbMsgType.ENTITY_CREATED.name(), e); + } + } + + public ListenableFuture processRuleChainMetadataMsgFromEdge(TenantId tenantId, Edge edge, RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg) { + log.trace("[{}] executing processRuleChainMetadataMsgFromEdge [{}] from edge [{}]", tenantId, ruleChainMetadataUpdateMsg, edge.getName()); + try { + edgeSynchronizationManager.getEdgeId().set(edge.getId()); + + switch (ruleChainMetadataUpdateMsg.getMsgType()) { + case ENTITY_CREATED_RPC_MESSAGE: + case ENTITY_UPDATED_RPC_MESSAGE: + saveOrUpdateRuleChainMetadata(tenantId, ruleChainMetadataUpdateMsg); + return Futures.immediateFuture(null); + case UNRECOGNIZED: + default: + return handleUnsupportedMsgType(ruleChainMetadataUpdateMsg.getMsgType()); + } + } catch (Exception e) { + String errMsg = String.format("Can't process rule chain metadata update msg %s", ruleChainMetadataUpdateMsg); + log.error(errMsg, e); + return Futures.immediateFailedFuture(new RuntimeException(errMsg, e)); + } finally { + edgeSynchronizationManager.getEdgeId().remove(); + } + } @Override public DownlinkMsg convertEdgeEventToDownlink(EdgeEvent edgeEvent) { diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java index 9c9ceb2b34..03bd7b4e66 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java @@ -52,6 +52,7 @@ import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.edge.EdgeEventUpdateMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceCredentialsUpdateNotificationMsg; +import org.thingsboard.server.dao.edge.EdgeSynchronizationManager; import org.thingsboard.server.dao.eventsourcing.ActionEntityEvent; import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; @@ -67,6 +68,7 @@ public class EntityStateSourcingListener { private final TenantService tenantService; private final TbClusterService tbClusterService; + private final EdgeSynchronizationManager edgeSynchronizationManager; @PostConstruct public void init() { @@ -270,6 +272,9 @@ public class EntityStateSourcingListener { private void onEdgeEvent(TenantId tenantId, EntityId entityId, Object entity, ComponentLifecycleEvent lifecycleEvent) { if (entity instanceof Edge) { + if (entityId.equals(edgeSynchronizationManager.getEdgeId().get())) { + return; + } tbClusterService.onEdgeStateChangeEvent(new ComponentLifecycleMsg(tenantId, entityId, lifecycleEvent)); } else if (entity instanceof EdgeEvent edgeEvent) { tbClusterService.onEdgeEventUpdate(new EdgeEventUpdateMsg(tenantId, edgeEvent.getEdgeId())); diff --git a/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java index 9212afa7c5..feac7adb04 100644 --- a/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java @@ -90,14 +90,12 @@ import org.thingsboard.server.gen.edge.v1.EdgeConfiguration; import org.thingsboard.server.gen.edge.v1.OAuth2ClientUpdateMsg; import org.thingsboard.server.gen.edge.v1.OAuth2DomainUpdateMsg; import org.thingsboard.server.gen.edge.v1.QueueUpdateMsg; -import org.thingsboard.server.gen.edge.v1.RuleChainMetadataRequestMsg; import org.thingsboard.server.gen.edge.v1.RuleChainMetadataUpdateMsg; import org.thingsboard.server.gen.edge.v1.RuleChainUpdateMsg; import org.thingsboard.server.gen.edge.v1.SyncCompletedMsg; import org.thingsboard.server.gen.edge.v1.TenantProfileUpdateMsg; import org.thingsboard.server.gen.edge.v1.TenantUpdateMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; -import org.thingsboard.server.gen.edge.v1.UplinkMsg; import org.thingsboard.server.gen.edge.v1.UserCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.v1.UserUpdateMsg; @@ -142,35 +140,14 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest { installation(); edgeImitator = new EdgeImitator("localhost", 7070, edge.getRoutingKey(), edge.getSecret()); + edgeImitator.expectMessageAmount(25); edgeImitator.ignoreType(OAuth2ClientUpdateMsg.class); edgeImitator.ignoreType(OAuth2DomainUpdateMsg.class); - edgeImitator.expectMessageAmount(26); edgeImitator.connect(); - requestEdgeRuleChainMetadata(); - verifyEdgeConnectionAndInitialData(); } - private void requestEdgeRuleChainMetadata() throws Exception { - RuleChainId rootRuleChainId = getEdgeRootRuleChainId(); - RuleChainMetadataRequestMsg.Builder builder = RuleChainMetadataRequestMsg.newBuilder() - .setRuleChainIdMSB(rootRuleChainId.getId().getMostSignificantBits()) - .setRuleChainIdLSB(rootRuleChainId.getId().getLeastSignificantBits()); - testAutoGeneratedCodeByProtobuf(builder); - UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder() - .addRuleChainMetadataRequestMsg(builder.build()); - edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); - } - - private RuleChainId getEdgeRootRuleChainId() throws Exception { - return doGetTypedWithPageLink("/api/ruleChains?type={type}&", new TypeReference>() { - }, - new PageLink(100, 0, "Edge Root Rule Chain"), - "EDGE") - .getData().get(0).getId(); - } - @After public void teardownEdgeTest() { try { @@ -213,6 +190,19 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest { doPost("/api/ruleChain/metadata", rootRuleChainMetadata, RuleChainMetaData.class); } + private RuleChainId getEdgeRootRuleChainId() throws Exception { + List edgeRuleChains = doGetTypedWithPageLink("/api/ruleChains?type={type}&", + new TypeReference>() {}, + new PageLink(100, 0, "Edge Root Rule Chain"), + "EDGE").getData(); + for (RuleChain edgeRuleChain : edgeRuleChains) { + if (edgeRuleChain.isRoot()) { + return edgeRuleChain.getId(); + } + } + throw new RuntimeException("Root rule chain not found"); + } + protected void extendDeviceProfileData(DeviceProfile deviceProfile) { DeviceProfileData profileData = deviceProfile.getProfileData(); List alarms = new ArrayList<>(); @@ -255,8 +245,8 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest { validateMsgsCnt(RuleChainUpdateMsg.class, 1); UUID ruleChainUUID = validateRuleChains(); - // 1 from request message - validateMsgsCnt(RuleChainMetadataUpdateMsg.class, 2); + // 1 from rule chain fetcher + validateMsgsCnt(RuleChainMetadataUpdateMsg.class, 1); validateRuleChainMetadataUpdates(ruleChainUUID); // 4 messages ('general', 'mail', 'connectivity', 'jwt') @@ -438,12 +428,11 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest { } private void validateRuleChainMetadataUpdates(UUID expectedRuleChainUUID) { - Optional ruleChainMetadataUpdateOpt = edgeImitator.findMessageByType(RuleChainMetadataUpdateMsg.class); - Assert.assertTrue(ruleChainMetadataUpdateOpt.isPresent()); - RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = ruleChainMetadataUpdateOpt.get(); + Optional ruleChainMetadataUpdateMsgOpt = edgeImitator.findMessageByType(RuleChainMetadataUpdateMsg.class); + Assert.assertTrue(ruleChainMetadataUpdateMsgOpt.isPresent()); + RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = ruleChainMetadataUpdateMsgOpt.get(); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, ruleChainMetadataUpdateMsg.getMsgType()); RuleChainMetaData ruleChainMetaData = JacksonUtil.fromString(ruleChainMetadataUpdateMsg.getEntity(), RuleChainMetaData.class, true); - Assert.assertNotNull(ruleChainMetaData); Assert.assertEquals(expectedRuleChainUUID, ruleChainMetaData.getRuleChainId().getId()); } diff --git a/application/src/test/java/org/thingsboard/server/edge/AssetEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/AssetEdgeTest.java index 79ec2f3d87..4d9840b936 100644 --- a/application/src/test/java/org/thingsboard/server/edge/AssetEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/AssetEdgeTest.java @@ -169,6 +169,7 @@ public class AssetEdgeTest extends AbstractEdgeTest { public void testSendAssetToCloud() throws Exception { Asset asset = buildAssetForUplinkMsg("Asset Edge 2"); + // created asset on edge UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); AssetUpdateMsg.Builder assetUpdateMsgBuilder = AssetUpdateMsg.newBuilder(); assetUpdateMsgBuilder.setIdMSB(asset.getUuidId().getMostSignificantBits()); @@ -191,6 +192,32 @@ public class AssetEdgeTest extends AbstractEdgeTest { Asset foundAsset = doGet("/api/asset/" + asset.getUuidId(), Asset.class); Assert.assertNotNull(foundAsset); Assert.assertEquals("Asset Edge 2", foundAsset.getName()); + + // update asset on edge + asset.setName("Asset Edge 2 Updated"); + + uplinkMsgBuilder = UplinkMsg.newBuilder(); + assetUpdateMsgBuilder = AssetUpdateMsg.newBuilder(); + assetUpdateMsgBuilder.setIdMSB(asset.getUuidId().getMostSignificantBits()); + assetUpdateMsgBuilder.setIdLSB(asset.getUuidId().getLeastSignificantBits()); + assetUpdateMsgBuilder.setEntity(JacksonUtil.toString(asset)); + assetUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE); + testAutoGeneratedCodeByProtobuf(assetUpdateMsgBuilder); + uplinkMsgBuilder.addAssetUpdateMsg(assetUpdateMsgBuilder.build()); + + testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); + + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + + Assert.assertTrue(edgeImitator.waitForResponses()); + + latestResponseMsg = edgeImitator.getLatestResponseMsg(); + Assert.assertTrue(latestResponseMsg.getSuccess()); + + foundAsset = doGet("/api/asset/" + asset.getUuidId(), Asset.class); + Assert.assertNotNull(foundAsset); + Assert.assertEquals("Asset Edge 2 Updated", foundAsset.getName()); } @Test diff --git a/application/src/test/java/org/thingsboard/server/edge/DashboardEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/DashboardEdgeTest.java index b848178c58..bbf3d17f0d 100644 --- a/application/src/test/java/org/thingsboard/server/edge/DashboardEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/DashboardEdgeTest.java @@ -184,6 +184,7 @@ public class DashboardEdgeTest extends AbstractEdgeTest { Dashboard dashboard = buildDashboardForUplinkMsg(savedCustomer); + // create dashboard on edge UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); DashboardUpdateMsg.Builder dashboardUpdateMsgBuilder = DashboardUpdateMsg.newBuilder(); dashboardUpdateMsgBuilder.setIdMSB(dashboard.getUuidId().getMostSignificantBits()); diff --git a/application/src/test/java/org/thingsboard/server/edge/DeviceEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/DeviceEdgeTest.java index a6cd127097..da8e7563e0 100644 --- a/application/src/test/java/org/thingsboard/server/edge/DeviceEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/DeviceEdgeTest.java @@ -593,8 +593,10 @@ public class DeviceEdgeTest extends AbstractEdgeTest { @Test public void testSendDeviceToCloud() throws Exception { - Device deviceMsg = buildDeviceForUplinkMsg("Edge Device 2", "test"); + String deviceName = "Edge Device 2"; + Device deviceMsg = buildDeviceForUplinkMsg(deviceName, "test"); + // create device on edge UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); DeviceUpdateMsg.Builder deviceUpdateMsgBuilder = DeviceUpdateMsg.newBuilder(); deviceUpdateMsgBuilder.setIdMSB(deviceMsg.getUuidId().getMostSignificantBits()); @@ -609,7 +611,25 @@ public class DeviceEdgeTest extends AbstractEdgeTest { Device device = doGet("/api/device/" + deviceMsg.getId().getId(), Device.class); Assert.assertNotNull(device); - Assert.assertEquals("Edge Device 2", device.getName()); + Assert.assertEquals(deviceName, device.getName()); + + // update device on edge + deviceMsg.setName(deviceName + " Updated"); + uplinkMsgBuilder = UplinkMsg.newBuilder(); + deviceUpdateMsgBuilder = DeviceUpdateMsg.newBuilder(); + deviceUpdateMsgBuilder.setIdMSB(deviceMsg.getUuidId().getMostSignificantBits()); + deviceUpdateMsgBuilder.setIdLSB(deviceMsg.getUuidId().getLeastSignificantBits()); + deviceUpdateMsgBuilder.setEntity(JacksonUtil.toString(deviceMsg)); + deviceUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE); + uplinkMsgBuilder.addDeviceUpdateMsg(deviceUpdateMsgBuilder.build()); + + edgeImitator.expectResponsesAmount(1); + edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + Assert.assertTrue(edgeImitator.waitForResponses()); + + device = doGet("/api/device/" + deviceMsg.getId().getId(), Device.class); + Assert.assertNotNull(device); + Assert.assertEquals(deviceName + " Updated", device.getName()); } @Test diff --git a/application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java index a5fea1694a..7c029e0471 100644 --- a/application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/RuleChainEdgeTest.java @@ -15,7 +15,7 @@ */ package org.thingsboard.server.edge; -import com.google.protobuf.AbstractMessage; +import com.datastax.oss.driver.api.core.uuid.Uuids; import org.junit.Assert; import org.junit.Test; import org.thingsboard.common.util.JacksonUtil; @@ -29,16 +29,17 @@ import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.common.data.rule.RuleChainType; import org.thingsboard.server.common.data.rule.RuleNode; import org.thingsboard.server.dao.service.DaoSqlTest; -import org.thingsboard.server.gen.edge.v1.RuleChainMetadataRequestMsg; import org.thingsboard.server.gen.edge.v1.RuleChainMetadataUpdateMsg; import org.thingsboard.server.gen.edge.v1.RuleChainUpdateMsg; import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.gen.edge.v1.UplinkMsg; +import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg; import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.Optional; +import java.util.UUID; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @@ -74,8 +75,9 @@ public class RuleChainEdgeTest extends AbstractEdgeTest { RuleChainMetaData ruleChainMetaData = JacksonUtil.fromString(ruleChainMetadataUpdateMsg.getEntity(), RuleChainMetaData.class, true); Assert.assertNotNull(ruleChainMetaData); Assert.assertEquals(ruleChainMetaData.getRuleChainId(), savedRuleChain.getId()); - - testRuleChainMetadataRequestMsg(savedRuleChain.getId()); + for (RuleNode ruleNode : ruleChainMetaData.getNodes()) { + Assert.assertEquals(CONFIGURATION_VERSION, ruleNode.getConfigurationVersion()); + } // unassign rule chain from edge edgeImitator.expectMessageAmount(1); @@ -97,60 +99,62 @@ public class RuleChainEdgeTest extends AbstractEdgeTest { } @Test - public void testSendRuleChainMetadataRequestToCloud() throws Exception { - RuleChainId edgeRootRuleChainId = edge.getRootRuleChainId(); - + public void testRuleChainToCloud() throws Exception { + String ruleChainName = "Rule Chain Edge"; + UUID uuid = Uuids.timeBased(); + + // create rule chain on edge + RuleChain edgeRuleChain = new RuleChain(); + edgeRuleChain.setTenantId(tenantId); + edgeRuleChain.setId(new RuleChainId(uuid)); + edgeRuleChain.setName(ruleChainName); UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); - RuleChainMetadataRequestMsg.Builder ruleChainMetadataRequestMsgBuilder = RuleChainMetadataRequestMsg.newBuilder(); - ruleChainMetadataRequestMsgBuilder.setRuleChainIdMSB(edgeRootRuleChainId.getId().getMostSignificantBits()); - ruleChainMetadataRequestMsgBuilder.setRuleChainIdLSB(edgeRootRuleChainId.getId().getLeastSignificantBits()); - testAutoGeneratedCodeByProtobuf(ruleChainMetadataRequestMsgBuilder); - uplinkMsgBuilder.addRuleChainMetadataRequestMsg(ruleChainMetadataRequestMsgBuilder.build()); + RuleChainUpdateMsg.Builder ruleChainUpdateMsgBuilder = RuleChainUpdateMsg.newBuilder(); + ruleChainUpdateMsgBuilder.setIdMSB(uuid.getMostSignificantBits()); + ruleChainUpdateMsgBuilder.setIdLSB(uuid.getLeastSignificantBits()); + ruleChainUpdateMsgBuilder.setEntity(JacksonUtil.toString(edgeRuleChain)); + ruleChainUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE); + testAutoGeneratedCodeByProtobuf(ruleChainUpdateMsgBuilder); + uplinkMsgBuilder.addRuleChainUpdateMsg(ruleChainUpdateMsgBuilder.build()); testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); edgeImitator.expectResponsesAmount(1); - edgeImitator.expectMessageAmount(1); edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + Assert.assertTrue(edgeImitator.waitForResponses()); - Assert.assertTrue(edgeImitator.waitForMessages()); - AbstractMessage latestMessage = edgeImitator.getLatestMessage(); - Assert.assertTrue(latestMessage instanceof RuleChainMetadataUpdateMsg); - RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = (RuleChainMetadataUpdateMsg) latestMessage; - RuleChainMetaData ruleChainMetadataMsg = JacksonUtil.fromString(ruleChainMetadataUpdateMsg.getEntity(), RuleChainMetaData.class, true); - Assert.assertNotNull(ruleChainMetadataMsg); - Assert.assertEquals(edgeRootRuleChainId, ruleChainMetadataMsg.getRuleChainId()); + UplinkResponseMsg latestResponseMsg = edgeImitator.getLatestResponseMsg(); + Assert.assertTrue(latestResponseMsg.getSuccess()); - testAutoGeneratedCodeByProtobuf(ruleChainMetadataUpdateMsg); - } + RuleChain ruleChain = doGet("/api/ruleChain/" + uuid, RuleChain.class); + Assert.assertNotNull(ruleChain); + Assert.assertEquals("Rule Chain Edge", ruleChain.getName()); - private void testRuleChainMetadataRequestMsg(RuleChainId ruleChainId) throws Exception { - RuleChainMetadataRequestMsg.Builder ruleChainMetadataRequestMsgBuilder = RuleChainMetadataRequestMsg.newBuilder() - .setRuleChainIdMSB(ruleChainId.getId().getMostSignificantBits()) - .setRuleChainIdLSB(ruleChainId.getId().getLeastSignificantBits()); - testAutoGeneratedCodeByProtobuf(ruleChainMetadataRequestMsgBuilder); + // update rule chain on edge + edgeRuleChain.setName(ruleChainName + " Updated"); + uplinkMsgBuilder = UplinkMsg.newBuilder(); + ruleChainUpdateMsgBuilder = RuleChainUpdateMsg.newBuilder(); + ruleChainUpdateMsgBuilder.setIdMSB(uuid.getMostSignificantBits()); + ruleChainUpdateMsgBuilder.setIdLSB(uuid.getLeastSignificantBits()); + ruleChainUpdateMsgBuilder.setEntity(JacksonUtil.toString(edgeRuleChain)); + ruleChainUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE); + testAutoGeneratedCodeByProtobuf(ruleChainUpdateMsgBuilder); + uplinkMsgBuilder.addRuleChainUpdateMsg(ruleChainUpdateMsgBuilder.build()); - UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder() - .addRuleChainMetadataRequestMsg(ruleChainMetadataRequestMsgBuilder.build()); testAutoGeneratedCodeByProtobuf(uplinkMsgBuilder); edgeImitator.expectResponsesAmount(1); - edgeImitator.expectMessageAmount(1); edgeImitator.sendUplinkMsg(uplinkMsgBuilder.build()); + Assert.assertTrue(edgeImitator.waitForResponses()); - Assert.assertTrue(edgeImitator.waitForMessages()); - AbstractMessage latestMessage = edgeImitator.getLatestMessage(); - Assert.assertTrue(latestMessage instanceof RuleChainMetadataUpdateMsg); - RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = (RuleChainMetadataUpdateMsg) latestMessage; - RuleChainMetaData ruleChainMetadataMsg = JacksonUtil.fromString(ruleChainMetadataUpdateMsg.getEntity(), RuleChainMetaData.class, true); - Assert.assertNotNull(ruleChainMetadataMsg); - Assert.assertEquals(ruleChainId, ruleChainMetadataMsg.getRuleChainId()); + latestResponseMsg = edgeImitator.getLatestResponseMsg(); + Assert.assertTrue(latestResponseMsg.getSuccess()); - for (RuleNode ruleNode : ruleChainMetadataMsg.getNodes()) { - Assert.assertEquals(CONFIGURATION_VERSION, ruleNode.getConfigurationVersion()); - } + ruleChain = doGet("/api/ruleChain/" + uuid, RuleChain.class); + Assert.assertNotNull(ruleChain); + Assert.assertEquals(ruleChainName + " Updated", ruleChain.getName()); } private RuleChainMetaData createRuleChainMetadata(RuleChain ruleChain) { diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java index b96a984f71..a2356ee149 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/rule/RuleChainService.java @@ -46,6 +46,8 @@ public interface RuleChainService extends EntityDaoService { RuleChain saveRuleChain(RuleChain ruleChain, boolean publishSaveEvent); + RuleChain saveRuleChain(RuleChain ruleChain, boolean publishSaveEvent, boolean doValidate); + boolean setRootRuleChain(TenantId tenantId, RuleChainId ruleChainId); RuleChainUpdateResult saveRuleChainMetaData(TenantId tenantId, RuleChainMetaData ruleChainMetaData, Function ruleNodeUpdater); diff --git a/common/edge-api/src/main/proto/edge.proto b/common/edge-api/src/main/proto/edge.proto index 6a7482a3c1..023ac00634 100644 --- a/common/edge-api/src/main/proto/edge.proto +++ b/common/edge-api/src/main/proto/edge.proto @@ -42,6 +42,8 @@ enum EdgeVersion { V_3_8_0 = 8; V_3_9_0 = 9; V_4_0_0 = 10; + + V_LATEST = 999; } /** @@ -303,7 +305,9 @@ message NotificationTemplateUpdateMsg { optional string entity = 4; } +// DEPRECATED. FOR REMOVAL message RuleChainMetadataRequestMsg { + option deprecated = true; int64 ruleChainIdMSB = 1; int64 ruleChainIdLSB = 2; } @@ -321,22 +325,30 @@ message RelationRequestMsg { string entityType = 3; } +// DEPRECATED. FOR REMOVAL message UserCredentialsRequestMsg { + option deprecated = true; int64 userIdMSB = 1; int64 userIdLSB = 2; } +// DEPRECATED. FOR REMOVAL message DeviceCredentialsRequestMsg { + option deprecated = true; int64 deviceIdMSB = 1; int64 deviceIdLSB = 2; } +// DEPRECATED. FOR REMOVAL message WidgetBundleTypesRequestMsg { + option deprecated = true; int64 widgetBundleIdMSB = 1; int64 widgetBundleIdLSB = 2; } +// DEPRECATED. FOR REMOVAL message EntityViewsRequestMsg { + option deprecated = true; int64 entityIdMSB = 1; int64 entityIdLSB = 2; string entityType = 3; @@ -394,14 +406,14 @@ message UplinkMsg { repeated DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg = 4; repeated AlarmUpdateMsg alarmUpdateMsg = 5; repeated RelationUpdateMsg relationUpdateMsg = 6; - repeated RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg = 7; + repeated RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg = 7 [deprecated = true]; repeated AttributesRequestMsg attributesRequestMsg = 8; repeated RelationRequestMsg relationRequestMsg = 9; - repeated UserCredentialsRequestMsg userCredentialsRequestMsg = 10; - repeated DeviceCredentialsRequestMsg deviceCredentialsRequestMsg = 11; + repeated UserCredentialsRequestMsg userCredentialsRequestMsg = 10 [deprecated = true]; + repeated DeviceCredentialsRequestMsg deviceCredentialsRequestMsg = 11 [deprecated = true]; repeated DeviceRpcCallMsg deviceRpcCallMsg = 12; - repeated WidgetBundleTypesRequestMsg widgetBundleTypesRequestMsg = 14; - repeated EntityViewsRequestMsg entityViewsRequestMsg = 15; + repeated WidgetBundleTypesRequestMsg widgetBundleTypesRequestMsg = 14 [deprecated = true]; + repeated EntityViewsRequestMsg entityViewsRequestMsg = 15 [deprecated = true]; repeated AssetUpdateMsg assetUpdateMsg = 16; repeated DashboardUpdateMsg dashboardUpdateMsg = 17; repeated EntityViewUpdateMsg entityViewUpdateMsg = 18; @@ -409,6 +421,8 @@ message UplinkMsg { repeated DeviceProfileUpdateMsg deviceProfileUpdateMsg = 20; repeated ResourceUpdateMsg resourceUpdateMsg = 21; repeated AlarmCommentUpdateMsg alarmCommentUpdateMsg = 22; + repeated RuleChainUpdateMsg ruleChainUpdateMsg = 23; + repeated RuleChainMetadataUpdateMsg ruleChainMetadataUpdateMsg = 24; } message UplinkResponseMsg { @@ -427,7 +441,7 @@ message DownlinkMsg { int32 downlinkMsgId = 1; SyncCompletedMsg syncCompletedMsg = 2; repeated EntityDataProto entityData = 3; - repeated DeviceCredentialsRequestMsg deviceCredentialsRequestMsg = 4; + repeated DeviceCredentialsRequestMsg deviceCredentialsRequestMsg = 4 [deprecated = true]; repeated DeviceUpdateMsg deviceUpdateMsg = 5; repeated DeviceProfileUpdateMsg deviceProfileUpdateMsg = 6; repeated DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg = 7; diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java index 3fb80332b5..b5fde88a0b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java @@ -118,7 +118,16 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC @Override @Transactional public RuleChain saveRuleChain(RuleChain ruleChain, boolean publishSaveEvent) { - ruleChainValidator.validate(ruleChain, RuleChain::getTenantId); + return saveRuleChain(ruleChain, publishSaveEvent, true); + } + + @Override + @Transactional + public RuleChain saveRuleChain(RuleChain ruleChain, boolean publishSaveEvent, boolean doValidate) { + log.trace("Executing doSaveRuleChain [{}]", ruleChain); + if (doValidate) { + ruleChainValidator.validate(ruleChain, RuleChain::getTenantId); + } try { RuleChain savedRuleChain = ruleChainDao.saveAndFlush(ruleChain.getTenantId(), ruleChain); if (ruleChain.getId() == null) {