Browse Source

Merge pull request #14537 from volodymyr-babak/support-ai-model-sync

Added AiModelEdgeFetcher and support deletion of AI Models in both directions
pull/14557/head
Viacheslav Klimov 10 months ago
committed by GitHub
parent
commit
7ca635e1b4
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 12
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java
  2. 48
      application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AiModelEdgeEventFetcher.java
  3. 53
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/ai/AiModelEdgeProcessor.java
  4. 10
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/ai/BaseAiModelProcessor.java
  5. 42
      application/src/test/java/org/thingsboard/server/edge/AiModelEdgeTest.java

12
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java

@ -21,6 +21,7 @@ import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.service.edge.EdgeContextComponent;
import org.thingsboard.server.service.edge.rpc.fetch.AdminSettingsEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.AiModelEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.AssetProfilesEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.AssetsEdgeEventFetcher;
import org.thingsboard.server.service.edge.rpc.fetch.CustomerEdgeEventFetcher;
@ -64,6 +65,12 @@ public class EdgeSyncCursor {
fetchers.add(new RuleChainsEdgeEventFetcher(ctx.getRuleChainService()));
fetchers.add(new AdminSettingsEdgeEventFetcher(ctx.getAdminSettingsService()));
fetchers.add(new TenantAdminUsersEdgeEventFetcher(ctx.getUserService()));
fetchers.add(new OAuth2EdgeEventFetcher(ctx.getDomainService()));
fetchers.add(new SystemWidgetTypesEdgeEventFetcher(ctx.getWidgetTypeService()));
fetchers.add(new TenantWidgetTypesEdgeEventFetcher(ctx.getWidgetTypeService()));
fetchers.add(new SystemWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService()));
fetchers.add(new TenantWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService()));
fetchers.add(new AiModelEdgeEventFetcher(ctx.getAiModelService()));
}
Customer publicCustomer = ctx.getCustomerService().findPublicCustomer(edge.getTenantId());
if (publicCustomer != null) {
@ -84,14 +91,9 @@ public class EdgeSyncCursor {
fetchers.add(new NotificationTemplateEdgeEventFetcher(ctx.getNotificationTemplateService()));
fetchers.add(new NotificationTargetEdgeEventFetcher(ctx.getNotificationTargetService()));
fetchers.add(new NotificationRuleEdgeEventFetcher(ctx.getNotificationRuleService()));
fetchers.add(new SystemWidgetTypesEdgeEventFetcher(ctx.getWidgetTypeService()));
fetchers.add(new TenantWidgetTypesEdgeEventFetcher(ctx.getWidgetTypeService()));
fetchers.add(new SystemWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService()));
fetchers.add(new TenantWidgetsBundlesEdgeEventFetcher(ctx.getWidgetsBundleService()));
fetchers.add(new OtaPackagesEdgeEventFetcher(ctx.getOtaPackageService()));
fetchers.add(new DeviceProfilesEdgeEventFetcher(ctx.getDeviceProfileService()));
fetchers.add(new TenantResourcesEdgeEventFetcher(ctx.getResourceService()));
fetchers.add(new OAuth2EdgeEventFetcher(ctx.getDomainService()));
}
}

48
application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/AiModelEdgeEventFetcher.java

@ -0,0 +1,48 @@
/**
* 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.fetch;
import lombok.AllArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.EdgeUtils;
import org.thingsboard.server.common.data.ai.AiModel;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.ai.AiModelService;
@AllArgsConstructor
@Slf4j
public class AiModelEdgeEventFetcher extends BasePageableEdgeEventFetcher<AiModel> {
private final AiModelService aiModelService;
@Override
PageData<AiModel> fetchEntities(TenantId tenantId, Edge edge, PageLink pageLink) {
return aiModelService.findAiModelsByTenantId(tenantId, pageLink);
}
@Override
EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, AiModel aiModel) {
return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.AI_MODEL,
EdgeEventActionType.ADDED, aiModel.getId(), null);
}
}

53
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/ai/AiModelEdgeProcessor.java

@ -20,7 +20,6 @@ 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.common.util.JacksonUtil;
import org.thingsboard.server.common.data.EdgeUtils;
import org.thingsboard.server.common.data.ai.AiModel;
import org.thingsboard.server.common.data.edge.Edge;
@ -30,7 +29,6 @@ import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.id.AiModelId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.gen.edge.v1.AiModelUpdateMsg;
import org.thingsboard.server.gen.edge.v1.DownlinkMsg;
@ -53,15 +51,17 @@ public class AiModelEdgeProcessor extends BaseAiModelProcessor implements AiMode
try {
edgeSynchronizationManager.getEdgeId().set(edge.getId());
switch (aiModelUpdateMsg.getMsgType()) {
case ENTITY_CREATED_RPC_MESSAGE:
case ENTITY_UPDATED_RPC_MESSAGE:
return switch (aiModelUpdateMsg.getMsgType()) {
case ENTITY_CREATED_RPC_MESSAGE, ENTITY_UPDATED_RPC_MESSAGE -> {
processAiModel(tenantId, aiModelId, aiModelUpdateMsg, edge);
return Futures.immediateFuture(null);
case UNRECOGNIZED:
default:
return handleUnsupportedMsgType(aiModelUpdateMsg.getMsgType());
}
yield Futures.immediateFuture(null);
}
case ENTITY_DELETED_RPC_MESSAGE -> {
deleteAiModel(tenantId, edge, aiModelId);
yield Futures.immediateFuture(null);
}
default -> handleUnsupportedMsgType(aiModelUpdateMsg.getMsgType());
};
} catch (DataValidationException e) {
return Futures.immediateFailedFuture(e);
} finally {
@ -95,36 +95,21 @@ public class AiModelEdgeProcessor extends BaseAiModelProcessor implements AiMode
return null;
}
@Override
public EdgeEventType getEdgeEventType() {
return EdgeEventType.AI_MODEL;
}
private void processAiModel(TenantId tenantId, AiModelId aiModelId, AiModelUpdateMsg aiModelUpdateMsg, Edge edge) {
Pair<Boolean, Boolean> resultPair = super.saveOrUpdateAiModel(tenantId, aiModelId, aiModelUpdateMsg);
Boolean wasCreated = resultPair.getFirst();
if (wasCreated) {
pushAiModelCreatedEventToRuleEngine(tenantId, edge, aiModelId);
Boolean created = resultPair.getFirst();
if (created) {
Optional<AiModel> aiModel = edgeCtx.getAiModelService().findAiModelById(tenantId, aiModelId);
aiModel.ifPresent(model -> pushEntityEventToRuleEngine(tenantId, edge, model, TbMsgType.ENTITY_CREATED));
}
Boolean nameWasUpdated = resultPair.getSecond();
if (nameWasUpdated) {
Boolean aiModelNameUpdated = resultPair.getSecond();
if (aiModelNameUpdated) {
saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.AI_MODEL, EdgeEventActionType.UPDATED, aiModelId, null);
}
}
private void pushAiModelCreatedEventToRuleEngine(TenantId tenantId, Edge edge, AiModelId aiModelId) {
try {
Optional<AiModel> aiModel = edgeCtx.getAiModelService().findAiModelById(tenantId, aiModelId);
if (aiModel.isPresent()) {
String aiModelAsString = JacksonUtil.toString(aiModel.get());
TbMsgMetaData msgMetaData = getEdgeActionTbMsgMetaData(edge, edge.getCustomerId());
pushEntityEventToRuleEngine(tenantId, aiModelId, edge.getCustomerId(), TbMsgType.ENTITY_CREATED, aiModelAsString, msgMetaData);
} else {
log.warn("[{}][{}] Failed to find aiModel", tenantId, aiModelId);
}
} catch (Exception e) {
log.warn("[{}][{}] Failed to push aiModel action to rule engine: {}", tenantId, aiModelId, TbMsgType.ENTITY_CREATED.name(), e);
}
@Override
public EdgeEventType getEdgeEventType() {
return EdgeEventType.AI_MODEL;
}
}

10
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/ai/BaseAiModelProcessor.java

@ -22,8 +22,10 @@ import org.springframework.data.util.Pair;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.ai.AiModel;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.id.AiModelId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.dao.service.DataValidator;
import org.thingsboard.server.gen.edge.v1.AiModelUpdateMsg;
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor;
@ -78,4 +80,12 @@ public abstract class BaseAiModelProcessor extends BaseEdgeProcessor {
return Pair.of(isCreated, isNameUpdated);
}
protected void deleteAiModel(TenantId tenantId, Edge edge, AiModelId aiModelId) {
Optional<AiModel> aiModel = edgeCtx.getAiModelService().findAiModelById(tenantId, aiModelId);
if (aiModel.isPresent()) {
edgeCtx.getAiModelService().deleteByTenantIdAndId(tenantId, aiModelId);
pushEntityEventToRuleEngine(tenantId, edge, aiModel.get(), TbMsgType.ENTITY_DELETED);
}
}
}

42
application/src/test/java/org/thingsboard/server/edge/AiModelEdgeTest.java

@ -32,8 +32,11 @@ import org.thingsboard.server.gen.edge.v1.UplinkResponseMsg;
import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import static org.awaitility.Awaitility.await;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
import static org.thingsboard.server.gen.edge.v1.UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE;
@DaoSqlTest
public class AiModelEdgeTest extends AbstractEdgeTest {
@ -42,7 +45,7 @@ public class AiModelEdgeTest extends AbstractEdgeTest {
private static final String UPDATED_AI_MODEL_NAME = "Updated Edge Test AiModel";
@Test
public void testAiModel_create_update_delete() throws Exception {
public void testAiModel_create_update_delete_fromCloud() throws Exception {
// create AiModel
AiModel aiModel = createSimpleAiModel(DEFAULT_AI_MODEL_NAME);
@ -91,26 +94,29 @@ public class AiModelEdgeTest extends AbstractEdgeTest {
}
@Test
public void testSendAiModelToCloud() throws Exception {
AiModel aiModel = createSimpleAiModel(DEFAULT_AI_MODEL_NAME);
UUID uuid = Uuids.timeBased();
UplinkMsg uplinkMsg = getUplinkMsg(uuid, aiModel, UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE);
checkAiModelOnCloud(uplinkMsg, uuid, aiModel.getName());
}
@Test
public void testUpdateAiModelNameOnCloud() throws Exception {
public void testAiModel_create_update_delete_toCloud() throws Exception {
// create
AiModel aiModel = createSimpleAiModel(DEFAULT_AI_MODEL_NAME);
UUID uuid = Uuids.timeBased();
UplinkMsg uplinkMsg = getUplinkMsg(uuid, aiModel, UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE);
checkAiModelOnCloud(uplinkMsg, uuid, aiModel.getName());
// update
aiModel.setName(UPDATED_AI_MODEL_NAME);
UplinkMsg updatedUplinkMsg = getUplinkMsg(uuid, aiModel, UpdateMsgType.ENTITY_UPDATED_RPC_MESSAGE);
checkAiModelOnCloud(updatedUplinkMsg, uuid, aiModel.getName());
// delete
UplinkMsg deleteUplinkMsg = getDeleteUplinkMsg(uuid);
edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(deleteUplinkMsg);
Assert.assertTrue(edgeImitator.waitForResponses());
await().atMost(30, TimeUnit.SECONDS).untilAsserted(() ->
doGet("/api/ai/model/" + uuid, AiModel.class, status().isNotFound())
);
}
@Test
@ -164,6 +170,20 @@ public class AiModelEdgeTest extends AbstractEdgeTest {
return aiModel;
}
private UplinkMsg getDeleteUplinkMsg(UUID uuid) throws InvalidProtocolBufferException {
UplinkMsg.Builder upLinkMsgBuilder = UplinkMsg.newBuilder();
AiModelUpdateMsg.Builder aiModelDeleteMsgBuilder = AiModelUpdateMsg.newBuilder();
aiModelDeleteMsgBuilder.setMsgType(ENTITY_DELETED_RPC_MESSAGE);
aiModelDeleteMsgBuilder.setIdMSB(uuid.getMostSignificantBits());
aiModelDeleteMsgBuilder.setIdLSB(uuid.getLeastSignificantBits());
testAutoGeneratedCodeByProtobuf(aiModelDeleteMsgBuilder);
upLinkMsgBuilder.addAiModelUpdateMsg(aiModelDeleteMsgBuilder.build());
testAutoGeneratedCodeByProtobuf(upLinkMsgBuilder);
return upLinkMsgBuilder.build();
}
private UplinkMsg getUplinkMsg(UUID uuid, AiModel aiModel, UpdateMsgType updateMsgType) throws InvalidProtocolBufferException {
UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
AiModelUpdateMsg.Builder aiModelUpdateMsgBuilder = AiModelUpdateMsg.newBuilder();

Loading…
Cancel
Save